SeaTunnel OssFile Sink 连接器全解析:从零构建写入阿里云 OSS 的数据同步作业
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
本文围绕 SeaTunnel(Apache SeaTunnel)连接器体系中用于将数据写入阿里云对象存储 OSS(Object Storage Service)的OssFile Sink插件展开,完整覆盖其支持引擎、依赖装配、文件格式与数据类型映射、全部可配置参数、分区与自定义文件名机制、基于 2PC 的精确一次语义,以及 text / parquet / orc / 多表等真实可运行的作业配置示例。读完本文,你将能够独立完成 OSS 读写作业的配置、排错与调优。
一、插件定位与核心能力
OssFile 是 SeaTunnel 的 V2 文件类 Sink 连接器,用于把上游(Source)流入的数据以指定文件格式写入阿里云 OSS。它基于 Hadoop 的AliyunOSSFileSystem实现,从源码结构看,其入口 OssFileSink 继承自文件 Sink 体系的通用基类BaseMultipleTableFileSink,因此天然支持多表写入能力。
支持引擎
- Spark
- Flink
- SeaTunnel Zeta(SeaTunnel 自研引擎)
二、使用依赖与 Jar 装配
OSS 通过 Hadoop 的oss://协议访问,因此对 Hadoop 相关依赖有硬性要求,且不同引擎装配位置不同。
Spark / Flink 引擎
- 必须确保 Spark / Flink 集群已集成 Hadoop,SeaTunnel 官方测试使用的 Hadoop 版本为 2.x。
- 必须确保
${SEATUNNEL_HOME}/plugins/目录下的hadoop-aliyun-xx.jar、aliyun-sdk-oss-xx.jar与jdom-xx.jar版本与集群 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/目录中存在以下四个 Jar:
seatunnel-shade-hadoop3-uber-3.1.4-3.0.0.jaraliyun-sdk-oss-3.4.1.jarhadoop-aliyun-3.1.4.jarjdom-1.1.jar
这三个版本号并非随意指定,而是由 connector 的 Maven 工程声明:在 connector-file-oss/pom.xml 中可以看到aliyun.sdk.oss.version=3.4.1、hadoop-aliyun.version=3.1.4、jdom.version=1.1,且hadoop-aliyun、aliyun-sdk-oss、jdom均以provided作用域引入——这意味着运行时必须由用户自行把这些依赖放到正确位置,这正是本文上面两步装配要求的来源。
源码层面的依赖映射
在 OssHadoopConf.java 中可以看到,连接器将配置中的bucket、access_key、access_secret、endpoint翻译为 Hadoop OSS 文件系统所需的参数:
- 文件系统实现类:
org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystem - URL Scheme:
oss access_key→ Hadoop 的ACCESS_KEY_IDaccess_secret→ Hadoop 的ACCESS_KEY_SECRETendpoint→ Hadoop 的ENDPOINT_KEY
理解这一层映射,有助于在 Hadoop 与 OSS 集成出现认证或协议错误时快速定位问题根源。
三、关键特性
- 多模态(Multimodal):支持以二进制文件格式读写任何格式的文件,例如视频、图片等。简而言之,任何文件都可以同步到目标位置。
- 精确一次(Exactly-Once):默认通过 2PC(两阶段提交)Commit 机制保证数据写入不丢不重。
- 支持多表写入:可从上游提取多张表的元数据,将不同表写入不同目录。
- 文件格式类型:
text、csv、parquet、orc、json、excel、xml、binary、canal_json、debezium_json、maxwell_json。
关于这些特性的概念性说明可参考 Connector V2 特性说明。
精确一次的实现原理(源码级)
文件类 Sink 的精确一次并不是把数据直接写到最终目录,而是采用“临时事务目录 + 提交期移动文件”的策略。在 AbstractWriteStrategy.java 中可以看到完整链路:
beginTransaction(checkpointId):每个 Checkpoint 开始时生成新的transactionId与对应的临时事务目录,后续数据先写入该临时目录;prepareCommit():Checkpoint 快照阶段关闭当前文件,产出FileCommitInfo(记录了需要移动的文件与分区信息),完成 2PC 的“准备”阶段;abortPrepare()/abortPrepare(transactionId):若提交失败,则直接删除整个临时事务目录,实现回滚;- 提交成功后,文件才通过
mv操作移动到目标目录,从而保证最终目录中不会出现半成品或重复数据。
默认情况下is_enable_transaction = true,此时文件名会自动加上${transactionId}_前缀,其前缀拼接逻辑同样可以在AbstractWriteStrategy.generateFileName()中看到。
四、数据类型映射
写入csv、text文件类型时,所有列都会被转换为字符串。对于orc与parquet这类列式格式,SeaTunnel 数据类型与文件格式类型之间的映射关系如下。
Orc 文件类型
| SeaTunnel 数据类型 | Orc 数据类型 |
|---|---|
| STRING | STRING |
| BOOLEAN | BOOLEAN |
| TINYINT | BYTE |
| SMALLINT | SHORT |
| INT | INT |
| BIGINT | LONG |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| DECIMAL | DECIMAL |
| BYTES | BINARY |
| DATE | DATE |
| TIME / TIMESTAMP | TIMESTAMP |
| ROW | STRUCT |
| NULL | 不支持的数据类型 |
| ARRAY | LIST |
| Map | Map |
Parquet 文件类型
| SeaTunnel 数据类型 | Parquet 数据类型 |
|---|---|
| STRING | STRING |
| BOOLEAN | BOOLEAN |
| TINYINT | INT_8 |
| SMALLINT | INT_16 |
| INT | INT32 |
| BIGINT | INT64 |
| FLOAT | FLOAT |
| DOUBLE | DOUBLE |
| DECIMAL | DECIMAL |
| BYTES | BINARY |
| DATE | DATE |
| TIME / TIMESTAMP | TIMESTAMP_MILLIS |
| ROW | GroupType |
| NULL | 不支持的数据类型 |
| ARRAY | LIST |
| Map | Map |
五、选项总览
下表为 OssFile Sink 的全部选项(与官方文档及 OssFileSinkFactory#optionRule 中声明的必填/可选约束保持一致):
| 名称 | 类型 | 必需 | 默认值 | 描述 |
|---|---|---|---|---|
| path | string | 是 | - | Sink 写入的 OSS 路径。配合bucket,实际位置为oss://<bucket><path> |
| tmp_path | string | 否 | /tmp/seatunnel | 结果文件先写入 tmp 路径,之后用mv将 tmp 目录提交到目标目录,因此需要一个 OSS 目录 |
| bucket | string | 是 | - | OSS 文件系统的桶地址,例如oss://tyrantlucifer-image-bed |
| access_key | string | 是 | - | OSS 桶的访问密钥 |
| access_secret | string | 是 | - | OSS 桶的访问密钥(密钥) |
| endpoint | string | 是 | - | OSS 端点,例如oss-cn-beijing.aliyuncs.com |
| custom_filename | boolean | 否 | false | 是否需要自定义文件名 |
| file_name_expression | string | 否 | "${transactionId}" | 仅在custom_filename为 true 时使用 |
| filename_time_format | string | 否 | "yyyy.MM.dd" | 仅在custom_filename为 true 时使用 |
| file_format_type | string | 否 | "csv" | 文件格式类型,支持text、csv、parquet、orc、json、excel、xml、binary、canal_json、debezium_json、maxwell_json |
| field_delimiter | string | 否 | '\001' | 仅当file_format_type为 text 时使用 |
| row_delimiter | string | 否 | "\n" | 仅当file_format_type为 text、csv、json 时使用 |
| have_partition | boolean | 否 | false | 是否需要处理分区 |
| partition_by | array | 否 | - | 只有在have_partition为 true 时才使用 |
| partition_dir_expression | string | 否 | "${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/" | 只有在have_partition为 true 时才使用 |
| is_partition_field_write_in_file | boolean | 否 | false | 只有在have_partition为 true 时才使用 |
| sink_columns | array | 否 | (空) | 当此参数为空时,所有字段都是接收列 |
| is_enable_transaction | boolean | 否 | true | 若为true,写入目标目录的数据不会丢失或重复;为true时自动在文件名前缀添加${transactionId}_ |
| batch_size | int | 否 | 1000000 | 单个文件的最大行数。对于 SeaTunnel Engine,文件中的行数由batch_size和checkpoint.interval共同决定 |
| compress_codec | string | 否 | none | 文件的压缩编解码器。Excel 格式不支持任何压缩格式 |
| common-options | object | 否 | - | Sink 插件通用参数,详见 Sink 常用选项 |
| max_rows_in_memory | int | 否 | - | 仅当file_format_type为 excel 时使用 |
| sheet_max_rows | int | 否 | 1048576 | 仅当file_format_type为 excel 时使用;每个工作表允许写入的最大行数 |
| sheet_name | string | 否 | Sheet${Random number} | 仅当file_format_type为 excel 时使用 |
| csv_string_quote_mode | enum | 否 | MINIMAL | 仅在 file_format 为 csv 时使用 |
| xml_root_tag | string | 否 | RECORDS | 仅在 file_format 为 xml 时使用 |
| xml_row_tag | string | 否 | RECORD | 仅在 file_format 为 xml 时使用 |
| xml_use_attr_format | boolean | 否 | - | 仅在 file_format 为 xml 时使用 |
| single_file_mode | boolean | 否 | false | 每个并行处理只会输出一个文件。启用此参数后batch_size将不再生效,输出文件名没有文件块后缀 |
| create_empty_file_when_no_data | boolean | 否 | false | 当上游没有数据同步时,仍然会生成相应的数据文件 |
| parquet_avro_write_timestamp_as_int96 | boolean | 否 | false | 仅在 file_format 为 parquet 时使用 |
| parquet_avro_write_fixed_as_int96 | array | 否 | - | 仅在 file_format 为 parquet 时使用 |
| enable_header_write | boolean | 否 | false | 仅当file_format_type为 text、csv 时使用。false:不写标头,true:写标头 |
| encoding | string | 否 | "UTF-8" | 仅当file_format_type为 json、text、csv、xml 时使用 |
| schema_save_mode | Enum | 否 | CREATE_SCHEMA_WHEN_NOT_EXIST | 在开启同步任务之前,对目标路径进行不同的处理 |
| data_save_mode | Enum | 否 | APPEND_DATA | 在开启同步任务之前,对目标路径中的数据文件进行不同的处理 |
| merge_update_event | boolean | 否 | false | 仅当file_format_type为 canal_json、debezium_json、maxwell_json 时使用 |
| schema_evolution_enabled | boolean | 否 | false | 开启 Schema 演变支持,适用于 CDC 管道。为 true 时,来自上游的 ADD/DROP/RENAME/MODIFY 列事件无需重启作业即可应用到 Sink。不支持 binary 格式 |
说明:在
OssFileSinkFactory#optionRule()中,path、bucket、access_key、access_secret、endpoint均被声明为required,与上表一致;其余选项按file_format_type、custom_filename、have_partition等前置条件进行conditional校验,配置工具(如 Web 控制台)会根据该规则动态展示可用选项。
核心参数详解
path [string]
目标目录路径,必填。注意path与bucket是拼接关系而非覆盖关系:例如配置bucket = "oss://seatunnel-test"且path = "/warehouse/events"时,文件实际写入oss://seatunnel-test/warehouse/events。
bucket [string]
OSS 文件系统的桶地址,例如oss://tyrantlucifer-image-bed。
access_key / access_secret [string]
OSS 桶的访问密钥与密钥,对应阿里云 AccessKey 体系,最终会透传给 Hadoop OSS 文件系统(fs.oss.accessKeyId/fs.oss.accessKeySecret),建议通过环境变量或密钥管理平台注入,避免明文落入作业配置。
endpoint [string]
OSS 端点,例如oss-cn-beijing.aliyuncs.com。选择与 bucket 所在地域一致的 endpoint 可降低延迟与流量费用。
custom_filename [boolean] 与 file_name_expression [string]
是否自定义文件名。file_name_expression描述了在path中创建的文件表达式,其中可以引入变量${now}或${uuid},例如test_${uuid}_${now};${now}表示当前时间,其格式由filename_time_format决定。
需要注意:如果is_enable_transaction为true,连接器会在文件名开头自动添加${transactionId}_前缀(源码中generateFileName()会先做变量替换,再拼上事务前缀与后缀)。
filename_time_format [String]
当file_name_expression中包含${now}时,此参数指定时间部分的格式,默认值为yyyy.MM.dd。常用时间符号:
| Symbol | Description |
|---|---|
| y | Year |
| M | Month |
| d | Day of month |
| H | Hour in day (0-23) |
| m | Minute in hour |
| s | Second in minute |
file_format_type [string]
支持text、csv、parquet、orc、json、excel、xml、binary、canal_json、debezium_json、maxwell_json。最终文件名以文件格式类型的后缀结尾,其中文本文件后缀为txt。
field_delimiter / row_delimiter [string]
field_delimiter是数据行中列之间的分隔符,仅用于 text 格式,默认\001(即 ASCII 单位分隔符,可避免与数据内容冲突);row_delimiter是文件中行之间的分隔符,用于 text、csv、json 格式,默认\n。
have_partition / partition_by / partition_dir_expression / is_partition_field_write_in_file
have_partition决定是否启用分区目录。当为true时:
partition_by(array)指定按哪些字段分区;partition_dir_expression指定分区目录表达式,默认"${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/",其中k0是第一个分区字段名,v0是第一个分区字段的值;is_partition_field_write_in_file为true时,分区字段及其值会一并写入数据文件。如果需要生成 Hive 可直接识别的数据文件,该值应设为false(Hive 通过目录结构识别分区)。
sink_columns [array]
指定哪些列需要写入文件,默认取 Transform 或 Source 输出的所有列。字段在数组中的顺序决定了文件实际写入的列顺序。
is_enable_transaction [boolean]
默认true,通过 2PC 保证写入目标目录的数据不丢失、不重复。为true时文件名自动添加${transactionId}_前缀。当前版本仅支持true。
batch_size [int]
单个文件的最大行数,默认 1000000。对于 SeaTunnel Engine,文件中的行数由batch_size和checkpoint.interval共同决定:如果checkpoint.interval足够大,writer 会持续写入直到文件行数超过batch_size再滚动新文件;如果checkpoint.interval较小,则每个 Checkpoint 触发时都会创建新文件(每个 Checkpoint 对应一个新事务)。
compress_codec [string]
各格式支持的压缩编解码器:
- txt:
lzo、none - json:
lzo、none - csv:
lzo、none - orc:
lzo、snappy、lz4、zlib、none - parquet:
lzo、snappy、lz4、gzip、brotli、zstd、none
提示:excel 类型不支持任何压缩格式。
single_file_mode / create_empty_file_when_no_data
single_file_mode(默认 false):每个并行子任务只输出一个文件;启用后batch_size不生效,输出文件名不带文件块后缀。注意源码 BaseFileSink 中对该模式有限制:开启 Checkpoint 或流式模式下不支持该模式。create_empty_file_when_no_data(默认 false):上游无数据时也生成对应数据文件(在prepareCommit()中通过提前创建输出流实现)。
csv_string_quote_mode [enum]
CSV 字符串引用模式,可选值:
ALL:所有字符串字段都会被引用。MINIMAL:仅对包含特殊字符(字段分隔符、引号字符或行分隔符)的字段加引号。NONE:从不引用字段;当分隔符出现在数据中时,打印器会用转义符作为前缀,若未设置转义符,格式校验会抛出异常。
xml_root_tag / xml_row_tag / xml_use_attr_format
xml_root_tag:XML 文件中根元素的标签名,默认RECORDS;xml_row_tag:XML 文件中数据行的标签名,默认RECORD;xml_use_attr_format:是否使用标签属性格式处理数据。
max_rows_in_memory / sheet_max_rows / sheet_name(Excel 专属)
max_rows_in_memory:Excel 格式下内存中可缓存的最大数据项数;sheet_max_rows:每个工作表允许写入的最大行数,默认1048576;sheet_name:写入的工作表名称,默认Sheet${Random number}。
parquet_avro_write_timestamp_as_int96 / parquet_avro_write_fixed_as_int96
两者均仅适用于 parquet 文件:前者支持将时间戳写入 Parquet INT96;后者支持将 12 字节字段写入 Parquet INT96。
encoding [string]
仅当file_format_type为 json、text、csv、xml 时使用,指定写入文件的编码,默认UTF-8。该参数最终由Charset.forName(encoding)解析,因此传入非标准字符集名称会在此处抛异常。
schema_save_mode [Enum]
同步任务开启前对目标路径的处理策略:
RECREATE_SCHEMA:路径不存在时创建;路径已存在时删除并重新创建。CREATE_SCHEMA_WHEN_NOT_EXIST:路径不存在时创建,存在时直接复用。ERROR_WHEN_SCHEMA_NOT_EXIST:路径不存在时报错。IGNORE:忽略路径的处理。
data_save_mode [Enum]
同步任务开启前对目标路径中数据文件的处理策略:
DROP_DATA:使用路径,但删除路径中已有的数据文件。APPEND_DATA:使用路径,并在路径中追加新文件写入数据。ERROR_WHEN_DATA_EXISTS:路径中已存在数据文件时报错。
merge_update_event [boolean]
仅当file_format_type为 canal_json、debezium_json、maxwell_json 时使用。设为true时,序列化数据时UPDATE_AFTER与UPDATE_BEFORE会合并为UPDATE;设为false时两者不合并。
enable_header_write [boolean]
仅当file_format_type为 text、csv 时使用。false:不写标头;true:写标头。
通用选项(common-options)
plugin_input、parallelism、metadata_datasource_id等 Sink 插件通用参数,详见 Sink 常用选项。其中:
plugin_input:当不指定时,当前插件处理配置文件中上一个插件输出的数据集;指定后则处理该参数对应的数据集;parallelism:未指定时继承env中的并行度,指定时覆盖之;metadata_datasource_id:从外部元数据服务获取连接配置的数据源 ID。
六、schema_evolution_enabled:CDC 场景下的 Schema 演变
设置为true时,文件 Sink 可在运行时处理 CDC Schema 变更事件(ADD COLUMN、DROP COLUMN、RENAME COLUMN、MODIFY COLUMN),无需重启作业。每次 Schema 变更时,当前输出文件会被关闭,并以新 Schema 打开一个新文件。
支持的格式:除binary外的所有文件格式。将schema_evolution_enabled与file_format_type = binary组合使用时,作业启动时会抛出配置校验错误。
分区约束:当have_partition = true时,不允许删除partition_by中列出的分区列,违反时会立即抛出异常——分区列在 Schema 变更过程中必须保持稳定。
当schema_evolution_enabled = false(默认值)时:若上游 CDC Source 配置了schema-changes.enabled = true且 Sink 收到AlterTableEvent,作业会立即抛出如下错误:
Received AlterTableEvent but schema_evolution_enabled=false at this sink. Either set schema_evolution_enabled=true to handle schema changes, or set schema-changes.enabled=false at the CDC source to suppress them.
使用默认 CDC Source 配置(schema-changes.enabled = false)的用户不受影响。
已知限制:Schema 变更与 Checkpoint 不是原子操作。若作业在文件轮转与 Schema 元数据更新之间的窗口期崩溃,恢复后写入的数据行可能使用变更前的 Schema。这是与其他 SeaTunnel Sink 共同存在的已知架构限制,完整的重启后 DDL 正确性支持需要配套的 CDC Source 修复。
CDC 管道中的使用示例:
LocalFile { path = "/tmp/cdc/${table_name}" file_format_type = "parquet" schema_evolution_enabled = true have_partition = true partition_by = ["updated_at_month"] }七、实战:创建 OSS 数据同步作业
下面四个示例均以FakeSource作为上游数据源,演示不同文件格式与特性组合下的完整作业配置。作业配置使用 SeaTunnel 的 HOCON 风格配置文件,可直接复制替换凭证后运行。
示例一:text 格式 + 分区 + 自定义文件名 + 指定列
# 设置要执行的任务的基本配置 env { parallelism = 1 job.mode = "BATCH" } # 创建产品数据源 source { FakeSource { schema = { fields { name = string age = int } } } } # 将数据写入 Oss sink { OssFile { 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 schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode="APPEND_DATA" } }执行后,文件将写入形如oss://tyrantlucifer-image-bed/seatunnel/sink/age=18/<transactionId>_2026.09.18.txt的目录结构中,且由于is_partition_field_write_in_file = true,age字段值会同时出现在数据行中。
示例二:parquet 格式 + 分区 + 指定列
# 设置要执行的任务的基本配置 env { parallelism = 1 job.mode = "BATCH" } # Create a source to product data source { FakeSource { schema = { fields { name = string age = int } } } } # 将数据写入 Oss sink { OssFile { path = "/seatunnel/sink" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" have_partition = true partition_by = ["age"] partition_dir_expression = "${k0}=${v0}" is_partition_field_write_in_file = true file_format_type = "parquet" sink_columns = ["name","age"] schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode="APPEND_DATA" } }示例三:orc 格式的简单配置
# 设置要执行的任务的基本配置 env { parallelism = 1 job.mode = "BATCH" } # Create a source to product data source { FakeSource { schema = { fields { name = string age = int } } } } # 将数据写入 Oss sink { OssFile { path="/seatunnel/sink" bucket = "oss://tyrantlucifer-image-bed" access_key = "xxxxxxxxxxx" access_secret = "xxxxxxxxxxx" endpoint = "oss-cn-beijing.aliyuncs.com" file_format_type = "orc" schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode="APPEND_DATA" } }示例四:多表写入
OssFile Sink 支持多表写入:从上游提取 source 元数据后,可在path中使用${database_name}、${table_name}和${schema_name}三个占位符,将不同表的数据落到不同目录。下面的配置让一个FakeSource产出fake1、fake2两张表,Sink 端通过path = "/tmp/fake_empty/text/${table_name}"自动按表名分流:
env { 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 { FakeSource { tables_configs = [ { schema = { table = "fake1" fields { c_map = "map<string, string>" c_array = "array<int>" 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 = "map<string, string>" c_array = "array<int>" 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 } } } }, { schema = { table = "fake2" fields { c_map = "map<string, string>" c_array = "array<int>" 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 = "map<string, string>" c_array = "array<int>" 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 } } } } ] } } sink { OssFile { bucket = "oss://whale-ops" access_key = "xxxxxxxxxxxxxxxxxxx" access_secret = "xxxxxxxxxxxxxxxxxxx" endpoint = "https://oss-accelerate.aliyuncs.com" path = "/tmp/fake_empty/text/${table_name}" row_delimiter = "\n" partition_dir_expression = "${k0}=${v0}" is_partition_field_write_in_file = true file_name_expression = "${transactionId}_${now}" file_format_type = "text" filename_time_format = "yyyy.MM.dd" is_enable_transaction = true compress_codec = "lzo" schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode="APPEND_DATA" } }多表写入的实现同样有源码支撑:OssFileSink 直接继承BaseMultipleTableFileSink(见 BaseMultipleTableFileSink.java),通过getWriteCatalogTable()返回当前表元数据,路径中的${table_name}等占位符在写入阶段按表展开。
八、运行提示
- 作业提交前请先参考 SeaTunnel 部署方案 完成环境准备,并确认第二节中的依赖 Jar 已正确放置到
plugins/(Spark/Flink)或lib/(Zeta)目录。 - 建议在本地先用
FakeSource配合最小配置做连通性验证,确认endpoint、bucket与凭证无误后,再逐步叠加分区、自定义文件名、压缩等特性。 - 连接器单测
OssFileFactoryTest会校验 Sink/Source 工厂的optionRule()非空(见 OssFileFactoryTest.java),若你基于该插件做二次开发或封装,可通过同样的方式守护配置规则。 - 该连接器的完整变更记录见 connector-file-oss 变更日志。
九、小结
OssFile Sink 通过 Hadoop AliyunOSSFileSystem 屏蔽了 OSS 底层协议差异,向上提供了一套覆盖 text / csv / parquet / orc / json / excel / xml / binary / 三类 CDC JSON 共 11 种文件格式、分区目录、自定义文件名、压缩、多表写入与 2PC 精确一次语义的完整能力。掌握本文的选项语义与源码链路,即可在生产环境中快速构建稳定、可维护的「数据 → OSS」同步管道。
【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考