跳到主要内容
版本:Next

Hudi

Hudi 接收器连接器

描述​

用于将数据写入 Hudi。

主要特性​

Hive Metastore 同步

SeaTunnel Hudi sink 会写入 Hudi 数据文件和 .hoodie 元数据,但不会在 Hive Metastore 中注册表或将表同步到 Hive Metastore。hoodie.datasource.hive_sync.* 配置不是受支持的 sink 选项,也不会传递给 Hudi 写入客户端。需要注册到 Hive Metastore 时,请单独运行 Apache Hudi HiveSyncTool 或其他表注册流程。

选项​

基础配置:

名称类型是否必填默认值
table_dfs_pathstring是-
conf_files_pathstring否-
table_listarray否-
schema_save_modeenum否CREATE_SCHEMA_WHEN_NOT_EXIST
data_save_modeenum否APPEND_DATA
multi_table_sink_replicaint否1
common-optionsconfig否-

表清单配置:

名称类型是否必填默认值
table_namestring是-
databasestring否default
table_typeenum否COPY_ON_WRITE
op_typeenum否INSERT
record_key_fieldsstring否-
partition_fieldsstring否-
precombine_fieldstring否-
batch_interval_msint否1000
batch_sizeint否1000
insert_shuffle_parallelismint否2
upsert_shuffle_parallelismint否2
min_commits_to_keepint否20
max_commits_to_keepint否30
index_typeenum否BLOOM
index_class_namestring否-
record_byte_sizeint否1024
cdc_enabledboolean否false

注意:写入单表时,可以把 table_list 中的表配置项平铺到外层。多表作业中,表级配置需放在各自的 table_list 条目内;table_dfs_path、conf_files_path、schema_save_mode 和 data_save_mode 保持在 sink 层级。

record_key_fields 在 UPSERT 模式下必填(启动时校验),在 BULK_INSERT 模式下也必填(当前未校验——缺少该配置会在写入时抛出 NullPointerException,而不是配置期错误)。对于 CDC 输入,上游记录必须包含 record_key_fields 引用的字段;仅当需要 Hudi CDC 变更日志时才设置 cdc_enabled = true。

table_name [string]​

table_name Hudi 表的名称。

database [string]​

database Hudi 表所属的数据库。

table_dfs_path [string]​

table_dfs_path Hudi 表的 DFS 根路径,例如 "hdfs://nameservice/data/hudi/"。

table_type [enum]​

table_type Hudi 表的类型,可选值为 COPY_ON_WRITE 和 MERGE_ON_READ。

record_key_fields [string]​

record_key_fields Hudi 表的记录键字段。当 op_type 为 UPSERT 时,必须配置该项。

partition_fields [string]​

partition_fields Hudi 表的分区字段.

precombine_field [string]​

precombine_field Hudi 表的预合并字段,它用于在写入前进行预合并.

index_type [string]​

index_type Hudi 表的索引类型。当前支持 BLOOM、SIMPLE、GLOBAL_BLOOM。

index_class_name [string]​

index_class_name Hudi 表自定义索引名称,例如: org.apache.seatunnel.connectors.seatunnel.hudi.index.CustomHudiIndex.

record_byte_size [Int]​

record_byte_size Hudi 表单行记录的大小, 该值可用于预估每个hudi数据文件中记录的大致数量。调整此参数与batch_size可以有效减少hudi数据文件写放大次数.

conf_files_path [string]​

conf_files_path 环境配置文件路径列表(本地路径),用于初始化 HDFS 客户端以读取 Hudi 表文件。示例:"/home/test/hdfs-site.xml;/home/test/core-site.xml;/home/test/yarn-site.xml"。

op_type [enum]​

op_type Hudi 表的操作类型。值可以是 insert、upsert 或 bulk_insert。

batch_interval_ms [Int]​

batch_interval_ms 为兼容性保留。在 Zeta 上需要定时刷新时,请在作业 env 中配置 sink.flush.interval。

batch_size [Int]​

batch_size 单次刷新到 Hudi 前最多缓存的记录数。

insert_shuffle_parallelism [Int]​

insert_shuffle_parallelism 插入数据到 Hudi 表的并行度。

upsert_shuffle_parallelism [Int]​

upsert_shuffle_parallelism 更新插入数据到 Hudi 表的并行度。

min_commits_to_keep [Int]​

min_commits_to_keep Hudi 表保留的最少提交数。

max_commits_to_keep [Int]​

max_commits_to_keep Hudi 表保留的最多提交数。

cdc_enabled [boolean]​

cdc_enabled 是否持久化Hudi表的CDC变更日志。启用后,在必要时持久化更改数据,表可以作为CDC模式进行查询.

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:当已有数据存在时报错

通用选项​

Sink插件通用参数,请参考 Sink Common Options 了解详细信息。

定时刷新​

定时刷新是仅由 Zeta 支持的引擎级能力。在作业的 env 中配置 sink.flush.interval 后,即使尚未达到 batch_size,Hudi Sink 也会写出待处理的记录。Spark 和 Flink 不会注入 FlushSignal,因此不会触发这种 定时刷新。

env {
sink.flush.interval = 5000
}

Hudi 定时刷新复用连接器现有的同步批量刷新和 Hudi 客户端 auto-commit 行为。Hudi Sink 没有 2PC 精确一次 写入器,因此定时刷新提供的是至少一次语义,重试可能产生额外的 commit。使用 INSERT 时,自动生成的 record key 还可能在恢复后产生重复行;使用具有稳定 record_key_fields 的 UPSERT 可以减少逻辑记录重复。

示例​

单表 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"
}
}

多表​

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"
}
]
}
}

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
}
}

S3 存储​

sink 可以写入 S3 兼容路径。connector-hudi 模块不依赖 hadoop-aws/aws-java-sdk,因此要解析 s3a:// 协议,需要先将 hadoop-aws 和匹配的 AWS SDK 包(或 SeaTunnel 的 seatunnel-hadoop-aws jar)放入 $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"
}
}

变更日志​