You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

使用reduce合并300个PySpark DataFrame后行计数极慢的排查与解决

解决方案:合并后调整分区解决性能问题

核心原因

用reduce结合unionAll逐个合并300个DataFrame时,会保留每个小DataFrame的原有分区。如果每个小DF有默认的200个分区,合并后总分区数会达到300*200=6000个,2万行数据分散到几千个分区,每个分区仅3-4行。Spark调度几千个小Task的开销远大于实际计算时间,导致行计数等操作耗时极长。

便捷解决办法

方法一:合并后直接压缩分区(推荐)

合并完成后立即用coalesce减少分区数,该操作不会触发Shuffle,性能开销极低:

from functools import reduce
from pyspark.sql import DataFrame

Lst = [df1, df2, ..., df300]
# 合并后压缩到4个分区(可根据集群情况调整为2-8)
Final_df = reduce(DataFrame.unionAll, Lst).coalesce(4)

注:coalesce适合减少分区,若需要增加分区或重新均匀分布数据,可改用repartition(4),但会触发Shuffle,小数据量下差异不大。

方法二:合并前先统一小DF的分区

如果每个小DF本身分区数过多,可先将每个小DF压缩到1个分区再合并,进一步降低中间调度压力:

# 先把每个小DF压缩到1个分区
Lst_optimized = [df.coalesce(1) for df in Lst]
# 合并后再调整到合适分区数
Final_df = reduce(DataFrame.unionAll, Lst_optimized).coalesce(4)

验证分区数

可以先查看合并后的分区数确认问题:

print(Final_df.rdd.getNumPartitions())

正常情况下,2万行数据设置2-8个分区即可满足性能需求。

内容的提问来源于stack exchange,提问作者Gaurav Adakmol

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 04:32:14