SeaTunnel OssJindoFile 连接器实战指南:通过 Jindo SDK 对接阿里云 OSS 的 Source/Sink 配置、原理与示例
2026/9/16 14:04:26 网站建设 项目流程

SeaTunnel OssJindoFile 连接器实战指南:通过 Jindo SDK 对接阿里云 OSS 的 Source/Sink 配置、原理与示例

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

SeaTunnel 的 OssJindoFile 连接器(connector-file-jindo-oss)基于阿里云 EMR 的 Jindo SDK,通过 HDFS 协议访问阿里云 OSS 对象存储,同时提供 Source(读取)与 Sink(写出)两种能力,支持 text、csv、parquet、orc、json、excel、xml、binary 等多种文件格式,并内置 exactly-once 语义。本文以 changelog 文档 与 Source 官方文档、Sink 官方文档 为主线,结合 connector-file-jindo-oss 模块源码 讲解其原理与完整实战配置,读完即可在 Spark、Flink 或 SeaTunnel Zeta 引擎上完成 OSS 数据入湖、文件迁移与二进制文件同步。

连接器概览与演进历史

OssJindoFile 连接器最早以 Add oss jindo source & sink connector (#3456)(插件历史变更)与本仓库的 connector-file-jindo-oss.md 变更日志 中,几个关键里程碑如下:

版本关键变更
2.3.0新增 OssJindo Source 与 Sink 连接器
2.3.2修复 file-oss 配置检查 bug,修正 file-oss-jindo 的 factoryIdentifier
2.3.3优化 jindo oss 连接器;新增file_filter_pattern文件过滤参数
2.3.4统一文件 Source/Sink 选项并更新文档;支持自定义行分隔符写入文本文件
2.3.5为 SFTP、FTP、LocalFile、HdfsFile 等文件连接器统一支持 XML 文件类型
2.3.6支持任意文件的传输(binary 模式);parquet 支持 timestamp/fixed 写为 int96
2.3.9支持null_format文本空值配置;Read/WriteStrategy 由setSeaTunnelRowTypeInfo改为setCatalogTable
2.3.10支持filename_extension读写参数、文件 Sink 单文件模式、无数据时创建空文件
2.3.11文本文件 Sink 支持row_delimiter选项
2.3.12文本文件处理支持自定义行分隔符;maxcompute sink writer 支持 timestamp 字段类型

从 plugin-mapping.properties 可以看到,seatunnel.source.OssJindoFileseatunnel.sink.OssJindoFile两个插件标识统一映射到connector-file-jindo-oss这一个 Maven 模块,即同一模块同时提供读取与写出能力。

支持引擎与核心特性

根据官方文档,该连接器同时支持 Spark、Flink 与 SeaTunnel Zeta 三种引擎。

Source 侧特性

  • batch 批处理:支持批式读取
  • multimodal 多模态:使用 binary 文件格式可读取任意格式文件(视频、图片、压缩包等)
  • exactly-once:在单次 pollNext 调用中读取一个 split 的全部数据,split 读取情况保存于 snapshot 中,保证精确一次
  • parallelism 并行度:支持并行读取
  • 支持的文件格式:textcsvparquetorcjsonexcelxmlbinarymarkdownpdf

Sink 侧特性

  • multimodal:以 binary 格式写出任意格式文件
  • exactly-once:默认使用 2PC(两阶段提交)保证精确一次
  • 多表写入:支持 multiple table write
  • 支持的文件格式:textcsvparquetorcjsonexcelxmlbinarycanal_jsondebezium_jsonmaxwell_json

注意 Sink 侧额外支持canal_jsondebezium_jsonmaxwell_json三种 CDC 事件格式;Source 侧则额外支持markdownpdf两种文档解析格式(Markdown 只支持读取、不支持写出)。

环境依赖与 Jindo SDK 安装

该连接器通过 Jindo(阿里云 EMR)SDK 走 HDFS 协议访问 OSS,因此依赖 Jindo 与 Hadoop 相关 jar 包,官方文档给出了明确的安装前提:

  1. 下载jindosdk-4.6.1.tar.gz,解压后将其lib目录下的jindo-sdk-4.6.1.jarjindo-core-4.6.1.jar复制到${SEATUNNEL_HOME}/lib下,且每个运行作业的节点都需要放置。
  2. 如果使用 Spark/Flink,需要确保集群已集成 Hadoop,官方测试的 Hadoop 版本为 2.x。
  3. 如果使用 SeaTunnel Engine,安装 SeaTunnel Engine 时已自动集成 Hadoop jar,可在${SEATUNNEL_HOME}/lib下确认。
  4. 由于为支持更多文件类型而内部走 HDFS 协议访问 OSS,连接器只支持 Hadoop 版本2.9.X+

从 pom.xml 可以看出,该模块声明了hadoop-common 2.9.2依赖且作用域为provided(运行时由运行环境提供),并依赖同级的connector-file-base基础模块。

核心原理:OssConf 如何把 Jindo 桥接到 Hadoop 抽象

从源码层面看,整个连接器只做了两件事:把 OSS 配置翻译成 Hadoop 文件系统配置,以及复用connector-file-base中通用的文件读写实现。

OssConf.java 继承自HadoopConf,核心逻辑如下:

  • HDFS_IMPL = "com.aliyun.emr.fs.oss.JindoOssFileSystem",即 Jindo 提供的 OSS Hadoop 文件系统实现类;
  • SCHEMA = "oss",即路径前缀oss://
  • buildWithReadonlyConfig()方法把配置项翻译成 Hadoop 配置项:
SeaTunnel 配置项Hadoop 配置项说明
bucket(作为hdfsNameKey传入构造器)OSS bucket 地址,如oss://tyrantlucifer-image-bed
access_keyfs.oss.accessKeyIdOSS 访问密钥 ID
access_secretfs.oss.accessKeySecretOSS 访问密钥 Secret
endpointfs.oss.endpointOSS 服务端点
(固定值)fs.AbstractFileSystem.oss.implcom.aliyun.emr.fs.oss.OSS
(固定值)fs.oss.implcom.aliyun.emr.fs.oss.JindoOssFileSystem
(固定值)fs.oss.upload.thread.concurrency上传线程并发数,固定 20
(固定值)fs.oss.upload.queue.size上传队列大小,固定 100

而 OssFileBaseOptions.java 中声明了四个必填的 OSS 专属选项(均无默认值):

ACCESS_KEY // access_key: OSS bucket access key ACCESS_SECRET // access_secret: OSS bucket access secret ENDPOINT // endpoint: OSS endpoint BUCKET // bucket: OSS bucket

OssFileSource.java 与 OssFileSink.java 分别继承BaseFileSourceBaseFileSink,仅实现initHadoopConf()getPluginName()两个方法,插件名来自FileSystemType.OSS_JINDO(值为OssJindoFile,见 FileSystemType.java),其余文件格式解析、分区处理、事务提交等全部复用 connector-file-base 的通用实现。

另外,OssJindoFactoryTest.java 验证了OssFileSourceFactoryOssFileSinkFactoryoptionRule()均能正常构造,可作为连接器注册与选项规则的自检入口。

Source 配置详解:从 OSS 读取文件

必选参数

参数类型必填默认值说明
pathstring-源文件路径
file_format_typestring-文件类型
bucketstring-OSS bucket 地址,如oss://tyrantlucifer-image-bed
access_keystring-OSS 访问密钥
access_secretstring-OSS 访问密钥 Secret
endpointstring-OSS 服务端点,如oss-cn-beijing.aliyuncs.com

常用可选参数

参数类型默认值说明
read_columnslist-读取列清单,实现字段投影
delimiter / field_delimiterstringtext 为\001,csv 为,字段分隔符(delimiter将在 2.3.5 之后废弃,请改用field_delimiter
row_delimiterstring\n行分隔符
parse_partition_from_pathbooleantrue是否从文件路径解析分区键与分区值
date_formatstringyyyy-MM-dd字符串转 date 的格式
datetime_formatstringyyyy-MM-dd HH:mm:ss字符串转 datetime 的格式
time_formatstringHH:mm:ss字符串转 time 的格式
skip_header_row_numberlong0跳过前 N 行(仅 text 与 csv)
schemaconfig-上游数据结构定义
file_filter_patternstring-基于正则的文件过滤模式
filename_extensionstring-按扩展名过滤文件
compress_codecstringnone压缩编解码器(lzo/none,orc/parquet 自动识别)
archive_compress_codecstringnone归档压缩格式(ZIP/TAR/TAR_GZ/GZ/NONE)
encodingstringUTF-8文件编码
null_formatstring-表示 null 的字符串(如\N
file_filter_modified_start / file_filter_modified_endstring-按修改时间过滤(yyyy-MM-dd HH:mm:ss
quote_charstring"CSV 字段包围字符
escape_charstring-CSV 转义字符
recursive_file_scanbooleantrue是否递归扫描子目录
sort_files_by_modification_timebooleanfalse是否按修改时间倒序排序文件

其中关键选项的用法细节如下:

  • schemafield_delimiter配合:当file_format_type为 text 且上游是tyrantlucifer#26#male这类定界文本时,不配 schema 会整行作为content字段;配置field_delimiter = "#"与 schema 后即可拆成name/age/gender三列。json 类型必须配 schema 才能解析,parquet/orc 则能自动识别元数据中的 schema。
  • parse_partition_from_path:读取形如oss://hadoop-cluster/tmp/seatunnel/parquet/name=tyrantlucifer/age=26的路径时,会自动为每条数据附加name=tyrantluciferage=26两个分区字段,注意不要在 schema 中重复定义分区字段。该选项的默认值为true,见 FileBaseSourceOptions.java。
  • file_filter_pattern:标准正则表达式;仅按文件名过滤时直接写文件名正则(如abc.*),如需同时匹配目录则表达式需以path开头(如/data/seatunnel/20241007/abc[h,g].*)。
  • XML 安全限制:出于 XXE 加固考虑,file_format_type = xml的文件若包含<!DOCTYPE ...>声明(即使是仅定义内部实体的良性声明),会以FILE_READ_FAILED错误拒绝读取,且没有配置项可以恢复旧行为。若 XML 文件由导出工具生成了 DOCTYPE 头,需先移除或预处理再接入。
  • markdown/pdf 文档解析:Source 侧可将 markdown 与 pdf 解析为结构化文档元素行(element_idelement_typeheading_leveltext等字段),并通过markdown_rag_metadata_enabled/pdf_rag_metadata_enabled追加source_uridocument_idchunk_idchunk_indexcontent_hash等 RAG 元数据列;注意仅支持单栏(自上而下)排版的 PDF。

Source 配置示例

读取 ORC 文件:

OssJindoFile { path = "/seatunnel/orc" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" file_format_type = "orc" }

读取 JSON 文件并指定 schema:

OssJindoFile { 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 } } }

二进制文件同步(multimodal)

源与 Sink 同时使用binary格式,即可把图片、压缩包等任意格式文件从 OSS 迁移到 S3、HDFS 等其他存储:

env { parallelism = 1 job.mode = "BATCH" } source { OssJindoFile { bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" path = "/seatunnel/read/binary/" file_format_type = "binary" } } sink { OssJindoFile { bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" path = "/seatunnel/read/binary2/" file_format_type = "binary" } }

按文件名过滤读取

env { parallelism = 1 job.mode = "BATCH" } source { OssJindoFile { bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" path = "/seatunnel/read/binary/" file_format_type = "binary" // 文件示例 abcD2024.csv file_filter_pattern = "abc[DX]*.*" } } sink { Console { } }

Sink 配置详解:向 OSS 写出文件

必选参数

参数类型必填说明
pathstringSink 写入的目标目录(不存在会自动创建)
bucketstringOSS bucket 地址
access_keystringOSS 访问密钥
access_secretstringOSS 访问密钥 Secret
endpointstringOSS 服务端点

常用可选参数

参数类型默认值说明
tmp_pathstring/tmp/seatunnel结果文件先写临时路径,再用mv提交到目标目录,需为 OSS 目录
custom_filenamebooleanfalse是否自定义文件名
file_name_expressionstring${transactionId}custom_filename=true时生效,支持${now}${uuid}变量
filename_time_formatstringyyyy.MM.ddcustom_filename=true时生效,${now}的时间格式
file_format_typestringcsv支持 text/csv/parquet/orc/json/excel/xml/binary/canal_json/debezium_json/maxwell_json
filename_extensionstring-覆盖默认文件扩展名,如.xml.jsondat
field_delimiterstringtext 为\001,csv 为,仅 text/csv 格式生效
row_delimiterstring\n仅 text/csv/json 格式生效
have_partitionbooleanfalse是否处理分区
partition_byarray-have_partition=true时生效,按选定字段分区
partition_dir_expressionstring${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/分区目录表达式
is_partition_field_write_in_filebooleanfalse分区字段是否同时写入数据文件(写 Hive 数据文件时应为 false)
sink_columnsarray写入文件的列,字段顺序决定实际写入顺序;为空则写全部列
is_enable_transactionbooleantrue开启后保证数据不丢失不重复,文件名自动加${transactionId}_前缀
batch_sizeint1000000单个文件的最大行数,与checkpoint.interval共同决定文件切分
compress_codecstringnone压缩编解码器
max_rows_in_memoryint-仅 excel 格式,内存中缓存的最大数据条数
sheet_max_rowsint1048576仅 excel 格式,每 sheet 最大行数
sheet_namestringSheet+随机数仅 excel 格式
csv_string_quote_modeenumMINIMAL仅 csv 格式(ALL/MINIMAL/NONE)
xml_root_tag / xml_row_tag / xml_use_attr_format-RECORDS / RECORD / -仅 xml 格式
single_file_modebooleanfalse每个并行度只输出一个文件,开启后batch_size失效,输出文件名无文件块后缀
create_empty_file_when_no_databooleanfalse上游无数据时仍生成对应数据文件
parquet_avro_write_timestamp_as_int96 / parquet_avro_write_fixed_as_int96-false / -仅 parquet 格式
encodingstringUTF-8仅 json/text/csv/xml 格式
merge_update_eventbooleanfalse仅 canal_json/debezium_json/maxwell_json,将 UPDATE_AFTER 与 UPDATE_BEFORE 合并为 UPDATE 事件
schema_evolution_enabledbooleanfalse为 CDC 管道启用 schema 演进(binary 格式不支持)

关键选项深入说明

file_name_expression与事务前缀${now}表示当前时间(格式由filename_time_format控制),${uuid}表示随机 UUID;当is_enable_transaction=true时,会自动在文件名头部追加${transactionId}_filename_time_format支持的常用符号包括y(年)、M(月)、d(日)、H(时)、m(分)、s(秒)。

is_enable_transaction:为true时通过 2PC 保证写入目标目录的数据不丢失、不重复。官方文档明确该参数当前仅支持true。注意文件最终扩展名取决于file_format_type,其中 text 文件的后缀为txt

batch_size与 checkpoint 的关系:在 SeaTunnel Engine 中,文件行数由batch_sizecheckpoint.interval共同决定——若 checkpoint 间隔足够大,sink writer 会持续写入直到文件行数超过batch_size;若 checkpoint 间隔很小,则每次 checkpoint 触发时都会新建文件。

compress_codec支持范围

  • txt / json / csv:lzonone
  • orc:lzosnappylz4zlibnone
  • parquet:lzosnappylz4gzipbrotlizstdnone
  • excel:不支持任何压缩格式

schema_evolution_enabled(CDC 场景):置为true后,文件 Sink 可在运行期处理 CDC schema 变更事件(ADD/DROP/RENAME/MODIFY COLUMN),每次 schema 变更会关闭当前输出文件并按新 schema 打开新文件,无需重启作业。支持除binary外的所有格式(binary 会在作业启动时报配置校验错误);当have_partition=true时,不允许删除partition_by中的分区列;若保持默认false而上游 CDC 源开启了schema-changes.enabled=true,则AlterTableEvent到达 Sink 时会立即抛出可操作的错误提示(Received AlterTableEvent but schema_evolution_enabled=false at this sink. ...),默认的 CDC 源配置(schema-changes.enabled = false)不受影响。已知限制:schema 变更与 checkpoint 不是原子的,作业恰好在文件轮转与 schema 元数据更新之间的窄窗口崩溃时,恢复后可能按变更前 schema 写行。

CDC 管道示例:

OssJindoFile { path = "/tmp/cdc/${table_name}" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" file_format_type = "parquet" schema_evolution_enabled = true have_partition = true partition_by = ["updated_at_month"] }

Sink 配置示例

text 格式 + 分区 + 自定义文件名

env { parallelism = 1 job.mode = "BATCH" } sink { OssJindoFile { path="/seatunnel/sink" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxx" access_secret = "xxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" 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}" filename_time_format = "yyyy.MM.dd" sink_columns = ["name","age"] is_enable_transaction = true } }

parquet 格式 + 指定写出列

env { parallelism = 1 job.mode = "BATCH" } sink { OssJindoFile { path = "/seatunnel/sink" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" file_format_type = "parquet" sink_columns = ["name","age"] } }

orc 格式最小配置

env { parallelism = 1 job.mode = "BATCH" } sink { OssJindoFile { path="/seatunnel/sink" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxx" access_secret = "xxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" file_format_type = "orc" } }

canal_json 格式 + CDC 事件合并

env { parallelism = 1 job.mode = "BATCH" } sink { OssJindoFile { path = "/seatunnel/sink" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxx" access_secret = "xxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" file_format_type = "canal_json" merge_update_event = true } }

常见问题与使用建议

  • Hadoop 版本限制:连接器内部走 HDFS 协议访问 OSS,仅支持 Hadoop 2.9.X+;使用 Spark/Flink 时需自行保证集群已集成对应版本的 Hadoop,使用 SeaTunnel Zeta 时lib目录已内置 Hadoop jar。
  • Jindo jar 缺失:作业启动报类加载失败时,优先检查每个节点${SEATUNNEL_HOME}/lib下是否存在jindo-sdk-4.6.1.jarjindo-core-4.6.1.jar
  • 事务与文件前缀:开启is_enable_transaction后文件名会带${transactionId}_前缀,若对下游文件名有严格要求,可结合file_name_expression设计命名规则。
  • 分区字段重复定义:Source 侧开启parse_partition_from_path后,不要在 schema 中重复定义从路径解析出的分区字段。
  • XML 文件被拒:若 XML 含<!DOCTYPE ...>声明会读取失败,属于 XXE 加固的安全设计,需在接入前去除 DOCTYPE。

如需完整的参数列表与更多示例,可直接查阅仓库中的 OssJindoFile Source 文档、OssJindoFile Sink 文档、Sink Common Options 与 Source Common Options,并在 connector-file-jindo-oss 模块中查看具体实现。

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

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

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

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

立即咨询