【Hudi】 Flink → Hudi 实时入湖实战

22 阅读 2188 字 · 约 8 分钟

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]
  1. Flink Source 从 Kafka 拉消息,反序列化,提取 event_date,按 recordKey=order_id 做 shuffle
  2. 写入阶段:Hudi 的 Flink writer 把记录写进 Avro log 文件(MOR),写进 event_date=2024-08-14/ 分区下的对应 file group
  3. Index 定位:写之前查 Bloom Index,找到 order_id=1001 在哪个 file group 的哪个 log 文件里。如果是第一次出现就建新 log,如果是更新就追加到已有的 log
  4. Commit 提交:Flink checkpoint 触发,Hudi 在 .hoodie/ 下写一个 20240814100500.deltacommit,标记这批数据已提交。读端看到这个 commit 就能读到新数据
  5. 异步 Compact:后台定时把 log 文件和 data file 合并成新的 Parquet,减少读放大。compact 本身也是个 commit,合并完的数据文件通过 Timeline 被标记为正式数据
  6. 异步 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 的设计:每个组件单独看都不复杂,串起来就成了一套完整的数据湖解决方案。