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');
环境截图

优化建议
一、数据读取优化:替换循环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优化策略
- 调整Shuffle分区数
默认spark.sql.shuffle.partitions为200,对于365-730GB的总数据量,该值过小会导致分区过大,引发性能问题。建议设置为:
# 按每分区100-200MB估算,或按总核心数的2-3倍设置 spark.conf.set("spark.sql.shuffle.partitions", 1000)
- 优化Shuffle内存配置
- 提升Shuffle可用内存比例:
spark.conf.set("spark.shuffle.memoryFraction", 0.3) - 开启小分区跳过合并排序:
spark.conf.set("spark.shuffle.sort.bypassMergeThreshold", 500)
- 避免不必要的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核,最大化并行度。
四、写入阶段优化
- 控制单文件大小
设置每个文件的最大记录数,避免生成过大文件:
spark.conf.set("spark.sql.files.maxRecordsPerFile", 1000000)
- 合并小分区
SQL执行后若分区数过多,写入前用coalesce合并,避免生成大量小文件:
dfResults.coalesce(365).write.partitionBy("Year", "Month", "Day").parquet(DESTINATION, mode='overwrite')
内容的提问来源于stack exchange,提问作者mateoc15
相关产品推荐
相关产品推荐

