Spark ETL复用变量时persist与unpersist的工作机制及相关疑问
Spark中复用变量名时persist()与unpersist()的行为解析
先看你给出的示例代码,核心问题是变量复用导致的DataFrame对象引用丢失,下面逐个解答你的疑问:
1. 究竟会持久化和取消持久化什么内容?
- 持久化的是第一次赋值的原始DataFrame(即读JSON生成的那个):调用
df.persist()时,变量df指向的是这个初始DF,此时Spark会标记该DF需要被缓存。 - 后续每次执行
df = df.withColumn(...)这类操作,都会生成全新的DataFrame对象,变量df只是转而指向新对象——但最初被标记持久化的DF并没有被覆盖,只是失去了直接的变量引用。 - 最后一行的
df.unpersist()毫无作用:此时df指向的是最后一次转换后的新DF,这个DF从未被调用过persist(),所以取消持久化操作没有目标。而最初被持久化的原始DF,会因为没有引用,等待Spark的自动缓存回收机制处理,但不会被这次unpersist()清理。
2. persist()和unpersist()方法属于Action还是Transformation?
两者都不属于,它们是缓存控制类方法:
persist():只是给目标DF/RDD打上“需要缓存”的标记,不会触发任何计算。只有当后续遇到Action类操作(比如count()、show())时,才会实际执行计算并将结果存入指定缓存介质(内存、磁盘等)。unpersist():会立刻触发缓存清理操作,直接移除目标DF/RDD的缓存数据,不需要等待Action触发。
3. 在每次persist()之后,使用spark.catalog.clearCache()和unpersist()哪个更合适?
优先用unpersist(),原因如下:
spark.catalog.clearCache()是全局清理操作,会清空整个Spark应用中所有被缓存的DF/RDD数据,完全没必要这么做,还会破坏其他依赖缓存的任务性能。unpersist()是精准清理,只针对指定的DF/RDD进行缓存移除,不会影响其他缓存数据,效率更高也更安全。- 注意:如果要复用变量名,一定要保留住被持久化的DF引用,比如改成:
df = spark.read.json(...) original_df = df original_df.persist() # 后续转换操作复用df变量 df = df.withColumn(...) df = df.another_transformation() # 清理时调用原始DF的unpersist() original_df.unpersist()
内容的提问来源于stack exchange,提问作者Vlad Vlad
相关产品推荐
相关产品推荐

