好的,所以上面给出的堆栈跟踪不足以理解根本原因,但正如您所提到的,您使用的连接很可能是因为这个原因而发生的。我在加入时遇到了同样的问题,如果您深入了解堆栈跟踪,您会看到类似 -
+- *HashAggregate(keys=[], functions=[partial_count(1)], output=[count#73300L])
+- *Project
+- *BroadcastHashJoin
...
Caused by: java.util.concurrent.TimeoutException: Futures timed out after [300 seconds]
这提示了它失败的原因,Spark 尝试使用“Broadcast Hash Join”加入,它具有超时和广播大小阈值,其中任何一个都会导致上述错误。根据基础错误修复此问题 -
增加“spark.sql.broadcastTimeout”,默认为300秒-
spark = SparkSession
.builder
.appName("AppName")
.config("spark.sql.broadcastTimeout", "1800")
.getOrCreate()
或者增加广播阈值,默认是10MB -
spark = SparkSession
.builder
.appName("AppName")
.config("spark.sql.autoBroadcastJoinThreshold", "20485760 ")
.getOrCreate()
或者通过将值设置为-1来禁用广播连接
spark = SparkSession
.builder
.appName("AppName")
.config("spark.sql.autoBroadcastJoinThreshold", "-1")
.getOrCreate()
更多细节可以在这里找到 - https://spark.apache.org/docs/latest/sql-performance-tuning.html