SeaTunnel FTP 文件连接器(connector-file-ftp)能力全景与演进史:从 2.2.0-beta 到 dev 的源码级解读
2026/9/18 17:09:26 网站建设 项目流程

SeaTunnel FTP 文件连接器(connector-file-ftp)能力全景与演进史:从 2.2.0-beta 到 dev 的源码级解读

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

SeaTunnel(Apache SeaTunnel)的 FTP 文件连接器位于seatunnel-connectors-v2/connector-file/connector-file-ftp,它通过自研的 Hadoop FileSystem 适配层(SeaTunnelFTPFileSystem)对接 Apache Commons Net 的 FTPClient,为上层统一的 File Source/Sink 框架提供 FTP 协议的读写能力。本文以仓库内 FTP 连接器 Change Log 为主线,结合连接器源码与配置工厂实现,系统梳理该连接器从 2.2.0-beta 到 dev 版本的能力演进、连接模式与可靠性机制、多表与保存模式,并给出可直接落地的完整配置示例。读完本文,你将掌握 FTP 连接器的全部核心参数、底层实现原理与版本能力边界,能够独立完成基于 FTP 的数据同步任务设计与排障。

连接器定位:File 连接器家族中的 FTP 一员

FTP 连接器是 SeaTunnel 文件连接器家族(LocalFile、HdfsFile、SftpFile、OssFile、S3File 等)中基于HDFS 协议模拟实现的一员。从源码结构看(FtpConf.java),它把 FTP 服务器封装成一个ftp://host:port的 Hadoop 文件系统:

  • 协议 schema 固定为ftp
  • 底层实现类为org.apache.seatunnel.connectors.seatunnel.file.ftp.system.SeaTunnelFTPFileSystem
  • 用户名、密码以fs.ftp.user.<host>fs.ftp.password.<host>的形式写入 Hadoop Configuration。

SeaTunnelFTPFileSystem继承自 HadoopFileSystem并实现StreamingFileSystem接口(SeaTunnelFTPFileSystem.java),在其上实现了opencreatelistStatusmkdirsdeleterenamegetFileStatus等完整文件操作,因此上层的 File Source / Sink(FtpFileSourceFtpFileSink)只需复用 connector-file-base 的通用实现即可工作。这一设计使得 FTP 连接器天然继承了文件连接器家族的全部能力(多格式、多表、分区、压缩等)。

与 SFTP 类似,FTP 连接器不维护工作目录状态setWorkingDirectory为空实现,getWorkingDirectory恒返回主目录),每次文件操作都会通过connect()建立独立连接、操作结束后disconnect()释放,源码注释明确说明这是为了避免每次 API 调用都承担 TCP 连接开销的权衡。

核心配置参数:基础连接与连接模式

FTP 专属参数定义在 FtpFileBaseOptions.java 中,Source 与 Sink 通过FtpFileSourceOptionsFtpFileSinkOptions直接继承。

参数类型必填默认值说明
hostString-FTP 服务器地址
portInteger-FTP 服务器端口(底层未显式指定时使用 FTP 标准端口 21)
userString-FTP 登录用户名
passwordString-FTP 登录密码
connection_modeEnumactive_localFTP 连接模式,可选active_localpassive_local
remote_verification_enabledBooleantrue是否开启 FTP 数据通道的远程主机校验(2.3.11 新增)
control_encodingStringUTF-8FTP 控制连接字符编码,用于支持文件路径中的特殊字符

其中hostportuserpassword在 Sink 侧被 FtpFileSinkFactory.java 标记为required,即写 FTP 时必须显式配置;Source 侧(FtpFileSourceFactory.java)则标记为optional,因为可以结合 Catalog / 多表配置注入。

connection_mode:主动/被动模式与自动降级

连接模式枚举定义在 FtpConnectionMode.java:

  • active_local:FTP 主动模式(服务器主动回连客户端数据端口);
  • passive_local:FTP 被动模式(客户端主动连接服务器数据端口),更适用于客户端位于 NAT/防火墙之后的场景。

