Spark缓存重分区GlobalTempView异常:重复执行重分区操作排查
Spark缓存重分区数据集的踩坑问题
场景与代码
我有个Spark作业要把一个数据集和多个数据集做关联,为了避免重复读HDFS的开销,我创建了GlobalTempView。原数据集在HDFS里是700个分区,但Spark默认分区数是200,我想缓存重分区后的数据集,结果没达到预期,代码如下:
Dataset<Row> accountToUserIdDataset = hdfsClient.loadDataSet("UserDataSetPath"); // 原数据700分区 accountToUserIdDataset = accountToUserIdDataset.repartition(200); // 重分区,避免每次shuffle都要调分区 accountToUserDetailsDataset.createOrReplaceGlobalTempView("USERID_TO_DETAILS"); // 给多会话用的全局视图 accountToUserIdDataset.cache(); accountToUserDetailsDataset.count(); // 触发缓存物化 // 在不同Spark会话里关联这个数据集
执行后的异常现象
- 多轮关联时,HDFS读取阶段(700分区的那个)确实被跳过了,但重分区到200的操作每次都重复执行;
- Spark存储页面显示
accountToUserDetailsDataset确实是以200分区缓存的; - 作业截图里,HDFS读取的stage7被跳过,但重分区的stage28没被跳过,我本来以为这个阶段也该被跳过。
我的疑问
这是我对Spark缓存的理解错了,还是本来就可以缓存重分区后的数据集?另外看到相关讨论提到,要是关联时只用到数据集的部分列,Spark会重新洗牌,会不会是这个原因?
问题分析与解决方案
核心bug:缓存对象与全局视图未绑定
你代码里有个关键错误:你缓存的是accountToUserIdDataset对象,但跨会话关联时用的是USERID_TO_DETAILS全局视图。全局视图存储的是数据集的逻辑执行计划,不是你已经物化好的缓存数据。当其他会话查询这个视图时,Spark会重新解析执行计划,自然会重复执行重分区步骤,不会复用你之前缓存的实例。
重分区重复执行的两个原因
- 视图不关联缓存:全局视图本质是「逻辑计划的快照」,不是指向缓存数据的引用。跨会话查询时,Spark会从头构建物理计划,重分区操作会被重新执行,除非你基于缓存后的数据集创建视图。
- 列裁剪的优化干扰:如果关联时只用到数据集的部分列,Spark优化器会计算成本:「只读取需要的列+重新分区」的开销,可能比「读取缓存里的全量列」更低,这时候它会直接跳过缓存,重新执行重分区(因为缓存的全量列里有你用不上的部分,重新计算反而更省资源)。
解决办法
- 直接传递缓存的数据集实例:如果允许在会话间传递对象,别用全局视图,直接把
accountToUserIdDataset传到其他会话使用,这样肯定会复用缓存。 - 先缓存物化,再创建视图:调整代码顺序,先把重分区后的数据集缓存并物化,再基于它创建全局视图:
这样全局视图的逻辑计划会指向缓存的物化数据,跨会话查询时直接读缓存,重分区和HDFS读取都会被跳过。Dataset<Row> accountToUserIdDataset = hdfsClient.loadDataSet("UserDataSetPath"); accountToUserIdDataset = accountToUserIdDataset.repartition(200); accountToUserIdDataset.cache(); accountToUserIdDataset.count(); // 先把缓存落地 accountToUserIdDataset.createOrReplaceGlobalTempView("USERID_TO_DETAILS"); // 基于缓存数据建视图 - 禁用列裁剪(不推荐):如果你确定全量列缓存更高效,可以添加配置
spark.sql.optimizer.excludedRules=org.apache.spark.sql.catalyst.optimizer.ColumnPruning关掉列裁剪,但这会影响整个Spark作业的优化效率,不到万不得已别用。
验证方法
查看跨会话查询视图的执行计划,若里面包含InMemoryTableScan说明缓存生效;如果还是有Repartition和HadoopRDD,则缓存未被复用。
内容的提问来源于stack exchange,提问作者best wishes
相关产品推荐
相关产品推荐

