【问题标题】:DataFrame join optimization - Broadcast Hash JoinDataFrame 连接优化 - 广播哈希连接
【发布时间】:2015-12-02 19:24:09
【问题描述】:

我正在尝试有效地连接两个 DataFrame,其中一个较大,第二个较小。

有没有办法避免所有这些洗牌?我无法设置autoBroadCastJoinThreshold,因为它只支持整数——而且我尝试广播的表略大于整数字节数。

有没有办法强制广播忽略这个变量?

【问题讨论】:

  • 感觉你的实际问题是“有没有办法强制广播忽略这个变量?”它可以通过我下面提到的属性来控制。由于没有人提到,为了使其相关,我给出了这个迟到的答案。希望有帮助!
  • 请接受一次作为已接受的答案。它也将指向其他人。谢谢!

标签: apache-spark dataframe apache-spark-sql apache-spark-1.4


【解决方案1】:

Broadcast Hash Joins(类似于 map side join 或 Mapreduce 中的 map-side combine):

在 SparkSQL 中,您可以通过调用 queryExecution.executedPlan 查看正在执行的连接类型。与核心 Spark 一样,如果其中一个表比另一个小得多,您可能需要广播哈希连接。您可以在加入之前通过在 DataFrame 上调用方法 broadcast 来提示 Spark SQL 应该广播给定的 DF 以进行加入

示例: largedataframe.join(broadcast(smalldataframe), "key")

在 DWH 术语中,largedataframe 可能类似于 fact
smalldataframe 可能类似于 dimension

正如我最喜欢的书 (HPS) 所描述的那样。看下面有更好的理解..

注意:以上broadcast 来自import org.apache.spark.sql.functions.broadcast 而不是来自SparkContext

Spark 也会自动使用spark.sql.conf.autoBroadcastJoinThreshold 来确定是否应该广播一个表。

提示:参见 DataFrame.explain() 方法

def
explain(): Unit
Prints the physical plan to the console for debugging purposes.

有没有办法强制广播忽略这个变量?

sqlContext.sql("SET spark.sql.autoBroadcastJoinThreshold = -1")


注意:

另一个类似的开箱即用笔记 w.r.t.蜂巢(不是火花):类似 可以使用 hive 提示 MAPJOIN 来实现,如下所示...

Select /*+ MAPJOIN(b) */ a.key, a.value from a join b on a.key = b.key

hive> set hive.auto.convert.join=true;
hive> set hive.auto.convert.join.noconditionaltask.size=20971520
hive> set hive.auto.convert.join.noconditionaltask=true;
hive> set hive.auto.convert.join.use.nonstaged=true;
hive> set hive.mapjoin.smalltable.filesize = 30000000; // default 25 mb made it as 30mb

延伸阅读:请参考我的article on BHJ, SHJ, SMJ

【讨论】:

  • 当我想做 smallDF.join(broadcast(largeDF, "left_outer") 时做 largeDF.join(broadcast(smallDF), "right_outer") 有意义吗?因为 smallDF 应该被保存在内存中而不是 largeDF
  • 请。关注SparkStrategies.scala 这里提到了所有案例。在这里写它不是单一的衬里......将 left_outer 更改为 right_outer 结果将被更改。
  • 但正常情况下Table1 LEFT OUTER JOIN Table2, Table2 RIGHT OUTER JOIN Table1 是相等的
【解决方案2】:

您可以使用left.join(broadcast(right), ...) 提示要广播的数据帧

【讨论】:

  • 此广播的正确导入是什么?我知道这个符号 broadcast 是不可解析的。
  • 在org.apache.spark.sql.functions下,需要spark 1.5.0或更新版本
  • 有机会提示广播连接到 SQL 语句吗?
  • 您确实可以在 SQL 语句中使用该提示,但不确定它的效果如何。例如我像left inner join broadcast(right) 一样使用它,不确定它是否适用于子查询
  • 要在 SparkSQL 中广播,可以例如:val df = broadcast(spark.table("tableA")).createTempView("tableAView")spark.sql("SELECT ... FROM tableAView a JOIN tableB b")
【解决方案3】:

设置spark.sql.autoBroadcastJoinThreshold = -1 将完全禁用广播。看 Other Configuration Options in Spark SQL, DataFrames and Datasets Guide.

【讨论】:

    【解决方案4】:

    这是 spark 的当前限制,请参阅 SPARK-6235。 2GB 限制也适用于广播变量。

    您确定没有其他好方法可以做到这一点,例如不同的分区?

    否则,您可以通过手动创建多个广播变量来破解它,每个变量

    【讨论】:

    • 我已经设法将一个较小的表的大小减小到略低于 2 GB,但似乎广播无论如何都没有发生。 (autoBroadcast 只是不会选择它)。不过,它适用于小表(100 MB)。
    【解决方案5】:

    我发现此代码适用于 Spark 2.11 版本 2.0.0 中的广播加入。

    import org.apache.spark.sql.functions.broadcast  
    
    val employeesDF = employeesRDD.toDF
    val departmentsDF = departmentsRDD.toDF
    
    // materializing the department data
    val tmpDepartments = broadcast(departmentsDF.as("departments"))
    
    import context.implicits._
    
    employeesDF.join(broadcast(tmpDepartments), 
       $"depId" === $"id",  // join by employees.depID == departments.id 
       "inner").show()
    

    这里是上面代码Henning Kropp Blog, Broadcast Join with Spark的参考

    【讨论】:

      【解决方案6】:

      使用连接提示将优先于配置autoBroadCastJoinThreshold,因此使用提示将始终忽略该阈值。

      此外,当使用连接提示时,Adaptive Query Execution(自 Spark 3.x 起)也不会更改提示中给出的策略。

      Spark SQL 中,您可以应用连接提示,如下所示:

      SELECT /*+ BROADCAST */ a.id, a.value FROM a JOIN b ON a.id = b.id
      
      SELECT /*+ BROADCASTJOIN */ a.id, a.value FROM a JOIN b ON a.id = b.id
      
      SELECT /*+ MAPJOIN */ a.id, a.value FROM a JOIN b ON a.id = b.id
      

      注意,关键字BROADCAST、BROADCASTJOIN和MAPJOIN都是hints.scala代码中写的别名。

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2018-06-17
        • 1970-01-01
        • 1970-01-01
        • 2019-02-10
        • 2012-09-16
        • 1970-01-01
        • 2020-10-08
        • 1970-01-01
        相关资源
        最近更新 更多