- 数据湖
- 湖仓一体
- 大数据
- 数据存储
【免费下载链接】hudi
Upserts, Deletes And Incremental Processing on Big Data.
Hudi 的CustomKeyGenerator允许在一张表上按字段粒度混用不同的分区键生成策略:既可以用字段原值直接作为分区(SIMPLE),也可以把时间戳字段格式化为yyyy-MM-dd之类的日期目录(TIMESTAMP),从而满足"国家 + 日期"这类复合分区场景。本文以仓库中hudi-trino连接器的测试数据集hudi_custom_keygen_pt_v8_mor为骨架,完整还原其建表脚本与配置,并结合CustomKeyGenerator源码和 Trino 测试初始化器,讲解如何创建、写入并验证一张带自定义键生成器的 MOR 表。
一、这张测试表要解决什么问题
在 hudi-trino/src/test/resources/hudi-testing-data/hudi_custom_keygen_pt_v8_mor.md 中,该测试数据的创建脚本说明了表的形态:
- 表类型:MOR(读时合并)表,使用CustomKeyGenerator生成记录键与分区路径;
- 数据版本:Hudi Table Revision
f7157954a9819137446d8e3a1d331d003f069414(仓库里同时存在对应的 hudi_custom_keygen_pt_v8_mor.zip 归档,供 Trino 集成测试直接解压使用); - 文件形态:无 log 文件(未触发增量日志写入/compaction),测试时通过关闭 inline compaction 保证表保持纯 base file 状态。
它属于 hudi-trino/src/test/resources/hudi-testing-data 目录下"带分区路径的 MOR 表"系列测试数据之一,同系列还包括hudi_timestamp_keygen_pt_epoch_to_yyyy_mm_dd_hh_v8_mor、hudi_multi_pt_v8_mor等,专门用于验证 Trino Hudi 连接器对各类 keygen/分区形态的读取兼容性。
二、建表脚本:混合 SIMPLE 与 TIMESTAMP 分区类型
以下是该文档提供的完整建表脚本(保留原始注释),它将两个分区字段分别声明为不同分区键类型:
test("Create MOR table with custom keygen partition field") { withTempDir { tmp => val tableName = "hudi_custom_keygen_pt_v8_mor" spark.sql( s""" |CREATE TABLE $tableName ( | id INT, | name STRING, | price DOUBLE, | ts LONG, | -- Partition Source Fields -- | partition_field_country STRING, | partition_field_date BIGINT |) USING hudi | LOCATION '${tmp.getCanonicalPath}' | TBLPROPERTIES ( | primaryKey = 'id', | type = 'mor', | preCombineField = 'ts', | -- Timestamp Keygen and Partition Configs -- | hoodie.table.keygenerator.class = 'org.apache.hudi.keygen.CustomKeyGenerator', | hoodie.datasource.write.partitionpath.field = 'partition_field_country:SIMPLE,partition_field_date:TIMESTAMP', | hoodie.keygen.timebased.timestamp.type = 'EPOCHMILLISECONDS', | hoodie.keygen.timebased.output.dateformat = 'yyyy-MM-dd', | hoodie.keygen.timebased.timezone = 'UTC' | ) PARTITIONED BY (partition_field_country, partition_field_date) """.stripMargin) // To not trigger compaction scheduling, and compaction spark.sql(s"set hoodie.compact.inline.max.delta.commits=9999") spark.sql(s"set hoodie.compact.inline=false") // Configure Hudi properties spark.sql(s"SET hoodie.metadata.enable=true") spark.sql(s"SET hoodie.metadata.index.column.stats.enable=true") // Insert data with new partition values spark.sql(s"INSERT INTO $tableName VALUES(1, 'a1', 100.0, 1000, 'SG', 1749284360000)") spark.sql(s"INSERT INTO $tableName VALUES(2, 'a2', 200.0, 1000, 'SG', 1749204000000)") spark.sql(s"INSERT INTO $tableName VALUES(3, 'a3', 101.0, 1001, 'US', 1749202000000)") spark.sql(s"INSERT INTO $tableName VALUES(4, 'a4', 201.0, 1001, 'CN', 1749102000000)") spark.sql(s"INSERT INTO $tableName VALUES(5, 'a5', 300.0, 1002, 'MY', 1747102000000)") spark.sql(s"INSERT INTO $tableName VALUES(6, 'a6', 301.0, 1000, 'SG', 1749284360000)") spark.sql(s"INSERT INTO $tableName VALUES(7, 'a7', 401.0, 1000, 'SG', 1749204000000)") // Generate logs through updates // NOTE: The query below will throw an error // spark.sql(s"UPDATE $tableName SET price = ROUND(price * 1.02, 2)") // NOTE: The query below will throw an error // spark.sql(s"SELECT * FROM $tableName").show(false) } }脚本同时给出了三条关键运行期配置:
| 配置项 | 值 | 作用 |
|---|---|---|
hoodie.compact.inline.max.delta.commits | 9999 | 抬高触发 inline compaction 的 delta commit 阈值,防止测试过程中自动触发合并,保持"无 log 文件"的期望形态 |
hoodie.compact.inline=false | false | 显式关闭 inline compaction |
hoodie.metadata.enable | true | 开启 Hudi 元数据表 |
hoodie.metadata.index.column.stats.enable | true | 在元数据表中构建列统计索引,供查询裁剪(column stats index)使用 |
关于末尾两行注释需要说明:UPDATE与SELECT *在脚本中被注释并标注"会抛错",这是测试数据生成脚本在特定 Spark 会话/环境下的行为约束,读者在自己的环境里复现时不一定复现同样的错误,但生成数据文件本身并不依赖这两条语句。
三、CustomKeyGenerator 的分区键解析原理
hoodie.table.keygenerator.class = 'org.apache.hudi.keygen.CustomKeyGenerator'指向 hudi-client/hudi-spark-client/src/main/java/org/apache/hudi/keygen/CustomKeyGenerator.java。该类是一个"通用键生成器",核心设计是:
- 分区字段配置格式:
hoodie.datasource.write.partitionpath.field按字段名:分区键类型,字段名:分区键类型的逗号分隔格式声明,每个字段可以挂不同的类型; - 逐字段分配生成器:构造时对每个分区字段调用
getPartitionKeyGenerators,解析出PartitionKeyType后分发到对应的内置生成器——SIMPLE类型交给SimpleKeyGenerator,TIMESTAMP类型交给TimestampBasedKeyGenerator(见 CustomKeyGenerator.java); - 记录键(record key):根据
primaryKey字段个数决定——单字段用SimpleKeyGenerator,多字段组合用ComplexKeyGenerator; - 分区路径拼接:多个分区字段生成的值按顺序用分隔符
/拼接成完整分区路径。
在 Avro 数据写入路径上,有对应的 hudi-client/hudi-client-common/src/main/java/org/apache/hudi/keygen/CustomAvroKeyGenerator.java 实现同样的逻辑,其PartitionKeyType枚举仅支持SIMPLE与TIMESTAMP两种(见 CustomAvroKeyGenerator.java),并提供getPartitionFieldAndKeyType做field:Type的解析与合法性校验。所以本表中partition_field_country:SIMPLE,partition_field_date:TIMESTAMP的含义是:
partition_field_country按SIMPLE处理:分区目录直接用字段原值(如SG、US、CN、MY);partition_field_date按TIMESTAMP处理:把 BIGINT 类型的时间戳(毫秒)按配置格式化为日期字符串(如2025-06-06)作为分区目录。
两个字段组合后产生的分区路径形如SG/2025-06-06、US/2025-06-06。
四、TIMESTAMP 类型分区的三项关键配置
时间戳分区行为由hoodie.keygen.timebased.*系列配置控制,本例中的三项配置直接决定了分区目录长什么样:
| 配置项 | 本例值 | 说明 |
|---|---|---|
hoodie.keygen.timebased.timestamp.type | EPOCHMILLISECONDS | 输入时间戳的格式类型。EPOCHMILLISECONDS表示输入值为 Unix 毫秒时间戳,另支持EPOCHSECONDS、SCALAR以及TIMESTAMP_MILLIS/TIMESTAMP_MICROS等类型 |
hoodie.keygen.timebased.output.dateformat | yyyy-MM-dd | 输出分区目录的日期格式。本例输出到"天"粒度,若要细分到小时可写成yyyy-MM-dd HH(同目录下的hudi_timestamp_keygen_pt_epoch_to_yyyy_mm_dd_hh_v8_mor即采用小时粒度) |
hoodie.keygen.timebased.timezone | UTC | 时间戳换算时使用的时区。时区不同会导致同一毫秒值落在不同日历日期,生产环境务必与业务语义对齐 |
这些配置在KeyGeneratorOptions中均有对应的配置键定义,例如hoodie.datasource.write.recordkey.field、hoodie.datasource.write.partitionpath.field定义在 hudi-common/src/main/java/org/apache/hudi/keygen/constant/KeyGeneratorOptions.java,其中PARTITIONPATH_FIELD_NAME明确指出"实际取值通过对字段值调用 toString() 获得"。
时区换算验证:7 条数据如何落成 5 个分区
脚本中 7 条 INSERT 的毫秒时间戳在UTC、yyyy-MM-dd格式下换算结果如下:
| 记录 | 国家字段 | 毫秒时间戳 | 换算出的日期分区 | 最终分区路径 |
|---|---|---|---|---|
| id=1 | SG | 1749284360000 | 2025-06-07 | SG/2025-06-07 |
| id=2 | SG | 1749204000000 | 2025-06-06 | SG/2025-06-06 |
| id=3 | US | 1749202000000 | 2025-06-06 | US/2025-06-06 |
| id=4 | CN | 1749102000000 | 2025-06-05 | CN/2025-06-05 |
| id=5 | MY | 1747102000000 | 2025-05-13 | MY/2025-05-13 |
| id=6 | SG | 1749284360000 | 2025-06-07 | SG/2025-06-07 |
| id=7 | SG | 1749204000000 | 2025-06-06 | SG/2025-06-06 |
7 条记录最终收敛为 5 个唯一分区:US/2025-06-06、CN/2025-06-05、MY/2025-05-13、SG/2025-06-06、SG/2025-06-07。这与 Trino 测试初始化器中注册的分区集合完全一致(见下文第五节),可据此验证建表配置的正确性。
五、测试数据在 Trino 连接器中的注册与使用
该 zip 归档在 Trino 集成测试中由 hudi-trino/src/test/java/io/trino/plugin/hudi/testing/ResourceHudiTablesInitializer.java 消费。其工作流程为:
initializeTables把hudi-testing-data资源目录下所有 zip 解压到临时目录,再通过copyDir拷贝到测试文件系统,并对每个文件做 SHA-256 校验(拷贝后重读比对哈希,防止损坏);TestingTable枚举中的HUDI_CUSTOM_KEYGEN_PT_V8_MOR条目注册了该表的数据列、分区列与分区路径(见 ResourceHudiTablesInitializer.java):- 数据列:
id INT、name STRING、price DOUBLE、ts LONG(与建表脚本一致,另叠加_hoodie_*元数据列); - 分区列:
partition_field_country STRING、partition_field_date STRING——注意存储为 Hive 分区值时日期已被格式化为字符串,即 keygen 的输出值同步进了 metastore; - 分区映射:
partition_field_country=US/partition_field_date=2025-06-06等 5 个分区,恰好对应建表脚本插入数据后生成的 5 个分区路径;
- 数据列:
- 该条目
isCreateRtTable=false,即只为这张 MOR 表创建默认的 ro(Read Optimized)外部表,不额外创建_rt实时表。
另外注意copyDir中有一条值得留意的过滤逻辑:跳过所有.crc文件,注释明确说明"Hudi 遇到 crc 文件会出问题",这也是在本地/云存储上部署 Hudi 表时容易踩的坑之一。
六、复现与验证要点总结
要把这张表从文档脚本复现成可查询的表,建议按以下顺序操作:
- 建表:在 Spark 3.x + Hudi 写端环境中按上文
CREATE TABLE脚本执行,LOCATION指向你的存储路径;primaryKey、preCombineField、type是必填的 Spark SQL Hudi 表属性; - 关闭 inline compaction:
hoodie.compact.inline=false+hoodie.compact.inline.max.delta.commits=9999,保证测试期间不产生 log 文件合并,维持"MOR 但无 log"的状态,便于后续以纯 base file 形态验证查询; - 写入:执行 7 条
INSERT INTO ... VALUES(...),字段顺序为(id, name, price, ts, partition_field_country, partition_field_date); - 验证分区:到表路径下检查目录结构,应出现
partition_field_country=SG/partition_field_date=2025-06-06之类的 5 个分区(若开启 hive_style_partitioning 则为key=value形式,否则为裸值路径),对照第五节的分区映射确认 keygen 行为; - 接入 Trino:将该表目录作为外部表挂载后,用 ResourceHudiTablesInitializer.java 同款方式(解压 zip → 建外部表 → 注册分区)即可让 Trino 通过 Hudi 连接器读取。
七、结语
hudi_custom_keygen_pt_v8_mor测试数据是理解 HudiCustomKeyGenerator的一手教材:它同时演示了 SIMPLE 与 TIMESTAMP 两种分区键类型如何在同一条hoodie.datasource.write.partitionpath.field中混用、时间戳分区配置如何影响目录粒度、以及这种表形态在 Trino 连接器中如何被注册和校验。结合 CustomKeyGenerator.java 与 CustomAvroKeyGenerator.java 的源码,可以完整串联"建表配置 → keygen 解析 → 分区目录生成 → 查询端读取"的全链路。如果你的业务同样需要"按国家分区 + 按天分区"的复合布局,直接复用本文的配置组合即可快速落地。
- 数据湖
- 湖仓一体
- 大数据
- 数据存储
【免费下载链接】hudi
Upserts, Deletes And Incremental Processing on Big Data.
相关推荐
Hudi 多分区字段 MOR 表的 Spark SQL 建表与 Trino 分区裁剪实战(基于 hudi-trino 测试数据集)
Hudi 多分区字段 MOR 表的 Spark SQL 建表与 Trino 分区裁剪实战(基于 hudi trino 测试数据集) 导读 本文以 Apache
数据湖湖仓一体大数据数据存储构建多文件组分区 MOR 测试数据集:Apache Hudi Trino 连接器索引跳过与分区剪枝验证实战
构建多文件组分区 MOR 测试数据集:Apache Hudi Trino 连接器索引跳过与分区剪枝验证实战 这篇技术指南围绕 Apache Hudi 仓库中面向
数据湖湖仓一体大数据数据存储Apache Hudi 非分区 MOR 测试表生成指南:基于 hudi_non_part_mor 的 Trino 连接器测试实践
Apache Hudi 非分区 MOR 测试表生成指南:基于 hudi_non_part_mor 的 Trino 连接器测试实践 本文以 Hudi 仓库中 hu
数据湖湖仓一体大数据数据存储
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考