8.1 Join Expressions Join表达式
判断是否应该连接两个数据集.
通过汇集两组数据进行联接计算,类似于SQL的join, 但是Spark能够过滤不匹配的值从而进行操作.
8.2 Join Types Join类型
经过Join表达式判断是否连接后,该步奏判断如何连接,分为如下几种情况:
- Inner joins 内联(包含左右两边共有的值)
- Outer joins 外联(包含左右两边所有的值)
- Left outer joins 左联(包含左边数据集所有的值)
- Right outer joins 右联(包含右边数据集所有的值)
- Left semi joins 左半开联(左边数据集出现在右边的所有的值)
- Left anti joins (左边数据集未出现在右边的所有的值)
- Natural joins (左右两边通过隐式匹配具有相同name的连接)
- Cross (or Cartesian) joins (对左右两边具有相同的row进行匹配连接)
这一章基本就是SQL join差不多
下面是这章所需示例数据集的创建Person, graduateProgram, sparkStatus 三个DF:
// Scala
val person = Seq(
(0, "Bill Chambers", 0, Seq(100)),
(1, "Matei Zaharia", 1, Seq(500, 250, 100)),
(2, "Michael Armbrust", 1, Seq(250, 100)))
.toDF("id", "name", "graduate_program", "spark_status")
val graduateProgram = Seq(
(0, "Masters", "School of Information", "UC Berkeley"),
(2, "Masters", "EECS", "UC Berkeley"),
(1, "Ph.D.", "EECS", "UC Berkeley"))
.toDF("id", "degree", "department", "school")
val sparkStatus = Seq(
(500, "Vice President"),
(250, "PMC Member"),
(100, "Contributor"))
.toDF("id", "status")
// Scala
person.createOrReplaceTempView("person")
graduateProgram.createOrReplaceTempView("graduateProgram")
sparkStatus.createOrReplaceTempView("sparkStatus")
8.3 Inner Joins 内联
这一章基本就是SQL join差不多
没有什么好说的,和SQL的内联相似,如果两边都有匹配值,则返回新的DF:
// Scala
val joinExpression = person.col("graduate_program") === graduateProgram.col("id")
没有则不返回:
//Scala
val wrongJoinExpression = person.col("name") === graduateProgram.col("school")
还可以写成这样:
// Scala
var joinType = "inner"
person.join(graduateProgram, joinExpression, joinType).show()
8.4 Outer Joins 外联
同样和SQL一样:
// Scala
joinType = "outer"
person.join(graduateProgram, joinExpression, joinType).show()
8.5 Left Outer Joins 左联
同SQL:
// Scala
joinType = "left_outer"
graduateProgram.join(person, joinExpression, joinType).show()
8.6 Right Outer Joins 右联
同SQL:
joinType = "right_outer"
person.join(graduateProgram, joinExpression, joinType).show()
8.7 Left Semi Joins 左半开联
并不是传统的连接,只是比较左边数据集出现在右边的所有的值,如果存在就保留,不存在就不保留:
// Scala
joinType = "left_semi"
graduateProgram.join(person, joinExpression, joinType).show()
val gradProgram2 = graduateProgram.union(Seq(
(0, "Masters", "Duplicated Row", "Duplicated School")).toDF())
gradProgram2.createOrReplaceTempView("gradProgram2")
gradProgram2.createOrReplaceTempView("gradProgram2")
gradProgram2.join(person, joinExpression, joinType).show()
8.8 Left Anti Joins
为上面的反义:
joinType = "left_anti"
graduateProgram.join(person, joinExpression, joinType).show()
8.9 Natural Joins
据说谨慎使用,那先不学习
8.10 Cross (Cartesian) Joins
类似于SQL的Cross暂时本人水平有限,先放SQL于Scala的比较代码:
// Scala
joinType = "cross"
graduateProgram.join(person, joinExpression, joinType).show()
-- SQL
SELECT * FROM graduateProgram CROSS JOIN person
ON graduateProgram.id = person.graduate_program
8.11 使用Joins的挑战
复杂类型的Joins
暂时空置
处理重复列名的Joins
暂时空置
8.12 Spark如何处理Joins
暂时空置