SeaTunnel OssFile Sink 连接器全解析从零构建写入阿里云 OSS 的数据同步作业【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文围绕 SeaTunnelApache SeaTunnel连接器体系中用于将数据写入阿里云对象存储 OSSObject Storage Service的OssFile Sink插件展开完整覆盖其支持引擎、依赖装配、文件格式与数据类型映射、全部可配置参数、分区与自定义文件名机制、基于 2PC 的精确一次语义以及 text / parquet / orc / 多表等真实可运行的作业配置示例。读完本文你将能够独立完成 OSS 读写作业的配置、排错与调优。一、插件定位与核心能力OssFile 是 SeaTunnel 的 V2 文件类 Sink 连接器用于把上游Source流入的数据以指定文件格式写入阿里云 OSS。它基于 Hadoop 的AliyunOSSFileSystem实现从源码结构看其入口 OssFileSink 继承自文件 Sink 体系的通用基类BaseMultipleTableFileSink因此天然支持多表写入能力。支持引擎SparkFlinkSeaTunnel ZetaSeaTunnel 自研引擎二、使用依赖与 Jar 装配OSS 通过 Hadoop 的oss://协议访问因此对 Hadoop 相关依赖有硬性要求且不同引擎装配位置不同。Spark / Flink 引擎必须确保 Spark / Flink 集群已集成 HadoopSeaTunnel 官方测试使用的 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/目录中存在以下四个 Jarseatunnel-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.version3.4.1、hadoop-aliyun.version3.1.4、jdom.version1.1且hadoop-aliyun、aliyun-sdk-oss、jdom均以provided作用域引入——这意味着运行时必须由用户自行把这些依赖放到正确位置这正是本文上面两步装配要求的来源。源码层面的依赖映射在 OssHadoopConf.java 中可以看到连接器将配置中的bucket、access_key、access_secret、endpoint翻译为 Hadoop OSS 文件系统所需的参数文件系统实现类org.apache.hadoop.fs.aliyun.oss.AliyunOSSFileSystemURL Schemeossaccess_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 数据类型STRINGSTRINGBOOLEANBOOLEANTINYINTBYTESMALLINTSHORTINTINTBIGINTLONGFLOATFLOATDOUBLEDOUBLEDECIMALDECIMALBYTESBINARYDATEDATETIME / TIMESTAMPTIMESTAMPROWSTRUCTNULL不支持的数据类型ARRAYLISTMapMapParquet 文件类型SeaTunnel 数据类型Parquet 数据类型STRINGSTRINGBOOLEANBOOLEANTINYINTINT_8SMALLINTINT_16INTINT32BIGINTINT64FLOATFLOATDOUBLEDOUBLEDECIMALDECIMALBYTESBINARYDATEDATETIME / TIMESTAMPTIMESTAMP_MILLISROWGroupTypeNULL不支持的数据类型ARRAYLISTMapMap五、选项总览下表为 OssFile Sink 的全部选项与官方文档及 OssFileSinkFactory#optionRule 中声明的必填/可选约束保持一致名称类型必需默认值描述pathstring是-Sink 写入的 OSS 路径。配合bucket实际位置为oss://bucketpathtmp_pathstring否/tmp/seatunnel结果文件先写入 tmp 路径之后用mv将 tmp 目录提交到目标目录因此需要一个 OSS 目录bucketstring是-OSS 文件系统的桶地址例如oss://tyrantlucifer-image-bedaccess_keystring是-OSS 桶的访问密钥access_secretstring是-OSS 桶的访问密钥密钥endpointstring是-OSS 端点例如oss-cn-beijing.aliyuncs.comcustom_filenameboolean否false是否需要自定义文件名file_name_expressionstring否${transactionId}仅在custom_filename为 true 时使用filename_time_formatstring否yyyy.MM.dd仅在custom_filename为 true 时使用file_format_typestring否csv文件格式类型支持text、csv、parquet、orc、json、excel、xml、binary、canal_json、debezium_json、maxwell_jsonfield_delimiterstring否\001仅当file_format_type为 text 时使用row_delimiterstring否\n仅当file_format_type为 text、csv、json 时使用have_partitionboolean否false是否需要处理分区partition_byarray否-只有在have_partition为 true 时才使用partition_dir_expressionstring否${k0}${v0}/${k1}${v1}/.../${kn}${vn}/只有在have_partition为 true 时才使用is_partition_field_write_in_fileboolean否false只有在have_partition为 true 时才使用sink_columnsarray否空当此参数为空时所有字段都是接收列is_enable_transactionboolean否true若为true写入目标目录的数据不会丢失或重复为true时自动在文件名前缀添加${transactionId}_batch_sizeint否1000000单个文件的最大行数。对于 SeaTunnel Engine文件中的行数由batch_size和checkpoint.interval共同决定compress_codecstring否none文件的压缩编解码器。Excel 格式不支持任何压缩格式common-optionsobject否-Sink 插件通用参数详见 Sink 常用选项max_rows_in_memoryint否-仅当file_format_type为 excel 时使用sheet_max_rowsint否1048576仅当file_format_type为 excel 时使用每个工作表允许写入的最大行数sheet_namestring否Sheet${Random number}仅当file_format_type为 excel 时使用csv_string_quote_modeenum否MINIMAL仅在 file_format 为 csv 时使用xml_root_tagstring否RECORDS仅在 file_format 为 xml 时使用xml_row_tagstring否RECORD仅在 file_format 为 xml 时使用xml_use_attr_formatboolean否-仅在 file_format 为 xml 时使用single_file_modeboolean否false每个并行处理只会输出一个文件。启用此参数后batch_size将不再生效输出文件名没有文件块后缀create_empty_file_when_no_databoolean否false当上游没有数据同步时仍然会生成相应的数据文件parquet_avro_write_timestamp_as_int96boolean否false仅在 file_format 为 parquet 时使用parquet_avro_write_fixed_as_int96array否-仅在 file_format 为 parquet 时使用enable_header_writeboolean否false仅当file_format_type为 text、csv 时使用。false不写标头true写标头encodingstring否UTF-8仅当file_format_type为 json、text、csv、xml 时使用schema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST在开启同步任务之前对目标路径进行不同的处理data_save_modeEnum否APPEND_DATA在开启同步任务之前对目标路径中的数据文件进行不同的处理merge_update_eventboolean否false仅当file_format_type为 canal_json、debezium_json、maxwell_json 时使用schema_evolution_enabledboolean否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。常用时间符号SymbolDescriptionyYearMMonthdDay of monthHHour in day (0-23)mMinute in hoursSecond in minutefile_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_filehave_partition决定是否启用分区目录。当为true时partition_byarray指定按哪些字段分区partition_dir_expression指定分区目录表达式默认${k0}${v0}/${k1}${v1}/.../${kn}${vn}/其中k0是第一个分区字段名v0是第一个分区字段的值is_partition_field_write_in_file为true时分区字段及其值会一并写入数据文件。如果需要生成 Hive 可直接识别的数据文件该值应设为falseHive 通过目录结构识别分区。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]各格式支持的压缩编解码器txtlzo、nonejsonlzo、nonecsvlzo、noneorclzo、snappy、lz4、zlib、noneparquetlzo、snappy、lz4、gzip、brotli、zstd、none提示excel 类型不支持任何压缩格式。single_file_mode / create_empty_file_when_no_datasingle_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_formatxml_root_tagXML 文件中根元素的标签名默认RECORDSxml_row_tagXML 文件中数据行的标签名默认RECORDxml_use_attr_format是否使用标签属性格式处理数据。max_rows_in_memory / sheet_max_rows / sheet_nameExcel 专属max_rows_in_memoryExcel 格式下内存中可缓存的最大数据项数sheet_max_rows每个工作表允许写入的最大行数默认1048576sheet_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-optionsplugin_input、parallelism、metadata_datasource_id等 Sink 插件通用参数详见 Sink 常用选项。其中plugin_input当不指定时当前插件处理配置文件中上一个插件输出的数据集指定后则处理该参数对应的数据集parallelism未指定时继承env中的并行度指定时覆盖之metadata_datasource_id从外部元数据服务获取连接配置的数据源 ID。六、schema_evolution_enabledCDC 场景下的 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_enabledfalse at this sink. Either set schema_evolution_enabledtrue to handle schema changes, or set schema-changes.enabledfalse 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_modeAPPEND_DATA } }执行后文件将写入形如oss://tyrantlucifer-image-bed/seatunnel/sink/age18/transactionId_2026.09.18.txt的目录结构中且由于is_partition_field_write_in_file trueage字段值会同时出现在数据行中。示例二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_modeAPPEND_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_modeAPPEND_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 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 } } } }, { schema { table fake2 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 } } } } ] } } 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_modeAPPEND_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),仅供参考
