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