SeaTunnel Hudi Sink 连接器详解:配置参数、多表写入与 CDC 实战指南 📅 发布时间:2026/9/16 16:07:41 👁 浏览次数: 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 如何写 HudiHudi Sink 连接器插件名Hudi用于将 SeaTunnel 作业中的数据写入 Hudi 表。它基于 Hudi 官方提供的HoodieJavaWriteClientJava 引擎写入客户端实现支持写入 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_modeenum否CREATE_SCHEMA_WHEN_NOT_EXIST作业启动前如何处理目标表结构data_save_modeenum否APPEND_DATA作业启动前如何处理已存在的表数据common-optionsConfig否-Sink 公共参数表列表配置table_list 内名称类型必填默认值说明table_namestring是-目标 Hudi 表名databasestring否default目标 Hudi 数据库名table_typeenum否COPY_ON_WRITEHudi 表类型COPY_ON_WRITE或MERGE_ON_READop_typeenum否INSERT写操作INSERT、UPSERT或BULK_INSERTrecord_key_fieldsstring否-用于构建 Hudi record key 的字段UPSERT必填partition_fieldsstring否-用于构建分区路径的字段precombine_fieldstring否-用于解决同一记录多次更新的字段batch_interval_msInt否1000当前未使用刷写仅由batch_size或 checkpoint 触发batch_sizeInt否1000刷写前缓冲的最大行数insert_shuffle_parallelismInt否2insert 操作的 shuffle 并行度upsert_shuffle_parallelismInt否2upsert 操作的 shuffle 并行度min_commits_to_keepInt否20清理时保留的最少 commit 数max_commits_to_keepInt否30清理时保留的最多 commit 数index_typeenum否BLOOMHudi 索引类型BLOOM、SIMPLE或GLOBAL_BLOOMindex_class_namestring否-自定义 Hudi 索引类的全限定类名record_byte_sizeInt否1024每条记录的平均字节数估算值cdc_enabledboolean否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 行的RowKindINSERT/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 optionsSink 插件的公共参数详见 Sink Common Options。底层写入流程与一致性语义理解连接器的工作方式有助于正确设置参数。以 HudiSinkWriter 与 HudiRecordWriter 为主线写入链路如下缓冲每条SeaTunnelRow经 HudiRecordConverter 转换为HoodieRecordHoodieAvroPayloadSeaTunnel 行类型通过 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 ServerwithEmbeddedTimelineServerEnabled(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-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交付重试可能产生额外 commitINSERT下恢复后可能产生重复行而UPSERT 稳定的record_key_fields可以限制重复的逻辑记录。实战示例单表 UPSERTop_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 层各表的databasetable_name组合决定了其实际存储路径。连接器通过 HudiClientManager 按tableName 队列索引缓存和复用HoodieJavaWriteClient实例避免每张表重复创建客户端。CDC 到 Hudi当目标 Hudi 表需要持久化 CDC 变更日志信息时开启cdc_enabledsink { 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 Metastorehoodie.datasource.hive_sync.*选项不会生效需要单独运行HiveSyncTool精确一次连接器不具备 exactly-once 能力交付语义为 at-least-once关键业务去重请依赖UPSERT 稳定的record_key_fieldsbatch_interval_ms已弃用定时刷写请改用 Zeta 的env.sink.flush.intervalBULK_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 Sink2.3.8优化 Hudi Sink2.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),仅供参考