值得强调的是 Change Log 2.3.9 中“Fix FTP connector connection_mode is not effective”(#7865)这一修复。在 SeaTunnelFTPFileSystem.connect() 中,连接模式通过fs.ftp.connection.mode配置项读取;而 setFsFtpConnectionMode() 实现了主动模式的自动降级逻辑:

  1. 若配置为active_local,先enterLocalActiveMode(),随后尝试创建一个测试目录/.ftptest<timestamp>
  2. 若创建失败,捕获 IOException 并自动切换为被动模式,同时将fs.ftp.connection.mode更新为passive_local
  3. 无论成败,finally中都会清理测试目录。

也就是说,即使显式配置了主动模式,在网络环境不支持时连接器也会自动回退到被动模式,而不是直接失败。

remote_verification_enabled:数据通道远程主机校验(2.3.11)

2.3.11 新增的“FTP data channels remote host verification”选项(#9324)对应remote_verification_enabled参数。在 connect() 中,该开关被透传到 Commons Net 的FTPClient.setRemoteVerificationEnabled()。默认开启(true),用于校验数据连接返回的 IP 是否与控制连接一致,提升安全性;在部分代理/NAT 环境(服务器回连地址与预期不一致)下若连接异常,可显式设置为false关闭校验。

control_encoding:控制连接编码

control_encoding默认UTF-8,在FTPClient连接前通过setControlEncoding()设置(见 SeaTunnelFTPFileSystem.java#L144-L146),用于支持文件路径中的中文等特殊字符。源码注释特别强调该设置必须在连接建立之前完成。

配置解析链路

FtpConf.buildWithConfig() 负责将 SeaTunnel 配置转换为 Hadoop 配置:

  • 构造ftp://<host>:<port>作为 defaultFS;
  • 写入fs.ftp.user.<host>fs.ftp.password.<host>fs.ftp.connection.modefs.ftp.remote.verification.enabledfs.ftp.control.encoding等键;
  • 通过setExtraOptions挂载到 HadoopConf。

支持的引擎、格式与数据类型

FTP 连接器复用了 File 连接器家族的完整能力矩阵,适用于 Spark / Flink / SeaTunnel Zeta 三种引擎(具体集成方式参见 文件连接器通用说明)。

文件格式能力

从 FtpFileSourceFactory.optionRule() 与 FtpFileSinkFactory.optionRule() 可见,file_format_type支持:textcsvjsonexcelxmlmarkdownpdfbinaryparquetorc等。其中:

  • binary格式(2.3.6 “Supports the transfer of any file”#6826)支持以二进制块读取任意文件,实现视频、图片等任意文件的同步;
  • markdownpdf是 dev 分支最新加入的解析能力(Add markdown parser #9714);
  • 各格式还有专属参数,例如excelsheet_nameexcel_enginePOI/EasyExcel)、poi_excel_max_file_sizexmlxml_row_tagxml_use_attr_formattext/csvfield_delimiterrow_delimiterskip_header_row_number等。

数据类型映射

文件本身没有强类型定义,通过配置schema显式声明每列的目标 SeaTunnel 数据类型,支持 STRING、SHORT、INT、BIGINT、BOOLEAN、DOUBLE、DECIMAL、FLOAT、DATE、TIME、TIMESTAMP、BYTES、ARRAY、MAP 等类型(与 SFTP Source 文档 中的数据类型映射一致)。

压缩支持

compress_codec按格式区分:text/json/csv支持lzoorc支持lzosnappylz4zlibparquet支持lzosnappylz4gzipbrotlizstd;excel 不支持压缩。2.3.4 还加入了 LZO 读取支持,2.3.8 新增archive_compress_codec支持读取归档压缩文件。

版本演进主线:从 2.2.0-beta 到 dev

Change Log(docs/zh/connectors/changelog/connector-file-ftp.md)完整记录了连接器的演进轨迹,按主题可归纳为以下几条主线:

起步:FTP Sink 与 Source 诞生(2.2.0-beta)

  • FTP 文件 Sink 支持(#2483):最初仅支持将数据写入 FTP;
  • FTP Sink 重构并新增 FTP Source(#2774):补齐了读端能力,连接器成为完整的读写双端组件;
  • 同期还改进了 parquet 读取(#2841)并修复了 Hive ORC 读取问题(#2845)。

选项体系与工厂机制成型(2.3.0 ~ 2.3.1)

2.3.0 是一次结构性大版本:为文件连接器引入统一的Option 与 Factory 机制(#3375)、重构代码结构(#3238)、统一文件连接器异常处理(#3525)、补充 Hadoop3 uber 包(#3755)。2.3.1 阶段继续完善:将file type统一更名为file_format_type(#4249)、重构 schema 解析(#4157)、为文件读写加入压缩支持(#3899)、改进文件连接器 option rule 与文档(#3812)、增加 get source 方法(#3846)。

格式与过滤能力扩充(2.3.2 ~ 2.3.6)

  • 2.3.2:新增 Excel Sink 与 Source(#4164);
  • 2.3.3:新增file_filter_pattern文件过滤配置(#5153),支持按文件名、目录(以path开头)正则过滤;
  • 2.3.4:能力密集版本——Source/Sink 增加connection_mode(#6077/#6099)、多表 File API 下沉到 File Base 模块(#6033)、支持多 Hadoop 账号(#5903)、引入新错误定义规则(#5793)、统一文件连接器选项与文档(#5680)、支持 LZO 读取压缩(#5083)、支持读取空目录(#5591)、schema 支持 column/primaryKey/constraintKey(#5564)、text/csv 格式新增enable_header_write(#5567);
  • 2.3.5:为 FTP/SFTP/LocalFile/HdfsFile 等文件连接器增加XML 文件类型支持(#6327);
  • 2.3.6:parquet 支持将 fixed/timestamp 以 int96 写入(#6971)、支持任意文件的二进制传输(#6826)。

多表与保存模式(2.3.8 ~ 2.3.9)

  • 2.3.8:FTP Sink 支持多表与 save mode(#7665),支持读取归档压缩文件(#7633);
  • 2.3.9:FTP Source 支持多表(#7795),修复connection_mode不生效问题(#7865),text 读取支持null_format自定义空值格式(#8109),(S)FTP 创建目录补充 debug 日志(#8286)。

文件管理与可靠性增强(2.3.10 ~ 2.3.12)

  • 2.3.10:FTP 连接器目录操作可靠性修复(#8959)、新增filename_extension读写参数(#8769)、重构 connector common options(#8634)、Sink 支持无数据时创建空文件(#8543)、Sink 支持单文件模式(#8518);
  • 2.3.11:新增 FTP 数据通道远程主机校验选项remote_verification_enabled(#9324)、更新文件连接器配置(#9034)、text Sink 增加row_delimiter(#9017);
  • 2.3.12:text 文件处理支持自定义行分隔符(#9608)。

dev 分支最新动态

当前 dev 分支正在加入Markdown 解析器(#9714),用于文本类文件的 Markdown 解析与 RAG 元数据处理(markdown_rag_metadata_enabled参数已出现在 Source option rule 中)。

多表(Multiple Table)与表管理能力

从 2.3.8(Sink)与 2.3.9(Source)开始,FTP 连接器支持多表配置。Source 侧table_configs与单表path互斥(exclusive规则见 FtpFileSourceFactory.java#L58),多表配置通过MultipleTableFTPFileSourceConfig解析,每个表可独立指定pathschemafile_format_type

此外,连接器还提供 FtpFileCatalog.java 与对应 Factory,支持以 Catalog 方式管理 FTP 文件表元数据,配合schema_save_modedata_save_mode(Sink 侧)实现表结构/数据的安全落库策略。

实战配置示例

以下示例均以当前仓库为基准编写,可直接套用。

示例一:FTP → Console(text 格式读取)

env { parallelism = 2 job.mode = "BATCH" } source { FTP { host = "192.168.1.100" port = 21 user = "seatunnel" password = "your_password" path = "/data/input" file_format_type = "text" field_delimiter = "\001" row_delimiter = "\n" skip_header_row_number = 1 connection_mode = "passive_local" # 读取 FTP 后删除源文件 post_sync_action = "delete" schema = { fields { id = INT name = STRING ts = TIMESTAMP } } } } sink { Console { parallelism = 1 } }

示例二:多表 Source(2.3.9+)

source { FTP { host = "192.168.1.100" port = 21 user = "seatunnel" password = "your_password" table_configs = [ { table_path = "/data/orders" table_name = "orders" file_format_type = "csv" schema = { fields { order_id = BIGINT, amount = DOUBLE } } }, { table_path = "/data/users" table_name = "users" file_format_type = "json" schema = { fields { user_id = BIGINT, name = STRING } } } ] } }

示例三:FTP Sink(CSV + 分区 + 单文件模式)

sink { FTP { host = "192.168.1.100" port = 21 user = "seatunnel" password = "your_password" path = "/data/output" file_format_type = "csv" field_delimiter = "," row_delimiter = "\n" enable_header_write = true # 分区能力 have_partition = true partition_by = ["dt"] partition_dir_expression = "${dt}" is_partition_field_write_in_file = true # 单文件模式(2.3.10+):所有数据写入一个文件 single_file_mode = true # 无数据时也创建空文件(2.3.10+) create_empty_file_when_no_data = true filename_extension = ".csv" schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode = "APPEND_DATA" } }

示例四:任意文件同步(binary 格式,2.3.6+)

source { FTP { host = "192.168.1.100" port = 21 user = "seatunnel" password = "your_password" path = "/data/images" file_format_type = "binary" binary_chunk_size = 4096 } } sink { LocalFile { path = "/tmp/images_backup" file_format_type = "binary" } }

环境依赖说明

与 SFTP 文档 中说明一致:若使用 Spark/Flink 引擎,需确保集群已集成 Hadoop(实测版本 2.x,连接器要求Hadoop 2.9.X+);若使用 SeaTunnel Engine,安装包已自动集成 Hadoop jar,可在${SEATUNNEL_HOME}/lib下确认。连接器通过 HDFS 协议模拟访问 FTP,因此这类依赖是必需的。

测试保障与工程规范

仓库为 FTP 连接器提供了对应的单元测试与工厂测试:

  • SeaTunnelFTPFileSystemTest.java:验证 FTP 文件系统的核心操作逻辑;
  • FtpFileFactoryTest.java:验证 Source/Sink Factory 的 option rule 与插件发现机制。

Change Log 中的 2.3.1 条目还记录了 spotless 代码格式化、模块可读命名(#4114)等工程化改进;2.3.9 的 metrics 关联改进(#7786)使指标信息可关联到逻辑计划节点,属于 SeaTunnel 引擎层通用能力。

总结

SeaTunnel FTP 文件连接器通过“Hadoop FileSystem 适配层 + 文件连接器家族通用框架”的架构,用较小成本实现了与 LocalFile、HdfsFile 等一致的丰富能力:十余种文件格式、多表、save mode、压缩、分区、文件过滤、任意文件二进制传输,以及主动/被动连接模式与自动降级、远程主机校验、控制连接编码等 FTP 专属可靠性机制。从 2.2.0-beta 的 Sink 起步,到 2.3.x 的格式扩充、多表支持与工程重构,再到 dev 的 Markdown/PDF 解析探索,其演进主线清晰可循。读者可按需组合本文章节中的配置示例,快速落地 FTP 数据集成任务;更深入的实现细节可直接阅读 connector-file-ftp 源码目录。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询