SeaTunnel FtpFile Sink Connector 使用指南:将数据输出到 FTP 服务器
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
SeaTunnel 的FtpFileSink 插件用于将上游 Source 或 Transform 产出的数据写入 FTP 服务器,支持 text、csv、parquet、orc、json、excel、xml、binary 八种文件格式,并内置 2PC 事务提交、分区目录、自定义文件名、schema 演进(CDC 场景)等能力。读完本文,你将掌握 FtpFile Sink 的全部配置项语义、事务与临时目录提交机制,以及文本、分区、多表写入和 SFTP 四类可直接落地的作业配置。
插件定位与能力总览
FtpFile是 SeaTunnel 文件类 Sink 家族中的一员,其核心职责是把数据以文件形式"投递"到 FTP 服务器。从源码看,该插件的 Sink 主体非常轻量——FtpFileSink.java 直接继承文件类 Sink 的公共基类BaseMultipleTableFileSink,并通过FtpConf.buildWithConfig将用户配置转换为 Hadoop 兼容的文件系统参数,最终由 SeaTunnelFTPFileSystem.java(基于 Apache Commons Net 的 FTPClient 实现)完成真实的读写。也就是说,FTP 只被当作一个"文件系统"看待,文件类 Sink 的通用能力(事务、分区、格式、压缩等)全部被继承下来。
插件支持的关键特性如下:
- 多模态(multimodal):以二进制文件格式读写任意格式的文件,例如视频、图片等,任何文件都可以同步到目标位置;
- exactly-once(精确一次):默认使用 2PC 提交保证数据不丢不重;
- 多表写入(multiple table write):一个作业可将多个表分别写入各自目录;
- 文件格式:text、csv、parquet、orc、json、excel、xml、binary。
运行前提(来自官方文档提示)
- 如果使用 Spark/Flink 作为引擎,必须保证你的 Spark/Flink 集群已集成 Hadoop(官方测试过的 Hadoop 版本为 2.x);
- 如果使用 SeaTunnel Engine,则安装包已自动集成 hadoop jar,可通过检查
${SEATUNNEL_HOME}/lib目录下的 jar 包确认。
相关特性说明可参阅 connector-v2-features,通用 Sink 参数见 sink-common-options,该插件的变更记录见 connector-file-ftp changelog。
配置项总览
下表为 FtpFile Sink 支持的全部参数(与官方文档保持一致):
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
| 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 连接模式 |
| remote_verification_enabled | boolean | no | true | 是否启用 FTP 数据通道的远程主机校验 |
| control_encoding | string | no | UTF-8 | FTP 控制连接字符编码,对含空格或非 ASCII 字符的路径很有用 |
| 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" | 输出文件格式 |
| filename_extension | string | no | - | 用自定义扩展名覆盖默认扩展名,如.xml、.json、dat、.customtype |
| field_delimiter | string | no | text 为 '\001',csv 为 ',' | 仅当 file_format_type 为 text 和 csv 时生效 |
| row_delimiter | string | no | "\n" | 仅当 file_format_type 为 text、csv 和 json 时生效 |
| have_partition | boolean | no | false | 是否需要分区处理 |
| partition_by | array | no | - | 仅当 have_partition 为 true 时生效 |
| partition_dir_expression | string | no | "${k0}=${v0}/${k1}=${v1}/.../${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 | - | 仅当 file_format_type 为 excel 时生效 |
| sheet_max_rows | int | no | 1048576 | 仅当 file_format_type 为 excel 时生效 |
| sheet_name | string | no | Sheet${随机数} | 仅当 file_format_type 为 excel 时生效 |
| csv_string_quote_mode | enum | no | MINIMAL | 仅当 file_format 为 csv 时生效 |
| xml_root_tag | string | no | RECORDS | 仅当 file_format 为 xml 时生效 |
| xml_row_tag | string | no | RECORD | 仅当 file_format 为 xml 时生效 |
| xml_use_attr_format | boolean | no | - | 仅当 file_format 为 xml 时生效 |
| single_file_mode | boolean | no | false | 每个并行度只输出一个文件;开启后 batch_size 不生效,输出文件名不带文件块后缀 |
| create_empty_file_when_no_data | boolean | no | false | 上游无数据同步时仍生成对应数据文件 |
| parquet_avro_write_timestamp_as_int96 | boolean | no | false | 仅当 file_format 为 parquet 时生效 |
| parquet_avro_write_fixed_as_int96 | array | no | - | 仅当 file_format 为 parquet 时生效 |
| enable_header_write | boolean | no | false | 仅当 file_format_type 为 text、csv 时生效:false 不写表头,true 写表头 |
| encoding | string | no | "UTF-8" | 仅当 file_format_type 为 json、text、csv、xml 时生效 |
| schema_evolution_enabled | boolean | no | false | 为 CDC 管道启用 schema 演进,支持 ADD/DROP/RENAME/MODIFY 列事件,无需重启作业;binary 格式不支持 |
| schema_save_mode | string | no | CREATE_SCHEMA_WHEN_NOT_EXIST | 已有目录的处理方式 |
| data_save_mode | string | no | APPEND_DATA | 已有数据的处理方式 |
| multi_table_sink_replica | int | no | 1 | 多表 Sink 作业中每个表使用的 writer 副本数 |
连接参数详解
host / port / user / password
这四个参数是必填的连接基础信息:FTP 服务器主机、端口、登录用户名和密码。在源码 FtpFileBaseOptions.java 中,它们被定义为无默认值的必填 Option;FtpConf.java 会据此构造ftp://host:port的默认文件系统地址,并将用户名、密码写入fs.ftp.user.<host>、fs.ftp.password.<host>两个 Hadoop 配置键,供底层SeaTunnelFTPFileSystem在建立连接时使用。
path
目标目录路径(必填)。支持多表场景下在路径中嵌入${table_name}占位符,让每个上游表写入独立的 FTP 目录,例如/data/ftp/job1/${table_name}。
tmp_path
临时目录(必填),默认值为/tmp/seatunnel。写入流程为:结果文件先写入临时目录,提交时通过mv(底层对应 FTPrename)将临时目录下的文件移动到目标目录。因此该目录必须是一个真实存在的 FTP 目录。这样做的好处是配合 2PC 提交,保证最终目录中不会出现"半成品"文件。
connection_mode
FTP 数据通道连接模式,默认active_local,支持active_local和passive_local两种取值,对应源码中的枚举 FtpConnectionMode.java。在 SeaTunnelFTPFileSystem.java 的连接建立逻辑中可以看到:
active_local:进入本地主动模式,并会创建一个/ .ftptest<时间戳>测试目录来验证主动模式是否可用;若失败则自动降级切换为被动模式,并同步更新配置;passive_local:直接进入本地被动模式。
无论哪种模式,连接成功后都会设置二进制文件类型(BINARY_FILE_TYPE)、1MB 缓冲区(DEFAULT_BUFFER_SIZE)以及块传输模式(BLOCK_TRANSFER_MODE),因此 FTP 通道本身始终以二进制方式传输数据,文本/CSV 等格式化的差异由上层 Writer 处理。
remote_verification_enabled
是否启用 FTP 数据通道的远程主机校验,默认true,对应FTPClient.setRemoteVerificationEnabled。当 FTP 服务器位于 NAT 之后或数据连接地址与实际主机不一致时,可以关闭该校验。
control_encoding
FTP 控制连接的字符编码,默认UTF-8。源码在建立连接前调用client.setControlEncoding(controlEncoding)(见 SeaTunnelFTPFileSystem.java),该设置对路径中包含空格、特殊字符或非 ASCII 字符的场景至关重要。除非你的 FTP 服务器要求其他控制通道编码,否则保持UTF-8即可。
文件命名与格式配置
custom_filename / file_name_expression / filename_time_format
custom_filename:是否自定义文件名,默认false;file_name_expression:仅当custom_filename为true时生效,描述将要写入path的文件名表达式。可以在表达式中使用变量${now}或${uuid},例如test_${uuid}_${now}。${now}表示当前时间,其格式由filename_time_format定义;- 注意:如果
is_enable_transaction为true,插件会自动在文件名头部加上${transactionId}_前缀。
filename_time_format默认值为yyyy.MM.dd。常用时间格式符号如下:
| Symbol | Description |
|---|---|
| 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
支持的文件类型:text、csv、parquet、orc、json、excel、xml、binary。最终文件名的后缀与文件格式类型对应,其中 text 文件的默认后缀为txt。默认值为csv。可以使用filename_extension覆盖默认扩展名(例如.xml、.json、dat、.customtype)。
field_delimiter / row_delimiter
field_delimiter:一行数据中列之间的分隔符,仅 text 和 csv 格式需要。默认值:text 为'\001',csv 为',';row_delimiter:文件中行与行之间的分隔符,仅 text、csv 和 json 格式需要,默认"\n"。
enable_header_write / encoding
enable_header_write:仅 text、csv 格式生效,false不写表头,true写表头,默认false;encoding:仅 json、text、csv、xml 格式生效,指定输出文件的字符编码(如 UTF-8、ISO-8859-1),默认"UTF-8"。该参数会通过Charset.forName(encoding)解析(对应 FileBaseOptions.java 中的ENCODING定义)。
分区写入配置
have_partition:是否启用分区处理,默认false;partition_by:仅当have_partition为true时生效,按所选字段对数据进行分区;partition_dir_expression:仅当have_partition为true时生效。指定partition_by后,插件会根据分区信息生成对应的分区目录,最终文件写入分区目录内。默认表达式为${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/,其中k0是第一个分区字段,v0是其取值;is_partition_field_write_in_file:仅当have_partition为true时生效。若为true,分区字段及其值也会写入数据文件;若想写出 Hive 数据文件,该值应设为false。
列、事务与文件拆分
sink_columns
指定需要写入文件的列,默认值为从Transform或Source获取的全部列。字段的排列顺序决定文件实际写入的顺序。
is_enable_transaction
若为true(默认值),插件保证数据写入目标目录时不丢失、不重复。注意:开启后会自动在文件名头部添加${transactionId}_。目前仅支持true。
其底层实现依托文件类 Sink 公共基类提供的 2PC 提交:Writer 将文件先写到临时目录(tmp_path),checkpoint 触发后通过 FileSinkAggregatedCommitter.java 对每个事务的临时文件执行移动/重命名(FTP 层面对应rename)到目标目录,失败的文件进入重试列表,从而保证精确一次语义。
batch_size
单个文件的最大行数,默认1000000。对于 SeaTunnel Engine,文件行数由batch_size和checkpoint.interval共同决定:如果checkpoint.interval足够大,writer 会持续写入直到文件行数超过batch_size;如果checkpoint.interval较小,则每次新 checkpoint 触发时都会创建新文件。
single_file_mode / create_empty_file_when_no_data
single_file_mode:默认false。开启后每个并行度只输出一个文件,此时batch_size不再生效,输出文件名不带文件块后缀;create_empty_file_when_no_data:默认false。开启后,即使上游没有数据同步,也仍会生成对应的数据文件。
压缩、CSV/XML/Excel/Parquet 专属参数
compress_codec
文件压缩编码,默认none,按格式支持如下:
- txt:
lzo、none - json:
lzo、none - csv:
lzo、none - orc:
lzo、snappy、lz4、zlib、none - parquet:
lzo、snappy、lz4、gzip、brotli、zstd、none
提示:excel 格式不支持任何压缩格式。
Excel 相关(file_format_type = excel)
max_rows_in_memory:内存中可缓存的最大数据条数;sheet_max_rows:每个 sheet 的最大行数,默认1048576;sheet_name:工作簿中写入的 sheet 名称,默认Sheet${随机数}。
CSV 引号模式(file_format_type = csv)
csv_string_quote_mode的可选值及语义:
ALL:所有 String 字段都加引号;MINIMAL:仅对包含特殊字符(如字段分隔符、引号字符或行分隔符串中的任意字符)的字段加引号;NONE:从不加引号。当数据中出现分隔符时,printer 会在其前面加上转义字符;若未设置转义字符,格式校验将抛出异常。
XML 相关(file_format_type = xml)
xml_root_tag:指定 XML 文件根元素标签名,默认RECORDS;xml_row_tag:指定 XML 文件数据行标签名,默认RECORD;xml_use_attr_format:指定是否使用标签属性格式处理数据。
Parquet 相关(file_format_type = parquet)
parquet_avro_write_timestamp_as_int96:支持将时间戳写入 Parquet INT96,默认false;parquet_avro_write_fixed_as_int96:支持从 12 字节字段写入 Parquet INT96。
目录与数据的 Save Mode
schema_save_mode(已有目录处理方式)
RECREATE_SCHEMA:目录不存在则创建,目录已存在则删除后重建;CREATE_SCHEMA_WHEN_NOT_EXIST(默认):目录不存在则创建,目录已存在则跳过;ERROR_WHEN_SCHEMA_NOT_EXIST:目录不存在时报错;IGNORE:忽略对该表的处理。
data_save_mode(已有数据处理方式)
DROP_DATA:保留目录但删除数据文件;APPEND_DATA(默认):保留目录、保留数据文件;ERROR_WHEN_DATA_EXISTS:存在数据文件时报错。
Schema 演进(CDC 管道场景)
schema_evolution_enabled默认false。设为true后,文件 Sink 可以在运行期处理 CDC 的 schema 变更事件(ADD COLUMN、DROP COLUMN、RENAME COLUMN、MODIFY COLUMN),无需重启作业;每次 schema 变更时,当前输出文件会被关闭,并以更新后的 schema 打开新文件。
使用约束与限制:
支持的格式:除
binary外的所有文件格式。若file_format_type = binary且开启该选项,作业启动时会在配置校验阶段直接失败;分区约束:当
have_partition = true时,不允许删除partition_by中列出的列,否则会快速失败。分区列必须在 schema 变更期间保持稳定;当
schema_evolution_enabled = false(默认)时:若上游 CDC Source 开启了schema-changes.enabled = true且AlterTableEvent到达 Sink,作业会立即抛出如下可操作的错误:Received AlterTableEvent but schema_evolution_enabled=false at this sink. Either set schema_evolution_enabled=true to handle schema changes, or set schema-changes.enabled=false at the CDC source to suppress them.使用默认 CDC Source 配置(
schema-changes.enabled = false)的用户完全不受影响;已知限制:schema 变更与 checkpoint 并非原子操作。如果作业在文件轮转与 schema 元数据更新之间的极窄窗口内崩溃,恢复后写入的行可能仍使用变更前的 schema。这是 SeaTunnel 其他 Sink 共有的架构性差距,若要获得"重启+DDL 完全正确"的语义,需要后续 CDC Source 侧的修复配合(另行跟踪)。
CDC 管道中的示例配置:
FtpFile { host = "xxx.xxx.xxx.xxx" port = 21 user = "username" password = "password" path = "/data/ftp/cdc/${table_name}" file_format_type = "parquet" schema_evolution_enabled = true }多表写入与 writer 副本
multi_table_sink_replica指定多表 Sink 作业中每个表使用的 writer 副本数,默认1;仅当每个表需要更高的 Sink writer 并行度时才调大。配合path中的${table_name}占位符,即可实现"上游多表、各自落盘到独立 FTP 目录"的目标。
完整示例
示例一: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"] }示例二:text 格式 + 分区 + 自定义文件名 + 指定列
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" }示例三:多表写入 + Save Mode
当上游 Source 有多个表、且每个表需要写入各自的 FTP 目录时,在path中嵌入${table_name}。schema_save_mode与data_save_mode决定写入前如何处理已有目录和文件:
FtpFile { host = "xxx.xxx.xxx.xxx" port = 21 user = "username" password = "password" path = "/data/ftp/seatunnel/job1/${table_name}" 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" schema_save_mode=RECREATE_SCHEMA data_save_mode=DROP_DATA }通过 SFTP 写入
FtpFileSink 除了ftp://之外还支持sftp://URI。认证与主机信任(host-key trust)的配置方式与 Source 侧一致——使用 SSH 密钥或密码,外加一个known_hosts文件(连接器不会自动信任未知主机):
sink { FtpFile { fs.defaultFS = "sftp://sftp.example.example.com:22" path = "/upload/landing/" user = "seatunnel" file_format_type = "parquet" ftp_properties = { "fs.sftp.user." = "seatunnel" "fs.sftp.keyfile" = "/etc/seatunnel/id_rsa" "fs.sftp.host" = "sftp.example.example.com" "fs.sftp.port" = "22" "fs.sftp.knownHosts" = "/etc/seatunnel/known_hosts" } } }常见问题排查要点
- 连接失败或登录失败:优先核对
host、port、user、password。底层SeaTunnelFTPFileSystem在登录失败时会抛出包含 reply code 的异常信息("Login failed on server - %s, port - %d as user '%s'"),可据此向 FTP 服务器管理员确认账户权限; - 主动/被动模式问题:FTP 服务器位于 NAT 或防火墙之后时,主动模式(默认)的数据通道可能无法建立,此时显式设置
connection_mode = "passive_local";从源码看,主动模式失败时连接器会自动降级为被动模式并记录日志; - 路径含特殊字符或中文乱码:保持
control_encoding = "UTF-8",该设置在建立连接前即生效; - 文件没出现在目标目录:检查
tmp_path是否真实存在于 FTP 服务器,并确认提交阶段(2PC commit)是否成功——临时文件会先出现在tmp_path,提交成功后才mv到path; - CDC 作业报 AlterTableEvent 错误:若上游开启了
schema-changes.enabled = true,需要在本 Sink 设置schema_evolution_enabled = true,或在上游关闭 schema 变更事件。
如需进一步了解该插件的单元测试与工厂类实现,可查看 FtpFileFactoryTest.java 与 SeaTunnelFTPFileSystemTest.java。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考