
SeaTunnel FtpFile Sink Connector 使用指南将数据输出到 FTP 服务器【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnelSeaTunnel 的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 支持的全部参数与官方文档保持一致NameTypeRequiredDefaultDescriptionhoststringyes-FTP 服务器主机名portintyes-FTP 服务器端口userstringyes-FTP 登录用户名passwordstringyes-FTP 登录密码pathstringyes-目标目录路径tmp_pathstringyes/tmp/seatunnel结果文件先写入该临时目录随后通过mv提交到目标目录需要是 FTP 目录connection_modestringnoactive_localFTP 连接模式remote_verification_enabledbooleannotrue是否启用 FTP 数据通道的远程主机校验control_encodingstringnoUTF-8FTP 控制连接字符编码对含空格或非 ASCII 字符的路径很有用custom_filenamebooleannofalse是否需要自定义文件名file_name_expressionstringno${transactionId}仅当 custom_filename 为 true 时生效filename_time_formatstringnoyyyy.MM.dd仅当 custom_filename 为 true 时生效file_format_typestringnocsv输出文件格式filename_extensionstringno-用自定义扩展名覆盖默认扩展名如.xml、.json、dat、.customtypefield_delimiterstringnotext 为 \001csv 为 ,仅当 file_format_type 为 text 和 csv 时生效row_delimiterstringno\n仅当 file_format_type 为 text、csv 和 json 时生效have_partitionbooleannofalse是否需要分区处理partition_byarrayno-仅当 have_partition 为 true 时生效partition_dir_expressionstringno${k0}${v0}/${k1}${v1}/.../${kn}${vn}/仅当 have_partition 为 true 时生效is_partition_field_write_in_filebooleannofalse仅当 have_partition 为 true 时生效sink_columnsarrayno空为空时所有字段都作为输出列is_enable_transactionbooleannotrue是否启用事务batch_sizeintno1000000单个文件最大行数compress_codecstringnonone压缩编码common-optionsobjectno-Sink 公共参数max_rows_in_memoryintno-仅当 file_format_type 为 excel 时生效sheet_max_rowsintno1048576仅当 file_format_type 为 excel 时生效sheet_namestringnoSheet${随机数}仅当 file_format_type 为 excel 时生效csv_string_quote_modeenumnoMINIMAL仅当 file_format 为 csv 时生效xml_root_tagstringnoRECORDS仅当 file_format 为 xml 时生效xml_row_tagstringnoRECORD仅当 file_format 为 xml 时生效xml_use_attr_formatbooleanno-仅当 file_format 为 xml 时生效single_file_modebooleannofalse每个并行度只输出一个文件开启后 batch_size 不生效输出文件名不带文件块后缀create_empty_file_when_no_databooleannofalse上游无数据同步时仍生成对应数据文件parquet_avro_write_timestamp_as_int96booleannofalse仅当 file_format 为 parquet 时生效parquet_avro_write_fixed_as_int96arrayno-仅当 file_format 为 parquet 时生效enable_header_writebooleannofalse仅当 file_format_type 为 text、csv 时生效false 不写表头true 写表头encodingstringnoUTF-8仅当 file_format_type 为 json、text、csv、xml 时生效schema_evolution_enabledbooleannofalse为 CDC 管道启用 schema 演进支持 ADD/DROP/RENAME/MODIFY 列事件无需重启作业binary 格式不支持schema_save_modestringnoCREATE_SCHEMA_WHEN_NOT_EXIST已有目录的处理方式data_save_modestringnoAPPEND_DATA已有数据的处理方式multi_table_sink_replicaintno1多表 Sink 作业中每个表使用的 writer 副本数连接参数详解host / port / user / password这四个参数是必填的连接基础信息FTP 服务器主机、端口、登录用户名和密码。在源码 FtpFileBaseOptions.java 中它们被定义为无默认值的必填 OptionFtpConf.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_modeFTP 数据通道连接模式默认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_encodingFTP 控制连接的字符编码默认UTF-8。源码在建立连接前调用client.setControlEncoding(controlEncoding)见 SeaTunnelFTPFileSystem.java该设置对路径中包含空格、特殊字符或非 ASCII 字符的场景至关重要。除非你的 FTP 服务器要求其他控制通道编码否则保持UTF-8即可。文件命名与格式配置custom_filename / file_name_expression / filename_time_formatcustom_filename是否自定义文件名默认falsefile_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。常用时间格式符号如下SymbolDescriptionyYear年MMonth月dDay of month日HHour in day (0-23)时mMinute in hour分sSecond 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_delimiterfield_delimiter一行数据中列之间的分隔符仅 text 和 csv 格式需要。默认值text 为\001csv 为,row_delimiter文件中行与行之间的分隔符仅 text、csv 和 json 格式需要默认\n。enable_header_write / encodingenable_header_write仅 text、csv 格式生效false不写表头true写表头默认falseencoding仅 json、text、csv、xml 格式生效指定输出文件的字符编码如 UTF-8、ISO-8859-1默认UTF-8。该参数会通过Charset.forName(encoding)解析对应 FileBaseOptions.java 中的ENCODING定义。分区写入配置have_partition是否启用分区处理默认falsepartition_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_pathcheckpoint 触发后通过 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_datasingle_file_mode默认false。开启后每个并行度只输出一个文件此时batch_size不再生效输出文件名不带文件块后缀create_empty_file_when_no_data默认false。开启后即使上游没有数据同步也仍会生成对应的数据文件。压缩、CSV/XML/Excel/Parquet 专属参数compress_codec文件压缩编码默认none按格式支持如下txtlzo、nonejsonlzo、nonecsvlzo、noneorclzo、snappy、lz4、zlib、noneparquetlzo、snappy、lz4、gzip、brotli、zstd、none提示excel 格式不支持任何压缩格式。Excel 相关file_format_type excelmax_rows_in_memory内存中可缓存的最大数据条数sheet_max_rows每个 sheet 的最大行数默认1048576sheet_name工作簿中写入的 sheet 名称默认Sheet${随机数}。CSV 引号模式file_format_type csvcsv_string_quote_mode的可选值及语义ALL所有 String 字段都加引号MINIMAL仅对包含特殊字符如字段分隔符、引号字符或行分隔符串中的任意字符的字段加引号NONE从不加引号。当数据中出现分隔符时printer 会在其前面加上转义字符若未设置转义字符格式校验将抛出异常。XML 相关file_format_type xmlxml_root_tag指定 XML 文件根元素标签名默认RECORDSxml_row_tag指定 XML 文件数据行标签名默认RECORDxml_use_attr_format指定是否使用标签属性格式处理数据。Parquet 相关file_format_type parquetparquet_avro_write_timestamp_as_int96支持将时间戳写入 Parquet INT96默认falseparquet_avro_write_fixed_as_int96支持从 12 字节字段写入 Parquet INT96。目录与数据的 Save Modeschema_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_enabledfalse at this sink. Either set schema_evolution_enabledtrue to handle schema changes, or set schema-changes.enabledfalse 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_modeRECREATE_SCHEMA data_save_modeDROP_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到pathCDC 作业报 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),仅供参考