【问题标题】:Does Spark SQL do predicate pushdown on filtered equi-joins?Spark SQL 是否对过滤的 equi-join 进行谓词下推?
【发布时间】: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_udfgroup_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


    【解决方案1】:

    很大程度上取决于您所说的下推。如果你问|a.id_a - b.id_b| &lt; 2 是否作为join 逻辑的一部分在a.grp = b.grp 旁边执行,答案是否定的。不基于相等的谓词不直接包含在join 条件中。

    您可以说明的一种方法是使用 DAG 而不是执行计划 它应该或多或少像这样:

    如您所见,filter 是作为与SortMergeJoin 不同的转换执行的。另一种方法是在删除a.grp = b.grp 时分析执行计划。您会看到它将join 扩展为笛卡尔积,然后是filter,没有进行额外的优化:

    d = a.join(b,(F.abs(a['id_a'] - b['id_b']) < 2)).drop(b['grp'])
    
    ## == Physical Plan ==
    ## Project [id_a#2L,grp#1L,id_b#5L]
    ## +- Filter (abs((id_a#2L - id_b#5L)) < 2)
    ##    +- CartesianProduct
    ##       :- ConvertToSafe
    ##       :  +- Project [id#0L AS id_a#2L,(FLOOR((cast(id#0L as double) / 10.0)) * 10) AS grp#1L]
    ##       :     +- Scan ExistingRDD[id#0L] 
    ##       +- ConvertToSafe
    ##          +- Project [id#3L AS id_b#5L]
    ##             +- Scan ExistingRDD[id#3L]
    

    这是否意味着您的代码(不是带有笛卡尔的代码 - 您真的想在实践中避免这种情况)会生成一个巨大的中间表?

    不,它没有。 SortMergeJoinfilter 都作为单个阶段执行(参见 DAG)。虽然DataFrame 操作的一些细节可以在稍低的级别上应用,但它基本上只是 Scala 上的一系列转换Iteratorsas shown in a very illustrative way by Justin Pihony,不同的操作可以被压缩在一起,而无需添加任何特定于 Spark 的逻辑。一种或另一种两种过滤器将应用于单个任务。

    【讨论】:

    • 迭代器压缩很有意义,但它对我提出了另一个问题。几个月前,我在 Spark 1.2.1 中使用 RDD 进行了similar experiment,得到了相反的结果。似乎您的相同逻辑应该适用于该示例:join + filter 应该只是将迭代器压缩在一起,除非 Spark RDD join 不能那样工作。
    • 哦,对不起 :) python_join.dispatch 实际上返回了一个生成器表达式。
    • 明确一点 - 你觉得this 麻烦吗?因为您链接的问题/答案描述了
    • 啊,我看错了代码。这些对确实仅在 1.2 之后的生成器中生成。
    • python_full_outer_join 中有一个列表 comp 剩余但我今天修复了它。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-10
    • 2020-01-11
    • 2019-01-21
    • 2020-01-28
    • 1970-01-01
    • 2019-07-04
    相关资源
    最近更新 更多