数据湖中的小文件问题:成因、影响与自动化合并策略

本文深入分析数据湖中小文件问题的成因与性能影响,介绍基于 Spark、Iceberg、Delta Lake 和 Hudi 的自动化小文件合并策略,并给出触发条件、资源控制与常见误区等工程落地建议。适合数据平台工程师和数仓团队参考。

小文件问题是怎么来的

很多团队在数据湖上跑得越久,越容易被一个不算复杂的问题缠住:表里的文件越来越多,单个文件却越来越碎。一开始只是查询变慢,后来连元数据服务和调度器都开始报警,排查之后才发现,一张数据量不高的表,居然积累了几十万个几十 KB 的 Parquet 小文件。这就是数据湖里非常典型的小文件问题

AI technology illustration

小文件问题,本质上不是“文件多”三个字能概括的。它会消耗 NameNode 内存、放大查询扫描的开销、拖慢任务调度,甚至让写入任务变得更加脆弱。更麻烦的是,一旦形成,靠单纯加机器很难消除,必须在管道层面引入主动的合并策略。这篇文章会先从成因讲起,再拆解影响,最后重点聊聊自动化合并的思路和落地方式。

数据湖的小文件问题,和一次性大导入不同,它主要是在持续、零散的写入过程中积累出来的。最常见的来源有下面几种。

  • 频繁批处理写入:每个批任务写一批数据,如果分区多、并行度又高,一次写入就会产生几十甚至上百个文件。
  • 动态分区写入:Spark 或 Hive 按某个动态列分区时,每个分区都可能有多个并发任务向其中写文件,分区越多,文件碎片越多。
  • 流式微批写入:Flink、Spark Streaming 每隔几秒提交一个微批,每个微批至少生成一个新文件,一天下来新增文件数量非常可观。
  • 事务表的数据文件版本:Iceberg、Delta Lake 这类表在每次 INSERT、UPDATE、DELETE 后都会生成新的数据快照文件,旧文件需要保留一段时间才被清理,文件数会更快膨胀。
  • 分区粒度过细:把一张本来数据量不大的表按小时、按城市拆成大量分区,每个分区里的数据又不多,文件自然无法填满。

可以看到,小文件问题的共性在于写入频率快、单次数据量小、分区维度多。如果不做干预,文件增长速度通常会线性甚至加速增长,直到超过存储和计算体系的某种隐式上限。

小文件拖慢数据湖的四个路径

查询引擎的 Task 调度开销。Spark、Hive 在扫描数据时,会按照文件或文件块生成输入分片。小文件多意味着输入分片数量巨大,一个读取全表的查询可能被拆成几十万个 Task,调度、序列化、Executor 分配的开销远远超过实际执行时间,最终表现为集群很忙,查询却很慢。

元数据服务的内存压力。在 HDFS 上,每个文件、每个 block 都会在 NameNode 中保留一条元数据记录。虽然单条记录不大,但几百万小文件会显著增加 JVM 堆内存消耗,严重时甚至触发 Full GC。对象存储虽然没有 NameNode,但每次 List 请求的数量、MetaStore 里分区和文件的统计信息,也会成为新的瓶颈。

写入链路的脆弱性。写入端同样会被小文件反噬。最常见的场景是,动态分区写入时,由于下游输出文件过多,每个写入分片都对应一个文件写入器,Driver 需要维护的 Writer 数量暴涨,最终导致内存溢出或文件描述符耗尽。另一个场景是,如果底层是对象存储,短时间内大量小文件写入还可能触发请求频率限制。

下游消费成本和语义复杂化。许多计算引擎在打开表时会先拉取文件列表并推断 Schema,文件数越多,这一阶段消耗的时间越长。对于需要多次扫描的亚秒级查询,小文件造成的固定开销无法忽略,这也是很多人在引入数据湖后发现“查询还不如直接查原始文件快”的原因之一。

怎么判断表里已经有小文件问题

判断一张表是否存在小文件问题,最直接的办法是遍历底层目录,统计文件数量、总大小以及平均文件大小。以 HDFS 上的表为例,下面这条命令可以快速统计文件数和总字节数。

hdfs dfs -ls -R /data/warehouse/ods/order_table | awk '$1 ~ /^-/{count++; size+=$5} END{print count, size}'

文件数只能说明有多少,不能说明有多碎。更准确的判断要看平均文件大小是否远低于目标文件大小。在 Spark 或 Python 环境里,也可以写一个简短脚本遍历底层路径,计算每个分区的文件分布,再决定是否触发合并。

如果统计结果显示单分区文件数超过几千、平均大小低于 32MB,而目标文件块大小设置为 128MB,那么基本可以断定存在需要治理的小文件积累。当然,具体阈值要结合查询频率和存储系统来定,不应盲目照搬。

合并小文件的基础思路

合并的本质,是读取原来的小文件,以更合适的并行度重新写入,让每个输出文件的大小尽量接近底层存储和查询引擎期望的块大小。下面三种思路在实际工程中最常用。

写入时减少文件数

在数据写入阶段就控制文件数量,是最省成本的做法。对于 Spark SQL 批任务,可以通过调整 shuffle 分区数来间接控制输出文件数。

SET spark.sql.shuffle.partitions=256;
SET spark.sql.adaptive.enabled=true;
SET spark.sql.adaptive.advisoryPartitionSizeInBytes=128MB;

