Spark按date和country分区写入Delta时文件数量与任务数不符的疑问
你好,我来帮你理清楚这个问题——你的理解「任务数=分区数=文件数」其实是个简化的理想情况,但实际写入Delta时,因为Spark的shuffle机制和partitionBy的工作逻辑,两者并不会严格对应,具体原因如下:
Shuffle分区与输出分区并非一一映射
你用的partitionBy("date","country")是输出层面的目录分区,而Spark显示的607个任务对应的是计算阶段的shuffle分区(由spark.sql.shuffle.partitions参数控制,这里显然这个参数被设为了607)。
执行写入流程时,Spark会先按shuffle分区做数据分发(基于哈希算法):同一个date+country组合的数据会被哈希到同一个shuffle分区,但不同的组合可能因为哈希碰撞被塞进同一个shuffle分区。每个任务处理一个shuffle分区时,会把分区内的数据按date+country拆分到对应的目录下,这就导致:- 单个任务可能会向多个
date+country目录写入文件片段 - 单个
date+country目录会收到多个任务的写入,自然就生成了多个文件
- 单个任务可能会向多个
数据分布的影响
你的输入有968个不同的country,加上1个date,理论上最多有968个date+country组合,但你的shuffle分区数是607(小于组合数),这必然导致多个组合被挤到同一个shuffle分区。反过来,如果某个date+country组合的数据量特别大,Spark可能会自动拆分这个key的分区(比如开启了数据倾斜优化),也会让该目录下的文件数量变多。Delta Lake的写入细节
Delta Lake本身有后台小文件合并的优化,但如果你的写入任务刚完成,可能还没触发这个合并;另外,写入过程中的临时文件、任务重试也可能导致目录下出现多个文件,但Delta会自动清理无效文件,最终保留的都是有效数据文件。
如果你想更精准地控制每个date+country目录下的文件数量,可以试试这些方法:
- 把
spark.sql.shuffle.partitions调整为接近或略大于date+country的组合数(比如1000),减少哈希碰撞的概率 - 写入完成后执行
OPTIMIZE <path> ZORDER BY (...)命令合并小文件 - 用
repartition("date","country")配合partitionBy,强制让每个date+country组合对应一个RDD分区,这样每个目录下只会生成一个文件(但要注意数据倾斜问题,如果某个组合数据量过大,单个文件会非常庞大,反而影响后续的查询性能)
备注:内容来源于stack exchange,提问作者Matthew

