Spark Dataset unpersist异常:调用data.unpersist()为何删除extension缓存?
首先直接给结论:这种情况不是Spark的预期行为,大概率是你的代码写法或者对Spark缓存机制的理解有偏差导致的。下面详细拆解原因和解决方案:
为什么会出现这个问题?
先看你代码里的一个关键细节:
data.join(df1, "key") //etc, more transformations data.cache(); // 用于保存后避免重计算
Spark的Dataset是不可变的——所有转换操作(比如join)都会返回一个全新的Dataset实例,而你这里没有把join后的结果赋值给任何变量,导致后续缓存的data还是最初从spark.read(...)得到的原始Dataset,而不是经过转换后的版本。这可能是第一个核心误区。
另外,关于unpersist()的作用:它只会删除当前调用方法的Dataset实例对应的缓存数据,不会主动删除其他Dataset的缓存——除非这些Dataset的缓存依赖于父Dataset的缓存?其实不是的,每个Dataset的缓存都是独立存储的,缓存的是该Dataset计算后的最终结果,和父Dataset的缓存没有直接关联。那你遇到的情况可能是:
- 你调用
data.unpersist()时,extension的缓存还没完全写入存储(比如extension.count()触发的缓存还在异步执行),导致后续操作误以为缓存被删除; - 代码中
extension的血统和data的实例有意外的关联(比如转换操作没有生成新的Dataset,但join肯定会生成新实例,这种概率很低)。
怎么解决?
要确保unpersist旧Dataset时不影响后续链式Dataset的缓存,你需要做到这几点:
1. 正确处理Dataset的不可变性,保存转换结果
所有转换操作的结果必须赋值给新变量,确保你缓存、操作的都是经过转换后的正确Dataset:
val data = spark.read(...) // 将join后的结果赋值给新变量 val transformedData = data.join(df1, "key") // 这里可以添加其他转换操作 transformedData.cache() // 缓存的是转换后的数据集 transformedData.write.parquet() // 基于转换后的数据集生成extension val extension = transformedData.join(...) // 添加你的join、select等操作 extension.cache() // 缓存extension自己的计算结果 extension.count() // 触发缓存计算 // 现在unpersist旧数据集,不会影响extension的缓存 transformedData.unpersist() // 后续操作正常使用extension即可 extension.write.csv() extension.groupBy("key").agg(...) extension.write.parquet()
2. 确保extension的缓存完全完成后再unpersist
虽然count()会同步触发缓存计算,但如果你担心异步写入的问题,可以明确指定存储级别并确保计算完成,比如:
import org.apache.spark.storage.StorageLevel extension.persist(StorageLevel.MEMORY_ONLY) // 明确存储级别 extension.count() // 同步完成计算和缓存 data.unpersist() // 此时extension的缓存已经稳定存在
3. 明确区分不同Dataset的缓存实例
记住:每个Dataset的缓存都是独立的,只要extension是一个独立的实例(通过转换操作生成),并且已经完成缓存计算,那么data.unpersist()只会删除data的缓存,不会碰extension的缓存。
另外补充:Spark的Dataset缓存机制和RDD是完全一致的,都是基于实例的缓存,血统只是用来在缓存丢失时重新计算,缓存的数据本身是独立存储的——所以只要你确保extension已经被计算并缓存,它就和父Dataset的缓存状态无关。
内容的提问来源于stack exchange,提问作者T. Gawęda

