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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 10:47:04