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

PySpark写入Parquet性能不佳,遇Shuffle问题求资源配置与调优建议

PySpark Shuffle优化与资源配置建议

问题背景

我能正常运行PySpark,但对性能相关文档理解不清。现有按日期分区的源数据,每日对应1-2GB的Parquet文件,本次需读取365天的文件,通过unionAll拼接DataFrame,选取字段后执行SQL,最终写入新的分区Parquet文件。目前遇到Shuffle性能错误,可用的内存、Executors、Cores资源充足,但不清楚该场景下的最优配置组合。附上代码与环境截图,求Shuffle优化、Executors及Cores配置的建议。

现有代码

# 源文件按年、月、日分区 - 每个分区对应一个1-2GB的文件
SOURCE ="abfss://...my_files/Year={}/Month={}/Day={}";
DESTINATION = "abfss://...another_place";

# 获取过去DAYS_TO_LOAD天的日期列表,用于遍历需要读取的源文件
DAYS_TO_LOAD = 365;
today = datetime.today();
start_date = today - timedelta(days = DAYS_TO_LOAD);
range_tuple = date_range_tuple(start_date, today);

print(range_tuple); # 示例输出:((2024, 5, 19), (2024, 5, 20), (2024, 5, 21), (2024, 5, 22), ...)

# 创建DataFrame并通过unionAll逐天追加数据
df = None;

for year, month, day in range_tuple:
    if year >= 2023:
        # right()用于截取后两位数字,实现月份、日期的补零对齐
        path = SOURCE.format(str(year), right("0" + str(month),2),right("0" + str(day),2)); 

    print(path);
    
    
    if df == None: # 初始化DataFrame
        df = spark.read.load(path);
    else: # 追加数据到已有DataFrame
        df = df.unionAll(spark.read.load(path));
        
df\
    .select(\
        "约20个字段",\
        ...
    )\
    .createOrReplaceTempView("my_source");
    
# 输出分区的年/月/日字段与源数据的日期分区不同
df = df.withColumn('utc_date', from_unixtime(df.event_time_epoch / 1000, "yyyy-MM-dd"))\

print(df.rdd.getNumPartitions());

# 尝试过coalesce和repartition,但效果不明显
# df = df.coalesce(DAYS_TO_LOAD);
# print(df.rdd.getNumPartitions());
df = df.repartition("utc_date");
print(df.rdd.getNumPartitions());

dfResults = spark.sql(SOURCE_QUERY); # 引用上面的"my_source"临时视图

dfResults\
    .write\
    .partitionBy("Year", "Month", "Day")\
    .parquet(DESTINATION, mode = 'overwrite');

环境截图

Spark环境配置

优化建议

一、数据读取优化:替换循环unionAll为批量读取

当前循环unionAll的方式会不断生成复杂的执行计划,且单文件读取会增加调度开销。直接利用Spark的分区发现能力批量读取:

# 直接读取整个时间范围的分区,Spark自动识别Year/Month/Day分区
df = spark.read.load("abfss://...my_files/")\
    .filter((col("Year") >= 2023) & 
            (to_date(concat_ws("-", col("Year"), col("Month"), col("Day"))) >= start_date) & 
            (to_date(concat_ws("-", col("Year"), col("Month"), col("Day"))) <= today))

如果源路径是Hive风格分区(路径含Year=xxxx键值对),Spark会自动将分区列加载为DataFrame字段,无需手动拼接路径。

二、Shuffle优化策略

  1. 调整Shuffle分区数
    默认spark.sql.shuffle.partitions为200,对于365-730GB的总数据量,该值过小会导致分区过大,引发性能问题。建议设置为:
# 按每分区100-200MB估算,或按总核心数的2-3倍设置
spark.conf.set("spark.sql.shuffle.partitions", 1000)
  1. 优化Shuffle内存配置
  • 提升Shuffle可用内存比例:
    spark.conf.set("spark.shuffle.memoryFraction", 0.3)
    
  • 开启小分区跳过合并排序:
    spark.conf.set("spark.shuffle.sort.bypassMergeThreshold", 500)
    
  1. 避免不必要的Shuffle
  • 若SQL中已有按utc_date或输出分区字段的聚合/分组,将repartition放到SQL执行之后,或直接利用write.partitionBy的分区逻辑,避免重复Shuffle。
  • 用broadcast join优化小表关联,减少Shuffle数据量:
    SELECT /*+ BROADCAST(small_table) */ * FROM my_source s JOIN small_table st ON s.id = st.id
    

三、Executors与Cores配置

结合环境截图(Driver内存64GB,最大Executor内存64GB,最大核心数8),建议:

  • Executor数量:按集群总核心数计算,比如集群有64核,每个Executor用8核,可设置8个Executor。
  • Executor内存:保持64GB不变,调整堆外内存避免溢出:
    spark.conf.set("spark.executor.memoryOverhead", "8g")
    
  • 核心数配置:每个Executor用8核(spark.executor.cores=8),设置spark.task.cpus=1,确保单任务占用1核,最大化并行度。

四、写入阶段优化

  1. 控制单文件大小
    设置每个文件的最大记录数,避免生成过大文件:
spark.conf.set("spark.sql.files.maxRecordsPerFile", 1000000)
  1. 合并小分区
    SQL执行后若分区数过多,写入前用coalesce合并,避免生成大量小文件:
dfResults.coalesce(365).write.partitionBy("Year", "Month", "Day").parquet(DESTINATION, mode='overwrite')

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 02:42:34