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
相关产品推荐
相关产品推荐

