SeaTunnel FtpFile Sink 连接器完全指南:将数据写入 FTP 服务器的配置、参数与实现原理
SeaTunnel FtpFile Sink 连接器完全指南:将数据写入 FTP 服务器的配置、参数与实现原理
FtpFile Sink 是 Apache SeaTunnel 内置的文件型输出连接器,负责把上游 Source 或 Transform 产出的数据以 text、csv、parquet、orc、json、excel、xml、binary 等格式写入 FTP 服务器指定目录。本文以官方文档 docs/en/connector-v2/sink/FtpFile.md 为主体,结合仓库中 connector-file-ftp 与 connector-file-base 模块源码,完整讲解该连接器的全部配置参数、运行环境要求、分区与事务机制、完整配置示例,以及 2PC 提交与 FTP 文件系统适配的底层实现。读完本文,你将能够独立完成一个从任意 Source 到 FTP 的 SeaTunnel 作业的配置、调优与排障。
连接器简介与适用场景
FtpFile 连接器的职责非常聚焦:Output data to Ftp——将 SeaTunnel 管道中的 SeaTunnelRow 数据流式写入 FTP 服务器。它与同模块的 FtpFileSource 形成一对读写组合,适用于数据落地 FTP、与外部系统通过 FTP 交换文件的场景。
从源码看,连接器插件名来自 FileSystemType.FTP,Sink 侧实现为 FtpFileSink,它继承自文件连接器公共基类 BaseFileSink,因此所有文件型 Sink 的通用能力(分区、事务、格式、压缩)它都天然具备。
运行环境前提
官方文档对运行环境有明确提示,务必先满足:
- 使用 Spark / Flink 运行时:集群必须已经集成 Hadoop,官方实测的 Hadoop 版本为 2.x。如果集群未集成 Hadoop,FTP 文件系统适配层将无法加载,作业会启动失败。
- 使用 SeaTunnel Engine 运行时:无需额外操作,下载并安装 SeaTunnel Engine 时已自动集成 Hadoop 相关 jar。可通过检查
${SEATUNNEL_HOME}/lib目录下的 jar 包来确认。
这一前提在源码中也得到印证:FtpFileSink 通过 Hadoop 配置体系来桥接 FTP 协议(见下文“FTP 文件系统适配”),本质上仍依赖 Hadoop FileSystem 抽象。
核心特性(Key features)
| 特性 | 说明 |
|---|---|
| exactly-once | 默认通过两阶段提交(2PC)保证 exactly-once,详见 connector-v2-features |
| 文件格式 | text、csv、parquet、orc、json、excel、xml、binary 共 8 种,均完整支持 |
其中 exactly-once 承诺的落地实现可追溯到 BaseFileSink:它实现了 SeaTunnelSink<SeaTunnelRow, FileSinkState, FileCommitInfo, FileAggregatedCommitInfo> 泛型签名,分别提供了 createWriter、restoreWriter(从 checkpoint 状态恢复写)以及 createAggregatedCommitter 返回 FileSinkAggregatedCommitter,这正是 2PC 中预提交与最终提交的载体(详见下文“源码级原理”)。
完整 Options 参数表
下表完整列出官方文档定义的连接器参数。标注 “Required = yes” 的 6 个参数(host、port、user、password、path、tmp_path)在配置缺失时,连接器会直接抛出配置校验异常。
| 名称 | 类型 | 是否必填 | 默认值 | 说明 |
|---|---|---|---|---|
| host | string | yes | - | FTP 服务器地址 |
| port | int | yes | - | FTP 服务器端口 |
| user | string | yes | - | FTP 用户名 |
| password | string | yes | - | FTP 密码 |
| path | string | yes | - | 目标目录路径 |
| tmp_path | string | yes | /tmp/seatunnel | 结果文件先写入临时目录,再通过 mv 提交到目标目录;必须是 FTP 目录 |
| connection_mode | string | no | active_local | FTP 连接模式 |
| custom_filename | boolean | no | false | 是否自定义文件名 |
| file_name_expression | string | no | "${transactionId}" | 仅当 custom_filename 为 true 时生效 |
| filename_time_format | string | no | "yyyy.MM.dd" | 仅当 custom_filename 为 true 时生效 |
| file_format_type | string | no | "csv" | 输出文件格式 |
| field_delimiter | string | no | '\001' | 仅 text 格式生效,列分隔符 |
| row_delimiter | string | no | "\n" | 仅 text 格式生效,行分隔符 |
| have_partition | boolean | no | false | 是否开启分区处理 |
| partition_by | array | no | - | 仅 have_partition 为 true 时生效,按所选字段分区 |
| partition_dir_expression | string | no | "{v0}/{v1}/.../{vn}/" | 仅 have_partition 为 true 时生效 |
| is_partition_field_write_in_file | boolean | no | false | 仅 have_partition 为 true 时生效 |
| sink_columns | array | no | (空) | 为空时写入全部字段;顺序决定文件实际写入顺序 |
| is_enable_transaction | boolean | no | true | 是否开启事务保护 |
| batch_size | int | no | 1000000 | 单个文件最大行数 |
| compress_codec | string | no | none | 压缩编解码器 |
| common-options | object | no | - | Sink 公共参数 |
| max_rows_in_memory | int | no | - | 仅 excel 格式生效,内存中缓存的最大数据条数 |
| sheet_name | string | no | Sheet${随机数} | 仅 excel 格式生效,工作簿工作表名 |
| xml_root_tag | string | no | RECORDS | 仅 xml 格式生效,XML 根元素标签 |
| xml_row_tag | string | no | RECORD | 仅 xml 格式生效,XML 数据行标签 |
| xml_use_attr_format | boolean | no | - | 仅 xml 格式生效,是否使用标签属性格式 |
| parquet_avro_write_timestamp_as_int96 | boolean | no | false | 仅 parquet 格式生效 |
| parquet_avro_write_fixed_as_int96 | array | no | - | 仅 parquet 格式生效 |
| encoding | string | no | "UTF-8" | 仅 json、text、csv、xml 格式生效 |
上表默认值均可在 BaseSinkConfig 中逐一找到对应常量(如 DEFAULT_TMP_PATH = "/tmp/seatunnel"、DEFAULT_FILE_NAME_EXPRESSION = "${transactionId}"、DEFAULT_BATCH_SIZE = 1000000),参数定义与文档完全一致,可作为排障时的权威依据。
连接参数:host / port / user / password / path / tmp_path / connection_mode
连接五要素与目标路径
host[string]:FTP 服务器主机地址,必填。port[int]:FTP 服务端口,必填(通常为 21)。user[string]:FTP 登录用户名,必填。password[string]:FTP 登录密码,必填。path[string]:目标写入目录,必填。最终文件会生成在该目录下。tmp_path[string]:临时目录,必填,默认/tmp/seatunnel。写入机制是:结果文件先写进临时目录,完成后用mv命令整体提交到目标目录,这保证了目标目录不会出现半成品文件。由于mv发生在 FTP 服务器侧,tmp_path必须是一个 FTP 可达的目录,且建议与path位于同一文件系统/挂载范围内。
在 FtpFileSink.java 的 prepare 阶段,会通过 CheckConfigUtil.checkAllExists 强制校验 host、port、user、password 四项(path 的校验由工厂的 OptionRule 完成),任一缺失都会抛出 FileConnectorException(错误码 CONFIG_VALIDATION_FAILED)。
connection_mode:主动/被动连接模式
- 默认值
active_local,可选值:active_local、passive_local。 - 对应枚举定义在 FtpConnectionMode,其模式字符串与 Apache Commons Net 的连接模式一致:
active_local:FTP 主动模式(服务器主动连接客户端的数据端口);passive_local:FTP 被动模式(客户端主动连接服务器数据端口)。
生产环境中,客户端位于 NAT/防火墙之后时通常需要切换为 passive_local,否则数据通道可能建立失败。
FTP 文件系统适配(源码级)
FtpFile 连接器通过 Hadoop FileSystem 抽象访问 FTP 协议,这一桥接逻辑在 FtpConf 中实现:
- 文件系统实现类:
org.apache.seatunnel.connectors.seatunnel.file.ftp.system.SeaTunnelFTPFileSystem; - 协议 scheme:
ftp; - 默认文件系统地址
defaultFS由配置拼接为ftp://{host}:{port}; - 认证与连接参数被映射为 Hadoop 配置项写入
hadoopConf:fs.ftp.user.{host}→ FTP 用户名;fs.ftp.password.{host}→ FTP 密码;fs.ftp.connection.mode→ 连接模式(仅在配置了connection_mode时写入)。
也就是说,只要让 BaseFileSinkWriter 拿到这份 HadoopConf,写数据、建目录、mv 提交等所有文件操作都复用了文件连接器公共模块的通用实现,FTP 特有逻辑被收敛在配置映射这一层。这也是“先临时目录、后 mv 提交”机制能够透明工作的原因。
文件名与时间格式控制
custom_filename [boolean]
是否自定义文件名,默认 false。
file_name_expression [string]
仅当 custom_filename = true 时生效。描述要写入 path 的文件名表达式,支持变量 ${now} 与 ${uuid},例如 test_${uuid}_${now}:
${now}表示当前时间,其格式由filename_time_format定义;- 需要注意:若
is_enable_transaction = true,会自动在文件名头部追加${transactionId}_前缀,用于事务标识(例如transactionId_20250101_xxxx这类形态)。
默认表达式为 ${transactionId}(见 BaseSinkConfig 的 DEFAULT_FILE_NAME_EXPRESSION)。
filename_time_format [string]
仅当 custom_filename = true 时生效。当 file_name_expression 中包含 xxxx-${now} 时,该参数指定时间格式,默认 yyyy.MM.dd。常用符号如下:
| 符号 | 含义 |
|---|---|
| y | 年(Year) |
| M | 月(Month) |
| d | 日(Day of month) |
| H | 小时(Hour in day 0-23) |
| m | 分钟(Minute in hour) |
| s | 秒(Second in minute) |
输出文件格式与行列分隔
file_format_type [string]
支持的格式:text、csv、parquet、orc、json、excel、xml、binary。默认值为 csv。
请特别注意:最终文件名会以格式后缀结尾,其中 text 格式的后缀是 txt,即 file_format_type = "text" 时生成 xxx.txt 文件。
field_delimiter [string]
行内列之间的分隔符,仅 text 格式需要。默认 '\001'(SOH 控制符,即 \u0001),这是 Hadoop 生态 text 文件常见的默认列分隔。
row_delimiter [string]
文件内行与行之间的分隔符,仅 text 格式需要,默认 "\n"。
历史上(2.3.0-beta 之前)存在无法从配置文件解析
'\t'作为分隔符的缺陷,相关修复可见文档 Changelog 部分,配置时建议直接书写真实转义字符。
encoding [string]
仅 json、text、csv、xml 四种格式生效。写入文件的编码,默认 "UTF-8",该参数最终通过 Charset.forName(encoding) 解析。
压缩配置:compress_codec
compress_codec 指定输出文件的压缩编解码器,默认 none。不同格式支持的压缩算法不同,务必按下表对照使用(excel 格式不支持任何压缩):
| 文件格式 | 支持压缩算法 |
|---|---|
| txt | lzo、none |
| json | lzo、none |
| csv | lzo、none |
| orc | lzo、snappy、lz4、zlib、none |
| parquet | lzo、snappy、lz4、gzip、brotli、zstd、none |
源码 BaseSinkConfig 中以 TXT_COMPRESS、PARQUET_COMPRESS、ORC_COMPRESS 三个单选项分别约束了 txt/json/csv、parquet、orc 的可选压缩集,与文档表格完全一致;同时 FtpFileSinkFactory 的 OptionRule 也按 file_format_type 做了条件绑定,配置不合法时工厂校验即会报错。
分区写入
have_partition [boolean]
是否开启分区处理,默认 false。
partition_by [array]
仅当 have_partition = true 时生效。按所选字段进行分区,支持多字段组合。
partition_dir_expression [string]
仅当 have_partition = true 时生效。指定 partition_by 后,连接器会根据分区信息生成对应的分区目录,最终文件落在分区目录内。默认表达式为 ${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/,其中 k0 是第一个分区字段名,v0 是第一个分区字段的值,以此类推。例如按 age 分区时,典型结果是 age=25/ 这样的目录层级。
is_partition_field_write_in_file [boolean]
仅当 have_partition = true 时生效。若为 true,分区字段及其值也会写入数据文件内部;若为 false,分区字段只体现在目录结构中。官方文档特别举例:如果要产出 Hive 数据文件,该值应设为 false(Hive 分区约定分区列不在数据文件内重复出现)。
写入列与事务保护
sink_columns [array]
需要写入文件的列,默认值为 Source 或 Transform 输出的全部列。数组中字段的顺序即文件实际写入顺序,可用于调整列序或裁剪列。
is_enable_transaction [boolean]
- 默认
true。开启后,连接器保证写入目标目录的数据不丢失、不重复。 - 开启时,文件名头部会自动追加
${transactionId}_前缀。 - 官方文档明确:当前只支持
true。
batch_size [int]
单个文件的最大行数,默认 1000000。在 SeaTunnel Engine 中,一个文件的实际行数由 batch_size 与 checkpoint.interval 共同决定:
- 若
checkpoint.interval足够大,sink writer 会持续写入当前文件,直到行数超过batch_size才切换新文件; - 若
checkpoint.interval较小,则每次 checkpoint 触发时,sink writer 就会创建一个新文件(以 checkpoint 边界切分文件,便于与事务提交对齐)。
因此,想要控制文件大小,应同时调大 checkpoint.interval 并设置合适的 batch_size。
Excel / XML / Parquet 专属参数
max_rows_in_memory [int]
仅 excel 格式生效。内存中可缓存的最大数据条数,用于控制写 Excel 时的内存占用。
sheet_name [string]
仅 excel 格式生效。写入的工作簿工作表名,默认值为 Sheet${随机数}。
xml_root_tag [string]
仅 xml 格式生效。指定 XML 文件的根元素标签名,默认 RECORDS。
xml_row_tag [string]
仅 xml 格式生效。指定 XML 文件中数据行的标签名,默认 RECORD。
xml_use_attr_format [boolean]
仅 xml 格式生效。是否使用标签属性格式处理数据。
parquet_avro_write_timestamp_as_int96 [boolean]
仅 parquet 格式生效,默认 false。支持将 timestamp 类型写为 Parquet INT96。
parquet_avro_write_fixed_as_int96 [array]
仅 parquet 格式生效。支持将 12 字节字段写为 Parquet INT96。
完整配置示例
示例一:text 格式的基础配置
最简单的 text 落地配置(官方文档原始示例):
FtpFile {
host = "xxx.xxx.xxx.xxx"
port = 21
user = "username"
password = "password"
path = "/data/ftp"
file_format_type = "text"
field_delimiter = "\t"
row_delimiter = "\n"
sink_columns = ["name","age"]
}
该配置将上游数据以 name、age 两列顺序写入 /data/ftp 下的 *.txt 文件,列间以 \t 分隔,行间以 \n 分隔,未开启自定义文件名(默认使用 ${transactionId} 表达式并带事务前缀),未开启分区。
示例二:分区 + 自定义文件名 + 列裁剪的组合配置
同时启用 have_partition、custom_filename 与 sink_columns 的完整配置(官方文档原始示例):
FtpFile {
host = "xxx.xxx.xxx.xxx"
port = 21
user = "username"
password = "password"
path = "/data/ftp/seatunnel/job1"
tmp_path = "/data/ftp/seatunnel/tmp"
file_format_type = "text"
field_delimiter = "\t"
row_delimiter = "\n"
have_partition = true
partition_by = ["age"]
partition_dir_expression = "${k0}=${v0}"
is_partition_field_write_in_file = true
custom_filename = true
file_name_expression = "${transactionId}_${now}"
sink_columns = ["name","age"]
filename_time_format = "yyyy.MM.dd"
}
行为推演(结合前文参数语义):
- 数据先写入临时目录
/data/ftp/seatunnel/tmp,提交时mv到目标目录/data/ftp/seatunnel/job1; - 按
age字段分区,分区目录表达式为${k0}=${v0},例如age=25/; is_partition_field_write_in_file = true,分区列age也会写入数据文件;- 文件名形如
${transactionId}_2025.01.01.txt(${now}按yyyy.MM.dd格式化,text 后缀为 txt); - 仅写入
name、age两列,顺序按数组声明。
源码级工作原理与验证
Sink 主类与配置校验
FtpFileSink 通过 @AutoService(SeaTunnelSink.class) 注册,插件名为 FileSystemType.FTP 对应的字符串。其 prepare 方法强制校验 host/port/user/password,随后调用 FtpConf.buildWithConfig 构造 Hadoop 配置,其余逻辑全部继承自 BaseFileSink。
2PC exactly-once 的实现路径
从 BaseFileSink 的接口实现可以看出完整的提交链路:
createWriter/restoreWriter:创建或恢复BaseFileSinkWriter,写入发生在临时路径;createAggregatedCommitter:返回FileSinkAggregatedCommitter,在 checkpoint 提交阶段统一把各 writer 的临时文件mv到目标目录;- 三个序列化器(
FileCommitInfo、FileAggregatedCommitInfo、FileSinkState)保证 checkpoint 状态可持久化、可恢复。
这套“写临时文件 → checkpoint 预提交 → 聚合提交(mv 到目标目录)”的机制,正是“数据不丢失、不重复”与文件名带 transactionId 前缀的原因。文档中的“Next version”Changelog 提到修复了“从 states 恢复 writer 时获取 transaction 直接失败”的问题,正是这条链路中的状态恢复环节。
工厂参数约束与单元测试
FtpFileSinkFactory 通过 OptionRule 声明了参数依赖关系:path、host、port、user、password 为必填;field_delimiter、row_delimiter、压缩等参数仅在对应 file_format_type 下生效;file_name_expression、filename_time_format 仅在 custom_filename = true 时生效;partition_by 等仅在 have_partition = true 时生效。仓库中的单元测试 FtpFileFactoryTest 对 Sink 与 Source 工厂的 optionRule() 均做了非空断言,确保参数规则可用。
Sink 与 Source 的能力差异(注意)
同一 FTP 模块下的 FtpFileSource 在读取侧有明确限制:Ftp 文件 Source 仅支持读取 text、csv、json 三种格式,配置为 orc、parquet 会直接抛出 ILLEGAL_ARGUMENT 异常;而本篇文章所讲的 FtpFile Sink 则支持 8 种格式的写入。在规划“FTP 读 → FTP 写”的自同步任务时,需注意两侧格式能力并不对称。
版本演进 Changelog
- 2.2.0-beta(2022-09-26):新增 Ftp File Sink 连接器。
- 2.3.0-beta(2022-10-20):三处缺陷修复——Windows 环境下路径不正确(PR 2980);文件系统获取报错(PR 3117);配置文件无法解析
'\t'作为分隔符(PR 3083)。 - Next version:
- 缺陷修复(PR 3258):上游字段为 null 时抛 NullPointerException;sink 列映射失败;从 states 恢复 writer 时直接获取 transaction 失败;
- 能力增强:支持为每个文件设置 batch_size(PR 3625);支持文件压缩(PR 3899)。
常见问题速查
| 现象 | 排查方向 |
|---|---|
| 作业启动报配置校验失败 | 检查 host、port、user、password、path 是否齐全,对应 FtpFileSink.prepare 与 OptionRule 的必填声明 |
| Spark/Flink 下无法写入 | 确认集群已集成 Hadoop 2.x,FTP 文件系统适配依赖 Hadoop FileSystem 抽象 |
| 数据通道建立失败 | 客户端处于 NAT/防火墙后时,将 connection_mode 改为 passive_local |
| 文件一直不落到目标目录 | 确认 tmp_path 为 FTP 可达目录,mv 提交依赖该目录可写;同时理解文件切分由 batch_size 与 checkpoint.interval 共同决定 |
| 目标目录出现未知前缀文件名 | is_enable_transaction = true 时文件名自动带 ${transactionId}_ 前缀,属预期行为 |
| 写 Hive 数据文件列重复 | 分区写入时确认 is_partition_field_write_in_file = false |
| 压缩参数报错 | 对照“压缩配置”表格,excel 不支持压缩,orc/parquet 各有允许的算法白名单 |
更多 Sink 公共参数(如 date_format、datetime_format、time_format、表头写入等)可查阅 Sink Common Options,FtpFile 专属配置之外的能力均可复用公共参数体系。