【问题标题】:Un-persisting all dataframes in (py)spark取消持久化(pyspark)中的所有数据帧
【发布时间】: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


    【解决方案1】:

    Spark 2.x

    你可以使用Catalog.clearCache:

    from pyspark.sql import SparkSession
    
    spark = SparkSession.builder.getOrCreate
    ...
    spark.catalog.clearCache()
    

    Spark 1.x

    你可以使用SQLContext.clearCache方法

    从内存缓存中删除所有缓存的表。

    from pyspark.sql import SQLContext
    from pyspark import SparkContext
    
    sqlContext = SQLContext.getOrCreate(SparkContext.getOrCreate())
    ...
    sqlContext.clearCache()
    

    【讨论】:

    • 目前这是一个很好的解决方案,因为它允许我在合理的断点处清除完整的缓存。我将合并这一点,但我担心当我扩大规模并开始使用更大的数据集时,我的旧缓存开始开始失控。如果我想随时清除旧缓存,建议创建一个新变量(或临时变量),并显式取消保留旧对象。像:df.cache()df_new = df.withColumn('C1+C2', df['C1'] + df['C2']) ; df_new.cache() ; df.unpersist()。如果这是唯一的方法,这似乎有点麻烦......
    • 我担心我做错了什么。在我的完整应用程序中,我的作业最终会由于内存不足错误而崩溃。数据帧的每个单独副本都相当小(小于 100 MB),但缓存似乎永远存在;即使在将输出写入文件并继续下一步之后。我会看看我是否可以生成一个更小的工作示例来展示这一点。
    • 我认为您对缓存的工作方式有一个错误的印象。它可能会发生也可能不会发生,数据只能部分缓存,即使缓存了,也可以在用户不知情的情况下从缓存中逐出。
    • 感谢您的澄清。我不确定我的观察是否与您描述的行为一致。在我的测试中,除非我使用clearCache() 明确释放它们,否则丢失的小缓存会持续存在。这会导致内存不足错误。如果缓存在幕后被释放,我不希望我分配的内存饱和,即使我不使用clearCache。您是否知道为什么即使执行程序内存不足,缓存也可能不会被驱逐?
    • 我正在尝试更新 OP 中的示例以显示我所看到的行为。我在运行内存少于 500 MB 的示例时遇到问题(读取 4 kb 文件!)。我将发布一个新问题来解决这个问题,一旦我可以生成一个最小的工作示例,就会回到这个问题。
    【解决方案2】:

    我们经常使用这个

    for (id, rdd) in sc._jsc.getPersistentRDDs().items():
        rdd.unpersist()
        print("Unpersisted {} rdd".format(id))
    

    sc 是 sparkContext 变量。

    【讨论】:

      【解决方案3】:

      当您在数据帧上使用缓存时,它是一种转换,当您对其执行任何操作(如 count()、show() 等)时,它会被懒惰地评估。

      在您的情况下,在进行第一次缓存之后,您正在调用 show(),这就是数据帧缓存在内存中的原因。现在,您再次对数据帧执行转换以添加其他列并再次缓存新数据帧,然后再次调用操作命令 show,这会将第二个数据帧缓存在内存中。如果您的数据帧的大小足以容纳一个数据帧,那么当您缓存第二个数据帧时,它将从内存中删除第一个数据帧,因为它没有足够的空间容纳第二个数据帧。

      要记住:除非您在多个操作中使用数据帧,否则不应缓存数据帧,否则在性能方面会超载,因为缓存本身是更昂贵的操作。

      【讨论】:

        猜你喜欢
        • 2022-01-07
        • 2016-10-27
        • 1970-01-01
        • 2021-03-07
        • 1970-01-01
        • 2022-10-24
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多