ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Apache Hudi MOR 表全面类型写入与 Trino 读取验证实战指南

Apache Hudi MOR 表全面类型写入与 Trino 读取验证实战指南 数据湖湖仓一体大数据数据存储【免费下载链接】hudiUpserts, Deletes And Incremental Processing on Big Data.项目地址https://gitcode.com/gh_mirrors/hud/hudi点击查看免费下载导读本文围绕 hudi-trino/src/test/resources/hudi-testing-data/hudi_comprehensive_types_v8_mor.md 这份测试数据生成脚本完整剖析如何在 Spark SQL 中创建一个覆盖数值、字符串、二进制、日期时间以及多层嵌套复杂类型Array / Map / Struct 及其任意组合的 Hudi MOR 表并通过 UPDATE 生成日志文件后交由 Trino Hudi 连接器完成类型级读取验证。读完本文你将掌握Hudi 支持的全部 Spark 类型到 Parquet 的写入方式、MOR 表配合 Metadata TableMDT多索引的建表参数组合、不同版本V6 / V8测试表的差异以及 Trino 侧对每一列类型映射的验证手段。一、文档背景测试数据生成脚本的定位1.1 它是 Trino Hudi 连接器的测试数据说明书该 Markdown 文件本身是一段 Scala 测试代码的完整转写对应的可执行实现位于 hudi-trino/src/test/java/io/trino/plugin/hudi/TestHudiSmokeTest.java 中testComprehensiveTypes(ResourceHudiTablesInitializer.TestingTable table, boolean isRtTable)方法约 L1140-L1329。其作用是通过 Spark SQL 生成一张全类型的 MOR 表并将生成的表数据打包进同目录下的 hudi_comprehensive_types_v8_mor.zip供 Trino 连接器的单元/集成测试直接加载。在 ResourceHudiTablesInitializer.java 的TestingTable枚举L339-L363中该表被注册为HUDI_COMPREHENSIVE_TYPES_V8_MOR(hudiComprehensiveTypesColumns(), hudiComprehensiveTypesPartitionColumns(), hudiComprehensiveTypesPartitions(), true)其中最后一个布尔参数isCreateRtTable true表示测试框架同时会生成/注册 Read OptimizedRO与 Real-TimeRT两张视图从而分别验证只读基础文件和合并日志文件两条读取路径。1.2 V8 与 V6 的关系同一目录下还提供了 hudi_comprehensive_types_v6_mor.md 与对应 zip二者表结构与生成逻辑一致区别仅在于 Hudi 表版本Table Version。在测试参数工厂comprehensiveTestParameters()TestHudiSmokeTest.java L1468-L1480中V6 与 V8 两张表会与true/false两个 RT 参数组合出 4 组测试用例确保新旧表版本都能通过类型验证。文档中标注的Revision: 444cac26cb1077fd2b7deefc7b3713bacb270f9c用于锁定生成这批测试数据的 Hudi 源码提交保证数据可复现。二、建表脚本全类型 Schema 逐列解读2.1 完整建表语句以下为文档中的核心建表脚本略去 Scala 字符串转义它一次性覆盖了 Hudi/Spark 支持的几乎所有标量与嵌套类型CREATE TABLE hudi_type_test_mor ( uuid STRING, precombine_field LONG, -- Numeric Types col_boolean BOOLEAN, col_tinyint TINYINT, col_smallint SMALLINT, col_int INT, col_bigint BIGINT, col_float FLOAT, col_double DOUBLE, col_decimal DECIMAL(10, 2), -- String Types col_string STRING, col_varchar VARCHAR(50), col_char CHAR(10), -- Binary Type col_binary BINARY, -- Datetime Types col_date DATE, col_timestamp TIMESTAMP, -- col_timestamp_ntz TIMESTAMP_NTZ, (No support on Hudi for now) -- Complex types col_array_int ARRAYINT, col_array_string ARRAYSTRING, col_map_string_int MAPSTRING, INT, col_struct STRUCTf1: STRING, f2: INT, f3: BOOLEAN, col_array_struct ARRAYSTRUCTnested_f1: DOUBLE, nested_f2: ARRAYSTRING, col_map_string_struct MAPSTRING, STRUCTnested_f3: DATE, nested_f4: DECIMAL(5,2), col_array_struct_with_map ARRAYSTRUCTf_arr_struct_str: STRING, f_arr_struct_map: MAPSTRING, INT, col_map_struct_with_array MAPSTRING, STRUCTf_map_struct_arr: ARRAYBOOLEAN, f_map_struct_ts: TIMESTAMP, col_struct_nested_struct STRUCTouter_f1: INT, nested_struct: STRUCTinner_f1: STRING, inner_f2: BOOLEAN, col_array_array_int ARRAYARRAYINT, col_map_string_array_double MAPSTRING, ARRAYDOUBLE, col_map_string_map_string_date MAPSTRING, MAPSTRING, DATE, -- Array of structs with single (inner) fields do not work with parquet.version 1.13.1 col_struct_array_struct STRUCTouter_f2: STRING, struct_array: ARRAYSTRUCTinner_f3: TIMESTAMP, inner_f4: STRING, col_struct_map STRUCTouter_f3: BOOLEAN, struct_map: MAPSTRING, BIGINT, part_col STRING ) USING hudi LOCATION ${tmp.getCanonicalPath} TBLPROPERTIES ( primaryKey uuid, type mor, preCombineField precombine_field ) PARTITIONED BY (part_col)2.2 标量类型分组说明分组列名Spark 类型说明主键/预合并uuid/precombine_fieldSTRING / LONG分别对应primaryKey与preCombineField数值col_booleanBOOLEAN布尔数值col_tinyint/col_smallintTINYINT / SMALLINT8/16 位整数数值col_int/col_bigintINT / BIGINT32/64 位整数数值col_float/col_doubleFLOAT / DOUBLE单/双精度浮点数值col_decimalDECIMAL(10, 2)高精度定点数字符串col_string/col_varchar/col_charSTRING / VARCHAR(50) / CHAR(10)长度语义不同二进制col_binaryBINARY以 VARBINARY 落盘日期时间col_date/col_timestampDATE / TIMESTAMP日期与时间戳需要注意文档中明确标注的一行注释col_timestamp_ntz TIMESTAMP_NTZ目前 Hudi 尚不支持No support on Hudi for now因此在建表与插入语句中它都被注释掉。这说明在选择时间类型时当前 Hudi 版本应使用带时区语义的 TIMESTAMP而非 Spark 3.4 引入的 TIMESTAMP_NTZ。2.3 嵌套复杂类型任意组合覆盖复杂类型部分是整张表设计的重点用于检验 Parquet 嵌套列在 Hudi → Trino 链路上的完整映射单层组合col_array_int数组、col_map_string_intMap、col_structStructStruct 内嵌col_struct_nested_structStruct 里再套 Struct、col_struct_mapStruct 里带 Map、col_struct_array_structStruct 里带 Array Array 内嵌col_array_structArray Struct 内还有 Array、col_array_array_int二维数组Map 内嵌col_map_string_structMap 值为 Struct、col_map_string_array_doubleMap 值为 Array、col_map_string_map_string_dateMap 的值为 Map三层组合col_array_struct_with_mapArrayStruct…, Map、col_map_struct_with_arrayMap…, StructArray, TIMESTAMP。文档在col_struct_array_struct一列上专门注释Array of structs with single (inner) fields do not work with parquet.version 1.13.1。这是一个真实的工程约束——Parquet 1.13.1 在单字段内层 Struct 的数组场景存在兼容性问题测试数据特意用双字段inner_f3 TIMESTAMP, inner_f4 STRING的内层 Struct 绕开该缺陷也为读者在使用旧版 Parquet 时提供了可参考的规避经验。三、TBLPROPERTIES 与 Session 级配置MOR MDT 的完整组合3.1 表级属性TBLPROPERTIES ( primaryKey uuid, type mor, preCombineField precombine_field ) PARTITIONED BY (part_col)type mor声明 Merge-On-Read 表类型基础数据写入 Parquet 文件更新记录写入增量日志文件log file由读取端合并primaryKey uuid主键列用于记录去重与更新定位preCombineField precombine_field预合并字段同一主键多条记录按该字段取值决定最新记录值大者胜出本脚本中每次 UPDATE 都同步递增该字段保证更新语义正确PARTITIONED BY (part_col)按part_col分区脚本分别在分区A与B写入数据。3.2 会话级配置逐条解析建表后脚本通过spark.sql(set ...)设置了 6 个关键会话配置共同塑造了如何生成这张测试表spark.sql(sset hoodie.compact.inline.max.delta.commits9999) spark.sql(sset hoodie.compact.inlinefalse) spark.sql(sset hoodie.parquet.small.file.limit0) spark.sql(sset hoodie.metadata.compact.max.delta.commits1) spark.sql(sset hoodie.metadata.index.column.stats.enabletrue) spark.sql(sset hoodie.metadata.record.index.enabletrue) spark.sql(sset hoodie.metadata.index.secondary.enabletrue)配置项取值作用hoodie.compact.inline.max.delta.commits9999将内联触发 compaction 的累计 commit 阈值调大避免测试过程中自动触发 compactionhoodie.compact.inlinefalsefalse彻底关闭内联 compactionhoodie.parquet.small.file.limit0关闭小文件合并直接写入全新 Parquet 文件便于得到规整的文件布局hoodie.metadata.compact.max.delta.commits1MDTMetadata Table内部每 1 个 delta commit 就压缩一次保证后续读取 MDT 高效hoodie.metadata.index.column.stats.enabletrue开启 MDT 列统计索引Column Stats Indexhoodie.metadata.record.index.enabletrue开启 MDT 记录索引Record Index配合主键加速点查hoodie.metadata.index.secondary.enabletrue开启 MDT 二级索引Secondary Index支撑col_double上的CREATE INDEX其中三个 MDT 索引配置是文档强调的重点注释明确指出Partition stats index 会与 column stats index 一起启用Partition stats index is enabled together with column stats index即开启列统计索引的同时会自动带上分区统计索引。这些配置组合意味着这张测试表是一张MDT 全功能开启的 MOR 表Trino 连接器在测试中SessionBuilder.from(getSession()).withMdtEnabled(true)正是以 MDT 模式读取它。3.3 关闭 compaction 与 UPDATE 的关系关闭 compaction 是生成带日志文件测试数据的必要前提。脚本先写入 3 行数据然后执行两条 UPDATEUPDATE hudi_type_test_mor SET col_double col_double 100, precombine_field precombine_field 1 WHERE part_col A UPDATE hudi_type_test_mor SET col_string updated string, precombine_field precombine_field 1 WHERE part_col B由于 compaction 被禁用这两次 UPDATE 产生的变更会以log file 追加的形式保留下来而不会立即被 compaction 合并回 Parquet 基础文件。这正是验证 MOR 表 RO/RT 两条读取路径差异的关键RT 表合并了日志能读到更新后的值col_double 110.123、col_string updated string而 RO 表只读基础文件读到的是原始值10.123、NULL。四、插入数据标量 cast 与嵌套类型的 NULL 语义4.1 行 1 / 行 2分区 A 的完整数据行 1uuid1展示了每个类型的标准插入写法重点看几个 casttrue, cast(1 as tinyint), cast(100 as smallint), 1000, 100000L, 1.1, 10.123, cast(123.45 as decimal(10,2)), string val 1, cast(varchar val 1 as varchar(50)), cast(charval1 as char(10)), cast(binary1 as binary), cast(2025-01-15 as date), cast(2025-01-15 11:30:00 as timestamp), array(1, 2, 3), array(a, b, c), map(key1, 10, key2, 20), struct(struct_str1, 55, false), array(struct(1.1, array(n1,n2)), struct(2.2, array(n3))), map(mapkey1, struct(cast(2024-11-01 as date), cast(9.8 as decimal(5,2)))), ...要点BIGINT字面量用100000L后缀DECIMAL用cast(x as decimal(p,s))显式指定精度BINARY、VARCHAR、CHAR、DATE、TIMESTAMP都必须通过cast完成从字符串字面量到目标类型的转换ARRAY用array(...)构造MAP用map(k1, v1, k2, v2)构造STRUCT用struct(f1, f2, f3)按声明字段顺序构造多层嵌套用array(struct(...))、map(key, struct(...))逐层包裹与 Schema 中ARRAYSTRUCT...、MAPSTRING, STRUCT...的声明一一对应。4.2 行 3分区 B 的 NULL 全面覆盖行 3uuid3是 NULL 语义的压测用例除uuid、precombine_field与part_col外几乎所有列都显式插入 NULL并在个别嵌套列上保留了非 NULL 值null, null, null, null, null, null, null, null, -- 数值列全 NULL null, null, null, -- 字符串列全 NULL null, -- 二进制列 NULL null, null, -- 日期时间列 NULL null, null, null, -- 基础复杂类型 NULL null, -- col_struct NULL array(struct(3.3, array(n4))), -- col_array_struct 有值 null, ... -- 其余复杂类型 NULL这一行的价值在于验证两个关键语义标量列 NULLTrino 侧期望值统一为CAST(NULL AS 对应类型)如CAST(NULL AS DOUBLE)、CAST(NULL AS DECIMAL(10,2))嵌套列的部分空行 2 的col_array_struct整体为 NULL而行 3 的该列又有值行 2 的col_map_string_struct、col_array_struct_with_map、col_map_struct_with_array、col_struct_map中出现了Map 值为 NULL、Struct 内字段为 NULL、Array 元素为 NULL的细粒度 NULL用于验证 Trino 在嵌套路径上对 NULL 的逐层处理。4.3 二级索引的创建CREATE INDEX idx_double ON hudi_type_test_mor (col_double)这行语句依赖前面开启的hoodie.metadata.index.secondary.enabletrue在col_double上创建二级索引并写入 MDT使 Trino 侧可以对该列做索引裁剪加速。五、Trino 侧验证类型映射与读取断言5.1 测试执行方式testComprehensiveTypes是一个参数化测试参数来自comprehensiveTestParameters()对HUDI_COMPREHENSIVE_TYPES_V6_MOR/V8_MOR两张表分别跑isRtTable true / false共 4 组。测试会话通过.withMdtEnabled(true)强制以 MDT 模式读取且测试注释特别说明不使用#assertQuery()H2QueryRunner 无法定义 MAP 等类型改用#getQueryRunner()TrinoQueryRunner执行真实查询。这是测试编写上的重要约束Trino 的 H2 校验器不支持 MAP 类型因此复杂类型断言必须走真实 Trino 引擎。5.2 Spark 类型 → Trino 类型映射表依据测试期望值TestHudiSmokeTest.java L1150-L1268可以整理出完整映射Spark 类型Trino 类型示例期望值BOOLEANBOOLEANtrue/CAST(NULL AS BOOLEAN)TINYINTTINYINTTINYINT 1SMALLINTSMALLINTSMALLINT 100INTINTEGERINTEGER 1000BIGINT / LONGBIGINTBIGINT 100000FLOATREALREAL 1.1DOUBLEDOUBLEDOUBLE 110.123RT 表为更新后值DECIMAL(10,2)DECIMAL(10,2)DECIMAL 123.45STRINGVARCHARstring val 1VARCHAR(50)VARCHAR(50)CAST(varchar val 1 AS VARCHAR(50))CHAR(10)CHAR(10)CAST(charval1 AS CHAR(10))BINARYVARBINARYX62696e61727931binary1 的 UTF-8 字节DATEDATEDATE 2025-01-15TIMESTAMPTIMESTAMP(3)TIMESTAMP 2025-01-15 11:30:00.000ARRAYTARRAYTARRAY[1, 2, 3]MAPK,VMAP(K,V)MAP(ARRAY[key1, key2], ARRAY[10, 20])STRUCT...ROW(...)CAST(ROW(struct_str1, 55, false) AS ROW(f1 VARCHAR, f2 INTEGER, f3 BOOLEAN))几个值得注意的细节FLOAT → REALSpark 的FLOAT是单精度对应 Trino 的REAL而非DOUBLEBINARY → VARBINARYcast(binary1 as binary)的 UTF-8 字节62 69 6e 61 72 79 31恰好是字符串 binary1因此期望值为十六进制字面量X62696e61727931TIMESTAMP 精度Trino 默认将 TIMESTAMP 显示为毫秒精度TIMESTAMP(3)STRUCT → ROWStruct 字段在 Trino 中映射为 ROW 类型嵌套结构逐层转换为嵌套 ROW。5.3 断言逻辑逐列 全列 嵌套字段测试分三个层次验证L1276-L1328逐列验证对columnsToTest列表中的 30 个列逐一执行SELECT col FROM sourceTable与SELECT 期望值 UNION ALL SELECT 期望值 ...构造的期望查询比对便于定位具体哪一列的类型映射出错全列一次验证SELECT col1, col2, ... FROM sourceTable全列联合查询验证所有列同时读取的一致性嵌套字段提取验证对col_map_string_struct执行SELECT (map_values(col_map_string_struct))[1].nested_f4 AS extracted_nested_f4 FROM sourceTable验证Map → 值数组 → 取第一个 ROW → 点取嵌套字段这条嵌套访问链路Trino 数组索引从 1 开始确认深层次 Struct 字段能被正确解析。sourceTable的选择体现了 RO/RT 差异isRtTable true时读 RT 表合并日志false时读 RO 表仅基础文件因此precombine_field、col_double分区 A 被 UPDATE、col_string分区 B 被 UPDATE三列在不同模式下期望值不同——这正是对 MOR 表日志合并语义的端到端验证。5.4 测试数据加载测试时ResourceHudiTablesInitializer.java 负责解压 hudi_comprehensive_types_v8_mor.zip按hudiComprehensiveTypesColumns()、hudiComprehensiveTypesPartitionColumns()、hudiComprehensiveTypesPartitions()注册表元数据含_hoodie_*元数据列定义见该文件 L365-L370并为isCreateRtTable true的表同时注册 RO 与 RT 两张表。这也是该 md 文件与其 zip 配套使用的原因md 是如何生成zip 是已生成的结果。六、实战要点总结类型选型当前 Hudi 不支持TIMESTAMP_NTZ时间字段统一使用TIMESTAMPFLOAT在 Trino 侧读作REALBINARY读作VARBINARY十六进制展示STRUCT读作ROW。嵌套类型Parquet 1.13.1 下单字段内层 Struct 的数组有兼容性问题应像本脚本一样让内层 Struct 至少包含两个字段或升级 Parquet 版本。MOR 验证技巧通过hoodie.compact.inlinefalse与hoodie.compact.inline.max.delta.commits9999关闭 compaction再用 UPDATE 生成日志文件即可在同一张表上对比 RO 表基础文件与 RT 表基础文件 日志合并的读取差异更新时务必同步递增preCombineField以维持预合并语义。MDT 索引组合hoodie.metadata.index.column.stats.enabletrue会连带启用分区统计索引若需二级索引CREATE INDEX必须同时开启hoodie.metadata.index.secondary.enabletrue配合hoodie.metadata.record.index.enabletrue可得到一张 MDT 全索引 MOR 表这也是 Trino 连接器withMdtEnabled(true)测试模式的标配。测试数据资产需要一张现成的全类型 MOR 表用于 Trino 验证时可直接解压 hudi_comprehensive_types_v8_mor.zip 使用其生成脚本即本文剖析的 md 文件V6 与 V8 两个版本的数据分别由 hudi_comprehensive_types_v6_mor.md 与本文档生成可在 TestHudiSmokeTest.java 的testComprehensiveTypes中直接复用。赞分享数据湖湖仓一体大数据数据存储【免费下载链接】hudiUpserts, Deletes And Incremental Processing on Big Data.项目地址https://gitcode.com/gh_mirrors/hud/hudi点击查看免费下载相关推荐Apache Hudi TimestampBasedKeyGenerator 实战SCALAR 时间戳生成 yyyy-MM-dd hh 分区的 MOR 表Trino 端到端验证Apache Hudi TimestampBasedKeyGenerator 实战SCALAR 时间戳生成 yyyy MM dd hh 分区的 MOR 表T数据湖湖仓一体大数据数据存储Hudi 多分区字段 MOR 表的 Spark SQL 建表与 Trino 分区裁剪实战基于 hudi-trino 测试数据集Hudi 多分区字段 MOR 表的 Spark SQL 建表与 Trino 分区裁剪实战基于 hudi trino 测试数据集 导读 本文以 Apache数据湖湖仓一体大数据数据存储DataHub Redshift 元数据采集实战指南从权限配置、Lineage 到 Usage 与 ProfilingDataHub Redshift 元数据采集实战指南从权限配置、Lineage 到 Usage 与 Profiling 导读 本文以 DataHub 官方 R数据湖湖仓一体大数据数据存储上一篇提前加载数据让应用启动速度更快下一篇探索与故障排查的利器pgCenter —— PostgreSQL 管理工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表