news 2026/9/16 16:07:38

SeaTunnel Hudi Sink 连接器详解:配置参数、多表写入与 CDC 实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SeaTunnel Hudi Sink 连接器详解:配置参数、多表写入与 CDC 实战指南

SeaTunnel Hudi Sink 连接器详解:配置参数、多表写入与 CDC 实战指南

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

本指南以 SeaTunnel 仓库中 Hudi Sink 官方文档 为主体,结合connector-hudi模块源码深入讲解:如何通过 SeaTunnel 将数据写入 Apache Hudi 表、全部配置参数的语义与默认值、多表写入与 CDC 变更日志的配置方式,以及 Timer Flush 定时刷写等 Zeta 引擎专属能力。读完本文,你将能够独立编写单表 UPSERT、多表同步、CDC 入湖及 S3 存储等完整的 SeaTunnel Hudi 作业配置,并理解其底层写入机制与一致性边界。

概述:SeaTunnel 如何写 Hudi

Hudi Sink 连接器(插件名Hudi)用于将 SeaTunnel 作业中的数据写入 Hudi 表。它基于 Hudi 官方提供的HoodieJavaWriteClient(Java 引擎写入客户端)实现,支持写入 HDFS、本地文件系统以及 S3 兼容对象存储上的 Hudi 表。

从 connector-v2-features 特性定义 看,该连接器支持以下核心能力:

特性支持情况
exactly-once(精确一次)不支持(-),连接器不提供 2PC 两阶段提交写入器
cdc(变更数据捕获)支持(),可基于主键写入 INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE 行类型
support multiple table write(多表写入)支持(),单个作业可同时写入多张 Hudi 表
timer flush(定时刷写)支持(),由 Zeta 引擎注入FlushSignal驱动定时刷写

注意:Hive Metastore 同步。Hudi Sink 只负责写 Hudi 数据文件和.hoodie元数据,不会在 Hive Metastore 中注册表或同步表结构。所有形如hoodie.datasource.hive_sync.*的选项都不是受支持的 Sink 选项,也不会传递给 Hudi 写客户端。当需要 Hive Metastore 注册时,请单独运行 Apache Hudi 的HiveSyncTool或其他表注册流程。

配置总览:两类配置项

Hudi Sink 的配置分为两层:基础配置(Sink 级)表列表配置(table_list 内)。其解析逻辑可在 HudiSinkConfig.java 与 HudiTableConfig.java 中看到。

基础配置

名称类型必填默认值说明
table_dfs_pathstring-Hudi 表数据与元数据的根路径
conf_files_pathstring-HDFS 客户端配置文件(分号分隔的本地路径列表)
table_listArray-多表作业中每张表的独立设置
schema_save_modeenumCREATE_SCHEMA_WHEN_NOT_EXIST作业启动前如何处理目标表结构
data_save_modeenumAPPEND_DATA作业启动前如何处理已存在的表数据
common-optionsConfig-Sink 公共参数

表列表配置(table_list 内)

名称类型必填默认值说明
table_namestring-目标 Hudi 表名
databasestringdefault目标 Hudi 数据库名
table_typeenumCOPY_ON_WRITEHudi 表类型:COPY_ON_WRITEMERGE_ON_READ
op_typeenumINSERT写操作:INSERTUPSERTBULK_INSERT
record_key_fieldsstring-用于构建 Hudi record key 的字段,UPSERT必填
partition_fieldsstring-用于构建分区路径的字段
precombine_fieldstring-用于解决同一记录多次更新的字段
batch_interval_msInt1000当前未使用;刷写仅由batch_size或 checkpoint 触发
batch_sizeInt1000刷写前缓冲的最大行数
insert_shuffle_parallelismInt2insert 操作的 shuffle 并行度
upsert_shuffle_parallelismInt2upsert 操作的 shuffle 并行度
min_commits_to_keepInt20清理时保留的最少 commit 数
max_commits_to_keepInt30清理时保留的最多 commit 数
index_typeenumBLOOMHudi 索引类型:BLOOMSIMPLEGLOBAL_BLOOM
index_class_namestring-自定义 Hudi 索引类的全限定类名
record_byte_sizeInt1024每条记录的平均字节数估算值
cdc_enabledbooleanfalse开启后持久化 Hudi CDC 变更日志数据

