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_path | string | 是 | - | Hudi 表数据与元数据的根路径 |
| conf_files_path | string | 否 | - | HDFS 客户端配置文件(分号分隔的本地路径列表) |
| table_list | Array | 否 | - | 多表作业中每张表的独立设置 |
| schema_save_mode | enum | 否 | CREATE_SCHEMA_WHEN_NOT_EXIST | 作业启动前如何处理目标表结构 |
| data_save_mode | enum | 否 | APPEND_DATA | 作业启动前如何处理已存在的表数据 |
| common-options | Config | 否 | - | Sink 公共参数 |
表列表配置(table_list 内)
| 名称 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|
| table_name | string | 是 | - | 目标 Hudi 表名 |
| database | string | 否 | default | 目标 Hudi 数据库名 |
| table_type | enum | 否 | COPY_ON_WRITE | Hudi 表类型:COPY_ON_WRITE或MERGE_ON_READ |
| op_type | enum | 否 | INSERT | 写操作:INSERT、UPSERT或BULK_INSERT |
| record_key_fields | string | 否 | - | 用于构建 Hudi record key 的字段,UPSERT必填 |
| partition_fields | string | 否 | - | 用于构建分区路径的字段 |
| precombine_field | string | 否 | - | 用于解决同一记录多次更新的字段 |
| batch_interval_ms | Int | 否 | 1000 | 当前未使用;刷写仅由batch_size或 checkpoint 触发 |
| batch_size | Int | 否 | 1000 | 刷写前缓冲的最大行数 |
| insert_shuffle_parallelism | Int | 否 | 2 | insert 操作的 shuffle 并行度 |
| upsert_shuffle_parallelism | Int | 否 | 2 | upsert 操作的 shuffle 并行度 |
| min_commits_to_keep | Int | 否 | 20 | 清理时保留的最少 commit 数 |
| max_commits_to_keep | Int | 否 | 30 | 清理时保留的最多 commit 数 |
| index_type | enum | 否 | BLOOM | Hudi 索引类型:BLOOM、SIMPLE或GLOBAL_BLOOM |
| index_class_name | string | 否 | - | 自定义 Hudi 索引类的全限定类名 |
| record_byte_size | Int | 否 | 1024 | 每条记录的平均字节数估算值 |
| cdc_enabled | boolean | 否 | false | 开启后持久化 Hudi CDC 变更日志数据 |
配置的组织方式:当作业只写一张表时,可以把table_list内的配置项平铺到外层(即直接写在Hudi {}块下);多表作业时,表级选项必须放在各自table_list条目内,而table_dfs_path、conf_files_path、schema_save_mode、data_save_mode仍保持在 Sink 层。这一设计在 HudiTableConfig.of() 中体现:当table_list未配置时,连接器会把外层平铺配置组装成一个单元素列表。
参数间约束:
record_key_fields在op_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_path、database、table_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_type为UPSERT时必须配置。底层实现中,该字段会参与构建HoodieKey与去重逻辑(见 HudiRecordWriter.prepareRecords()),同一HoodieKey的记录在缓冲区内会被后者覆盖。
partition_fields [string]
Hudi 表的分区键字段,用于生成分区路径。配置后 Hudi 会按该字段对数据进行分区存储。
precombine_field [string]
Hudi 表的 precombine 字段,在实际写入前用于 preCombining,即同一记录有多条更新时,据此字段决定最终保留哪一条。该选项由 PRECOMBINE_FIELD 定义,变更记录显示该参数在 2.3.12 版本中加入。
index_type [string]
Hudi 表的索引类型。当前支持BLOOM、SIMPLE和GLOBAL_BLOOM,默认BLOOM。索引用于在 upsert 时定位记录所属的文件组。
index_class_name [string]
自定义 Hudi 索引类的全限定类名,例如org.apache.seatunnel.connectors.seatunnel.hudi.index.CustomHudiIndex。从源码看,index_class_name与index_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用于构建HoodieJavaEngineContext和HoodieJavaWriteClient。
op_type [enum]
Hudi 表的操作类型,取值为insert、upsert或bulk_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.commits。min_commits_to_keep与max_commits_to_keep会一起传给HoodieArchivalConfig.archiveCommitsWith(),用于配置归档策略(见 HudiUtil)。
cdc_enabled [boolean]
是否持久化 CDC 变更日志,默认false。开启后,在必要时持久化变更数据,Hudi 表可以按 CDC 查询模式(CDC query mode)被查询。从源码看,连接器会依据 SeaTunnel 行的RowKind(INSERT/UPDATE_BEFORE/UPDATE_AFTER/DELETE)将变更标记为"写入"或"删除",DELETE与UPDATE_BEFORE走删除路径,INSERT与UPDATE_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 发现名为Hudi的CatalogFactory创建 Catalog,再交给DefaultSaveModeHandler按schema_save_mode/data_save_mode执行建表、删表、清数据等前置动作。
common options
Sink 插件的公共参数,详见 Sink Common Options。
底层写入流程与一致性语义
理解连接器的工作方式有助于正确设置参数。以 HudiSinkWriter 与 HudiRecordWriter 为主线,写入链路如下:
- 缓冲:每条
SeaTunnelRow经 HudiRecordConverter 转换为HoodieRecord<HoodieAvroPayload>(SeaTunnel 行类型通过 AvroSchemaConverter 转为 Avro Schema),按HoodieKey存入LinkedHashMap缓冲,同一 key 的新记录覆盖旧记录; - 刷写触发:
batchCount达到batch_size,或 checkpoint 到达(prepareCommit()调用 flush),或 Zeta 定时信号到达时,执行 flush; - 提交:flush 时调用
writeClient.startCommit()开启一次 commit,再根据op_type调用insert()/upsert()/bulkInsert(),删除记录则调用delete(); - 关闭:
close()时先做最后一次 flush 再关闭写客户端。
从源码结构看,Hudi Sink不提供 2PC 精确一次写入器,因此整体交付语义为at-least-once(至少一次)。这意味着:
- 作业失败重试可能产生额外的 commit;
- 使用
INSERT时,恢复后自动生成的 record key 可能产生重复行; - 使用
UPSERT且record_key_fields稳定时,可将重复的逻辑记录限制在可控范围内。
写客户端的构建细节集中在 HudiUtil.createHoodieJavaWriteClient(),可以看到连接器为每次写入固定了如下 Hudi 配置:
- 引擎类型
EngineType.JAVA,使用HoodieJavaWriteClient; - 关闭嵌入式 Timeline Server(
withEmbeddedTimelineServerEnabled(false)); - 开启自动清理、关闭异步清理(
withAutoClean(true).withAsyncClean(false)); - Parquet 压缩编码固定为
SNAPPY(HoodieStorageConfig.parquetCompressionCodec); - 归档阈值由
min_commits_to_keep/max_commits_to_keep决定; - 平均记录大小估算由
record_byte_size提供(approxRecordSize)。
另外,pom.xml 显示该模块基于 Hudi 0.15.0 的hudi-java-client与hudi-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_type为UPSERT时,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_path和table_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_path与conf_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_fields对BULK_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),仅供参考