ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

SeaTunnel OssFile Source Connector 完全指南:从阿里云 OSS 批量读取多格式文件

2026/9/29 6:22:10 拓冰建站 浏览量
SeaTunnel OssFile Source Connector 完全指南:从阿里云 OSS 批量读取多格式文件 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文以 OssFile.md 为核心骨架结合seatunnel-connectors-v2/connector-file/connector-file-oss源码完整讲解 SeaTunnel 如何通过OssFileSource 连接器读取阿里云 OSS 上的 text / csv / parquet / orc / json / excel / xml / binary 文件覆盖依赖准备、参数详解、数据类型映射、单表与多表任务配置以及底层基于 Hadoop Aliyun OSS 文件系统的实现原理。OssFile是 SeaTunnel 文件类 Source 家族connector-file 模块中面向阿里云 OSS 的实现与LocalFile、S3File、HdfsFile等共享同一套BaseMultipleTableFileSource架构。读完本文你可以独立完成OSS → SeaTunnel → 任意 Sink的批式数据同步作业并理解分区字段解析、列投影、多表读取等进阶能力背后的源码机制。支持的计算引擎OssFileSource 支持以下引擎文档声明SparkFlinkSeaTunnel Zeta自研引擎推荐由于读取动作本质上是通过 HadoopFileSystemAPI 完成的因此在 Spark / Flink 引擎下需要集群具备 Hadoop 能力文档说明已验证的 Hadoop 版本为 2.x而在 SeaTunnel Zeta 引擎下则依赖内置的 Hadoop 3 运行时。依赖准备两种引擎的不同要求针对 Spark / Flink 引擎确保 Spark / Flink 集群已集成 Hadoop文档验证版本为 Hadoop 2.x。在${SEATUNNEL_HOME}/plugins/目录下必须存在以下 jarhadoop-aliyun-xx.jaraliyun-sdk-oss-xx.jarjdom-xx.jar版本匹配要求hadoop-aliyun的版本必须与 Spark / Flink 使用的 Hadoop 版本一致aliyun-sdk-oss与jdom的版本必须与hadoop-aliyun版本对应。例如hadoop-aliyun-3.1.4.jar依赖aliyun-sdk-oss-3.4.1.jar与jdom-1.1.jar。针对 SeaTunnel Zeta 引擎在${SEATUNNEL_HOME}/lib/目录下必须存在以下 jarseatunnel-hadoop3-3.1.4-uber.jaraliyun-sdk-oss-3.4.1.jarhadoop-aliyun-3.1.4.jarjdom-1.1.jar从源码看OssHadoopConf.java 将底层FileSystem实现类硬编码为org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystemscheme 为oss这正是依赖hadoop-aliyunjar 的原因——该 jar 提供了oss://协议到阿里云 OSS 的映射。核心特性Key FeaturesOssFileSource 的能力矩阵如下与官方 connector-v2-features 对照能力支持batch批式✅stream流式❌exactly-once精确一次✅column projection列投影✅parallelism并行读取✅support user-defined split用户自定义分片❌file format type✅ text / csv / parquet / orc / json / excel / xml / binary关于 exactly-once文档说明Read all the data in a split in a pollNext call. What splits are read will be saved in snapshot.——即连接器在单次pollNext中完整消费一个 split分片并将已读取的 split 列表写入快照FileSourceState.java配合 checkpoint 机制实现故障恢复后不重不漏的精确一次语义。支持的文件格式与数据类型映射OssFile支持text、csv、parquet、orc、json、excel、xml、binary八种格式其中数据类型映射与具体文件类型强相关以下按格式分类说明。JSON 文件类型指定file_format_type json时必须同时配置schema选项连接器依赖 schema 将 JSON 数据解析为目标行结构。上游数据单条{code: 200, data: get success, success: true}同一文件可存放多条数据以换行分隔JSON Lines{code: 200, data: get success, success: true} {code: 300, data: get failed, success: false}对应的 schema 配置HOCONschema { fields { code int data string success boolean } }连接器生成的数据codedatasuccess200get successtrueText / CSV 文件类型指定file_format_type text或csv时schema 可选不配置 schema整行内容作为一个字段content输出。上游数据tyrantlucifer#26#male将被解析为contenttyrantlucifer#26#male配置 schema除 CSV 外CSV 自带列分隔符必须同时配置field_delimiter选项注意文档表格中的delimiter参数与源码中 BaseSourceConfigOptions.FIELD_DELIMITER 一致默认值为\001即 Hive 默认分隔符且delimiter是它的 fallback keyfield_delimiter # schema { fields { name string age int gender string } }解析结果nameagegendertyrantlucifer26maleOrc 文件类型指定file_format_type orc时schema不需要配置连接器自动从 ORC 文件元数据中读取字段类型。类型映射如下Orc Data typeSeaTunnel Data typeBOOLEANBOOLEANINTINTBYTEBYTESHORTSHORTLONGLONGFLOATFLOATDOUBLEDOUBLEBINARYBINARYSTRING / VARCHAR / CHARSTRINGDATELOCAL_DATE_TYPETIMESTAMPLOCAL_DATE_TIME_TYPEDECIMALDECIMALLIST(STRING)STRING_ARRAY_TYPELIST(BOOLEAN)BOOLEAN_ARRAY_TYPELIST(TINYINT)BYTE_ARRAY_TYPELIST(SMALLINT)SHORT_ARRAY_TYPELIST(INT)INT_ARRAY_TYPELIST(BIGINT)LONG_ARRAY_TYPELIST(FLOAT)FLOAT_ARRAY_TYPELIST(DOUBLE)DOUBLE_ARRAY_TYPEMapK,VMapTypeK、V 递归转换为 SeaTunnel 类型STRUCTSeaTunnelRowTypeParquet 文件类型指定file_format_type parquet时schema 同样可选自动探测映射关系Parquet Data typeSeaTunnel Data typeINT_8BYTEINT_16SHORTDATEDATETIMESTAMP_MILLISTIMESTAMPINT64LONGINT96TIMESTAMPBINARYBYTESFLOATFLOATDOUBLEDOUBLEBOOLEANBOOLEANFIXED_LEN_BYTE_ARRAYTIMESTAMP / DECIMALDECIMALDECIMALLIST(STRING)STRING_ARRAY_TYPELIST(BOOLEAN)BOOLEAN_ARRAY_TYPELIST(TINYINT)BYTE_ARRAY_TYPELIST(SMALLINT)SHORT_ARRAY_TYPELIST(INT)INT_ARRAY_TYPELIST(BIGINT)LONG_ARRAY_TYPELIST(FLOAT)FLOAT_ARRAY_TYPELIST(DOUBLE)DOUBLE_ARRAY_TYPEMapK,VMapTypeK、V 递归转换为 SeaTunnel 类型STRUCTSeaTunnelRowType从源码看格式探测与行类型解析由 ReadStrategyFactory.java 根据file_format_type分派到 OrcReadStrategy.java / ParquetReadStrategy.java 等实现类schema 的处理逻辑集中在 BaseFileSourceConfig.java配置了schema则通过CatalogTableUtil.buildWithConfig构建 CatalogTable否则buildSimpleTextTable对 text/csv/json/excel/xml 会调用readStrategy.setSeaTunnelRowTypeInfo校正实际行类型对 orc/parquet/binary 则调用getSeaTunnelRowTypeInfoWithUserConfigRowType从文件元数据读取。参数详解Options下表完整列出OssFileSource 的全部参数默认值与说明与官方文档一致标注部分补充了源码佐证nametyperequireddefaultDescriptionpathstringyes-需要读取的 OSS 路径可包含子路径但子路径需满足分区格式要求配合parse_partition_from_path使用file_format_typestringyes-文件类型textcsvparquetorcjsonexcelxmlbinarybucketstringyes-OSS bucket 地址例如oss://seatunnel-testendpointstringyes-OSS 文件系统 endpoint例如oss-cn-beijing.aliyuncs.comread_columnslistno-数据源要读取的列清单用于实现字段投影text/csv/parquet/orc/json/excel/xml 均支持读取 text/json/csv 时使用该功能必须配置schemaaccess_keystringno-OSS 访问密钥 AccessKey IDaccess_secretstringno-OSS 访问密钥 AccessKey Secretdelimiterstringno\001读取 text 文件时的字段分隔符默认\001与 Hive 默认分隔符一致源码中实际主键为field_delimiterdelimiter为 fallback keyparse_partition_from_pathbooleannotrue是否从文件路径解析分区键与值。例如读取oss://hadoop-cluster/tmp/seatunnel/parquet/nametyrantlucifer/age26每条记录都会附带nametyrantlucifer、age26两个字段date_formatstringnoyyyy-MM-dd日期格式支持yyyy-MM-dd、yyyy.MM.dd、yyyy/MM/dddatetime_formatstringnoyyyy-MM-dd HH:mm:ss日期时间格式支持yyyy-MM-dd HH:mm:ss、yyyy.MM.dd HH:mm:ss、yyyy/MM/dd HH:mm:ss、yyyyMMddHHmmsstime_formatstringnoHH:mm:ss时间格式支持HH:mm:ss、HH:mm:ss.SSSskip_header_row_numberlongno0跳过文件前 N 行仅对 txt 与 csv 生效例如skip_header_row_number 2schemaconfigno-上游数据的 schema 定义sheet_namestringno-读取 Excel 工作簿的 sheet 名仅file_format excel时使用xml_row_tagstringno-指定 XML 文件中数据行的标签名仅file_format xml时使用xml_use_attr_formatbooleanno-是否以标签属性格式处理数据仅file_format xml时使用compress_codecstringnonone文件的压缩编解码器encodingstringnoUTF-8文件编码file_filter_patternstringno-文件过滤模式例如*.txt表示只读取.txt结尾的文件common-optionsconfigno-Source 插件公共参数详见 Source Common Options从 OssFileSourceFactory.java 的optionRule()可以看到参数的条件生效逻辑file_format_type text时启用field_delimiterfile_format_type xml时启用xml_row_tag与xml_use_attr_formatfile_format_type为 text/json/excel/csv/xml 时启用schema。endpoint、access_key、access_secret、bucket在 OssConfigOptions.java 中定义为 OSS 专属必填项。compress_codec压缩编解码按文件格式区分支持范围text / json / csvlzo、noneorc / parquet自动识别压缩类型无需额外设置encoding编码仅file_format_type为 json / text / csv / xml 时生效。该参数最终通过Charset.forName(encoding)解析对应源码 BaseSourceConfigOptions.ENCODING默认UTF-8。file_filter_pattern文件过滤用于过滤待读取文件的正则/通配模式例如*.txt只读取以.txt结尾的文件。文件清单的获取发生在 BaseFileSourceConfig.parseFilePaths 中通过readStrategy.getFileNamesByPath(rootPath)完成若路径下文件枚举失败会抛出FILE_LIST_GET_FAILED异常。schemaschema 定义仅当file_format_type为 text、json、excel、xml 或 csv即无法从元数据读取 schema 的格式时必须配置。嵌套的fields子配置用于声明上游数据的字段名与类型。实战创建 OSS 数据同步任务示例一读取 ORC 文件并打印到控制台# Set the basic configuration of the task to be performed env { parallelism 1 job.mode BATCH } # Create a source to connect to Oss source { OssFile { path /seatunnel/orc bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type orc } } # Console printing of the read Oss data sink { Console { } }示例二读取 JSON 文件需配置 schema# Set the basic configuration of the task to be performed env { parallelism 1 job.mode BATCH } # Create a source to connect to Oss source { OssFile { path /seatunnel/json bucket oss://tyrantlucifer-image-bed access_key xxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxxxxx endpoint oss-cn-beijing.aliyuncs.com file_format_type json schema { fields { id int name string } } } } # Console printing of the read Oss data sink { Console { } }启动方式将上述配置保存为配置文件后使用bin/seatunnel.sh --config config_file提交任务详见 SeaTunnel 部署文档。运行前提是${SEATUNNEL_HOME}/lib/Zeta或${SEATUNNEL_HOME}/plugins/Spark/Flink已放置前文所述的 OSS 依赖 jar。多表读取Multiple Table当需要从多个 OSS 路径读取并注册为多张表时使用tables_configs数组通过schema.table指定表名并设置result_table_name。分为两种情况无需配置 schema 的格式如 orcenv { parallelism 1 spark.app.name SeaTunnel spark.executor.instances 2 spark.executor.cores 1 spark.executor.memory 1g spark.master local job.mode BATCH } source { OssFile { tables_configs [ { schema { table fake01 } bucket oss://whale-ops access_key xxxxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxx endpoint https://oss-accelerate.aliyuncs.com path /test/seatunnel/read/orc file_format_type orc }, { schema { table fake02 } bucket oss://whale-ops access_key xxxxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxx endpoint https://oss-accelerate.aliyuncs.com path /test/seatunnel/read/orc file_format_type orc } ] result_table_name fake } } sink { Assert { rules { table-names [fake01, fake02] } } }需要配置 schema 的格式如 json除schema.table外还需在schema.fields中声明每个表的完整字段结构支持mapstring, string、arrayint、decimal(38, 18)、嵌套c_row等复杂类型env { execution.parallelism 1 spark.app.name SeaTunnel spark.executor.instances 2 spark.executor.cores 1 spark.executor.memory 1g spark.master local job.mode BATCH } source { OssFile { tables_configs [ { bucket oss://whale-ops access_key xxxxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxx endpoint https://oss-accelerate.aliyuncs.com path /test/seatunnel/read/json file_format_type json schema { table fake01 fields { c_map mapstring, string c_array arrayint c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double c_bytes bytes c_date date c_decimal decimal(38, 18) c_timestamp timestamp c_row { C_MAP mapstring, string C_ARRAY arrayint C_STRING string C_BOOLEAN boolean C_TINYINT tinyint C_SMALLINT smallint C_INT int C_BIGINT bigint C_FLOAT float C_DOUBLE double C_BYTES bytes C_DATE date C_DECIMAL decimal(38, 18) C_TIMESTAMP timestamp } } } }, { bucket oss://whale-ops access_key xxxxxxxxxxxxxxxxxxx access_secret xxxxxxxxxxxxxxxxxxx endpoint https://oss-accelerate.aliyuncs.com path /test/seatunnel/read/json file_format_type json schema { table fake02 fields { c_map mapstring, string c_array arrayint c_string string c_boolean boolean c_tinyint tinyint c_smallint smallint c_int int c_bigint bigint c_float float c_double double c_bytes bytes c_date date c_decimal decimal(38, 18) c_timestamp timestamp c_row { C_MAP mapstring, string C_ARRAY arrayint C_STRING string C_BOOLEAN boolean C_TINYINT tinyint C_SMALLINT smallint C_INT int C_BIGINT bigint C_FLOAT float C_DOUBLE double C_BYTES bytes C_DATE date C_DECIMAL decimal(38, 18) C_TIMESTAMP timestamp } } } } ] result_table_name fake } } sink { Assert { rules { table-names [fake01, fake02] } } }多表能力由 MultipleTableOssFileSourceConfig.java 承载其底层复用BaseMultipleTableFileSourceBaseMultipleTableFileSource.java与 MultipleTableFileSourceSplitEnumerator.java 实现多表 split 的分配与枚举。源码视角OssFile Source 的底层实现OssFile连接器的类结构非常精简它通过组合文件连接器公共基座复用全部读取能力OssFileSource.java直接继承BaseMultipleTableFileSource插件名返回FileSystemType.OSS.getFileSystemPluginName()即OssFile。OssFileSourceFactory.java以AutoService(Factory.class)注册到 SPI负责插件识别、参数规则optionRule()与 Source 实例创建。OssHadoopConf.java核心对接层——将bucket、access_key、access_secret、endpoint映射为 HadoopAliyunOSSFileSystem的Constants.ACCESS_KEY_ID、ACCESS_KEY_SECRET、ENDPOINT_KEY配置构建出oss://协议的 HadoopConf。OssFileSourceConfig.javagetHadoopConfig()委托OssHadoopConf.buildWithConfig。由此可以推断出完整的数据链路OssFile Source → BaseMultipleTableFileSource多表拆分→ FileSourceSplitEnumerator按文件/分区生成 split支持并行→ 各 ReadStrategyText/Json/Orc/Parquet/Excel/Xml/Binary→ AliyunOSSFileSystem 读取 → SeaTunnelRow 输出到下游 Sink。并行读取能力即来源于 split 被分配给多个 reader 并行消费FileSourceSplitEnumerator.java 负责具体的 split 分配策略。版本演进Changelog2.2.0-beta2022-09-26新增 OSS File Source 连接器。2.3.0-beta2022-10-20[BugFix] 修复 Windows 环境下路径错误的问题PR 2980[Improve] 支持从 SeaTunnelRow 字段中提取分区信息PR 3085[Improve] 支持从文件路径解析字段PR 2985从这段演进可以看出从文件路径解析分区字段即parse_partition_from_path与从行字段提取分区是后续版本持续增强的亮点能力也是数据湖/数仓场景下按分区目录增量读取的关键支撑。相关文档与进一步阅读SeaTunnel 部署文档任务提交与引擎启动方式Source Common OptionsSource 插件公共参数如result_table_name、parallelism等Connector V2 Featuresbatch / stream / exactly-once 等能力定义源码参考connector-file-oss 模块、connector-file-base 模块赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐终极Unity游戏翻译方案3步实现多语言无障碍体验终极Unity游戏翻译方案3步实现多语言无障碍体验 还在为外语游戏中的对话一头雾水而烦恼吗是否曾经因为语言障碍而错过了精彩的剧情体验XUnity Auto数据工程大数据批处理流处理SeaTunnel GcsFile Source Connector 完全指南从 GCS 对象存储读取多格式文件SeaTunnel GcsFile Source Connector 完全指南从 GCS 对象存储读取多格式文件 Google Cloud Storage 文数据集成ETL大数据批处理流处理变更数据捕获OssFile 数据源连接器从 OSS 读取多格式文件到 SeaTunnel 的完整实战指南OssFile 数据源连接器从 OSS 读取多格式文件到 SeaTunnel 的完整实战指南 本文以 SeaTunnel 官方文档 OssFile.md ht数据集成ETL大数据批处理流处理变更数据捕获上一篇审查输出 Schemav6下一篇Home Assistant ISY994 集成使用 isy994.get_zwave_parameter 动作读取 Z-Wave 设备参数创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考