【发布时间】: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