相关疑难解决方法(0)

合并减少整个阶段的并行性(火花)

有时,Spark会以低效的方式“优化”数据框架计划。考虑以下Spark 2.1中的示例(也可以在Spark 1.6中复制):

val df = sparkContext.parallelize((1 to 500).map(i=> scala.util.Random.nextDouble),100).toDF("value")

val expensiveUDF = udf((d:Double) => {Thread.sleep(100);d})

val df_result = df
.withColumn("udfResult",expensiveUDF($"value"))

df_result
.coalesce(1)
.saveAsTable(tablename)
Run Code Online (Sandbox Code Playgroud)

在此示例中,我想在对数据帧进行昂贵的转换后写入1个文件(这只是一个演示此问题的示例)。Spark向上移动coalesce(1),使得UDF仅应用于包含1个分区的数据帧,从而破坏了并行性(有趣的repartition(1)是,这种行为不起作用)。

概括地说,当我想在转换的某个部分中增加并行度,但此后降低并行度时,就会发生此行为。

我发现了一种解决方法,包括缓存数据框,然后触发对数据框的完整评估:

val df = sparkContext.parallelize((1 to 500).map(i=> scala.util.Random.nextDouble),100).toDF("value")

val expensiveUDF = udf((d:Double) => {Thread.sleep(100);d})

val df_result = df
.withColumn("udfResult",expensiveUDF($"value"))
.cache

df_result.rdd.count // trigger computation

df_result
.coalesce(1)
.saveAsTable(tablename)
Run Code Online (Sandbox Code Playgroud)

我的问题是:在这种情况下,还有另一种方法可以告诉Spark不降低并行度吗?

scala apache-spark

6
推荐指数
1
解决办法
2066
查看次数

标签 统计

apache-spark ×1

scala ×1