配置的组织方式:当作业只写一张表时,可以把table_list内的配置项平铺到外层(即直接写在Hudi {}块下);多表作业时,表级选项必须放在各自table_list条目内,而table_dfs_pathconf_files_pathschema_save_modedata_save_mode仍保持在 Sink 层。这一设计在 HudiTableConfig.of() 中体现:当table_list未配置时,连接器会把外层平铺配置组装成一个单元素列表。

参数间约束

  • record_key_fieldsop_type = UPSERT必填,且是启动期校验的——缺失会在配置解析阶段直接抛出IllegalArgumentException(见 HudiTableConfig.of());
  • BULK_INSERT目前不校验record_key_fields,省略它会在写入阶段抛NullPointerException(而非配置时报错);
  • 多表作业中table_name不允许重复或为空,重复/空表名会在启动时被拒绝;
  • 对 CDC 输入,上游记录必须包含record_key_fields用到的字段;仅在确实需要 Hudi CDC 变更日志时才设置cdc_enabled = true

核心参数详解

table_name [string]

Hudi 表的名称。

database [string]

Hudi 表所属的数据库名,默认值为default。表路径由table_dfs_pathdatabasetable_name三者拼接得到:未指定 database 时为{table_dfs_path}/{table_name},指定后为{table_dfs_path}/{database}/{table_name},该推断逻辑实现在 HudiCatalogUtil.inferTablePath()。

table_dfs_path [string]

Hudi 表的 DFS 根路径,例如hdfs://nameserivce/data/hudi/。该路径下既存放 Hudi 数据文件,也存放.hoodie元数据。

table_type [enum]

Hudi 表类型,取值为COPY_ON_WRITE(写时复制)或MERGE_ON_READ(读时合并)。默认COPY_ON_WRITE

record_key_fields [string]

Hudi 表的记录键字段,用于生成 record key。当op_typeUPSERT必须配置。底层实现中,该字段会参与构建HoodieKey与去重逻辑(见 HudiRecordWriter.prepareRecords()),同一HoodieKey的记录在缓冲区内会被后者覆盖。

partition_fields [string]

Hudi 表的分区键字段,用于生成分区路径。配置后 Hudi 会按该字段对数据进行分区存储。

precombine_field [string]

Hudi 表的 precombine 字段,在实际写入前用于 preCombining,即同一记录有多条更新时,据此字段决定最终保留哪一条。该选项由 PRECOMBINE_FIELD 定义,变更记录显示该参数在 2.3.12 版本中加入。

index_type [string]

Hudi 表的索引类型。当前支持BLOOMSIMPLEGLOBAL_BLOOM,默认BLOOM。索引用于在 upsert 时定位记录所属的文件组。

index_class_name [string]

自定义 Hudi 索引类的全限定类名,例如org.apache.seatunnel.connectors.seatunnel.hudi.index.CustomHudiIndex。从源码看,index_class_nameindex_type是二选一的关系:配置了解类名则用自定义类构建索引,否则使用index_type指定的内置索引(见 HudiUtil.createHoodieJavaWriteClient())。

record_byte_size [Int]

每条记录的字节大小估算值。该值用于帮助估算每个 Hudi 数据文件中的大致记录数,调整它可以有效降低 Hudi 数据文件的写放大(write amplification)。默认1024

conf_files_path [string]

本地路径的 HDFS 配置文件列表(分号分隔),用于初始化 HDFS 客户端以读写 Hudi 表文件。示例:/home/test/hdfs-site.xml;/home/test/core-site.xml;/home/test/yarn-site.xml。底层实现通过Configuration.addResource()逐个加载这些文件(见 HudiUtil.getConfiguration()),加载后的Configuration用于构建HoodieJavaEngineContextHoodieJavaWriteClient

op_type [enum]

