Spark缓存DataFrame后仍重计算?Scala应用技术咨询
Spark缓存DataFrame后仍触发重新计算的问题
我正在开发一个基于Scala的Spark应用,包含6个函数,每个函数接收两个DataFrame作为输入,处理后输出一个结果DataFrame。我对每个函数的输出DataFrame进行缓存,并调用count()方法将其物化。在第6个函数中,我使用前5个输出的DataFrame执行Join操作,但发现Spark并未使用缓存的DataFrame,而是从第1个函数开始重新计算所需的DataFrame(根据查询计划和DAG)。请问这是预期行为吗?
代码示例
//functions 1 to 5 to generate joinedLeaseDF5 //Now after getting joinedLeaseDF5, caching it joinedLeaseDF5.cache() joinedLeaseDF5.count() val joinedLeaseDF6 = joinedLeaseDF5 .join(credDF, joinedLeaseDF5("cust_id") === credDF("cust_id"), "left_outer")
此时我能看到计算joinedLeaseDF5的所有逻辑正在执行。
回答
这不是预期行为,正常调用cache()后执行count()物化缓存,后续操作应该直接复用缓存的DataFrame,不会触发从头计算。出现这种情况通常有以下几种原因:
- 缓存实例未被正确引用:确保
joinedLeaseDF6使用的joinedLeaseDF5是你执行过cache()和count()的那个实例,没有被后续代码重新赋值为未缓存的新DataFrame。 - 缓存被自动清除:Spark会基于LRU策略在内存不足时自动清除缓存数据。可以通过Spark UI的「Storage」页面查看
joinedLeaseDF5的缓存状态,确认它是否还在缓存中。 - 转换操作未重新缓存:如果在
count()之后,你对joinedLeaseDF5执行了额外的转换(比如select、filter)但没有重新缓存转换后的DataFrame,后续使用的是未缓存的新实例,就会触发重新计算。 - 查询优化器的特殊决策:极少数情况下,Catalyst优化器会判断重新计算比读取缓存的开销更低(比如数据量极小),从而跳过缓存,但这种场景非常罕见。
排查建议
- 打开Spark UI的「Storage」标签,检查
joinedLeaseDF5的缓存状态:是否有对应的存储级别、缓存大小、是否被标记为「Cached」。 - 在定义
joinedLeaseDF6前,打印joinedLeaseDF5.queryExecution.simpleString,查看逻辑计划中是否包含缓存标记。 - 尝试用
persist(StorageLevel.MEMORY_AND_DISK)替代cache(),提升缓存的持久化优先级,降低被自动清除的概率。 - 确认在
count()之后没有对joinedLeaseDF5做任何未重新缓存的转换操作。
内容的提问来源于stack exchange,提问作者Karthik
相关产品推荐
相关产品推荐

