You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

DataFrame缓存时长及跨函数缓存作用域技术咨询

1. DataFrame的缓存时长是多久?

Spark里DataFrame的缓存时长没有固定的“过期时间”,主要看两种情况:

  • 只要你的Spark应用还在运行,缓存的DataFrame会一直留在内存(或磁盘,取决于你设置的缓存级别)里,直到你**主动调用df.unpersist()**来释放缓存,或者Spark因为内存不足,按照LRU(最近最少使用)策略自动淘汰掉一部分缓存数据——这时候如果后续再用到这个DataFrame,就得重新计算了。
  • 当整个Spark应用结束,SparkContext被销毁时,所有缓存的数据都会被清除。
2. 缓存DataFrame在其他函数中的作用域问题

咱们先拆解你的示例代码,再一步步说清楚:

首先看你的mydf()函数,每次调用它都会执行这几步:执行SQL生成新的DataFrame,调用cache()标记它需要缓存,然后返回这个DataFrame。这里要注意:cache()只是标记,实际的缓存写入是在第一次触发Action操作时才会执行(比如join()、show()这些都是Action)。

回到你的核心问题:joinWithDept()里调用mydf().join(...)会不会用缓存的数据集?
答案是:第一次调用会触发缓存,后续调用如果缓存还在就会复用,但这里有个容易踩的坑:你每次调用mydf()都会生成一个新的DataFrame对象,但只要这个新对象的逻辑执行计划(也就是背后的SQL计算逻辑)和之前缓存的那个一致,Spark会自动识别并复用缓存数据,不会重新计算。不过这种写法不够直观,很容易让人误以为每次都在创建新的未缓存的DataFrame。

更清晰、更稳妥的写法是把缓存的DataFrame抽出来作为一个共享变量,比如:

// 把缓存的DataFrame定义在函数外面,让所有需要的函数都能直接复用
val cachedEmpDF = sparkSession.sql("select * from emp").cache()

def joinWithDept(): Unit = {
  val deptdf1 = sparkSession.sql("select * from dept")
  val deptdf2 = cachedEmpDF.join(deptdf1, Seq("empid")) // 明确使用缓存好的数据集
  deptdf2.show()
}

def joinWithLocation(): Unit = {
  val locdf1 = sparkSession.sql("select * from location")
  val locdf2 = cachedEmpDF.join(locdf1, Seq("empid")) // 直接复用缓存,不用重新计算
  locdf2.show()
}

最后再总结一下作用域的关键点:

  • 缓存的生命周期和SparkSession绑定,和函数的调用栈无关——哪怕创建缓存DataFrame的函数已经执行完,只要SparkSession还在,缓存没被清除/淘汰,后续任何能拿到对应DataFrame引用(或者逻辑计划一致的DataFrame)的代码都能复用缓存。
  • 变量的作用域决定了你能不能拿到DataFrame的引用:如果缓存的DataFrame是函数内部的局部变量,只有在函数内部或者把它返回出去之后,其他代码才能拿到引用复用缓存;如果是全局/类级别的变量,所有函数都能直接用。

内容的提问来源于stack exchange,提问作者user5944508

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.21 07:30:39