Hudi 表的操作类型,取值为insertupsertbulk_insert,默认INSERT。三种操作分别映射到写客户端的insert()upsert()bulkInsert()调用(见 HudiRecordWriter.executeWrite())。其中BULK_INSERT适合大批量初始导入场景。

batch_interval_ms [Int]

为兼容性保留的参数,默认1000当前未使用;刷写仅由batch_size或 checkpoint 触发。若需要在 Zeta 引擎上做定时刷写,请在作业的env块中配置sink.flush.interval(详见下文"Timer Flush"章节)。

batch_size [Int]

刷写前缓冲的最大行数,默认1000。达到该阈值即触发一次 flush(见 HudiRecordWriter.writeRecord())。

insert_shuffle_parallelism [Int]

insert 数据到 Hudi 表时的 shuffle 并行度,默认2。对应写入配置中的.withParallelism(insert, upsert)(见 HudiUtil)。

upsert_shuffle_parallelism [Int]

upsert 数据到 Hudi 表时的 shuffle 并行度,默认2

min_commits_to_keep [Int]

清理(clean)时保留的最少 commit 数,默认20,对应 Hudi 的hoodie.keep.min.commits

max_commits_to_keep [Int]

清理时保留的最多 commit 数,默认30,对应 Hudi 的hoodie.keep.max.commitsmin_commits_to_keepmax_commits_to_keep会一起传给HoodieArchivalConfig.archiveCommitsWith(),用于配置归档策略(见 HudiUtil)。

cdc_enabled [boolean]

是否持久化 CDC 变更日志,默认false。开启后,在必要时持久化变更数据,Hudi 表可以按 CDC 查询模式(CDC query mode)被查询。从源码看,连接器会依据 SeaTunnel 行的RowKind(INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE)将变更标记为"写入"或"删除",DELETEUPDATE_BEFORE走删除路径,INSERTUPDATE_AFTER走写入路径(见 HudiRecordWriter.changeFlag())。

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:数据已存在时抛出错误。

上述 SaveMode 由 HudiSink.getSaveModeHandler() 实现:连接器通过 SPI 发现名为HudiCatalogFactory创建 Catalog,再交给DefaultSaveModeHandlerschema_save_mode/data_save_mode执行建表、删表、清数据等前置动作。

common options

Sink 插件的公共参数,详见 Sink Common Options。

底层写入流程与一致性语义

理解连接器的工作方式有助于正确设置参数。以 HudiSinkWriter 与 HudiRecordWriter 为主线,写入链路如下:

  1. 缓冲:每条SeaTunnelRow经 HudiRecordConverter 转换为HoodieRecord<HoodieAvroPayload>(SeaTunnel 行类型通过 AvroSchemaConverter 转为 Avro Schema),按HoodieKey存入LinkedHashMap缓冲,同一 key 的新记录覆盖旧记录;
  2. 刷写触发batchCount达到batch_size,或 checkpoint 到达(prepareCommit()调用 flush),或 Zeta 定时信号到达时,执行 flush;
  3. 提交:flush 时调用writeClient.startCommit()开启一次 commit,再根据op_type调用insert()/upsert()/bulkInsert(),删除记录则调用delete()
  4. 关闭close()时先做最后一次 flush 再关闭写客户端。

从源码结构看,Hudi Sink不提供 2PC 精确一次写入器,因此整体交付语义为at-least-once(至少一次)。这意味着:

  • 作业失败重试可能产生额外的 commit;
  • 使用INSERT时,恢复后自动生成的 record key 可能产生重复行;
  • 使用UPSERTrecord_key_fields稳定时,可将重复的逻辑记录限制在可控范围内。

写客户端的构建细节集中在 HudiUtil.createHoodieJavaWriteClient(),可以看到连接器为每次写入固定了如下 Hudi 配置:

  • 引擎类型EngineType.JAVA,使用HoodieJavaWriteClient
  • 关闭嵌入式 Timeline Server(withEmbeddedTimelineServerEnabled(false));
  • 开启自动清理、关闭异步清理(withAutoClean(true).withAsyncClean(false));
  • Parquet 压缩编码固定为SNAPPYHoodieStorageConfig.parquetCompressionCodec);
  • 归档阈值由min_commits_to_keep/max_commits_to_keep决定;
  • 平均记录大小估算由record_byte_size提供(approxRecordSize)。

