Apache Hudi 深入
概述
Apache Hudi 起源于 Uber,核心解决数据湖上的数据更新与增量处理:支持 UPSERT(更新插入)、增量查询、流批写。本文讲透 CoW/MoR 两种表类型、Incremental Query、Clustering 与 Cleaner。
一、Hudi 是什么
| 特性 | 说明 |
|---|---|
| 表格式 | 数据湖上可更新的表 |
| 核心 | UPSERT、增量、时间线 |
| 引擎 | Spark/Flink 支持 |
| 场景 | 近实时数仓、CDC 入湖 |
Hudi 表 = 数据文件(Parquet/Avro)
+ 时间线(Timeline,提交历史)
+ 索引(文件定位)二、两种表类型
2.1 Copy-on-Write(CoW)
写时复制:
更新 → 重写受影响的数据文件(整个文件)
读 → 读最新文件(快)
写 → 更新代价高(重写)2.2 Merge-on-Read(MoR)
读时合并:
更新 → 追加写增量(Log 文件)
读 → 合并基文件 + 日志(慢)
写 → 快(追加)| 对比 | CoW | MoR |
|---|---|---|
| 写性能 | 慢(重写) | 快(追加) |
| 读性能 | 快 | 慢(合并) |
| 更新延迟 | 高 | 低 |
| 适用 | 读多写少 | 写多读少 |
2.3 表类型 vs 查询类型
| 查询类型 | 说明 |
|---|---|
| Snapshot | 最新快照(CoW/MoR 均支持) |
| Incremental | 增量变更 |
| Read Optimized | 只读基文件(MoR 优化) |
三、时间线与索引
3.1 时间线(Timeline)
记录所有提交/写操作的元数据:
commit、clean、rollback、compaction| 作用 | 说明 |
|---|---|
| 版本管理 | 提交历史 |
| 恢复 | 未完成操作回滚 |
| 增量定位 | 提交时间戳 |
3.2 索引(Index)
定位记录所在文件:
更新时先查索引,找旧文件| 索引 | 说明 |
|---|---|
| 布隆过滤 | 按文件过滤 |
| 简单索引 | 遍历文件 |
| HBase 索引 | 外部存储 |
| 记录级索引 | 精确定位 |
四、UPSERT 写流程
4.1 流程
1. 接收增量(含 upsert 语义)
2. 索引定位:记录在哪个文件
3. CoW:重写文件(合并新旧)
MoR:写日志文件
4. 更新索引
5. 提交时间线Spark UPSERT:
df.write.format("hudi").option("hoodie.datasource.write.operation", "upsert")4.2 写操作类型
| 操作 | 说明 |
|---|---|
| upsert | 更新或插入 |
| insert | 只插入 |
| bulk_insert | 批量插入(优化) |
| delete | 删除记录 |
| delete_partition | 删分区 |
五、增量查询(Incremental Query)
5.1 概念
只读"某次提交之后的新变更"
→ 增量消费(类似流)5.2 使用
sql
-- 增量读(时间戳之后)
SELECT * FROM hudi_orders
WHERE _hoodie_commit_time > '20260804_100000';| 场景 | 说明 |
|---|---|
| CDC 增量入数仓 | 消费变更 |
| 增量同步 | 下游持续更新 |
| 审计 | 变更追踪 |
5.3 与时间线关系
增量查询基于时间线定位提交
→ 需保留足够提交历史六、Clustering 聚类
6.1 问题
频繁小写入 → 大量小文件 → 查询慢6.2 Clustering 原理
把小文件合并成大文件(异步)
重新布局数据,不改变语义| 策略 | 说明 |
|---|---|
| 按分区 | 分区内合并 |
| 按排序 | 合并时排序 |
| 自定义 | 条件触发 |
触发:
table service 定期执行
或手动 trigger| 效果 | 说明 |
|---|---|
| 减少文件数 | 查询快 |
| 布局优化 | 数据局部性 |
七、Cleaner 清理
7.1 作用
删除无用历史文件:
过期提交的数据文件
未完成的写残留7.2 清理策略
| 策略 | 说明 |
|---|---|
| 按版本 | 保留最近 N 个提交 |
| 按时间 | 保留最近 T 时间 |
| 触发 | 提交后定期执行 |
| 注意 | 说明 |
|---|---|
| 清理 vs 增量查询 | 清理过旧会丢增量 |
| 孤儿文件 | 未提交写的残留 |
八、Compaction(MoR)
8.1 作用
把 MoR 的日志文件合并回基文件
→ 恢复读性能| 方式 | 说明 |
|---|---|
| 同步 | 写入时合并 |
| 异步 | 后台任务合并 |
| 在线/离线 | 调度执行 |
压缩后:日志消失,基文件更新九、生态集成
| 引擎 | 支持 |
|---|---|
| Spark | 完整 |
| Flink | 流写 |
| Presto/Trino | 查询 |
| Hive | 表注册 |
| 场景 | 选型 |
|---|---|
| 更新密集 | MoR |
| 读密集 | CoW |
| CDC 入湖 | MoR + 增量 |
| 大表批量 | bulk_insert |
常见问题速查
| 问题 | 原因与处理 |
|---|---|
| 查询慢 | MoR 未压缩,执行 compaction |
| 小文件多 | Clustering 合并 |
| 增量丢 | Cleaner 清理过快 |
| 更新性能差 | 换 MoR 或检查索引 |
| 写失败残留 | 时间线回滚清理 |