【Hudi】 Flink → Hudi 实时入湖实战
Hudi 的存储、索引、Timeline、并发控制,每个组件单独看都不复杂。这篇把它们串起来,看一条真实链路:用 Flink 把数据写进 Hudi,在一个场景里把各个概念都过一遍。
场景:订单流水实时入湖
假设有一条 Kafka 流,主题是 orders,每秒几千条订单消息。消息体是 upsert 语义——同一个 order_id 可能出现多次,前面的先来,后面的更正。
目标是把这条流写进 Hudi,要求:
- 数据实时可见:新数据进湖后几分钟内能查到
- 能更新已有记录:前面先来的订单,后面更正时要覆盖,不要写重
- 宕机能恢复:任务挂了重启,不丢数据
表类型选 MOR
场景决定了选型:
- 写频率高,每秒几千条——COW 每次重写 Parquet 顶不住
- 读延迟要求不高,几分钟内能查到就行——MOR 的合并开销能接受
所以表类型选 MOR,写入走 Avro log,延迟低,compact 后面异步做。
recordKey 和分区
确定主键和分区:
recordKey = order_id
partitionPath = event_date # 按天分区
recordKey 是 order_id,同一个订单的多次更新改同一个 file group 的同一个 log file,Hudi 用 Index 定位到具体文件做 update。
分区选 event_date(按天),因为多数查询按天过滤。分区太细(比如按小时)文件太碎,分区太粗(按周、按月)读放大。
并发安全
Flink 任务并行度通常是多个 TaskManager 同时写,并发安全靠 Timeline + LockProvider:
- 每个 checkpoint 生成一个 commit,多个 TM 同时写,生成多个 commit
- 提交时 Hudi 做乐观并发检查:改了不同 file group 的 commit 并行成功,改了同一个 file group 的排队
- LockProvider(ZooKeeper)锁住提交动作本身,防止两个 TM 同时拿到同一个 commitTime 产生竞态
Flink checkpoint 机制和 Hudi 的 commit 机制是耦合的:checkpoint 成功 = commit 成功,checkpoint 失败 = 回滚。
端到端流程
一条 order_id=1001 的记录从 Kafka 到 Hudi 的完整路径:
flowchart LR
A[Kafka<br/>orders topic] --> B[Flink Source<br/>反序列化 + shuffle]
B --> C[Bloom Index<br/>定位 file group]
C --> D[写入 Avro log<br/>MOR 增量文件]
D --> E[Checkpoint 触发<br/>写 .deltacommit]
E --> F[.hoodie/ Timeline<br/>标记已提交]
G[后台 Compact<br/>log → Parquet] -.-> F
H[后台 Clean<br/>清理旧文件] -.-> F
F --> I[下游消费<br/>快照/增量/Time Travel]
- Flink Source 从 Kafka 拉消息,反序列化,提取
event_date,按recordKey=order_id做 shuffle - 写入阶段:Hudi 的 Flink writer 把记录写进 Avro log 文件(MOR),写进
event_date=2024-08-14/分区下的对应 file group - Index 定位:写之前查 Bloom Index,找到
order_id=1001在哪个 file group 的哪个 log 文件里。如果是第一次出现就建新 log,如果是更新就追加到已有的 log - Commit 提交:Flink checkpoint 触发,Hudi 在
.hoodie/下写一个20240814100500.deltacommit,标记这批数据已提交。读端看到这个 commit 就能读到新数据 - 异步 Compact:后台定时把 log 文件和 data file 合并成新的 Parquet,减少读放大。compact 本身也是个 commit,合并完的数据文件通过 Timeline 被标记为正式数据
- 异步 Clean:清理过期的旧版本文件和旧 log,释放存储空间
读端怎么用
数据写进 Hudi 后,下游可以用多种方式消费:
- 快照读:Spark/Flink 批读,每次扫
.hoodie/拿最新 commit,读对应的数据文件 - 增量读:只读某个 commitTime 之后的文件,做下游增量处理
- Time Travel:查历史快照,比如"昨天中午 12 点这批订单的金额是多少"
小结
这是一条典型的实时入湖链路:Kafka → Flink → Hudi MOR。每个概念在链路里都能找到对应:
| 概念 | 链路里的位置 |
|---|---|
| 存储格式(Parquet + Avro) | MOR 用 Avro 写 log,compact 后生成 Parquet |
| File Group + Index | 按 order_id 分桶,Bloom Index 定位到文件 |
| Timeline | 每次 checkpoint 一个 commit,靠 .commit 文件标记提交 |
| 并发控制 | 乐观锁 + LockProvider 保证多 TM 同时写不冲突 |
| 数据模型 | (recordKey, partitionPath, commitTime) 三元组贯穿全程 |
这就是 Hudi 的设计:每个组件单独看都不复杂,串起来就成了一套完整的数据湖解决方案。
寒蝉 Hancic

