一、引言
Hudi 的 Merge on Read(MOR)表通过将增量数据写入 Delta Log 文件来实现高吞吐写入,但随着 Log 文件累积,读取性能会逐步退化。Compaction 作为 Hudi 的核心 Table Service,负责将 Delta Log 与 Base File 合并,是平衡读写性能的关键机制。
二、Compaction 机制原理
Compaction 的本质就是将"读时合并"的开销前置为"写时合并"的后台任务,在读写性能之间取得可控的平衡点。
Hudi Compaction 采用计划与执行分离的两阶段模型,这是其设计的核心特征:
设计优势:
- 解耦:Schedule 可以由 Writer 在每次 commit 后触发,Execute 可以由独立进程(如独立 Spark 作业)异步执行
- 容错:若 Execute 失败,Plan 仍存在于 Timeline 上,下次可直接重试执行,无需重新调度
- 灵活性:支持 inline(同步)、async(异步同进程)、独立作业 三种执行模式
Hudi 提供多种内置策略(hoodie.compaction.strategy),决定"哪些 File Group 优先被 Compact":
策略 | 说明 | 适用场景 |
LogFileSizeBasedCompactionStrategy | 按 Log 文件总大小降序排列,优先 compact Log 最大的 File Group | 默认策略,适合通用场景 |
BoundedIOCompactionStrategy | 在 LogFileSize 策略基础上限制单次 Compaction 的总 IO 量 | IO 资源受限的环境 |
UnBoundedIOCompactionStrategy | 不做 IO 限制,compact 所有待合并的 File Group | 资源充足时快速追赶 |
DayBasedCompactionStrategy | 按分区日期排序,优先 compact 最新分区 | 时间分区表,关注最新数据读取性能 |
Compaction核心触发参数:
# 每隔多少次 commit 触发一次 compaction schedule hoodie.compact.inline.max.delta.commits=5 # 异步模式下的触发条件(Flink 场景) compaction.delta_commits=5 compaction.delta_seconds=3600单个 File Group 的 Compaction 执行过程:
合并过程中:
- 读取 Base File 中的所有记录
- 按顺序 replay 每个 Log Block 中的操作(insert/update/delete)
- 使用 RecordMerger(1.x)或 HoodieRecordPayload(0.x)进行记录级合并
- 输出最终结果到新的 Base File
三、Hudi 1.x vs 0.x:Compaction 的架构演进
维度 | Hudi 0.x | Hudi 1.x | 影响 |
记录合并模型 | HoodieRecordPayload | RecordMerger 接口 | 合并逻辑更灵活、可插拔 |
索引机制 | 文件级 BloomFilter / HBase 索引 | Record-Level Index(默认) | Compaction 时定位记录更高效 |
并发控制 | OCC (Optimistic Concurrency Control) | Non-Blocking Concurrency Control (NBCC) | Compaction 不阻塞 Writer |
Log Compaction | 0.14+ 引入 | 增强与完善 | 减少 Snapshot Query 合并开销 |
Table Services 调度 | 与 Writer 耦合较紧 | 统一 Table Service Manager | 调度更灵活、可观测 |
存储抽象 | 依赖 Hadoop FileSystem | 新 HoodieStorage 抽象 | Compaction 可运行于更多存储后端 |
文件格式 | 固定 Parquet + Avro Log | 可扩展的 Reader/Writer 抽象 | 为未来格式扩展做准备 |
四、最佳实践
Compaction 执行模式选择:
推荐方案:
- 生产环境首选:独立的 Compaction 作业(资源隔离、可独立扩缩容、失败不影响写入链路)
- 开发/测试环境:Inline Compaction(简单、无需额外作业)
- Flink 实时场景:异步 Compaction(同一 Flink Job 内)或独立 Compaction Job
关键参数调优:
# ===== 调度频率 ===== # 每 N 次 delta commit 触发一次 Compaction Schedule hoodie.compact.inline.max.delta.commits=5 # 建议值:根据写入频率调整,过于频繁会增加小文件合并开销, # 过于稀疏则 Log 堆积影响读性能 # ===== IO 控制 ===== # 使用 BoundedIO 策略时,单次 Compaction 最大处理数据量(MB) hoodie.compaction.target.io=512000 # 建议值:根据集群可用资源设定,避免 Compaction 占用过多 IO # ===== 并行度 ===== # Compaction 并行度(Spark 场景) hoodie.compaction.parallelism=200 # 建议值:≈ 待 compact 的 File Group 数量,避免过小导致长尾 # ===== Compaction 策略 ===== hoodie.compaction.strategy=org.apache.hudi.table.action.compact.strategy.LogFileSizeBasedCompactionStrategy # ===== 异步 Compaction (Flink) ===== compaction.async.enabled=true compaction.delta_commits=5 compaction.max_memory=512 # Flink compaction task memory (MB)Compaction排查思路:
Compaction 积压(pending 数持续增长) ├── 原因1:Compaction 执行资源不足 → 增加并行度/独立作业资源 ├── 原因2:写入速率 >> Compaction 速率 → 降低触发频率 or 增加资源 ├── 原因3:大 File Group 导致单次 Compaction 耗时过长 → 使用 BoundedIO 策略 └── 原因4:Compaction 失败重试 → 检查日志,排查 OOM/数据问题