PySpark连接过滤DataFrame时的报错排查及解决
PySpark DataFrame连接过滤报错排查与解决
场景与需求
- 操作对象:两个列结构完全一致的DataFrame,
df_dsp_f(约100,000条记录)、df_slv_f(约99,800条记录) - 计算环境:Databricks无服务器集群
- 核心需求:筛选仅存在于
df_dsp_f中的记录,为这些记录新增D_flag列并标记为"X"
初始实现代码
##Repartitioning on a specific column## df_dsp_f_p = df_dsp_f.repartition(2, "col_5") df_slv_f_p = df_slv_f.repartition(2, "col_5") ##Join and then select only records with X## df_del = df_dsp_f_p.alias("df1").join(df_slv_f_p.alias("df2"),\ on = (F.col("df1.col_1") == F.col("df2.col_1")) & \ (F.col("df1.col_2") == F.col("df2.col_2")), \ how = "left").\ select(F.when(F.col("df2.col_1").isNull() | F.col("df2.col_2").isNull(), "X").\ otherwise(F.lit("")).alias("D_flag")).filter(F.col("D_flag") == "X")
执行报错信息
执行df_del.count()时触发如下错误:
Job aborted due to stage failure: Task 0 in stage 2324.1 failed 4 times, most recent failure: Lost task 0.3 in stage 2324.1 (TID 4180) Reason: Command exited with code 134, sigabrt
后续优化尝试
已检查两个DataFrame的基数并针对性设置重分区,但问题未解决。随后简化连接逻辑,改用left_anti连接:
df_del = df_dsp_f_p.alias("df1").join(df_slv_f.alias("df2"), (F.col("df1.col_1") == F.col("df2.col_1")) & (F.col("df1.col_2") == F.col("df2.col_2")), how = "left_anti") df_del = df_del.withColumn("D_flag",F.lit("X"))
执行df_del.show()时仍出现相同报错。
问题根源与解决
最终确认问题源于Databricks无服务器集群性能不足,无法支撑该连接操作。更换集群后,问题彻底解决。
内容的提问来源于stack exchange,提问作者pythondumb
相关产品推荐
相关产品推荐

