Apache Hudi 想象成一个非常智能的“超级文件柜”。这个文件柜不仅要存东西还要能改东西、查东西。咱们今天要聊的就是怎么设计这个文件柜的核心机制。我会用最通俗的例子带你把 Copy On Write (COW)、Merge On Read (MOR) 以及主键分区这些概念彻底搞懂。一、 Copy On Write (COW) —— “完美主义者”模式想象一下你的文件柜里存着一页页的客户资料每一页都是一张纸。1. 写入流程与重写机制场景客户“张三”搬家了他的地址要从“北京”改成“上海”。这页资料目前锁在文件夹“第1号文件”里。COW 的做法取出把“第1号文件”整个拿出来。修改找到张三那行把地址改成“上海”。但是COW 不允许在原文件上涂改。它必须打印一张全新的、完美无缺的“第1号文件”里面包含所有原来的数据 张三的新地址。替换把旧的“第1号文件”扔进碎纸机把这张新的放回去。核心逻辑任何修改都伴随着包含该数据的整个文件的重写。2. 适用场景与读性能优势读性能极快因为文件柜里永远是最新的、整理好的、干干净净的文件。你伸手去拿拿出来就能看不用做任何额外处理。适用场景读多写少比如每天只更新一次数据但这一整天要被分析师查几百次。数据量适中因为每次写都要重写文件如果你总是一秒钟改一次数据光打印文件就累死了。二、 Merge On Read (MOR) —— “灵活便利贴”模式还是那个文件柜。MOR 模式觉得 COW 太死板改个地址还要重打整张纸太浪费了。1. Base File 与 Log File 的协作在 MOR 模式下文件柜里出现了两种东西Base File (基础文件)这就是那张整理好的正式资料纸Parquet 格式它是基础。Log File (日志文件)这就是“便利贴”Avro 格式。场景客户“李四”的电话号码变了。MOR 的做法它不去动那张正式资料纸。它快速拿出一张便利贴写上“李四的新电话是 138xxxx”然后啪地贴在文件夹旁边。完事儿写入速度极快。2. 读时合并与写时压缩那我要查“李四”的电话怎么办读时合并你先拿出那张正式的资料纸看到上面写着李四的旧电话。你发现旁边贴着一张便利贴。你必须在脑子里或者电脑里把这两部分拼起来哦旧电话作废用便利贴上的新电话。代价这就比直接读那张干净的纸要慢一点点因为你多了一步“合并”的动作。写时压缩便利贴贴多了会怎么样资料就乱了查起来越来越慢。所以MOR 有一个后台清洁工。2. 这个清洁工会定期比如每隔几个小时过来把所有便利贴撕下来真正地改写进那张正式资料纸里然后把便利贴扔掉。3. 这个过程就叫Compaction压缩/合并。3. 场景选择如何在 Latency延迟与 Throughput吞吐间权衡这个选择就像选快递模式Latency (数据延迟/多久能看到新数据)Throughput (写入吞吐量/每秒能写多少)适合场景COW高(因为要重写文件慢)低(重写太费劲写不动)追求秒级查询数据更新不频繁如日报表。MOR低(贴个便利贴就行极快)高(贴便利贴非常轻松能抗住大流量)追求秒级写入数据变化极快如订单状态、实时日志。一句话总结如果你要快写选 MOR如果你要稳读选 COW。三、 Primary Key 与 Partition 的关系现在我们知道怎么存文件了下一步是怎么整理这些文件。这就涉及到 Primary Key主键和 Partition分区。1. RecordKey 的设计原则RecordKey就是每一行数据的身份证号。唯一性最重要比如你的主键是“用户ID”。假如你来了两条数据用户ID1001说名字叫“Alice”过了一会儿又来一条用户ID1001说名字叫“Bob”。Hudi 怎么知道这两条数据其实是同一个人的不同状态全靠 RecordKey它看到 ID 都是 1001就知道这是“更新”操作而不是“新增”两个人。如果你的 RecordKey 选错了比如选了“性别”那所有男的记录都会打架所有女的记录也会打架数据就全乱套了。分布性如果你的身份证号前几位都是“110”比如北京户口那么“110”这个文件夹就会特别大挤死了而“330”浙江的文件夹可能没人去。好的主键应该让数据尽量均匀地散落在各个文件里这样大家干活才不累。2. 分区策略对查询性能的影响Partition分区就像是文件柜上的大抽屉标签。怎么分区通常按照时间天、小时或者业务属性城市、部门来分。对查询的影响场景 A好的分区你的数据按“省份”分区。老板问“给我看广东省的所有订单。”Hudi 太开心了它直接走到标着“广东省”的那个抽屉其他抽屉看都不看。速度极快场景 B坏的分区你的数据按“性别”分区只有男、女两个抽屉。老板问“给我看订单金额大于 100 元的数据。”Hudi 傻眼了。它不知道金额大于 100 的人是男是女啊它只能把“男抽屉”和“女抽屉”里的所有文件都翻一遍。速度极慢总结一下关系分区是为了让你一下子扔掉一大堆不需要的数据大范围搜索。主键是为了让你在剩下的数据里精确找到那一条要修改的数据精准定位。高阶口诀先分区缩小范围再主键精准定位。COW 适合读MOR 适合写。RecordKey 要唯一分区字段要常查。3.预聚合字段precombinehoodie.datasource.write.precombine.field我们可以从**“乱序数据”和“主键去重”**这两个核心场景入手用大白话解释它的作用、原理和实际意义。一、先搞懂为什么需要这个参数想象一下你在收快递你给商家下了两个订单同一商品不同地址但快递员先送了“旧地址”的后送了“新地址”的。如果快递系统没有“按最新地址派送”的规则可能就会把“旧地址”的快递放在“新地址”导致你收不到货。在Hudi中**“主键recordkey”就像“订单号”“更新数据”**就像“新地址”。当同一主键的数据比如同一用户的两次更新乱序到达比如先收到旧数据后收到新数据Hudi需要知道到底该保留哪条二、precombine.field的作用给“最新数据”加个“优先级标签”hoodie.datasource.write.precombine.field就是Hudi的“优先级规则”——指定一个字段用来判断哪条数据是“最新的”。比如你设置precombine.field tsts是时间戳字段Hudi就会当同一主键比如user_id1有两条数据时比较它们的ts值保留ts最大的那条也就是“最新”的丢弃旧的。三、举个“乱序更新”的例子假设你要更新用户信息表结构如下CREATE TABLE user_table ( user_id STRING, -- 主键recordkey name STRING, age INT, ts TIMESTAMP -- 预合并字段时间戳 ) USING hudi;场景1正常顺序更新先收到user_id1, nameAlice, age25, ts2024-08-01 10:00后收到user_id1, nameAlice, age26, ts2024-08-01 10:05→ Hudi保留ts10:05的数据最新。场景2乱序更新关键先收到user_id1, nameAlice, age26, ts2024-08-01 10:05新数据后收到user_id1, nameAlice, age25, ts2024-08-01 10:00旧数据→ 如果没有precombine.fieldHudi可能用“旧数据”覆盖“新数据”因为旧数据先到→ 但设置了precombine.fieldts后Hudi会比较ts保留ts10:05的数据正确。四、这个参数的本质解决“数据乱序”问题Hudi的写入通常是流式/批量的数据可能来自不同的源比如CDC、日志时间戳可能不一致甚至“旧数据”可能比“新数据”先到。precombine.field就是Hudi的“纠错机制”确保即使数据乱序最终存储的也是最新状态。五、实际操作中的注意事项字段必须存在precombine.field指定的字段比如ts必须在表结构中并且在写入数据时包含。字段类型要“可比较”通常用TIMESTAMP时间戳或BIGINT序列号不能是STRING字符串无法直接比较大小。场景适配流式写入如Kafka实时数据必须设置否则乱序数据会导致错误批量写入如每天全量同步可以不设置因为数据通常有序但建议保留避免未来数据乱序。总结一句话记住它hoodie.datasource.write.precombine.field“给数据加个‘最新标记’让Hudi知道该保留哪条”核心是解决同一主键的乱序更新问题确保数据最终状态正确。比如你在CDC场景中同步用户表一定要把数据库的update_time字段设为precombine.field这样即使数据延迟Hudi也能正确合并四、COW表案例演示1. 建表Spark SQL假设我们要建一个用户表user_id作为Primary Keydt作为Partition按天分区CREATE TABLE hudi_user_cow ( user_id STRING, name STRING, age INT, dt STRING -- Partition字段按天分区 ) USING hudi OPTIONS ( hoodie.table.name user_table_cow, hoodie.datasource.write.recordkey.field user_id, -- Primary Key hoodie.datasource.write.partitionpath.field dt, -- Partition字段 hoodie.datasource.write.table.type COPY_ON_WRITE, -- COW表 hoodie.datasource.write.precombine.field ts -- 预合并字段选最新数据 );2. 数据写入插入/更新插入新数据INSERT INTO hudi_user_cow VALUES (001, Alice, 25, 2024-08-01); INSERT INTO hudi_user_cow VALUES (002, Bob, 30, 2024-08-01);更新数据如user_id001的年龄从25改为26INSERT INTO hudi_user_cow VALUES (001, Alice, 26, 2024-08-01); -- COW会重写整个dt2024-08-01的文件3. 查询数据COW表读时无需合并直接读取最新数据SELECT * FROM hudi_user_cow WHERE dt 2024-08-01; -- 结果(001, Alice, 26, 2024-08-01), (002, Bob, 30, 2024-08-01)4. Compaction自动触发COW的Compaction由Hudi自动调度默认每5分钟检查一次无需手动操作。如果需要手动触发-- Spark SQL CALL hudi_compact(hudi_user_cow, 2024-08-01); -- 对dt2024-08-01的分区执行Compaction三、MOR表案例演示1. 建表Spark SQL同样建用户表但表类型改为MERGE_ON_READCREATE TABLE hudi_user_mor ( user_id STRING, name STRING, age INT, dt STRING ) USING hudi OPTIONS ( hoodie.table.name user_table_mor, hoodie.datasource.write.recordkey.field user_id, hoodie.datasource.write.partitionpath.field dt, hoodie.datasource.write.table.type MERGE_ON_READ, -- MOR表 hoodie.datasource.write.precombine.field ts );2. 数据写入插入/更新插入新数据INSERT INTO hudi_user_mor VALUES (001, Alice, 25, 2024-08-01); INSERT INTO hudi_user_mor VALUES (002, Bob, 30, 2024-08-01);更新数据如user_id001的年龄从25改为26INSERT INTO hudi_user_mor VALUES (001, Alice, 26, 2024-08-01); -- MOR只写增量log文件不重写基础文件3. 查询数据MOR表读时需要合并log和基础文件Hudi自动处理SELECT * FROM hudi_user_mor WHERE dt 2024-08-01; -- 结果(001, Alice, 26, 2024-08-01), (002, Bob, 30, 2024-08-01)4. Compaction必须定期执行MOR的Compaction需要手动或调度触发否则log文件越来越多读性能下降-- Spark SQL对dt2024-08-01的分区执行Compaction合并log和基础文件 CALL hudi_compact(hudi_user_mor, 2024-08-01); -- 或者通过Hudi的调度工具如Airflow定期执行四、实际操作注意事项Primary Key选择必须选择唯一且稳定的字段如user_id避免重复导致数据混乱。如果没有天然主键可以用组合字段如user_iddt。Partition选择优先选时间字段如dt方便按时间范围查询如WHERE dt BETWEEN 2024-08-01 AND 2024-08-07。避免选基数过低的字段如gender否则分区过多导致管理复杂。COW vs MOR选择COW适合读多写少场景如报表查询读性能优先。MOR适合写多读少场景如实时日志写入写性能优先。Compaction频率COW自动触发无需频繁手动操作。MOR必须定期执行如每天凌晨否则log文件堆积读性能下降。总结COW写时重写文件读快写慢Compaction自动。MOR写时写log读时合并写快读慢Compaction必须手动。Primary Key和Partition是Hudi的核心直接影响数据去重和查询性能。