【问题标题】:Does spark.sql.autoBroadcastJoinThreshold work for joins using Dataset's join operator?spark.sql.autoBroadcastJoinThreshold 是否适用于使用 Dataset 的连接运算符的连接?
【发布时间】:2017-05-15 16:01:53
【问题描述】:

我想知道spark.sql.autoBroadcastJoinThreshold 属性是否可用于在所有工作节点上广播较小的表(在进行连接时),即使连接方案使用 Dataset API 连接而不是使用 Spark SQL。

如果我的大表是 250 Gigs 而较小的表是 20 Gigs,我是否需要设置此配置:spark.sql.autoBroadcastJoinThreshold = 21 Gigs(可能)以便将整个表 / Dataset 发送到所有工作节点?

示例

  • 数据集 API 加入

    val result = rawBigger.as("b").join(
      broadcast(smaller).as("s"),
      rawBigger(FieldNames.CAMPAIGN_ID) === smaller(FieldNames.CAMPAIGN_ID), 
      "left_outer"
    )
    
  • SQL

    select * 
    from rawBigger_table b, smaller_table s
    where b.campign_id = s.campaign_id;
    

【问题讨论】:

    标签: apache-spark apache-spark-sql


    【解决方案1】:

    首先spark.sql.autoBroadcastJoinThresholdbroadcast 提示是独立的机制。即使autoBroadcastJoinThreshold 被禁用设置broadcast 提示将优先。使用默认设置:

    spark.conf.get("spark.sql.autoBroadcastJoinThreshold")
    
    String = 10485760
    
    val df1 = spark.range(100)
    val df2 = spark.range(100)
    

    Spark 将使用autoBroadcastJoinThreshold 并自动广播数据:

    df1.join(df2, Seq("id")).explain
    
    == Physical Plan ==
    *Project [id#0L]
    +- *BroadcastHashJoin [id#0L], [id#3L], Inner, BuildRight
       :- *Range (0, 100, step=1, splits=Some(8))
       +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, bigint, false]))
          +- *Range (0, 100, step=1, splits=Some(8))
    

    当我们禁用自动广播时,Spark 将使用标准 SortMergeJoin:

    spark.conf.set("spark.sql.autoBroadcastJoinThreshold", -1)
    df1.join(df2, Seq("id")).explain
    
    == Physical Plan ==
    *Project [id#0L]
    +- *SortMergeJoin [id#0L], [id#3L], Inner
       :- *Sort [id#0L ASC NULLS FIRST], false, 0
       :  +- Exchange hashpartitioning(id#0L, 200)
       :     +- *Range (0, 100, step=1, splits=Some(8))
       +- *Sort [id#3L ASC NULLS FIRST], false, 0
          +- ReusedExchange [id#3L], Exchange hashpartitioning(id#0L, 200)
    

    但可以强制使用BroadcastHashJoinbroadcast 提示:

    df1.join(broadcast(df2), Seq("id")).explain
    
    == Physical Plan ==
    *Project [id#0L]
    +- *BroadcastHashJoin [id#0L], [id#3L], Inner, BuildRight
       :- *Range (0, 100, step=1, splits=Some(8))
       +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, bigint, false]))
          +- *Range (0, 100, step=1, splits=Some(8))
    

    SQL 有自己的提示格式(类似于 Hive 中使用的那种):

    df1.createOrReplaceTempView("df1")
    df2.createOrReplaceTempView("df2")
    
    spark.sql(
     "SELECT  /*+ MAPJOIN(df2) */ * FROM df1 JOIN df2 ON df1.id = df2.id"
    ).explain
    
    == Physical Plan ==
    *BroadcastHashJoin [id#0L], [id#3L], Inner, BuildRight
    :- *Range (0, 100, step=1, splits=8)
    +- BroadcastExchange HashedRelationBroadcastMode(List(input[0, bigint, false]))
       +- *Range (0, 100, step=1, splits=8)
    

    所以回答您的问题 - autoBroadcastJoinThreshold 适用于使用 Dataset API 时,但在使用显式 broadcast 提示时不相关。

    此外,广播大型对象不太可能提供任何性能提升,并且在实践中通常会降低性能并导致稳定性问题。请记住,广播对象必须首先获取到驱动程序,然后发送到每个工作人员,最后加载到内存中。

    【讨论】:

      【解决方案2】:

      只是为了分享更多细节(从代码中)到@user6910411 的出色答案。


      引用source code(格式化我的):

      spark.sql.autoBroadcastJoinThreshold 配置在执行连接时将广播到所有工作节点的表的最大大小(以字节为单位)。

      通过将此值设置为 -1 可以禁用广播。

      请注意,目前仅支持运行了命令 ANALYZE TABLE COMPUTE STATISTICS noscan 的 Hive Metastore 表以及直接在数据文件上计算统计信息的基于文件的数据源表。

      spark.sql.autoBroadcastJoinThreshold 默认为 10M(即10L * 1024 * 1024),Spark 将检查要使用的连接(请参阅JoinSelection 执行计划策略)。

      6 种不同的加入选择,其中包括广播(使用 BroadcastHashJoinExecBroadcastNestedLoopJoinExec 物理运算符)。

      BroadcastHashJoinExec 将在有加入键和以下其中一项时被选中:

      • Join 是 CROSS, INNER, LEFT ANTI, LEFT OUTER, LEFT SEMI 和 right join side 之一可以广播,即尺寸小于spark.sql.autoBroadcastJoinThreshold
      • Join 是 CROSS、INNER 和 RIGHT OUTER 之一,left join 边可以广播,即大小小于spark.sql.autoBroadcastJoinThreshold

      BroadcastNestedLoopJoinExec 将在有 no 加入键且BroadcastHashJoinExec 的上述条件之一成立时被选中。

      换句话说,Spark 会自动选择正确的联接,包括基于 spark.sql.autoBroadcastJoinThreshold 属性(以及其他要求)的 BroadcastHashJoinExec 以及联接类型。

      【讨论】:

      • 当然!在这里:issues.apache.org/jira/browse/SPARK-26214 很高兴让 Spark 变得更好
      • no joining keys 是什么意思?这是否意味着不是equi-join 之类的?
      • no joining keys 是当您的查询在 SQL 中不使用 ON 或在 Dataset.join 运算符中没有键时。说得通?继续问,直到你对答案没问题。我最终可能会再次查看来源,以确保我没有问题:)
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2014-07-02
      • 2020-06-03
      • 1970-01-01
      • 2012-05-07
      • 1970-01-01
      • 1970-01-01
      • 2012-08-16
      相关资源
      最近更新 更多