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

