【问题标题】:Spark 1.6 DataFrame optimize join partitioningSpark 1.6 DataFrame 优化连接分区
【发布时间】:2019-02-10 11:30:43
【问题描述】:

我有一个关于 Spark DataFrame 分区的问题,我目前正在使用 Spark 1.6 来满足项目要求。这是我的代码摘录:

sqlContext.getConf("spark.sql.shuffle.partitions") // 6

val df = sc.parallelize(List(("A",1),("A",4),("A",2),("B",5),("C",2),("D",2),("E",2),("B",7),("C",9),("D",1))).toDF("id_1","val_1")
df.rdd.getNumPartitions // 4

val df2 = sc.parallelize(List(("B",1),("E",4),("H",2),("J",5),("C",2),("D",2),("F",2))).toDF("id_2","val_2")
df2.rdd.getNumPartitions // 4

val df3 = df.join(df2,$"id_1" === $"id_2")
df3.rdd.getNumPartitions // 6

val df4 = df3.repartition(3,$"id_1")
df4.rdd.getNumPartitions // 3

df4.explain(true)

以下是已创建的解释计划:

== Parsed Logical Plan ==
'RepartitionByExpression ['id_1], Some(3)
+- Join Inner, Some((id_1#42 = id_2#46))
   :- Project [_1#40 AS id_1#42,_2#41 AS val_1#43]
   :  +- LogicalRDD [_1#40,_2#41], MapPartitionsRDD[169] at rddToDataFrameHolder at <console>:26
   +- Project [_1#44 AS id_2#46,_2#45 AS val_2#47]
      +- LogicalRDD [_1#44,_2#45], MapPartitionsRDD[173] at rddToDataFrameHolder at <console>:26

== Analyzed Logical Plan ==
id_1: string, val_1: int, id_2: string, val_2: int
RepartitionByExpression [id_1#42], Some(3)
+- Join Inner, Some((id_1#42 = id_2#46))
   :- Project [_1#40 AS id_1#42,_2#41 AS val_1#43]
   :  +- LogicalRDD [_1#40,_2#41], MapPartitionsRDD[169] at rddToDataFrameHolder at <console>:26
   +- Project [_1#44 AS id_2#46,_2#45 AS val_2#47]
      +- LogicalRDD [_1#44,_2#45], MapPartitionsRDD[173] at rddToDataFrameHolder at <console>:26

== Optimized Logical Plan ==
RepartitionByExpression [id_1#42], Some(3)
+- Join Inner, Some((id_1#42 = id_2#46))
   :- Project [_1#40 AS id_1#42,_2#41 AS val_1#43]
   :  +- LogicalRDD [_1#40,_2#41], MapPartitionsRDD[169] at rddToDataFrameHolder at <console>:26
   +- Project [_1#44 AS id_2#46,_2#45 AS val_2#47]
      +- LogicalRDD [_1#44,_2#45], MapPartitionsRDD[173] at rddToDataFrameHolder at <console>:26

== Physical Plan ==
TungstenExchange hashpartitioning(id_1#42,3), None
+- SortMergeJoin [id_1#42], [id_2#46]
   :- Sort [id_1#42 ASC], false, 0
   :  +- TungstenExchange hashpartitioning(id_1#42,6), None
   :     +- Project [_1#40 AS id_1#42,_2#41 AS val_1#43]
   :        +- Scan ExistingRDD[_1#40,_2#41] 
   +- Sort [id_2#46 ASC], false, 0
      +- TungstenExchange hashpartitioning(id_2#46,6), None
         +- Project [_1#44 AS id_2#46,_2#45 AS val_2#47]
            +- Scan ExistingRDD[_1#44,_2#45]

据我所知,DataFrame 表示RDD 之上的抽象接口,因此应该将分区委托给 Catalyst 优化器。

事实上,与RDD 相比,许多转换接受多个分区参数,以便尽可能优化协同分区和协同定位,DataFrame 唯一改变分区的机会是调用方法重新分区,否则,使用配置参数spark.sql.shuffle.partitions 推断连接和聚合的分区数。

从我从上面的解释计划中可以看到和理解,似乎有一个 无用的重新分区(确实是随机播放) 到 6(默认值),然后再次重新分区到施加的最终值通过方法重新分区。

我相信优化器可以将连接的分区数更改为最终值 3。

有人可以帮我澄清一下吗?也许我错过了什么。

【问题讨论】:

  • 如果难以理解,体积较小。使用更大的卷和设置随机分区,您往往会获得实际使用的分区数量。 V 2.3.1 至少
  • join 过程中遵循的分区方案是 HashPartitioning,正如您在说明计划中看到的那样。连接后的分区数将等于spark.sql.shufle.partitions,但数据分布可能不同。例如,在您的示例中,您有 4 个不同的 ID,因此在 df3 拥有的 6 个分区中,2 个将是空的,不会导致任何作业/任务/阶段。

标签: apache-spark dataframe apache-spark-sql rdd


【解决方案1】:

如果你使用spark sql,你的shuffle partitions总是等于spark.sql.shufle.partitions。但是如果你启用这个spark.sql.adaptive.enabled就会添加EchangeCoordinator。现在这个coordinator的工作是确定需要从一个或多个阶段获取 shuffle 数据的阶段的 post-shuffle 分区数。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2016-04-12
    • 2015-12-02
    • 2016-03-06
    • 2022-01-23
    • 1970-01-01
    • 2018-01-19
    • 2017-01-15
    相关资源
    最近更新 更多