spark 出于某种未知原因决定使用 SortMergeJoin。做
有人知道如何解决这个问题吗?
原因: FullOuter(表示任何关键字outer、full、fullouter)不支持广播哈希连接(又名地图侧连接)
如何证明这一点?
举一个例子:
包 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 ... 之上