Spark 3.2 CacheManager.lookupCachedData无法识别已缓存数据问题
问题背景
测试环境为Spark 3.2 + spark-cassandra-connector 3.0.1,作业部署在Spark Job Server,使用长驻运行的Spark Context,通过多次连续执行同一作业验证Spark Cache Manager的缓存行为。
测试代码如下:
val manager = spark.sharedState.cacheManager val DF = collectData.retrieveDataFromCass(spark) // 从Cassandra成功加载2000行数据 val testCachedData = if (manager.lookupCachedData(DF.queryExecution.logical).isEmpty) 0 else 1 DF.createOrReplaceTempView(tempName1) spark.sqlContext.cacheTable(tempName1) DF.count() // 触发action,落地缓存 testCachedData
代码执行逻辑:
- 获取当前Spark会话的CacheManager实例
- 从Cassandra加载数据生成2000行规模的DataFrame
- 调用
manager.lookupCachedData传入当前DataFrame的逻辑计划判断缓存状态,未命中则testCachedData赋值为0,命中则赋值为1 - 创建临时视图,执行
cacheTable缓存表,调用count动作触发实际缓存计算,最终返回testCachedData结果
预期结果:首次执行作业返回0,后续重复执行同一作业返回1
实际结果:所有作业执行均返回0,即lookupCachedData始终判断无对应缓存,但Spark UI的STORAGE页面可明确看到对应缓存数据已存在。
核心疑问:同一长驻Spark Context(同个Spark应用)内,为何CacheManager无法识别已存在的缓存数据?
根因分析
出现该现象的核心原因有两点:
- 缓存绑定的逻辑计划和查询传入的逻辑计划不匹配
cacheTable(tempName1)缓存的是临时视图tempName1对应的CatalogRelation类型逻辑计划,而每次作业中通过retrieveDataFromCass新生成的DataFrame,对应的是直接读取Cassandra的CassandraSourceRelation逻辑计划,二者不属于同一个计划节点,无法完成匹配。 lookupCachedData是Spark内部API,匹配规则不做通用语义等价判断
该方法是Spark内部用于执行计划优化阶段匹配缓存的接口,不会主动判断两个逻辑计划是否读取同一份数据源、计算逻辑是否一致,仅对已登记的缓存计划做结构化匹配。每次新生成的Cassandra读取计划没有关联临时表元数据,也没有经过Analyzer阶段的解析替换,自然匹配不到已经登记的临时表缓存条目。
另外需要注意:每次调用retrieveDataFromCass都会生成全新的逻辑计划实例,哪怕读取逻辑完全一致,也不会直接和历史缓存绑定。
修正方案
- 停止直接调用内部API
CacheManager.lookupCachedData判断缓存状态,改用公开稳定APIspark.catalog.isCached(tempName1)判断临时表缓存状态,该接口直接读取Catalog元数据判断,不存在匹配偏差。 - 后续作业复用时,不要每次重新从Cassandra读取生成初始DataFrame,直接通过
spark.table(tempName1)读取临时视图生成DataFrame,此时生成的逻辑计划经过Analyzer解析后会自动匹配已有缓存条目,实现缓存复用。 - 若必须从DataFrame层面判断缓存命中,需先对DataFrame调用
queryExecution.assertAnalyzed()触发Analyze阶段,待计划中的数据源节点替换为可匹配缓存的节点后,再调用内部查询接口,该方式依赖Spark内部实现,不推荐在生产逻辑中使用。
内容的提问来源于stack exchange,提问作者Rand Abu Salim
相关产品推荐
相关产品推荐

