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

Docker中Spark2.3.0执行PySpark DataFrame内连接时遇awaitResult错误

解决Spark 2.3.0中小数据集内连接后的BroadcastExchangeExec错误

首先咱们得抓住问题的核心:你的两个DataFrame连接键的类型不匹配——df1的id是bigint,而df2的id是decimal(12,0)。虽然都是数值类型,但Spark 2.3.0在处理这种隐式类型转换的连接时,尤其是广播交换阶段,很容易触发内部执行异常。这跟广播超时完全没关系,所以你设置spark.sql.broadcastTimeout=-1自然解决不了问题。

直接有效的解决方案:显式统一连接键的类型

你只需要把其中一个DataFrame的id字段类型转换为和另一个一致,再执行连接即可,有两种方式可选:

方式1:把df2的decimal类型id转成bigint

# 转换df2的id字段类型
df2_cast = df2.withColumn("id", df2["id"].cast("bigint"))
# 执行内连接
df3 = df1.join(df2_cast, df1.id == df2_cast.id, "inner")
# 现在再执行show就正常了
df3.show(5)

方式2:把df1的bigint类型id转成decimal(12,0)

# 转换df1的id字段类型
df1_cast = df1.withColumn("id", df1["id"].cast("decimal(12,0)"))
# 执行内连接
df3 = df1_cast.join(df2, df1_cast.id == df2.id, "inner")
# 验证结果
df3.show(5)

为什么设置广播超时没用?

从你的报错栈可以看到,错误是org.apache.spark.SparkException: Exception thrown in awaitResult,发生在BroadcastExchangeExec.doExecuteBroadcast阶段——这不是广播超时(超时会有明确的TimeoutException提示),而是类型不匹配导致Spark在序列化/反序列化广播数据时出错,所以调整超时参数完全不对症。

额外提示

Spark 2.3.0是比较老的版本,这类类型处理的bug在后续版本(比如2.4+)中已经有修复。如果后续有条件升级Spark版本,也能避免这类问题。

内容的提问来源于stack exchange,提问作者Ajaykishan

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.28 07:29:47