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.OssJindoFile与seatunnel.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_json、debezium_json、maxwell_json三种 CDC 事件格式;Source 侧则额外支持markdown与pdf两种文档解析格式(Markdown 只支持读取、不支持写出)。
环境依赖与 Jindo SDK 安装
该连接器通过 Jindo(阿里云 EMR)SDK 走 HDFS 协议访问 OSS,因此依赖 Jindo 与 Hadoop 相关 jar 包,官方文档给出了明确的安装前提:
- 下载
jindosdk-4.6.1.tar.gz,解压后将其lib目录下的jindo-sdk-4.6.1.jar与jindo-core-4.6.1.jar复制到${SEATUNNEL_HOME}/lib下,且每个运行作业的节点都需要放置。 - 如果使用 Spark/Flink,需要确保集群已集成 Hadoop,官方测试的 Hadoop 版本为 2.x。
- 如果使用 SeaTunnel Engine,安装 SeaTunnel Engine 时已自动集成 Hadoop jar,可在
${SEATUNNEL_HOME}/lib下确认。 - 由于为支持更多文件类型而内部走 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_key | fs.oss.accessKeyId | OSS 访问密钥 ID |
access_secret | fs.oss.accessKeySecret | OSS 访问密钥 Secret |
endpoint | fs.oss.endpoint | OSS 服务端点 |
| (固定值) | fs.AbstractFileSystem.oss.impl | com.aliyun.emr.fs.oss.OSS |
| (固定值) | fs.oss.impl | com.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 bucketOssFileSource.java 与 OssFileSink.java 分别继承BaseFileSource与BaseFileSink,仅实现initHadoopConf()与getPluginName()两个方法,插件名来自FileSystemType.OSS_JINDO(值为OssJindoFile,见 FileSystemType.java),其余文件格式解析、分区处理、事务提交等全部复用 connector-file-base 的通用实现。
另外,OssJindoFactoryTest.java 验证了OssFileSourceFactory与OssFileSinkFactory的optionRule()均能正常构造,可作为连接器注册与选项规则的自检入口。
Source 配置详解:从 OSS 读取文件
必选参数
| 参数 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| path | string | 是 | - | 源文件路径 |
| file_format_type | string | 是 | - | 文件类型 |
| bucket | string | 是 | - | OSS bucket 地址,如oss://tyrantlucifer-image-bed |
| access_key | string | 是 | - | OSS 访问密钥 |
| access_secret | string | 是 | - | OSS 访问密钥 Secret |
| endpoint | string | 是 | - | OSS 服务端点,如oss-cn-beijing.aliyuncs.com |
常用可选参数
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| read_columns | list | - | 读取列清单,实现字段投影 |
| delimiter / field_delimiter | string | text 为\001,csv 为, | 字段分隔符(delimiter将在 2.3.5 之后废弃,请改用field_delimiter) |
| row_delimiter | string | \n | 行分隔符 |
| parse_partition_from_path | boolean | true | 是否从文件路径解析分区键与分区值 |
| date_format | string | yyyy-MM-dd | 字符串转 date 的格式 |
| datetime_format | string | yyyy-MM-dd HH:mm:ss | 字符串转 datetime 的格式 |
| time_format | string | HH:mm:ss | 字符串转 time 的格式 |
| skip_header_row_number | long | 0 | 跳过前 N 行(仅 text 与 csv) |
| schema | config | - | 上游数据结构定义 |
| file_filter_pattern | string | - | 基于正则的文件过滤模式 |
| filename_extension | string | - | 按扩展名过滤文件 |
| compress_codec | string | none | 压缩编解码器(lzo/none,orc/parquet 自动识别) |
| archive_compress_codec | string | none | 归档压缩格式(ZIP/TAR/TAR_GZ/GZ/NONE) |
| encoding | string | UTF-8 | 文件编码 |
| null_format | string | - | 表示 null 的字符串(如\N) |
| file_filter_modified_start / file_filter_modified_end | string | - | 按修改时间过滤(yyyy-MM-dd HH:mm:ss) |
| quote_char | string | " | CSV 字段包围字符 |
| escape_char | string | - | CSV 转义字符 |
| recursive_file_scan | boolean | true | 是否递归扫描子目录 |
| sort_files_by_modification_time | boolean | false | 是否按修改时间倒序排序文件 |
其中关键选项的用法细节如下:
schema与field_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=tyrantlucifer、age=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_id、element_type、heading_level、text等字段),并通过markdown_rag_metadata_enabled/pdf_rag_metadata_enabled追加source_uri、document_id、chunk_id、chunk_index、content_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 写出文件
必选参数
| 参数 | 类型 | 必填 | 说明 |
|---|---|---|---|
| path | string | 是 | Sink 写入的目标目录(不存在会自动创建) |
| bucket | string | 是 | OSS bucket 地址 |
| access_key | string | 是 | OSS 访问密钥 |
| access_secret | string | 是 | OSS 访问密钥 Secret |
| endpoint | string | 是 | OSS 服务端点 |
常用可选参数
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
| tmp_path | string | /tmp/seatunnel | 结果文件先写临时路径,再用mv提交到目标目录,需为 OSS 目录 |
| custom_filename | boolean | false | 是否自定义文件名 |
| file_name_expression | string | ${transactionId} | 仅custom_filename=true时生效,支持${now}、${uuid}变量 |
| filename_time_format | string | yyyy.MM.dd | 仅custom_filename=true时生效,${now}的时间格式 |
| file_format_type | string | csv | 支持 text/csv/parquet/orc/json/excel/xml/binary/canal_json/debezium_json/maxwell_json |
| filename_extension | string | - | 覆盖默认文件扩展名,如.xml、.json、dat |
| field_delimiter | string | text 为\001,csv 为, | 仅 text/csv 格式生效 |
| row_delimiter | string | \n | 仅 text/csv/json 格式生效 |
| have_partition | boolean | false | 是否处理分区 |
| partition_by | array | - | 仅have_partition=true时生效,按选定字段分区 |
| partition_dir_expression | string | ${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/ | 分区目录表达式 |
| is_partition_field_write_in_file | boolean | false | 分区字段是否同时写入数据文件(写 Hive 数据文件时应为 false) |
| sink_columns | array | 空 | 写入文件的列,字段顺序决定实际写入顺序;为空则写全部列 |
| is_enable_transaction | boolean | true | 开启后保证数据不丢失不重复,文件名自动加${transactionId}_前缀 |
| batch_size | int | 1000000 | 单个文件的最大行数,与checkpoint.interval共同决定文件切分 |
| compress_codec | string | none | 压缩编解码器 |
| max_rows_in_memory | int | - | 仅 excel 格式,内存中缓存的最大数据条数 |
| sheet_max_rows | int | 1048576 | 仅 excel 格式,每 sheet 最大行数 |
| sheet_name | string | Sheet+随机数 | 仅 excel 格式 |
| csv_string_quote_mode | enum | MINIMAL | 仅 csv 格式(ALL/MINIMAL/NONE) |
| xml_root_tag / xml_row_tag / xml_use_attr_format | - | RECORDS / RECORD / - | 仅 xml 格式 |
| single_file_mode | boolean | false | 每个并行度只输出一个文件,开启后batch_size失效,输出文件名无文件块后缀 |
| create_empty_file_when_no_data | boolean | false | 上游无数据时仍生成对应数据文件 |
| parquet_avro_write_timestamp_as_int96 / parquet_avro_write_fixed_as_int96 | - | false / - | 仅 parquet 格式 |
| encoding | string | UTF-8 | 仅 json/text/csv/xml 格式 |
| merge_update_event | boolean | false | 仅 canal_json/debezium_json/maxwell_json,将 UPDATE_AFTER 与 UPDATE_BEFORE 合并为 UPDATE 事件 |
| schema_evolution_enabled | boolean | false | 为 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_size与checkpoint.interval共同决定——若 checkpoint 间隔足够大,sink writer 会持续写入直到文件行数超过batch_size;若 checkpoint 间隔很小,则每次 checkpoint 触发时都会新建文件。
compress_codec支持范围:
- txt / json / csv:
lzo、none - orc:
lzo、snappy、lz4、zlib、none - parquet:
lzo、snappy、lz4、gzip、brotli、zstd、none - 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.jar与jindo-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),仅供参考