PySpark中.cache()是懒加载函数吗?缓存是否作用于最终DataFrame?
PySpark中.cache()的常见疑问解答
问题1:PySpark中的.cache()是否为懒加载函数?
是的,cache()是懒加载操作。它不会立即触发数据计算,只是在当前DataFrame对应的执行计划中添加一个“缓存标记”——告诉Spark:当这个DataFrame的计算被行动操作触发时,把计算结果缓存起来(默认存储级别为MEMORY_ONLY)。只有遇到show()、count()、write这类行动操作时,才会实际执行计算并完成缓存。
问题2:.cache()是否是一种前瞻性操作,会作用于一系列转换后生成的最终Spark DataFrame?
不是。cache()的作用范围仅限于调用它时的那个DataFrame对象,不会自动延续到后续转换生成的新DataFrame。每个转换操作(比如filter、withColumn)都会生成全新的DataFrame,这些新对象不会继承之前的缓存标记。
示例场景分析
看你给出的代码:
sdf = spark.read.format("delta").load("/path/to/my/data") sdf.cache() # 标记的是原始读取的Delta表DataFrame # 一系列转换生成新的DataFrame,覆盖了变量sdf sdf = sdf.filter("blah blah") sdf = sdf.withColumn("myvar", blah()) sdf = sdf.select("more blah") sdf.show() # 触发计算
执行show()后:
- 最终经过多次转换的
sdf不会被缓存,因为你从未给这个最终的DataFrame调用cache()。 - Spark并没有忽略最初的
cache()指令:如果执行计划中需要用到原始的未转换DataFrame(比如某些场景下的依赖链路),Spark会在计算原始DataFrame后将其缓存,但这和你最终的转换后DataFrame无关。 - 最初的
cache()不是无意义的请求,只是它的作用对象是原始DataFrame,而非后续转换后的版本。
如果想要缓存最终的结果,需要在最后一次转换完成后调用cache():
sdf = spark.read.format("delta").load("/path/to/my/data") sdf = sdf.filter("blah blah") sdf = sdf.withColumn("myvar", blah()) sdf = sdf.select("more blah") sdf.cache() # 标记最终DataFrame sdf.show() # 触发计算并缓存结果
内容的提问来源于stack exchange,提问作者pauljohn32
相关产品推荐
相关产品推荐

