【发布时间】:2016-07-14 19:46:22
【问题描述】:
我对使用 Spark SQL (1.6) 执行表单的“过滤等值连接”很感兴趣
A inner join B where A.group_id = B.group_id and pair_filter_udf(A[cols], B[cols])
这里的group_id 是粗略的:group_id 的单个值可以与 A 和 B 中的 10,000 条记录相关联。
如果 equi-join 自己执行,没有pair_filter_udf,group_id 的粗糙度会产生计算问题。例如,对于在 A 和 B 中都有 10,000 条记录的 group_id,连接中将有 1 亿个条目。如果我们有成千上万个这样的大组,我们将生成一个巨大的表,我们很容易耗尽内存。
因此,我们必须将pair_filter_udf 下推到连接中并让它在生成对时过滤它们,而不是等到所有对都生成之后。我的问题是 Spark SQL 是否这样做。
我设置了一个简单的过滤 equi-join 并询问 Spark 它的查询计划是什么:
# run in PySpark Shell
import pyspark.sql.functions as F
sq = sqlContext
n=100
g=10
a = sq.range(n)
a = a.withColumn('grp',F.floor(a['id']/g)*g)
a = a.withColumnRenamed('id','id_a')
b = sq.range(n)
b = b.withColumn('grp',F.floor(b['id']/g)*g)
b = b.withColumnRenamed('id','id_b')
c = a.join(b,(a.grp == b.grp) & (F.abs(a['id_a'] - b['id_b']) < 2)).drop(b['grp'])
c = c.sort('id_a')
c = c[['grp','id_a','id_b']]
c.explain()
结果:
== Physical Plan ==
Sort [id_a#21L ASC], true, 0
+- ConvertToUnsafe
+- Exchange rangepartitioning(id_a#21L ASC,200), None
+- ConvertToSafe
+- Project [grp#20L,id_a#21L,id_b#24L]
+- Filter (abs((id_a#21L - id_b#24L)) < 2)
+- SortMergeJoin [grp#20L], [grp#23L]
:- Sort [grp#20L ASC], false, 0
: +- TungstenExchange hashpartitioning(grp#20L,200), None
: +- Project [id#19L AS id_a#21L,(FLOOR((cast(id#19L as double) / 10.0)) * 10) AS grp#20L]
: +- Scan ExistingRDD[id#19L]
+- Sort [grp#23L ASC], false, 0
+- TungstenExchange hashpartitioning(grp#23L,200), None
+- Project [id#22L AS id_b#24L,(FLOOR((cast(id#22L as double) / 10.0)) * 10) AS grp#23L]
+- Scan ExistingRDD[id#22L]
这些是计划中的主要内容:
+- Filter (abs((id_a#21L - id_b#24L)) < 2)
+- SortMergeJoin [grp#20L], [grp#23L]
这些行给人的印象是过滤器将在连接后的单独阶段完成,这不是所需的行为。但也许它被隐式下推到连接中,而查询计划只是缺乏那种详细程度。
在这种情况下,我如何知道 Spark 正在做什么?
更新:
我正在运行 n=1e6 和 g=1e5 的实验,如果 Spark 不执行下推,这应该足以让我的笔记本电脑崩溃。由于它没有崩溃,我猜它正在下推。但如果想知道它是如何工作的,以及 Spark SQL 源代码的哪些部分负责这种出色的优化,那将会很有趣。
【问题讨论】:
标签: python apache-spark dataframe pyspark apache-spark-sql