另外,pom.xml 显示该模块基于 Hudi 0.15.0 的hudi-java-clienthudi-client-common,并将 Avro 包做了 shade 重定位以避免依赖冲突。在编写作业前请确保所用 SeaTunnel 发行版已包含connector-hudi插件。

Timer Flush 定时刷写(Zeta 专属)

Timer flush 是仅 Zeta 引擎支持的引擎级特性。在作业env块中配置sink.flush.interval,即可在batch_size尚未达到时也把待写入的 Hudi 记录刷写出去:

env { sink.flush.interval = 5000 }

Spark 与 Flink 引擎不会注入FlushSignal记录,因此不会触发这种定时刷写。在 Zeta 上,FlushSignal会与数据记录、checkpoint barrier 一起按顺序在 Sink 任务线程上被处理,触发回调context.registerFlushAction(this::timerFlush)中注册的 flush 动作(见 HudiSinkWriter)。

Hudi 定时刷写复用了连接器自身的同步批量刷写和 Hudi 客户端的 auto-commit 行为。由于 Hudi Sink 不提供 2PC 精确一次写入器,定时刷写提供的是at-least-once交付:重试可能产生额外 commit,INSERT下恢复后可能产生重复行,而UPSERT+ 稳定的record_key_fields可以限制重复的逻辑记录。

实战示例

单表 UPSERT

op_typeUPSERT时,record_key_fields必须配置:

sink { Hudi { table_dfs_path = "/tmp/seatunnel_mnt/hudi" database = "st" table_name = "st_test" table_type = "COPY_ON_WRITE" op_type = "UPSERT" record_key_fields = "c_bigint" batch_size = 1000 batch_interval_ms = 1000 } }

最简单表(追加写)

对于仅追加的写入,只需table_dfs_pathtable_name两个必填项:

sink { Hudi { table_dfs_path = "/tmp/seatunnel_mnt/hudi" table_name = "st_test" } }

多表写入

当上游数据源产出多张表时,使用table_list分别配置每张表。以下示例从 MySQL CDC 读取三张表并写入三张 Hudi 表(包含 INSERT、UPSERT、MERGE_ON_READ 三种形态):

env { parallelism = 1 job.mode = "STREAMING" checkpoint.interval = 5000 } source { Mysql-CDC { url = "jdbc:mysql://127.0.0.1:3306/seatunnel" username = "root" password = "******" table-names = ["seatunnel.role","seatunnel.user","galileo.Bucket"] } } transform { } sink { Hudi { table_dfs_path = "hdfs://nameserivce/data/" conf_files_path = "/home/test/hdfs-site.xml;/home/test/core-site.xml;/home/test/yarn-site.xml" table_list = [ { database = "st1" table_name = "role" table_type = "COPY_ON_WRITE" op_type = "INSERT" batch_size = 10000 }, { database = "st1" table_name = "user" table_type = "COPY_ON_WRITE" op_type = "UPSERT" record_key_fields = "user_id" batch_size = 10000 }, { database = "st1" table_name = "Bucket" table_type = "MERGE_ON_READ" } ] } }

注意:多表作业中table_dfs_pathconf_files_path位于 Sink 层,各表的database+table_name组合决定了其实际存储路径。连接器通过 HudiClientManager 按tableName + 队列索引缓存和复用HoodieJavaWriteClient实例,避免每张表重复创建客户端。

CDC 到 Hudi

当目标 Hudi 表需要持久化 CDC 变更日志信息时,开启cdc_enabled

sink { Hudi { table_dfs_path = "/tmp/seatunnel_mnt/hudi" database = "st" table_name = "st_test" table_type = "COPY_ON_WRITE" op_type = "UPSERT" record_key_fields = "id" cdc_enabled = true } }

此时上游 CDC 记录中的 INSERT/UPDATE_AFTER 会被写入,DELETE/UPDATE_BEFORE 会被转换为 Hudi 删除操作。

