ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

SeaTunnel File Connector 演进指南:从本地文件到多模态数据同步的完整技术图谱

SeaTunnel File Connector 演进指南:从本地文件到多模态数据同步的完整技术图谱 SeaTunnel File Connector 演进指南从本地文件到多模态数据同步的完整技术图谱【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本指南以 Apache SeaTunnel 仓库中 connector-file 变更日志 为核心骨架系统梳理 File 连接器家族的演进脉络与技术能力从 2.2.0-beta 的本地文件读写起步到覆盖 LocalFile、HdfsFile、S3、OSS、FTP、SFTP、COS、GCS、OBS 等十余种文件系统的统一架构再到多表同步、save mode、Excel 多引擎、CDC JSON 格式、二进制/多模态文件传输等高级特性。读完本文你将掌握 File 连接器的文件格式矩阵、核心参数语义、底层 ReadStrategy/WriteStrategy 机制以及每个重要特性对应的版本与源码位置可直接用于选型、排障与二次开发。一、File 连接器家族一份 changelog 背后的统一底座connector-file.md这份变更日志之所以能横跨十余个连接器是因为 SeaTunnel 采用了一个底座connector-file-base 多文件系统适配的架构。所有 File 系列连接器共享同一套读写策略、分片枚举、提交器与配置体系差异仅在于底层文件系统协议。从仓库目录结构可以清晰看到这一点公共底座connector-file-base 承载全部核心逻辑包括配置、Read/WriteStrategy、SplitEnumerator、Committer协议适配器connector-file-local、connector-file-hadoop、connector-file-s3、connector-file-oss、connector-file-ftp、connector-file-sftp、connector-file-cos、connector-file-gcs、connector-file-bos、connector-file-obs、connector-file-jindo-oss等每个模块仅实现各自文件系统FileSystem的接入。该体系从 2.2.0-beta 起逐步成型。日志中最早期的一批提交Add local file connector source #2419、Add hdfs file source connector #2420、Refactor the package of hdfs file connector #2402正是架构收敛的开始先是 LocalFile 与 HdfsFile 各自独立实现随后通过 Refactor the structure of file sink to reduce redundant codes #2555、Add base source connector code for connector-file-base #2399 等提交把公共逻辑抽入 base 模块后续 S3#3119、FTP#2483/#2774、OSS#2467/#2629、SFTP#3006等连接器都直接构建在该底座之上只需实现自己的 Catalog 与文件系统代理。二、文件格式支持矩阵从 text/csv 到多模态与 CDC JSONFile 连接器支持的文件格式在 2.x 系列版本中持续扩张这是 changelog 中最核心的一条演进主线。当前完整格式由 FileFormat.java 中的枚举定义格式读写方向引入/完善版本源码策略类text / csv读 写2.2.0-beta 起TextReadStrategy/TextWriteStrategy、CsvReadStrategy/CsvWriteStrategyjson读 写2.2.0-betaLocal file json support #2465JsonReadStrategy/JsonWriteStrategyparquet读 写2.2.0-betaAdd parquet writer #2273ParquetReadStrategy/ParquetWriteStrategyorc读 写2.2.0-betaSupport orc file format #2369OrcReadStrategy/OrcWriteStrategyexcelxlsx/xls读 写2.3.2 引入2.3.4 支持 .xls2.3.9 支持 EasyExcel 引擎ExcelReadStrategy/ExcelWriteStrategyxml读 写2.3.0/2.3.5扩展到 SFTP、FTP、LocalFile、HdfsFile 等XmlReadStrategy/XmlWriteStrategybinary读 写2.3.6Supports the transfer of any file #6826BinaryReadStrategy/BinaryWriteStrategymarkdown / pdf读dev多模态与 RAG 方向MarkdownReadStrategy/PdfReadStrategycanal_json / debezium_json / maxwell_json写2.3.12#9278/#9336CanalJsonWriteStrategy/DebeziumJsonWriteStrategy/MaxWellJsonWriteStrategy从源码看每种格式的读写逻辑被封装为独立的 Strategy 类并通过FileFormat.getReadStrategy()/getWriteStrategy(FileSinkConfig)工厂方法创建新增格式只需实现ReadStrategy/WriteStrategy接口并注册进枚举即可这也是格式快速扩张的结构性原因。值得注意的边界行为同样定义在FileFormat.java中canal_json、debezium_json、maxwell_json是只写格式读取会抛出UnsupportedOperationException用于将 CDC 事件落盘markdown、pdf是只读格式写入会抛出UnsupportedOperationException服务于多模态知识库/RAG 场景。此外supportFileSplit()方法明确 CSV、TEXT、JSON、PARQUET 四种格式支持文件切分enable_file_split。三、多表Multiple Table同步能力从单表到批量数据管道多表是 File 连接器 2.3.x 阶段的重大能力升级changelog 中有连续多条记录2.3.4LocalFileSource support multiple table、Add multiple table file sink to base #6049、Put Multiple Table File API to File Base Module #60332.3.8FTP/SFTP sink 支持多表与 save mode#7665/#76682.3.9SFTP/FTP source 支持多表#7824/#7795并引入Unified tables_configs and table_list #81002.3.10改进多表文件 source 的子任务分配算法#88782.3.12HDFS 文件 source/sink 支持多表#9816/#9651。多表能力的实现位于 base 模块的BaseMultipleTableFileSource/BaseMultipleTableFileSink以及MultipleTableFileSourceSplitEnumerator中。配置上通过tables_configs列表声明多张表的读取路径与 schema在 LocalFile source 文档 的选项表中可以看到tables_configs的说明used to define a multiple table task。结合日志中 Improved multiple table file source allocation algorithm for subtasks (#8878) 与 Improved file allocation algorithm for subtasks (#8453) 两条优化可以推断多表场景下分片如何在不同子任务间均衡分配是持续打磨的重点从源码结构看MultipleTableFileSplitStrategy与DefaultFileSplitStrategy分别承载了多表与单表的分片逻辑。四、源端核心参数深度解析File 连接器的源端参数集中定义在 FileBaseSourceOptions.java以下参数在 changelog 中均有对应演进记录4.1 文件发现与过滤file_filter_pattern2.3.3 引入 #51532.3.12 修复编译 #9658正则表达式用于按文件名过滤需要读取的文件是只同步部分文件的核心手段。filename_extension2.3.10 引入 #8769按扩展名过滤如csv、.txt、json、.xml与file_filter_pattern互为补充定义于公共类 FileBaseOptions.java。file_filter_modified_start/file_filter_modified_end2.3.12 引入 #9526 Support filtering files by last modified time按文件最后修改时间过滤时间格式默认yyyy-MM-dd HH:mm:ssstart 包含起始时刻end 不包含结束时刻。discovery_mode/scan_interval/start_modeonce默认与continuous两种发现模式后者支持增量监听新文件scan_interval默认 10Sstart_mode支持earliest/latest。recursive_file_scan是否递归扫描子目录默认 true。sort_files_by_modification_time按修改时间降序排序文件用于 schema 演化的场景下确保用最新文件做 schema 推断。4.2 解析控制field_delimiter/row_delimiter字段分隔符与行分隔符。字段分隔符对 text 默认\001行分隔符默认\n。2.3.1 支持将行分隔符设为空串#38542.3.11 为 text sink 增加row_delimiter#90172.3.12 支持可配置的 text 行分隔符#9608。csv_use_header_line是否使用 CSV 首行表头解析文件默认 false。skip_header_row_number跳过头部 N 行2.3.1 #3900。null_format定义代表 null 的字符串2.3.9 #8109。quote_char/escape_charCSV 引号字符默认与转义字符。parse_partition_from_path是否从文件路径解析分区字段默认 true2.3.0-beta #2985 Support parse field from file path。read_columns/read_partitions列投影与分区过滤。4.3 二进制与多模态binary_chunk_size2.3.12 #9391二进制文件读取的块大小默认 1024 字节增大可提升大文件性能但占用更多内存。binary_complete_file_mode是否将整个文件作为单个 chunk 一次性读入内存。sync_mode/target_path/update_strategy/compare_mode增量同步能力update模式仅支持binary格式compare_mode支持len_mtime长度修改时间与checksumHadoop FileSystem#getFileChecksum。4.4 Excel 读取excel_engine2.3.9 #8064POI默认与 EasyExcel 两种引擎可切换。poi_excel_max_file_sizePOI 引擎允许的最大 Excel 文件字节数默认 50MB52428800更大文件建议换 EasyExcel。sheet_name指定要读取的 sheet。五、汇端核心参数深度解析汇端参数集中在 FileBaseSinkOptions.java5.1 文件组织file_name_expression/custom_filename/filename_time_format自定义输出文件名支持${now}、${uuid}变量默认表达式为${transactionId}。single_file_mode2.3.10 #8518每个并行子任务将数据写入单个文件默认 false。create_empty_file_when_no_data2.3.10 #8543无数据时也生成空文件默认 false。batch_size每个切分文件的最大行数默认 1000000。enable_file_split/file_split_size大文件切分text 类格式会按row_delimiter对齐切分边界默认 128MB。5.2 分区写入have_partition/partition_by/partition_dir_expression/is_partition_field_write_in_file分区写入能力默认目录表达式为${k0}${v0}/${k1}${v1}/.../${kn}${vn}/即 Hive 风格分区目录2.3.1 支持为 Hive 连接器显式指定分区#3842。5.3 事务与提交is_enable_transaction默认 true是否启用写事务。事务由 FileSinkAggregatedCommitter 协调数据先写入tmp_path默认/tmp/seatunnelcheckpoint 完成后由 committer 将临时文件移动到目标目录并提交。changelog 中 Optimize files commit order #5045、Fix data file name will duplicate when use SeaTunnel Engine #3717、Fix WriteStrategy parallel writing thread unsafe issue #5546 均围绕此机制修复。5.4 表头与格式细节enable_header_write2.3.4 #5566/#5459text/csv 输出时是否写入列名表头默认 false。csv_string_quote_modeCSV 字符串引用模式默认 MINIMAL。sheet_max_rows2.3.12 #9668Excel sink 每个 sheet 的最大行数默认 1048576Excel 行数上限超出会分 sheet。max_rows_in_memory/sheet_nameExcel 写内存行数上限与目标 sheet 名。parquet_avro_write_timestamp_as_int96/parquet_avro_write_fixed_as_int962.3.6 #6971Parquet 以 INT96 类型写时间戳/12 字节字段。xml_root_tag/xml_row_tag/xml_use_attr_formatXML 根标签默认 RECORDS、行标签默认 RECORD与属性格式。5.5 CDC JSON 与 Schema 演化merge_update_event2.3.12canal_json/debezium_json/maxwell_json格式下将 UPDATE_AFTER 与 UPDATE_BEFORE 合并为一条 UPDATE 事件。schema_evolution_enabled默认 false从 CDC 数据源接收 ALTER TABLE 事件增/删/改/重名列并应用到 sink schema在每次 schema 变更边界轮转文件。encoding输出编码默认 UTF-82.3.5 #6489 支持为文件 source/sink 指定编码。六、Save Mode 与 Kerberos数据安全与认证6.1 Save Mode 演进save mode 是 File 连接器从机械写入走向语义化数据处理的标志性能力changelog 记录如下2.3.4S3 file save mode#6131、LocalFile save mode#70802.3.6LocalFile save mode 完善#7080 记录于 2.3.6 行2.3.8FTP/SFTP sink 支持多表与 save mode#7665/#76682.3.10OSS sink save mode 配置更新#9303。从FileBaseSinkOptions.java可以看到两类 mode 的枚举定义schema_save_mode默认CREATE_SCHEMA_WHEN_NOT_EXIST处理目标路径的建/删 schema 动作data_save_mode支持DROP_DATA、APPEND_DATA、ERROR_WHEN_DATA_EXISTS三种取值默认 APPEND_DATA决定同步开始前对目录中已有数据的处理方式。2.3.12 中 hive sink connector support overwrite mode #7891 则是 overwrite 语义在 Hive 上的落地。6.2 Kerberos 认证File 连接器对 Hadoop 系文件系统HDFS、Hive、S3、OSS、Iceberg 等提供完整的 Kerberos 支持changelog 中多次修复2.3.1hive 与 hdfs file connector 支持 kerberos#38402.3.3修复 Hadoop Kerberos 认证相关问题#51712.3.5OrcWriteStrategy/ParquetWriteStrategy 支持 kerberos 登录#64722.3.9修复 HdfsFile source 加载krb5_path配置#7870。参数定义在公共类FileBaseOptions/FileBaseSinkOptionskerberos_principal、krb5_path默认/etc/krb5.conf、kerberos_keytab_path。实现上由 HadoopLoginFactory.java 统一处理登录逻辑Hive 侧则通过 HiveMetaStoreProxy 启用。七、底层运行机制策略模式 分片枚举 事务提交7.1 读ReadStrategy 与 SplitEnumerator源端读取遵循枚举器发现文件 → 分配 split → 各 reader 用 ReadStrategy 解析的流水线BaseFileSource启动时通过FileDiscoveryScanner扫描路径把符合file_filter_pattern、filename_extension、修改时间过滤条件的文件封装为 FileSourceSplitFileSourceSplitEnumerator 负责将 split 分配给各个 reader并行度为 1 时全部分配给单个 reader源码第 118 行起的assignSplit逻辑Refactor file enumerator to prevent duplicate put split #8989 修复了 split 重复分配的问题各 reader 依据file_format_type调用对应 ReadStrategy如CsvReadStrategy、ParquetReadStrategy解析行数据并将已读 split 记录进快照实现 exactly-once文档中明确 Read all the data in a split in a pollNext call. What splits are read will be saved in snapshot。7.2 写WriteStrategy 与两阶段提交汇端写入由 WriteStrategyFactory 根据格式创建对应 WriteStrategy。2.3.0 之前文件命名易重复的问题#3717、并行写线程安全问题#5546以及 Json/Excel 写策略内存缓冲清理#5925均在这一层修复。提交流程走 SeaTunnel 两阶段提交写入临时目录 → checkpoint barrier →FileSinkAggregatedCommitter统一提交/回滚。八、压缩、编码与归档支持File 连接器在数据落地方面持续增强压缩2.3.1 #3899 Support compress通过compress_codec指定。text 支持 NONE/LZOparquet 支持 NONE/LZO/SNAPPY/LZ4/GZIP/BROTLI/ZSTDorc 支持 NONE/LZO/SNAPPY/LZ4/ZLIB。归档压缩archive_compress_codec2.3.8 #7633 支持读取归档压缩文件从源码看由ArchiveCompressFormat枚举管理。gz 读取2.3.9 #8025 LocalFile 支持读 gz#8181 gz 支持 excel。LZO 读取2.3.4 #5083 Support LZO compress on File Read。编码2.3.5 #6489source/sink 均可通过encoding指定字符集默认 UTF-8。九、版本演进总览与选型建议将 changelog 按大版本聚类可以清晰看到 File 连接器的能力里程碑版本核心变化2.2.0-betaLocalFile/HdfsFile sourcesink 落地json/parquet/orc 格式base 模块抽取文件切分#3625路径解析分区字段2.3.0-betaSFTP、S3 连接器hadoop3 uber 包批量切分支持用户自定义 schema 读 text2.3.0统一 Option 与 Factory#3375LZO 压缩编码支持XML 格式文件系统工具优化2.3.1file_format_type命名统一#4249压缩parquet INT96KerberosS3Catalog2.3.2Excel source/sink#4164OSS 配置修复临时文件读取修复2.3.3COS source/sinkfile_filter_pattern#5153提交顺序优化FTP e2e 测试2.3.4多表能力source/sinksave modeS3/LocalFile.xlsLZO 读空目录读取CSV 列顺序#9064表头写入2.3.5编码XML 扩展到多连接器parquet/ORC schema cast 与类型修复2.3.6任意文件传输binary#6826OBSS3 多表写Hive 多文件系统2.3.8FTP/SFTP 多表save mode归档压缩读S3/OSS FileCatalogIceberg Kerberos2.3.9EasyExcel 引擎Excel 公式/数字修复null_formatgz多表分片分配优化FTP connection_mode2.3.10common-csvfilename_extensionsingle_file_mode空文件创建SFTP 通配符转义2.3.11FTP 远程主机校验#9324text row_delimiter枚举器防重复CSV delimiter 修复2.3.12HDFS 多表 source/sink修改时间过滤#9526Excel sheet_max_rowsCDC JSON 三种格式 sinkbinary_chunk_sizeparquet 用户 schema#9596devHDFS 多表 source#9816多模态 embeddings#9673配合 markdown/pdf 的 RAG 方向选型建议基于仓库现状纯本地文件同步选connector-file-local对接 HDFS/Hive 生态选connector-file-hadoop可配套 Kerberos 与hdfs_site_path对象存储按云厂商选 S3/OSS/COS/GCS/OBS/BOS涉及增量文件监听与备份清理优先使用discovery_modecontinuouspost_sync_actiondelete/backupretention_max_age的组合CDC 数据落盘选择canal_json/debezium_json/maxwell_json写格式并配合merge_update_event知识库/RAG 场景关注 markdown/pdf 只读格式与markdown_rag_metadata_enabled、pdf_rag_metadata_enabled元数据开关。十、总结与延伸阅读SeaTunnel File 连接器通过一个 base 底座 策略模式 十余个文件系统适配器的架构在 2.2.x 到 2.3.x 的数个版本中完成了从简单本地文件读写到多表、多模态、CDC、增量同步的全面能力覆盖。这份 connector-file.md 变更日志 本身就是一张绝佳的地图每条记录都对应一段可以深入源码的功能实现。继续深入仓库的推荐路径连接器公共配置与语义FileBaseOptions.java、FileBaseSourceOptions.java、FileBaseSinkOptions.java格式与读写策略FileFormat.java、sink/writer与source/reader两个包下的全部 Strategy 类分片与提交source/split包下的 SplitEnumerator 家族、FileSinkAggregatedCommitter.java各连接器官方文档LocalFile source、HdfsFile source、以及 sink 文档 中对应各文件系统的页面均内含完整参数表与可直接运行的 HOCON 配置示例变更日志横向对比docs/en/connectors/changelog/目录下connector-file-local.md、connector-file-s3.md、connector-file-oss.md等文件记录了单个连接器的独立演进细节。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表