SeaTunnel FtpFile Sink 连接器完全指南:将数据写入 FTP 服务器的配置、参数与实现原理

原创2026-09-27 09:00:26263 阅读
文章标签:数据工程大数据批处理流处理

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 "k0={k0}={v0}/k1={k1}={v1}/.../kn={kn}={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"
}

行为推演(结合前文参数语义):

  1. 数据先写入临时目录 /data/ftp/seatunnel/tmp,提交时 mv 到目标目录 /data/ftp/seatunnel/job1;
  2. 按 age 字段分区,分区目录表达式为 ${k0}=${v0},例如 age=25/;
  3. is_partition_field_write_in_file = true,分区列 age 也会写入数据文件;
  4. 文件名形如 ${transactionId}_2025.01.01.txt(${now} 按 yyyy.MM.dd 格式化,text 后缀为 txt);
  5. 仅写入 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 专属配置之外的能力均可复用公共参数体系。

登录后查看全文
seatunnel