【问题标题】:Broadcast Hash Join (BHJ) in Spark for full outer join (outer, full, fullouter)Spark 中的广播哈希连接 (BHJ) 用于完全外连接(outer、full、fullouter)
【发布时间】:2019-05-12 02:44:26
【问题描述】:

如何强制 Spark 中的 Dataframes 完全外连接以使用 Boradcast Hash Join?这是代码sn-p:

sparkConfiguration.set("spark.sql.autoBroadcastJoinThreshold", "1000000000")
val Result = BigTable.join(
  org.apache.spark.sql.functions.broadcast(SmallTable),
  Seq("X", "Y", "Z", "W", "V"),
  "outer"
)

My SmallTable 的大小远小于上面指定的autoBroadcastJoinThreshold。此外,如果我使用内部、left_outerright_outer 连接,我从 DAG 可视化中看到连接按预期使用 BroadcastHashJoin

但是,当我使用“outer”作为连接类型时,spark 出于某种未知原因决定使用SortMergeJoin。有谁知道如何解决这个问题?根据我看到的左外连接的性能,BroadcastHashJoin 将有助于加速我的应用程序多倍。

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    spark 出于某种未知原因决定使用 SortMergeJoin。做 有人知道如何解决这个问题吗?

    原因: FullOuter(表示任何关键字outerfullfullouter)不支持广播哈希连接(又名地图侧连接)

    如何证明这一点?

    举一个例子:

    包 com.examples 导入 org.apache.log4j.{级别,记录器} 导入 org.apache.spark.internal.Logging 导入 org.apache.spark.sql.SparkSession 导入 org.apache.spark.sql.functions._ /** * 使用示例数据加入示例和一些基础演示。 * * @作者:拉姆·加迪亚拉姆 */ 对象 JoinExamples 扩展 Logging { // 关闭不必要的日志 Logger.getLogger("org").setLevel(Level.OFF) val spark: SparkSession = SparkSession.builder.config("spark.master", "local").getOrCreate; case class Person(name: String, age: Int, personid: Int) 案例类 Profile(名称:String,personId:Int,profileDescription:String) /** * 主要的 * * @param args 数组[字符串] */ def main(args: Array[String]): Unit = { spark.conf.set("spark.sql.join.preferSortMergeJoin", "false") 导入 spark.implicits._ spark.sparkContext.getConf.getAllWithPrefix("spark.sql").foreach(x => logInfo(x.toString())) /** * 在此处使用案例类创建 2 个数据框,一个是 Person df1,另一个是 profile df2 */ val df1 = spark.sqlContext.createDataFrame( spark.sparkContext.parallelize( 人(“萨拉特”,33,2) :: Person("KangarooWest", 30, 2) :: Person("Ravikumar Ramasamy", 34, 5) :: Person("Ram Ghadiyaram", 42, 9) :: Person("Ravi chandra Kancharla", 43, 9) :: 无)) val df2 = spark.sqlContext.createDataFrame( 配置文件(“Spark”,2,“SparkSQLMaster”) :: Profile("Spark", 5, "SparkGuru") :: Profile("Spark", 9, "DevHunter") :: 无 ) // 你可以用别名来引用列名来增加可读性 val df_asPerson = df1.as("dfperson") val df_asProfile = df2.as("dfprofile") /** * * 示例显示如何在数据框级别加入它们 * 下一个示例演示如何使用带有 createOrReplaceTempView 的 sql */ val join_df = df_asPerson.join( 广播(df_asProfile) , col("dfperson.personid") === col("dfprofile.personid") , "外") val加入=joined_df.select( col("dfperson.name") , col("dfperson.age") , col("dfprofile.name") , col("dfprofile.profileDescription")) join.explain(false) // 它将显示使用了哪个连接 加入.show } }

    我尝试对fullouter 加入使用广播提示,但框架忽略了它,下面的SortMergeJoin 是对此的解释计划。 结果:

    == 物理计划 == *项目 [name#4, age#5, name#11, profileDescription#13] +- SortMergeJoin [personid#6]、[personid#12]、FullOuter :- *排序 [personid#6 ASC NULLS FIRST], false, 0 : +- 交换散列分区(personid#6, 200) : +- *SerializeFromObject [staticinvoke(class org.apache.spark.unsafe.types.UTF8String, StringType, fromString, assertnotnull(input[0, com.examples.JoinExamples$Person, true]).name, true) AS 名称# 4、assertnotnull(input[0, com.examples.JoinExamples$Person, true]).age AS age#5, assertnotnull(input[0, com.examples.JoinExamples$Person, true]).personid AS personid#6] : +- 扫描 ExternalRDDScan[obj#3] +- *排序 [personid#12 ASC NULLS FIRST], false, 0 +- 交换哈希分区(personid#12, 200) +- LocalTableScan [name#11, personId#12, profileDescription#13] +--------------------+---+-----+------------------ + |姓名|年龄|姓名|简介描述| +--------------------+---+-----+------------------ + |拉维库马尔·拉马萨米| 34|火花|火花大师| |拉姆·加迪亚拉姆| 42|火花|开发者| |拉维·钱德拉·坎克...| 43|火花|开发者| |萨拉特| 33|火花| SparkSQLMaster| |袋鼠西| 30|火花| SparkSQLMaster| +--------------------+---+-----+------------------ +

    从 spark 2.3 开始 Merge-Sort join 是 spark 中默认的 join 算法。 但是,这可以通过使用内部参数来关闭 ‘spark.sql.join.preferSortMergeJoin’ 默认为真。

    fullouter join 以外的其他情况...如果您不想让 spark 在任何情况下使用 sortmergejoin,您可以设置以下属性。

    sparkSession.conf.set("spark.sql.join.preferSortMergeJoin", "false")
    

    这是您不想使用 sortmergejoin 的代码 SparkStrategies.scala (which is responsible & Converts a logical plan into zero or more SparkPlans) 的说明。

    展示。注意:

    此属性spark.sql.join.preferSortMergeJoin 如果为true,则更喜欢通过此PREFER_SORTMERGEJOIN 属性进行排序合并连接而不是随机散列连接。

    设置false意味着spark不能只选择broadcasthashjoin它也可以是其他任何东西(例如shuffle hash join)。

    以下文档位于SparkStrategies.scala 中,即在object JoinSelection extends Strategy with PredicateHelper ... 之上


    • 广播:如果连接的一侧的估计物理大小小于 用户可配置 [[SQLConf.AUTO_BROADCASTJOIN_THRESHOLD]] 阈值 或者如果该方有明确的广播提示(例如,用户应用了 [[org.apache.spark.sql.functions.broadcast()]] 函数到DataFrame),然后是那一边 加入的将被广播,另一端将被流式传输,没有洗牌 执行。如果加入的双方都有资格被广播,那么
    • Shuffle hash join:如果单个分区的平均大小足够小,可以构建散列表。

    • 排序合并:如果匹配的连接键是可排序的。

    【讨论】:

    【解决方案2】:

    广播连接不支持完全外连接。它仅支持以下类型:

    内心喜欢 |左外 |左半 |左反 |存在加入 |右外

    详情请见JoinStrategy

    【讨论】:

      猜你喜欢
      • 2015-12-02
      • 1970-01-01
      • 2012-11-06
      • 1970-01-01
      • 2018-06-17
      • 2018-11-12
      • 2015-10-08
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多