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

PySpark程序在添加广播块阶段停滞的问题排查求助

PySpark程序在添加广播块阶段停滞的问题排查求助

我现在遇到一个PySpark的棘手问题,想求助大家帮忙排查:

我需要处理一个1-20亿条记录的超大Parquet数据集,要基于另一个较小的参考DataFrame(约41.4万条记录)的条件进行过滤,采用左连接的方式,最后把结果写入Parquet文件。

奇怪的是,当参考DataFrame为空的时候,程序能正常运行完成;但当参考DataFrame有41.4万条记录时,程序会一直卡在日志 storage.BlockManagerInfo: [dispatcher-BlockManagerMaster]: Added broadcast_4_piece0 in memory on xxx:12345 (size:4.0 MiB, free: 10 GiB) 这一步,完全没有进展。

我的核心代码如下:

from pyspark.sql.functions import col

def extract_to_df(spark, ref_db):
    columns_to_drop = ["ColA", "ColB", "ColC"]
    
    # 定义连接条件
    join_cond_1 = (col("Col1") >= col("Col3a")) & (col("Col1") >= col("Col3b"))
    join_cond_2 = (col("Col2") >= col("Col3a")) & (col("Col2") >= col("Col3b"))
    
    # 读取主数据集
    df = spark.read.parquet(folder)
    
    # 过滤+UDF处理列
    df_2 = df.filter(df["Col4"]=="abc")\
             .withColumn("Col1", udf_col(col("Col1a")))\
             .withColumn("Col2", udf_col(col("Col2a")))
    
    # 第一次左连接+列处理
    df_tmp = df_2.join(ref_db, on=join_cond_1, how="left")\
                 .drop(*columns_to_drop)\
                 .withColumnRenamed("Col5", "Col5a")
    
    # 第二次左连接得到结果
    df_results = df_tmp.join(ref_db, on=join_cond_2, how="left")\
                       .drop(*columns_to_drop)\
                       .withColumnRenamed("Col6", "Col6a")
    
    # 过滤空值后写入
    df_final_results = df_results.dropna(subset=["Col5a", "Col6a"])
    df_final_results.write.mode("overwrite").parquet(output_folder)

def main():
    ref_db = spark.read.parquet("/ref_db.parquet")
    ref_db.persist()
    extract_to_df(spark, ref_db)

if __name__ == "__main__":
    main()

注:udf_col是我自定义的UDF,用于转换Col1a和Col2a的格式。

我自己的一些疑惑点:

  • 是不是因为非等值Join的条件,导致Spark无法正常优化广播Join,反而触发了低效的执行逻辑?
  • 自定义UDF会不会引入性能瓶颈,间接导致Join阶段卡住?
  • 连续两次左连接参考DF的写法有没有问题?或者列的删除、重命名操作会不会干扰执行计划?
  • 我对ref_db做了persist,但会不会存储级别不合适,或者这个操作反而起了反作用?

希望有经验的PySpark大佬能帮忙分析下问题所在,或者给一些排查的方向,谢谢!

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.07 06:43:06