有时,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不降低并行度吗?