【问题标题】:Understanding SparkSQL and its usage of partitioning了解 SparkSQL 及其对分区的使用
【发布时间】:2016-09-16 19:17:46
【问题描述】:

我正在尝试针对某些数据操作查询评估 Spark SQL。我对此感兴趣的场景:

table1: key, value1, value2
table2: key, value3, value4

create table table3 as
select * from table1 join table2 on table1.key = table2.key

听起来我应该能够创建 table1 和 table2 RDD(但我在文档中没有看到一个非常明显的例子)。 但更大的问题是——如果我已经成功地通过键对 2 个表 RDD 进行了分区,然后使用 Spark SQL 将它们连接起来,那么它是否足够聪明来利用分区?如果我因为加入而创建了一个新的 RDD,它也会被分区吗?换句话说,它会完全免洗牌吗? 我非常感谢有关这些主题的文档和/或示例的指针。

【问题讨论】:

标签: apache-spark apache-spark-sql


【解决方案1】:

如果您指的是RDDsDatasets 之间的转换,那么这两个问题的答案是否定的。

RDD 分区仅为RDD[(T, U)] 定义,在RDD 转换为Dataset 后将丢失。在某些情况下,您可以从现有数据布局中受益,但 join 不是其中之一,尤其是 RDDsDatasets 使用不同的散列技术(分别是标准的 hashCodeMurmurHash。你当然可以通过定义自定义分区器RDD 来模仿后者,但这并不是重点)。

Dataset 转换为RDD 时,有关分区的类似信息也会丢失。

您可以使用Dataset 分区来优化joins。例如,如果表已预先分区:

val n: Int = ??? 

val df1 =  Seq(
  ("key1", "val1", "val2"), ("key2", "val3", "val4")
).toDF("key", "val1", "val2").repartition(n, $"key").cache

val df2 = Seq(
  ("key1", "val5", "val6"), ("key2", "val7", "val8")
).toDF("key", "val3", "val4").repartition(n, $"key").cache

基于key 的后续join 将不需要额外的交换。

spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1

df1.explain
// == Physical Plan ==
// InMemoryTableScan [key#171, val1#172, val2#173]
//    +- InMemoryRelation [key#171, val1#172, val2#173], true, 10000, StorageLevel(disk, memory, deserialized, 1 replicas)
//          +- Exchange hashpartitioning(key#171, 3)
//             +- LocalTableScan [key#171, val1#172, val2#173]
df2.explain
// == Physical Plan ==
// InMemoryTableScan [key#201, val3#202, val4#203]
//    +- InMemoryRelation [key#201, val3#202, val4#203], true, 10000, StorageLevel(disk, memory, deserialized, 1 replicas)
//          +- Exchange hashpartitioning(key#201, 3)
//             +- LocalTableScan [key#201, val3#202, val4#203]
// 
df1.join(df3, Seq("key")).explain
// == Physical Plan ==
// *Project [key#171, val1#172, val2#173, val5#232, val6#233]
// +- *SortMergeJoin [key#171], [key#231], Inner
//    :- *Sort [key#171 ASC], false, 0
//    :  +- *Filter isnotnull(key#171)
//    :     +- InMemoryTableScan [key#171, val1#172, val2#173], [isnotnull(key#171)]
//    :           +- InMemoryRelation [key#171, val1#172, val2#173], true, 10000, StorageLevel(disk, memory, deserialized, 1 replicas)
//    :                 +- Exchange hashpartitioning(key#171, 3)
//    :                    +- LocalTableScan [key#171, val1#172, val2#173]
//    +- *Sort [key#231 ASC], false, 0
//       +- *Filter isnotnull(key#231)
//          +- InMemoryTableScan [key#231, val5#232, val6#233], [isnotnull(key#231)]
//                +- InMemoryRelation [key#231, val5#232, val6#233], true, 10000, StorageLevel(disk, memory, deserialized, 1 replicas)
//                      +- Exchange hashpartitioning(key#231, 3)
//                         +- LocalTableScan [key#231, val5#232, val6#233]

显然,我们并没有真正从单一连接中受益。因此,仅当单个表用于多个 joins 时才有意义。

Spark 也可以从join 创建的分区中受益,所以如果我们想执行另一个join

val df3 = Seq(
  ("key1", "val9", "val10"), ("key2", "val11", "val12")
).toDF("key", "val5", "val6")

df1.join(df3, Seq("key")).join(df3, Seq("key"))

我们将从第一个操作创建的结构中受益(注意ReusedExchange):

// == Physical Plan ==
// *Project [key#171, val1#172, val2#173, val5#682, val6#683, val5#712, val6#713]
// +- *SortMergeJoin [key#171], [key#711], Inner
//    :- *Project [key#171, val1#172, val2#173, val5#682, val6#683]
//    :  +- *SortMergeJoin [key#171], [key#681], Inner
//    :     :- *Sort [key#171 ASC], false, 0
//    :     :  +- Exchange hashpartitioning(key#171, 200)
//    :     :     +- *Filter isnotnull(key#171)
//    :     :        +- InMemoryTableScan [key#171, val1#172, val2#173], [isnotnull(key#171)]
//    :     :              +- InMemoryRelation [key#171, val1#172, val2#173], true, 10000, StorageLevel(disk, memory, deserialized, 1 replicas)
//    :     :                    +- Exchange hashpartitioning(key#171, 3)
//    :     :                       +- LocalTableScan [key#171, val1#172, val2#173]
//    :     +- *Sort [key#681 ASC], false, 0
//    :        +- Exchange hashpartitioning(key#681, 200)
//    :           +- *Project [_1#677 AS key#681, _2#678 AS val5#682, _3#679 AS val6#683]
//    :              +- *Filter isnotnull(_1#677)
//    :                 +- LocalTableScan [_1#677, _2#678, _3#679]
//    +- *Sort [key#711 ASC], false, 0
//       +- ReusedExchange [key#711, val5#712, val6#713], Exchange hashpartitioning(key#681, 200)

【讨论】:

  • 所以这是 DataSet 的例子,很好。这会映射到 SparkSql 吗?如果我通过加入在同一列上分区的 2 个 DF 使用 Spark SQL 创建一个新的 DF 会怎样?生成的 DF 会被分区吗?
  • 是的,SQL和DataFrame API在执行上没有区别。
  • 我的理解是它只会在 Spark 2.0 中变得聪明,对吧?
  • 据我所知,基本优化也应该在 1.6 中起作用,但要获得全部好处,您需要 2.0+
猜你喜欢
  • 2016-07-12
  • 2016-03-24
  • 2015-10-22
  • 2015-09-29
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-12-10
相关资源
最近更新 更多