【问题标题】:Apache Spark: broadcast join behaviour: filtering of joined tables and temp tablesApache Spark:广播连接行为:过滤连接表和临时表
【发布时间】:2021-07-09 16:06:32
【问题描述】:

我需要在 spark 中加入 2 个表。 但是我没有完全加入 2 个表,而是首先过滤掉了第二个表的一部分:

spark.sql("select * from a join b on a.key=b.key where b.value='xxx' ")

我想在这种情况下使用广播连接。

Spark 有一个参数定义广播连接的最大表大小:spark.sql.autoBroadcastJoinThreshold

配置表的最大大小(以字节为单位) 执行连接时向所有工作节点广播。通过设置这个 可以禁用值为 -1 的广播。请注意,目前 仅 Hive Metastore 表支持统计信息,其中 命令 ANALYZE TABLE COMPUTE STATISTICS noscan 已被 跑。 http://spark.apache.org/docs/2.4.0/sql-performance-tuning.html

我对此设置有以下疑问:

  1. spark 将与 autoBroadcastJoinThreshold 的值比较哪个表大小:FULL 大小,还是应用 where 子句后的大小?
  2. 我假设 spark 将在广播之前应用 where 子句,对吗?
  3. 文档说我需要事先运行 Hive 的分析表命令。当我将临时视图用作表格时,它将如何工作?据我了解,我无法针对通过 dataFrame.createorReplaceTempView("b") 创建的 spark 的临时视图运行分析表命令。我可以广播临时视图内容吗?

【问题讨论】:

    标签: apache-spark hive


    【解决方案1】:

    对选项 2 的理解是正确的。 您无法在 spark 中分析 TEMP 表。阅读here

    如果您想带头并想指定要广播的数据帧,则由spark决定,可以使用sn-p-

    df = df1.join(F.broadcast(df2),df1.some_col == df2.some_col, "left")
    

    【讨论】:

      【解决方案2】:

      我继续做了一些小实验来回答你的第一个问题。

      问题 1:

      • 创建了一个包含 3 行 [key,df_a_column] 的数据框 a
      • 创建了一个包含 10 行 [key,value] 的数据框 b
      • 跑:spark.sql("SELECT * FROM a JOIN b ON a.key = b.key").explain()
      == Physical Plan ==
      *(1) BroadcastHashJoin [key#122], [key#111], Inner, BuildLeft, false
      :- BroadcastExchange HashedRelationBroadcastMode(List(cast(input[0, int, false] as bigint)),false), [id=#168]
      :  +- LocalTableScan [key#122, df_a_column#123]
      +- *(1) LocalTableScan [key#111, value#112]
      

      正如预期的那样,广播了具有 3 行的 Smaller df a

      • 冉:spark.sql("SELECT * FROM a JOIN b ON a.key = b.key where b.value=\"bat\"").explain()
      == Physical Plan ==
      *(1) BroadcastHashJoin [key#122], [key#111], Inner, BuildRight, false
      :- *(1) LocalTableScan [key#122, df_a_column#123]
      +- BroadcastExchange HashedRelationBroadcastMode(List(cast(input[0, int, false] as bigint)),false), [id=#152]
         +- LocalTableScan [key#111, value#112]
      

      在这里您可以注意到数据帧b 已广播!意思是 spark 评估 size AFTER applying where 以选择要广播的那个。

      问题 2:

      是的,你是对的。从前面的输出中可以明显看出它首先应用的地方。

      问题 3: 不,您无法分析,但您可以通过在 SQL 中提示有关它的 spark 来广播 tempView 表。 ref

      例如:spark.sql("SELECT /*+ BROADCAST(b) */ * FROM a JOIN b ON a.key = b.key")

      如果你现在看到解释:

      == Physical Plan ==
      *(1) BroadcastHashJoin [key#122], [key#111], Inner, BuildRight, false
      :- *(1) LocalTableScan [key#122, df_a_column#123]
      +- BroadcastExchange HashedRelationBroadcastMode(List(cast(input[0, int, false] as bigint)),false), [id=#184]
         +- LocalTableScan [key#111, value#112]
      

      现在,如果您看到,即使数据帧 b 有 10 行,它也会被广播。 问题1,没有提示,a被广播了。

      注意:SQL spark 中的广播提示适用于 2.2


      了解物理计划的提示:

      • LocalTableScan[ list of columns ] 中找出数据框
      • 正在广播BroadcastExchange 的子树/列表下的数据帧。

      【讨论】:

        猜你喜欢
        • 2020-10-08
        • 2015-10-08
        • 1970-01-01
        • 2018-11-12
        • 1970-01-01
        • 1970-01-01
        • 2021-09-08
        • 2022-08-24
        • 1970-01-01
        相关资源
        最近更新 更多