- 数据湖
- 湖仓一体
- 大数据
- 数据存储
【免费下载链接】hudi
Upserts, Deletes And Incremental Processing on Big Data.
导读:本文围绕 Hudi 仓库中 Trino 连接器测试数据集
hudi_multi_fg_pt_v8_mor的构建文档展开,完整讲解如何用 Spark SQL 生成一张"多分区 × 多 FileGroup × 带日志文件 × 全量元数据索引"的 MOR 表,并剖析其中每个配置项对底层文件布局与索引的影响。读完本文,你将掌握一张可同时覆盖 Trino/Hudi 集成测试中 FileGroup、Log 文件、列统计索引、分区统计索引、记录级索引与二级索引六大场景的测试表的标准构建方法,并能独立复刻出同类测试数据。
数据集文档的定位:为 Trino 连接器测试而生的"表规格说明书"
在 Hudi 仓库中,hudi-trino/src/test/resources/hudi-testing-data/目录下的 Markdown 文件并不是普通的使用手册,而是每一份测试数据的"构建脚本 + 结构规格"说明。它们与同名.zip包一一对应:.zip中是已经生成好的真实 Hudi 表文件,.md则记录了生成这些数据所用的 Spark SQL 脚本和表结构约定。
hudi_multi_fg_pt_v8_mor.md就是其中的典型代表,它描述了一张名为hudi_multi_fg_pt_v8_mor的测试表:
- 这是一张MOR(Merge On Read)表,且MDT(Metadata Table,元数据表)已启用;
- 生成该数据集的 Hudi 版本修订号为
eb212c9dca876824b6c570665951777a772bc463; - Record Level Index(记录级索引)已启用,recordKey 为
id,name复合主键; - 在
price列上创建了二级索引; - 共2 个分区
[US, SG],每个分区2 个 FileGroup。
从文件命名可以推断,这是对应测试数据的v8 版本(_v8_),且在同一目录下存在同名的 v6 版本文档(hudi_multi_fg_pt_v6_mor.md,其 Revision 标注为release-0.15.0)。对照两份文档可以看到,v8 版本在 v6 的基础上新增了二级索引(secondary index)的创建步骤,这正对应了 Hudi 在 MDT 中对二级索引能力的持续演进。
表的最终结构一览
根据文档开头列出的结构清单,这张表在写入完成后呈现如下形态:
| 维度 | 规格 |
|---|---|
| 表类型 | MOR(Merge On Read) |
| MDT | 启用(Metadata Table) |
| 主键(recordKey) | id, name(复合主键) |
| 记录级索引 | 启用 |
| 二级索引 | 在price列上创建(索引名idx_price) |
| 分区 | 2 个:US、SG(分区列为country) |
| 每个分区的 FileGroup | 2 个 |
| Log 文件 | 存在(由 UPDATE 产生) |
数据列共 5 个:id int、name string、price double、ts long、country string,其中country同时充当分区列。
在 Trino 测试侧,这张表被注册在ResourceHudiTablesInitializer.java的TestingTable枚举中,条目为:
HUDI_MULTI_FG_PT_V8_MOR(hudiMultiFgRegularColumns(), hudiMultiFgPartitionsColumn(), hudiMultiFgPartitions(), false)其列定义(hudiMultiFgRegularColumns())与分区定义(hudiMultiFgPartitions())在源码中均有明确对应:
private static List<Column> hudiMultiFgRegularColumns() { return ImmutableList.of( column("id", HIVE_INT), column("name", HIVE_STRING), column("price", HIVE_DOUBLE), column("ts", HIVE_LONG)); } private static Map<String, String> hudiMultiFgPartitions() { return ImmutableMap.of( "country=SG", "country=SG", "country=US", "country=US"); }从该枚举的isCreateRtTable = false可知,这张表在 Trino 测试中只注册 RO 表(即只读表,直接读取 base 文件加日志合并后的结果),不额外创建带_rt后缀的实时表视图。这些信息印证了.md文档与测试代码之间严格的一致性——文档即规格,代码即实现。
完整构建脚本详解
文档给出了从零构建这张表的完整 Scala/Spark SQL 脚本。下面保留脚本全文,并逐段说明其作用。
test("Create table multi filegroup partitioned mor") { withTempDir { tmp => val tableName = "hudi_multi_fg_pt_mor" spark.sql( s""" |create table $tableName ( | id int, | name string, | price double, | ts long, | country string |) using hudi | location '${tmp.getCanonicalPath}' | tblproperties ( | primaryKey ='id,name', | type = 'mor', | preCombineField = 'ts' | ) partitioned by (country) """.stripMargin) // directly write to new parquet file spark.sql(s"set hoodie.parquet.small.file.limit=0") spark.sql(s"set hoodie.metadata.compact.max.delta.commits=1") // partition stats index is enabled together with column stats index spark.sql(s"set hoodie.metadata.index.column.stats.enable=true") spark.sql(s"set hoodie.metadata.record.index.enable=true") spark.sql(s"set hoodie.metadata.index.secondary.enable=true") spark.sql(s"set hoodie.metadata.index.column.stats.column.list=_hoodie_commit_time,_hoodie_partition_path,_hoodie_record_key,id,name,price,ts,country") // 2 filegroups per partition spark.sql(s"insert into $tableName values(1, 'a1', 100, 1000, 'SG'),(2, 'a2', 200, 1000, 'US')") spark.sql(s"insert into $tableName values(3, 'a3', 101, 1001, 'SG'),(4, 'a3', 201, 1001, 'US')") // create secondary index spark.sql(s"create index idx_price on $tableName (price)") // generate logs through updates spark.sql(s"update $tableName set price=price+1") } }第一步:建表 DDL
建表语句通过using hudi指定数据源为 Hudi,tblproperties中声明了三个核心表属性:
primaryKey = 'id,name':复合主键。两个字段共同构成 recordKey,写入时用于去重与更新定位,也是记录级索引的键基础;type = 'mor':表类型为 Merge On Read。基础文件按列式 Parquet 存储,后续更新先追加到基于行的 Log 文件,读取时再做合并;preCombineField = 'ts':预合并字段。当同一 recordKey 出现多条记录时,以ts值较大者为准,保证最终一致性。
分区列为country,因此表目录下会形成country=US/、country=SG/两个分区目录。
第二步:关键写入配置逐个解析
脚本中连续设置了 6 个 Spark 会话级 Hudi 配置,每个配置都对最终的文件布局与元数据索引产生直接影响:
hoodie.parquet.small.file.limit=0将"小文件合并阈值"设为 0,即关闭小文件合并行为。每次 commit 都会直接写入新的 Parquet 文件而不是填充到已有文件,这是保证"每个分区出现 2 个 FileGroup"的关键——两次 INSERT 各自生成独立的 base 文件,形成 2 个 FileGroup。
hoodie.metadata.compact.max.delta.commits=1控制 MDT 的 compaction 触发频率:每当 MDT 侧累积的 delta commit 达到 1 个即触发元数据表压缩。将阈值调小可以让元数据表在测试数据量很小的情况下也能形成压缩后的形态,确保测试能覆盖 MDT 压缩路径。
hoodie.metadata.index.column.stats.enable=true启用列统计索引(Column Stats Index)。该索引在 MDT 的column_stats分区中按列记录每个 FileSlice 的 min/max 等统计信息,供查询引擎做分区/文件裁剪。文档注释特别说明:分区统计索引(partition stats index)会随列统计索引一起启用,这是二者在实现上的联动关系。
hoodie.metadata.record.index.enable=true启用记录级索引(Record Level Index,RLI)。该索引提供 recordKey → FileGroup/FileSlice 的全局定位能力,是 Hudi 在 MOR 表上高效执行点查与更新定位的关键设施。
hoodie.metadata.index.secondary.enable=true启用二级索引(Secondary Index)。v8 版本相比 v6 版本的新增配置,配合后续的create index idx_price on $tableName (price)语句,在price列上建立二级索引,使得以price为过滤条件的查询可以直接借由 MDT 中二级索引分区快速定位目标记录。
hoodie.metadata.index.column.stats.column.list=...显式指定列统计索引要覆盖的列清单:_hoodie_commit_time, _hoodie_partition_path, _hoodie_record_key, id, name, price, ts, country。清单同时包含 5 个 Hudi 元数据列与全部 5 个业务列,意味着这张表的列统计索引是"全列覆盖"的,任何列的裁剪过滤都能命中统计信息。
这些配置项在源码层面均有对应定义。在 HoodieMetadataConfig.java 中可以看到SECONDARY_INDEX_ENABLE_PROP、SECONDARY_INDEX_PARALLELISM、SECONDARY_INDEX_NAME、SECONDARY_INDEX_COLUMN等二级索引相关属性的定义;而 TestHoodieMetadataConfig.java 中的测试用例也直接以hoodie.metadata.record.index.enable=true作为输入属性验证配置解析逻辑,证明这些配置键在 Hudi 通用配置层是稳定、可解析的。
第三步:两次 INSERT 制造多 FileGroup
insert into hudi_multi_fg_pt_mor values(1, 'a1', 100, 1000, 'SG'),(2, 'a2', 200, 1000, 'US') insert into hudi_multi_fg_pt_mor values(3, 'a3', 101, 1001, 'SG'),(4, 'a3', 201, 1001, 'US')两条 INSERT 各写一次 commit。由于hoodie.parquet.small.file.limit=0关闭了小文件合并,每次 commit 都会为每个分区生成新的 base 文件:
US分区:第一条插入记录(2, 'a2'),第二条插入(4, 'a3'),形成 2 个 FileGroup;SG分区:第一条插入(1, 'a1'),第二条插入(3, 'a3'),同样形成 2 个 FileGroup。
注意第二条插入中出现了两条name='a3'的记录(id=3在SG、id=4在US),由于 recordKey 是id,name复合键,这两条记录的 key 并不相同,因此它们被正常写入各自的文件组,不会被互相覆盖。
第四步:创建二级索引
create index idx_price on hudi_multi_fg_pt_mor (price)在price列上创建名为idx_price的二级索引。Hudi 会在 MDT 中为每个二级索引创建一个独立的索引分区(源码中对应PARTITION_NAME_SECONDARY_INDEX_PREFIX前缀的元数据分区),索引数据随后续写入增量维护。
第五步:UPDATE 生成 Log 文件
update hudi_multi_fg_pt_mor set price=price+1这是一条全表更新。对 MOR 表而言,UPDATE 不会直接改写已有的 Parquet base 文件,而是把更新后的记录追加写入对应的Log 文件(.log),这正是文档"生成 logs"注释的含义。更新完成后,每个 FileGroup 下形成"base 文件 + log 文件"的文件切片结构,Trino 在查询该表时需要完成 base 与 log 的实时合并(MOR 的 Read Optimized 语义之外的核心读取路径)。
这张表在 Trino 测试中的消费方式
.md文档定义了数据规格,而 Trino 连接器通过测试代码真正消费这些数据。可以从几个侧面看到它们的衔接:
测试数据装配链路。ResourceHudiTablesInitializer.java的initializeTables方法会:
- 把
hudi-testing-data资源目录下的.zip解压到临时目录; - 将解压出的 Hudi 表完整拷贝到 Trino 文件系统(拷贝过程会计算 SHA-256 哈希校验完整性,且跳过
.crc校验文件); - 通过
HudiConnector注入的TrinoFileSystemFactory与HiveMetastoreFactory,为每张测试表在 metastore 中注册外部表及分区(RO 存储格式使用HUDI_PARQUET_INPUT_FORMAT,RT 格式使用HUDI_PARQUET_REALTIME_INPUT_FORMAT); - 读取
HoodieTableMetaClient获取表版本并回填到TestingTable枚举。
索引能力在连接器侧的开关。HudiConfig.java 中默认开启isSecondaryIndexEnabled与isColumnStatsIndexEnabled,并定义了hudi.index.secondary-index-enabled、hudi.index.column-stats-index-enabled、hudi.index.column-stats.wait-timeout、hudi.index.record-index.wait-timeout、hudi.index.secondary-index.wait-timeout等内部配置(默认等待超时 2 秒),说明 Trino 侧在读取 Hudi 表时确实会加载并使用这些 MDT 索引。文档中这张表全量开启 MDT 索引,正是为了让这些读取路径在测试中被完整覆盖。
MDT 列统计数据的实际访问证据。在InlineSeekableDataInputStream.java的注释中保留了真实测试运行时的路径痕迹:inlinefs://.../hudi_multi_fg_pt_v8_mor/.hoodie/metadata/column_stats/,这直接证明 Trino 连接器在测试中确实访问了该表的 MDTcolumn_stats分区来读取列统计索引数据。
什么场景下应该使用这张表
文档在末尾明确列出了这张表的适用场景,这是选择测试数据时的直接决策依据:
- 需要分区内存在多个 FileGroup 的测试用例:两次 INSERT 配合关闭小文件合并,使
US、SG每个分区都有 2 个 FileGroup,可用于验证多文件组场景下的文件裁剪、并行读取与合并逻辑; - 需要 FileGroup 带 Log 文件的测试用例:全表 UPDATE 产生 log 文件,覆盖 MOR 表 base + log 的合并读取路径;
- 需要列统计索引的测试用例:全列(含 5 个元数据列 + 5 个业务列)启用了 column stats,可验证基于列统计的文件级/块级裁剪;
- 需要分区统计索引的测试用例:分区统计索引随列统计索引联动启用;
- 需要记录级索引的测试用例:
hoodie.metadata.record.index.enable=true开启 RLI,recordKey 为id,name; - 需要二级索引的测试用例:
price列上的idx_price二级索引可用于验证基于二级索引的查询加速路径。
如何复刻你自己的同构测试表
如果需要在本地复现或改造这张测试表,可以按以下步骤操作:
- 在 Spark 环境(需包含 Hudi Spark Bundle)中,将本文的 Scala 脚本中的
tableName、location替换为你自己的命名与路径; - 保持
tblproperties的primaryKey、type='mor'、preCombineField三个核心属性不变(MOR + 复合主键是这张表结构的根基); - 如需"多 FileGroup"形态,务必设置
hoodie.parquet.small.file.limit=0并分多条 INSERT 写入;每条 INSERT 产生一次 commit、每个分区新增一个 FileGroup; - 如需"带 Log 文件"形态,在 INSERT 之后追加 UPDATE 语句即可;MOR 表的更新天然落入 log 文件;
- 按需开启
hoodie.metadata.index.column.stats.enable、hoodie.metadata.record.index.enable、hoodie.metadata.index.secondary.enable三类索引,并通过hoodie.metadata.index.column.stats.column.list控制列统计的覆盖范围; - 生成完毕后,将表目录归档为 zip 放入测试资源目录,并在
TestingTable枚举中注册对应的列与分区定义,即可被 Trino 连接器测试框架加载。
需要说明的是:文档中给出的配置组合是为测试数据"定制"的(例如hoodie.metadata.compact.max.delta.commits=1会在每次增量后立即压缩元数据表),在生产环境中这些参数应按实际数据规模重新评估,不宜直接照搬。
小结
hudi_multi_fg_pt_v8_mor.md看似只是一份测试数据的构建脚本,但它实际上浓缩了 Hudi 一张"高端局"MOR 表的全部关键要素:复合主键、双分区、多 FileGroup、log 文件、MDT 下的列统计/分区统计/记录级/二级索引全量覆盖。理解这份文档,既能帮助你读懂 Trino 连接器测试数据集的构造逻辑,也能作为你在真实业务中组合配置 Hudi 元数据索引、控制文件布局的参考样板。
- 数据湖
- 湖仓一体
- 大数据
- 数据存储
【免费下载链接】hudi
Upserts, Deletes And Incremental Processing on Big Data.
相关推荐
Hudi 多分区字段 MOR 表的 Spark SQL 建表与 Trino 分区裁剪实战(基于 hudi-trino 测试数据集)
Hudi 多分区字段 MOR 表的 Spark SQL 建表与 Trino 分区裁剪实战(基于 hudi trino 测试数据集) 导读 本文以 Apache
数据湖湖仓一体大数据数据存储DataHub Redshift 元数据采集实战指南:从权限配置、Lineage 到 Usage 与 Profiling
DataHub Redshift 元数据采集实战指南:从权限配置、Lineage 到 Usage 与 Profiling 导读 本文以 DataHub 官方 R
数据湖湖仓一体大数据数据存储OceanBase表设计终极指南:分区与索引优化实战
OceanBase表设计终极指南:分区与索引优化实战 你是否还在为海量数据查询缓慢而困扰?作为企业级分布式关系型数据库,OceanBase凭借高可用性、高性能和
数据库分布式数据库关系型数据库后端高可用
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考