Delta Lake执行Merge操作时单分区生成过多文件的优化方案
Delta表Merge产生大量小文件问题优化
问题场景
日常通过Merge操作增量摄入新数据的表,正在从ORC格式迁移至Delta格式,执行如下Merge逻辑时遇到问题:
DeltaTable .forPath(sparkSession, deltaTablePath) .alias(SOURCE) .merge( rawDf.alias(RAW), joinClause // 关联逻辑使用主键,尽可能使用分区键,此处不展开 ) .whenNotMatched().insertAll() .whenMatched(dedupConditionUpdate).updateAll() .whenMatched(dedupConditionDelete).delete() .execute()
Merge执行完成后,每个受影响的分区都会生成数百个新文件。由于每日执行一次数据摄入,该问题会导致后续每次Merge操作的运行速度持续衰减。
当前组件版本
- Spark:2.4
- Delta Lake:0.6.1
核心疑问
- Delta是否支持在数据写入保存前执行重分区操作?
- 是否存在其他可优化该问题的方案?
解答
关于写入前重分区的支持说明
Delta 0.6.1版本没有在Merge事务流程中提供内置的写入前自动重分区配置项,但可以通过源端预处理实现写入前的数据重分布,从源头减少小文件生成:对参与Merge的源DataFrame(即代码中的rawDf),提前按照目标表的分区键+关联主键做repartition,避免源数据倾斜导致的写入文件碎片化:
// 按照目标表分区字段、主键重分区,分区数根据单批次数据量调整 val repartitionedRawDf = rawDf.repartition(24, $"partition_col", $"pk_col")
注意不要使用coalesce做重分区,会导致上游计算并行度下降,拖慢Merge整体执行速度。
可落地的优化方案
- 调整Shuffle并行度配置
Spark 2.4 + Delta 0.6.1 版本中,Merge写入阶段的文件数默认和spark.sql.shuffle.partitions参数值强相关,该参数默认值为200。如果单分区每日增量数据量不大,保留默认值会导致每个Shuffle分区仅写出KB/MB级的小文件,可根据单批次数据量下调该参数值:单分区增量在10G以内时,设置为2050即可,让每个Shuffle分区写出的文件大小保持在128MB1GB的合理区间。 - 分区级定期执行Optimize合并小文件
每日Merge任务执行完成后,针对当日受影响的分区执行Delta内置的Optimize命令,将碎片化小文件合并为默认1GB大小的标准文件,不要执行全表Optimize浪费计算资源:
如果后续Merge的关联性能较差,可以在Optimize时追加ZORDER BY主键的配置,进一步提升文件内数据有序性,加快Merge阶段的主键匹配速度。spark.sql(s"OPTIMIZE delta.`$deltaTablePath` WHERE partition_col = 'current_dt_part'") - 优化Merge条件避免全分区重写
检查dedupConditionUpdate、dedupConditionDelete两个匹配条件,如果条件中未携带分区字段,Delta会默认重写匹配到的整个分区的所有文件,产生大量冗余输出。需要将分区过滤条件追加到两个whenMatched的判断逻辑中,让Delta仅扫描、重写分区内真正匹配到行的文件,而非全分区覆盖重写。 - 版本升级(条件允许时)
当前使用的Delta 0.6.1属于非常早期的版本,Delta 0.8及后续版本针对Merge逻辑做了大量核心优化,包括动态分区裁剪、写入时自动按分区分布做Shuffle、小文件自动合并等能力,能从引擎层面大幅减少Merge产生的小文件数量。如果环境允许升级到Spark 3.x + Delta 2.x以上版本,该类小文件问题会得到明显改善。
注意:不要在Merge执行过程中直接修改目标表路径下的文件,会导致Delta事务日志不一致。所有重分区、文件合并操作,要么在源数据预处理阶段完成,要么在Merge事务提交后通过Delta官方提供的命令执行。
内容的提问来源于stack exchange,提问作者Ismail H
相关产品推荐
相关产品推荐

