数据湖的 Compaction 策略:什么时候该合并、合并粒度多大才合适

讲解数据湖Compaction的原理和核心问题,深入分析什么时候需要合并小文件、合并粒度如何选择,对比分区级、时间窗口级、表级等方案在不同场景下的优缺点,并给出工程实践建议,适合使用Delta Lake、Iceberg、Hudi的数据湖团队参考。

小文件问题不是新问题,但数据湖放大了它

在数据湖上跑任务,最烦人的问题往往不是查询引擎不给力,而是数据文件本身变得越来越碎。原本以为把数据都丢进 HDFS 或 S3 就完事,结果一次扫描要打开几万甚至几十万个小文件,元数据节点直接成了瓶颈,查询时间从分钟级变成小时级。于是 Compaction(合并小文件)成了数据湖维护里绕不开的话题。

AI technology illustration

但真正让团队拿不定主意的,不是要不要做 Compaction,而是到底什么时候触发合并、每次合并多大范围才算合理。合并太勤,计算资源消耗大;合并太少,小文件又持续蔓延。这个问题没有标准答案,只能按数据特征和业务场景去设计策略。这篇文章不打算给一个万能模板,而是把 Compaction 决策中真正需要权衡的因素讲清楚。

Compaction 的本质是空间与算力的互换

先明确一个概念:Compaction 不是后台任务里可有可无的“优化操作”,它是在用计算资源换未来的查询效率。每执行一次合并,都需要读取待合并的旧文件、排序或重写、再写入新文件。这个过程中,IO 和 CPU 都是实打实的成本。而收益是后续查询扫描的文件数变少,谓词下推和并发规划更有效。

所以 Compaction 策略的核心,是在“写入侧的生产速度”和“读取侧的消费能力”之间找一个平衡点。如果写入端每分钟都在追加新文件,而合并任务每天只跑一次,那文件数一定压不住。反过来,如果写入端本来就不太产生小文件,比如每天一个批任务写少量大文件,那就算不做 Compaction 也没问题。

真正麻烦的是合并任务和写任务的并发冲突

在 Hudi、Delta Lake 这类支持事务的数据湖格式中,合并意味着重写一批数据文件,写入新版本。如果此时上游还在对同一批文件进行增量提交,就会产生版本冲突。很多团队发现自动 Compaction 经常失败,不是资源不够,而是和写入端抢锁抢版本。要缓解这个问题,除了错峰执行,还要注意合并任务的粒度拆得足够细,让锁的持有时间变短。这是很多早期 Compaction 方案失败的主要原因,也是社区把“小文件合并”逐步拆成“文件组级别重写”的动机。

判断是否需要合并的几类信号

工程上不会有人天天盯着文件数,更实际的做法是通过几个指标来感知恶化程度。我见过比较靠谱的判断方式有这么几种:

  • 文件数量级:单个分区内小文件数量超过几千甚至上万,元数据服务(如 Hive Metastore、HDFS NameNode、甚至 S3 的 list 操作)就会开始拖后腿。
  • 查询计划耗时:执行 EXPLAIN 时发现扫描的 file chunks 数量异常多,实际读取数据量不大却花了很长时间。
  • ETL 调度时间膨胀:同一个增量处理任务,以前跑 20 分钟,现在要 40 分钟,且数据量并没有翻倍。

这些信号出现时,不一定要立刻启动全表级 Compaction,但至少要开始观察小文件的增长趋势。建议把文件统计做成每日指标,观察每个分区的文件数环比增长幅度。如果连续三天增长超过 30%,就说明写入侧的压缩逻辑可能出了问题。

合并粒度:分区级、时间窗口级、表级的权衡

合并粒度直接决定了一次 Compaction 的“破坏半径”。粒度太小,比如只合并几个文件,几乎没什么效果;粒度太大,比如整个表一把梭,可能在合并期间写入任务全堵住。不同粒度的核心区别是:你愿意为“减少文件数”付出多少并发牺牲和失败风险。

