Spark中cache()是否触发collect()至Driver?求助排查
首先直接点出核心原因:你的cache()本身不会主动把数据拉到Driver,但你的查询计划中包含了Spark自动触发的广播连接(Broadcast Join),而广播连接的实现必须先将被广播的数据集全量收集到Driver,这才是触发错误的根源。
具体错误溯源
从你的错误日志栈可以看到关键执行路径:
BroadcastExchangeExec.doExecuteBroadcast -> RDD.collect -> SparkException: Total size of serialized results exceeds spark.driver.maxResultSize
当Spark优化器判断JOIN中的某一张表足够小时,会自动选择广播连接策略:它会先把这张表的所有数据collect()到Driver节点,然后再将数据分发给所有Executor节点,这样每个Executor都能在本地完成JOIN操作,避免Shuffle开销。但如果这张表的实际大小超过了spark.driver.maxResultSize(你的配置是1024MB),就会触发这个报错。
那为什么cache()会和这个错误关联?因为cache()是懒加载的转换操作,当你后续调用writeParquet和writeTsv这两个action时,会触发整个查询计划的执行——包括广播连接的collect步骤。错误栈显示在cache()处,是因为Spark在执行持久化(cache)前需要先计算出完整的数据集,这个计算过程就包含了广播连接的执行。
你的场景分析
在你的LEFT JOIN查询中,Spark大概率把visits表(也就是你通过aggregateRows生成的聚合结果)判定为“小表”,选择了广播连接。但实际这个表的序列化后大小达到了1076MB,远超spark.driver.maxResultSize的限制,所以在广播时触发了错误。
你提到的cache()是将数据持久化在Executor节点的认知是对的,但前提是数据已经被正确计算出来——而计算过程中的广播步骤必须先把数据拉到Driver,这就和你的预期产生了冲突。
解决方案
针对这个问题,你可以从以下几个方向入手:
禁用或调整自动广播阈值
如果visits表实际并不小,修改Spark配置spark.sql.autoBroadcastJoinThreshold:- 设置为
-1可以完全禁用自动广播连接,让Spark选择SortMergeJoin(适合大表JOIN); - 或者调小阈值(比如默认是10MB,可根据实际情况调整),避免大表被误判为小表触发广播。
你可以在代码中设置:
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")- 设置为
调大Driver结果大小限制
如果确实需要使用广播连接,且Driver节点有足够内存,可以调大spark.driver.maxResultSize:spark.conf.set("spark.driver.maxResultSize", "2g")注意:这个值不能超过Driver的堆内存(
spark.driver.memory),否则会引发Driver OOM。修正groupBy的逻辑错误
你提到的groupBy不保留排序的问题确实存在:orderBy后执行groupBy,分组内的顺序无法被保证,所以first()取到的结果是不可预期的。正确的做法是使用窗口函数来获取每个分组的首行:import org.apache.spark.sql.expressions.Window def aggregateRows: sql.DataFrame = { val windowSpec = Window.partitionBy(groupBys.head, groupBys.tail: _*) .orderBy("headerTimestamp") projected .withColumn("rn", row_number().over(windowSpec)) .filter("rn = 1") .select( groupBys: _*, col("accountState"), col("userId"), col("subaffiliateId"), col("clientPlatform"), col("localTimestamp"), col("page").alias("firstPage") ) }这个逻辑能确保每个分组获取到
headerTimestamp最早的那一行数据。
内容的提问来源于stack exchange,提问作者hiroprotagonist

