【问题标题】:How to avoid re-evaluation of each transformation on pyspark data frame again and again如何避免一次又一次地重新评估pyspark数据框中的每个转换
【发布时间】:2019-10-30 16:56:36
【问题描述】:

我有一个 spark 数据框。我正在对数据框进行多次转换。我的代码如下所示:

df = df.withColumn ........
df2 = df.filter......
df = df.join(df1 ...
df = df.join(df2 ...

现在我有大约 30 多个这样的转换。我也知道数据框的持久化。因此,如果我有一些这样的转换:

df1 = df.filter.....some condition
df2 = df.filter.... some condtion
df3 = df.filter... some other conditon

在上述情况下,我正在持久化数据框“df”。

现在的问题是 spark 运行时间过长(8 + mts),或者有时它会因 Java 堆空间问题而失败。 但是,如果我保存到一个表(持久配置单元表)并从下一行的表中读取,经过大约 10 次以上的转换,大约需要 3 + mts 才能完成。即使我将它保存到内存表的中间,它也不起作用。 集群大小也不是问题。

# some transformations
df.write.mode("overwrite").saveAsTable("test")
df = spark.sql("select * from test")
# some transormations    ---------> 3 mts

# some transformations
df.createOrReplaceTempView("test")
df.count() #action statement for view to be created
df = spark.sql("select * from test")
# some more transformations  --------> 8 mts.    

看了spark sql plan(还是没完全看懂)。看起来 spark 正在一次又一次地重新评估相同的数据帧。

我做错了什么?我不必将其写入中间表。

编辑:我正在开发 azure databricks 5.3(包括 Apache Spark 2.4.0、Scala 2.11)

Edit2:问题是 rdd long lineage。如果 rdd 沿袭增加,我的 spark 应用程序看起来会越来越慢。

【问题讨论】:

  • 不能说没有看到实际代码、环境设置和数据

标签: apache-spark pyspark pyspark-sql pyspark-dataframes


【解决方案1】:

你应该使用缓存。

尝试使用

df.cache
df.count

使用 count 强制缓存所有信息。

另外我建议你看看thisthis

【讨论】:

    猜你喜欢
    • 2015-04-15
    • 1970-01-01
    • 2020-05-20
    • 1970-01-01
    • 2016-08-09
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多