-- 类似 Delta Lake 的 OPTIMIZE 语法,实际引擎有差异
OPTIMIZE dwd_orders
WHERE dt >= '2025-05-01' AND dt < '2025-06-01'
ZORDER BY (order_id)

上面这种按时间分区做局部合并,是目前最常见的做法。它对写入方的影响最小,因为合并只锁定一个时间范围内的文件,不会阻塞新写入的分区。如果你的表是按天分区的,那这种分区级合并几乎就是默认选项。

下面用一张表来对比不同合并粒度的特点,方便团队快速定位自己的场景:

合并粒度 目标范围 优点 风险/成本 适用场景
文件组级 几十到几百个文件 成本极低,实时性好 对整体文件数改善微弱 高频写入且只需缓解热点分区
分区级 单个或多个分区 不影响新写入,效果可控 需要选对分区策略 以日期/业务维度分区的典型数仓
时间窗口级 过去一段时间的数据 兼顾实时和批量,便于调度 窗口边界要仔细定义 有滑动窗口特征的流式写表
表级 整张表 文件数下降最明显 耗时长,可能引发写入冲突 离线批量更新后的一次性治理

如果数据量不大,表级合并还能接受;一旦表的大小到了 TB 级,表级合并的时间和失败回滚成本都会让人头疼。所以工程上通常用分区级或时间窗口级作为常规手段,表级只做最终兜底。

时间窗口级合并更适合流式数据湖

很多团队用 Kafka + Flink 往数据湖里写数据,文件是按 checkpoint 或批次产生的。每 1 分钟一个 checkpoint,一天就是 1440 个小文件。如果只按天合并,每次要处理的数据量太大,而且会积攒过多中间状态。

我见过一个可行的方案:每 15 分钟触发一次 Compaction,只合并最近 1 小时的数据,同时把昨天的数据拆成两个窗口滚动执行。这个做法对 Kafka 上游的尖峰流量也能扛住,因为单次合并运行时长不超过 5 分钟,即使失败重试也不影响新的写入。更重要的是,它让“文件数增长”和“合并消耗”都稳定在一个可控区间,不会出现某一天突然需要跑一个几小时的大合并。

Compaction 调度与执行的常见误区

误区一:合并次数越勤越好

理论上每次流式 checkpoint 后都可以做一次 Compaction,但实际是每次合并都要读一遍数据。如果合并频率和写入频率一样高,等于每写一份数据就重复读一次,开销直接翻倍。合理的做法是让 Compaction 的频率至少比写入频率低一个数量级,比如写入是分钟级,合并至少在 5 到 15 分钟以上。同时要为合并任务设置资源池上限,避免它抢走查询的资源。

误区二:粒度越大越能解决所有问题

表级合并虽然能让文件数量降到最低,但会引入两个问题:一是任务运行时间长,失败概率高;二是合并完成前,大量旧文件还活着,查询依然要扫描很多文件。而且大范围合并带来的排序重写,很可能把原本良好的数据局部性打乱,比如同一个订单的明细被分散到不同文件中,后续做 Z-order 聚类反而要花更长时间。

误区三:只盯着文件数,忽略文件大小均衡

文件数少了,但每个文件大小差异悬殊,同样会拖慢查询。比如一个 1GB 的大文件加上一堆 4KB 的小文件,任务调度粒度仍然不均衡。好的 Compaction 策略还应该控制合并后文件大小落在目标范围内,比如 64MB 到 256MB,避免出现极端倾斜。判断合并后文件大小的常用方式是:目标文件数 = 分区数据总量 / 目标大小,然后再决定分区内一次要啃掉多少文件。

设计一套可演进的 Compaction 策略

给团队做方案时,我建议走“先诊断、再分区、后渐进”的路径。别一开始就修改引擎参数,先搞清楚自己的表到底病在哪里。

第一步:用统计任务摸清家底

