如何动态选择spark.sql.shuffle.partitions

Nar*_*esh 5 apache-spark apache-spark-1.6

我当前正在使用spark处理数据,并且foreach分区打开与mysql的连接,并将其以1000的批数插入到数据库中。如SparkDocumentation所述,默认值为spark.sql.shuffle.partitions200,但我想保持其动态。因此,我该如何计算。因此,既不选择导致性能降低的非常高的值,也不选择导致性能降低的非常小的值OOM。

Raj*_*tti -4

您可以使用df.repartition(numPartitions)方法来执行此操作。您可以根据输入/中间输出做出决定,并将numPartitions 传递给 repartition() 方法。

df.repartition(numPartitions)   or rdd.repartition(numPartitions)
Run Code Online (Sandbox Code Playgroud)