0%

Spark SQL

1 Spark-SQL

1.1 join

Spark 中支持多种连接类型:

  • Inner Join : 内连接;
  • Full Outer Join : 全外连接;
  • Left Outer Join : 左外连接;
  • Right Outer Join : 右外连接;
  • Left Semi Join : 左半连接;
  • Left Anti Join : 左反连接;
  • Natural Join : 自然连接;
  • Cross (or Cartesian) Join : 交叉 (或笛卡尔) 连接

SQL JOINS

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
emp 员工表
|-- ENAME: 员工姓名
|-- DEPTNO: 部门编号
|-- EMPNO: 员工编号
|-- HIREDATE: 入职时间
|-- JOB: 职务
|-- MGR: 上级编号
|-- SAL: 薪资
|-- COMM: 奖金

dept 部门表
|-- DEPTNO: 部门编号
|-- DNAME: 部门名称
|-- LOC: 部门所在城市

-- LEFT SEMI JOIN
SELECT * FROM emp LEFT SEMI JOIN dept ON emp.deptno = dept.deptno
-- 等价于如下的 IN 语句
SELECT * FROM emp WHERE deptno IN (SELECT deptno FROM dept)

-- LEFT ANTI JOIN
SELECT * FROM emp LEFT ANTI JOIN dept ON emp.deptno = dept.deptno
-- 等价于如下的 IN 语句
SELECT * FROM emp WHERE deptno NOT IN (SELECT deptno FROM dept)

--CROSS JOIN
SELECT * FROM emp CROSS JOIN dept ON emp.deptno = dept.deptno

--自然连接是在两张表中寻找那些数据类型和列名都相同的字段,然后自动地将他们连接起来,并返回所有符合条件的结果。
SELECT * FROM emp NATURAL JOIN dept

--程序自动推断出使用两张表都存在的 dept 列进行连接
SELECT * FROM emp JOIN dept ON emp.deptno = dept.deptno

1.1.1 内部实现

broadcast join –> hash join –> sort-merge join

在对大表与大表之间进行连接操作时,通常都会触发 Shuffle Join,两表的所有分区节点会进行 All-to-All 的通讯,这种查询通常比较昂贵,会对网络 IO 会造成比较大的负担。

https://github.com/heibaiying

而对于大表和小表的连接操作,Spark 会在一定程度上进行优化,如果小表的数据量小于 Worker Node 的内存空间,Spark 会考虑将小表的数据广播到每一个 Worker Node,在每个工作节点内部执行连接计算,这可以降低网络的 IO,但会加大每个 Worker Node 的 CPU 负担。

是否采用广播方式进行 Join 取决于程序内部对小表的判断,如果想明确使用广播方式进行 Join,则可以在 DataFrame API 中使用 broadcast 方法指定需要广播的小表:

1
empDF.join(broadcast(deptDF), joinExpression).show()

1.2 Driver Collect Exec

数据必须收集到Driver的Exec:

Exec Statement 优化方向
CollectLimitExec 替换为GlobalLimitExec
TakeOrderedAndProjectExec
CollectTailExec

1.3 优化

优化规则 规则名称 简介
列裁剪 column_prune 对于上层算子不需要的列,不在下层算子输出该列,减少计算
子查询去关联 decorrelate 尝试对相关子查询进行改写,将其转换为普通 join 或 aggregation 计算
聚合消除 aggregation_eliminate 尝试消除执行计划中的某些不必要的聚合算子
投影消除 projection_eliminate 消除执行计划中不必要的投影算子
最大最小消除 max_min_eliminate 改写聚合中的 max/min 计算,转化为 order by + limit 1
谓词下推 predicate_push_down 尝试将执行计划中过滤条件下推到离数据源更近的算子上
外连接消除 outer_join_eliminate 尝试消除执行计划中不必要的 left join 或者 right join
分区裁剪 partition_processor 将分区表查询改成为用 union all,并裁剪掉不满足过滤条件的分区
聚合下推 aggregation_push_down 尝试将执行计划中的聚合算子下推到更底层的计算节点
TopN 下推 topn_push_down 尝试将执行计划中的 TopN 算子下推到离数据源更近的算子上
Join 重排序 join_reorder 对多表 join 确定连接顺序

1.4 逻辑优化

1.4.1 子查询相关的优化

关联子查询去关联

1.4.2 列裁剪

1.4.3 关联子查询去关联

1.4.4 Max/Min 消除

1.4.5 谓词下推

1.4.6 分区裁剪

1.4.7 TopN 和 Limit 下推

1.4.8 Join Reorder

1.5 物理优化

1.5.1 选择最优的索引进行表的访问

1.5.2 收集统计信息来获得表的数据分布情况

1.5.3 在错误索引的解决方案中会介绍当发现 TiDB 索引选错时,你应该使用那些手段来让它使用正确的索引

1.5.4 在 Distinct 优化中会介绍在物理优化中会做的一个有关 DISTINCT 关键字的优化,在这一小节中会介绍它的优缺点以及如何使用它。

1.6 [参考文献]

  1. The Business Intelligence for Hadoop Benchmark