先写一个简单的分析脚本,统计每个分区下文件数、平均大小、最大最小文件的差值。可以定期跑,输出到一个状态表。没有这份数据,任何 Compaction 参数都是拍脑袋。

# 伪代码:统计每个分区的文件分布
for partition in list_partitions(table):
    files = list_files(partition)
    record(partition, file_count, total_size, avg_size, max_file, min_file)

第二步:按数据热度和分区大小选择触发条件

一个相对稳妥的触发框架,是把数据分成冷、温、热三档,分别使用不同的合并节奏:

  • 冷分区(超过 30 天):一周做一次分区级合并,目标文件大小设为 256MB。
  • 温分区(过去 30 天内):每天做一次时间窗口合并,窗口覆盖最近 24 小时。
  • 热分区(今天内):每 15 分钟做一次小范围合并,只合并最近 1 小时的数据。

这个策略不是通用的,但它的好处是限制了一次合并的波及面,也让资源消耗能稳定在一个量级。具体阈值需要根据集群规模调整,但如果你的热分区有上千万行,那 15 分钟合并 1 小时的数据量也未必够,可能需要把窗口缩短为 15 分钟,或者提高合并任务并发度。

第三步:监控合并任务本身的开销

很多团队忽略这一点。Compaction 任务会占用资源池,如果和实时查询共用同一队列,很容易把查询延迟打上去。要单独划分一个维护队列,并监控合并任务的 CPU、内存和耗时。当合并任务的平均运行时长超过其预设窗口时,就该降低合并频率或缩小合并范围。还有一个容易忽略的指标是“合并失败率”。失败率超过 10% 时,多半是并发冲突,而不是资源问题,需要检查写入端提交频率和表锁机制。

经验提醒:在启用自动 Compaction 之前,先在测试环境用 10 倍小文件压力验证一下触发阈值。很多自动合并功能在文件数多到一定程度时会产生“合并风暴”,比不合并更可怕。

不同数据湖格式在 Compaction 上的处理差别

这里简单提一下三个主流格式的行为差异,帮助选择策略落地方式。

  • Delta Lake:原生 OPTIMIZE 命令,可以指定分区和 ZORDER,合并任务写回新版本事务日志。适合 Spark 生态,对批量合并的优化做得比较成熟。
  • Apache Iceberg:支持通过 Spark 或 Flink 执行 rewrite data files,可以细粒度控制文件组,也支持流式增量合并。如果你同时需要批读和流读,Iceberg 的 manifest 机制会更友好。
  • Apache Hudi:内置 Clustering 和 Archive,支持表服务自动异步执行,更贴近流式湖仓。它的 clustering 可以配置分区校验策略和文件组大小。

不过格式再不一样,底层的取舍逻辑是一样的:合并时机和粒度都取决于写入吞吐、查询读取模式、文件大小目标以及集群资源余量。把表格里的方案换成自己引擎的语法,就能开始落地。

最后:Compaction 不是孤立的最优解

如果你发现无论怎么调参数,文件还是会快速增长,大概率不是 Compaction 策略的问题,而是写入端的设计有问题。比如 Flink checkpoint 设置过小、Spark 写入分区数固定过高,或者下游 update 产生太多 delete gap。先把写入端的文件生成控制住,Compaction 才会轻松。这也是为什么很多团队最后会选择“写入端文件数预算 + 维护端定时 Compaction”的组合策略。

数据量再大,Compaction 也是一个可以渐进优化的工程问题。先跑起来,看到指标变化,再调整粒度,比一次性设计一个完美的框架要靠谱得多。二八法则在这里也很适用:先通过统计摸清小文件最集中的几个分区,把 80% 的收益用 20% 的成本拿下来,剩下的细节再慢慢打磨。

原创文章,作者:fudengji,如若转载,请注明出处:https://fudengji.cn/article/1038/

(0)
上一篇 54分钟前
下一篇 49分钟前

相关推荐