这里的关键思路是让每个输出分区对应一个目标文件,数据量除以分区数大致等于期望的文件大小。注意,分区数不能一味调小,否则单个文件过大,后续更新和读取都可能受影响。

定期批量重写

如果小文件已经形成,最直接的办法是使用 Spark 读取目标表,按分区或日期范围重分区,再覆盖写回。为了避免一次重写全表带来太大压力,通常限定最近 N 天或某个分区范围。这种方式的优点是灵活,不依赖特定表格式;缺点是并发写可能出现数据丢失,所以只适合离线批处理场景,重写作业本身也要控制输出文件大小,否则只是把一堆小文件搬成另一堆小文件。

利用表格式自带的合并能力

对于采用 Iceberg、Delta Lake 或 Hudi 的数据湖表,更推荐直接使用内置的文件合并能力。这些操作不仅会自动挑选需要重写的文件组,还能利用事务保证读写隔离,合并完成后还能触发快照清理。

-- Delta Lake 合并某一天分区的文件
OPTIMIZE db.table
WHERE dt = '2024-05-01'
ZORDER BY user_id;

-- Hudi 异步 clustering
CALL run_clustering(table => 'db.table', order => 'dt');

下面是三种主流表格式与自研 Spark 重写方案的简单对比。

实现方式 典型命令/Action 事务保障 适用场景
Iceberg rewrite_data_files 支持 大规模分析表、需要时间旅行
Delta Lake OPTIMIZE 支持 湖仓一体、带更新删除的表
Hudi Clustering 支持 流批一体、变更频繁的表
自研 Spark 重写 INSERT OVERWRITE 较弱 普通 Hive 表、无事务依赖

这些内置能力并不是“点了就立刻生效”的黑魔法。合并动作本身也是计算任务,需要消耗资源,并会产生新的数据版本。只有把调度和触发条件设计好,才能真正自动化。

自动化合并策略怎么设计

自动化合并,就是让系统根据表的状态自己决定何时、哪些范围需要合并。设计时需要解决三个问题:什么时候触发、合哪部分数据、怎么避免影响正常读写。

常见的触发条件有以下几类。

  • 按分区触发:当某个分区的文件数超过阈值,或平均文件大小低于目标大小的一半时,启动该分区的合并。
  • 按数据年龄触发:只对最近写入产生的小文件做合并,例如每 6 小时或每天检查一次新增分区。
  • 按查询成本触发:从执行计划中观察 scan 的文件数量,如果经常超过设定值,就把对应表加入合并队列。

一个典型的自动化流程是:由调度器启动检查任务,扫描表的文件统计信息,根据预设规则筛选出需要重写的分区,然后提交合并作业;合并完成后,再根据表格式清理过期快照。整个过程可以嵌入 Airflow、DolphinScheduler 或云平台的托管任务中。

def auto_compact(table):
    stats = scan_file_stats(table)
    for p in stats:
        if p.file_count > 10000 or p.avg_size < 64MB:
            submit_compact_job(table, p)

对于 Iceberg 和 Delta Lake,还可以把合并动作与写入任务解耦。Iceberg 的 RewriteDataFiles 能够选择一小部分文件组进行重写,不会锁表;Delta Lake 的 OPTIMIZE 也支持后台执行。这样即使合并期间有新的写入,也不会丢数据。

资源控制和频率同样重要。合并任务最好放在业务低峰期,并限制并发度,避免与查询任务抢资源。频率上,不是合并得越勤越好:每次合并都会带来额外的计算和快照开销,对于 T+1 批处理表,每天合并一次就够了;对于流式表,可能每隔几小时需要一次轻量合并。

几个容易踩的坑

  • 合并作业自己又造出小文件。如果输出分区数设置不当,或者数据在重写过程中发生了过多的 shuffle 溢写,合并完的文件平均大小可能仍然偏低。建议在合并任务后重新统计一次文件分布,验证是否达到预期。
  • 盲目追求大文件。文件大小不是越大越好。对点查、count 这类查询,文件块过大会造成不必要的扫描开销;对按主键更新的场景,大文件也会让索引或过滤效果变差。目标大小应该结合查询特征和底层块尺寸共同决定。
  • 合并期间不控并发。如果表格式没有 ACID 保障,例如直接对 Hive 表做 INSERT OVERWRITE,而同时业务任务还在写表,就可能覆盖新写入的数据。
  • 忽略版本文件清理。Iceberg、Delta Lake 合并后,被重写的旧数据文件可能还作为历史快照保留着,存储占用并没有下降,必须定期执行 expire_snapshots 或 VACUUM。

这些坑大多不是技术判断错误,而是对合并动作的边界理解不够。把合并当作一次性的补救,而不是持续的治理,是很多团队做完一次合并之后又快速反弹的主要原因。

写在最后

数据湖的小文件问题很难通过一次清洗永久解决,因为只要写入持续,碎裂就会重新出现。比较务实的做法,是把文件大小原则纳入到数据管道的设计和监控里:在写入时控制并行度,在表格式上引入自动合并,在底层建立文件规模的监控面板。当合并成了日常数据治理的一部分,小文件问题就不再是每一次告警都让你紧张的源头。

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

(0)
上一篇 23小时前
下一篇 9小时前

相关推荐