【发布时间】:2016-08-22 16:32:45
【问题描述】:
我是一个 spark 应用程序,有几个点我想保持当前状态。这通常是在一大步之后,或者缓存我想多次使用的状态。看来,当我第二次在我的数据帧上调用缓存时,一个新副本被缓存到内存中。在我的应用程序中,这会在扩展时导致内存问题。即使在我当前的测试中,给定的数据帧最大约为 100 MB,但中间结果的累积大小会超出执行程序上分配的内存。请参阅下面的一个小示例来显示此行为。
cache_test.py:
from pyspark import SparkContext, HiveContext
spark_context = SparkContext(appName='cache_test')
hive_context = HiveContext(spark_context)
df = (hive_context.read
.format('com.databricks.spark.csv')
.load('simple_data.csv')
)
df.cache()
df.show()
df = df.withColumn('C1+C2', df['C1'] + df['C2'])
df.cache()
df.show()
spark_context.stop()
simple_data.csv:
1,2,3
4,5,6
7,8,9
查看应用程序 UI,除了带有新列的数据框之外,还有原始数据框的副本。我可以通过在 withColumn 行之前调用 df.unpersist() 来删除原始副本。这是删除缓存中间结果的推荐方法吗(即在每个 cache() 之前调用 unpersist)。
另外,是否可以清除所有缓存的对象。在我的应用程序中,有一些自然断点,我可以简单地清除所有内存,然后转到下一个文件。我想在不为每个输入文件创建新的 spark 应用程序的情况下执行此操作。
提前谢谢你!
【问题讨论】:
标签: python caching apache-spark pyspark apache-spark-sql