S3 对象存储

Hudi Sink 可以写入 S3 兼容路径。需要注意:connector-hudi模块不依赖hadoop-aws/aws-java-sdk,因此要解析s3a://scheme,需要先把hadoop-aws及配套的 AWS SDK bundle(或 SeaTunnel 的seatunnel-hadoop-awsjar)放入$SEATUNNEL_HOME/lib(或连接器的插件 lib 目录),然后通过conf_files_path(或运行时 classpath)提供所需的 Hadoop 文件系统配置,最后使用s3a://表路径:

sink { Hudi { table_dfs_path = "s3a://hudi/" conf_files_path = "/etc/hadoop/core-site.xml;/etc/hadoop/hdfs-site.xml" table_name = "st_test" op_type = "UPSERT" record_key_fields = "id" } }

使用注意事项汇总

  • Hive Metastore:连接器不负责注册/同步 Hive Metastore,hoodie.datasource.hive_sync.*选项不会生效,需要单独运行HiveSyncTool
  • 精确一次:连接器不具备 exactly-once 能力,交付语义为 at-least-once;关键业务去重请依赖UPSERT+ 稳定的record_key_fields
  • batch_interval_ms已弃用:定时刷写请改用 Zeta 的env.sink.flush.interval
  • BULK_INSERT 校验缺口record_key_fieldsBULK_INSERT不做启动期校验,省略会在写入期抛 NPE,建议始终配置;
  • 多表唯一性table_name在多表作业中不可重复或为空;
  • S3 依赖:使用s3a://路径前需手动补充 AWS 相关依赖与 Hadoop 配置。

版本变更记录

Hudi Sink 连接器的完整变更历史见 connector-hudi 变更日志。其中与本主题直接相关的关键变更包括:

  • 2.3.12:为 Hudi Sink 新增precombine_field选项;
  • 2.3.9:支持将 CDC 变更日志事件写入 Hudi Sink;
  • 2.3.8:优化 Hudi Sink;
  • 2.3.7:新增多表 Sink 选项校验;
  • 2.3.6:新增 Hudi Sink 连接器。

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

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/16 16:05:29

老 Mac 升级后蓝牙灰了?OpenCore Legacy Patcher 点几下就能拉回来

老 Mac 升级后蓝牙灰了&#xff1f;OpenCore Legacy Patcher 点几下就能拉回来 【免费下载链接】OpenCore-Legacy-Patcher Experience macOS just like before 项目地址: https://gitcode.com/GitHub_Trending/op/OpenCore-Legacy-Patcher 2012 的 MacBook Pro 刚升到 M…

作者头像 李华
网站建设 2026/9/16 16:04:29

Vue3+Node.js+MySQL校园资产管理系统实战

简介&#xff1a;这是一套基于VueNode.jsMySQL技术栈开发的校园资产管理系统全栈源码&#xff0c;面向Web全栈初学者与高校课程设计者&#xff0c;解决校园场景下资产登记、借用归还、分类统计等核心管理需求。资源共61个文件&#xff0c;含25个JavaScript后端逻辑与API封装文件…

作者头像 李华
网站建设 2026/9/16 16:03:54

NSGA-II与NSGA-III算法解析:Python实现多目标优化与选择指南

简介&#xff1a;一套面向多目标优化学习与研究者的NSGA3与NSGA-II算法Python/Matlab实现代码包。针对工程设计、调度、投资组合等具有相互冲突目标的优化问题&#xff0c;提供从帕累托前沿构建、快速非支配排序到拥挤距离与分层选择的完整实现框架。压缩包共10个文件&#xff…

作者头像 李华
网站建设 2026/9/16 16:03:16

三菱PLC在药片装瓶机控制系统中的应用与优化

1. 自动药片装瓶机控制系统架构解析在制药行业的自动化生产线上&#xff0c;药片装瓶机堪称"精密舞者"&#xff0c;其控制系统需要协调机械臂、传送带、计数机构等十余个执行单元。这套由三菱FX3U PLC作为主控的系统&#xff0c;通过组态王上位机实现人机交互&#x…

作者头像 李华