Spark¶
一、Spark 基础概念¶
1.1 RDD 是什么¶
RDD(Resilient Distributed Dataset,弹性分布式数据集)是Spark的核心数据结构。 - 特点:不保存实际数据,仅封装计算逻辑 - 代码层面:是一个抽象类
1.2 RDD 的弹性体现在哪¶
- 自动的进行内存和磁盘的存储切换
- 基于血缘的高效容错
- task如果失败会自动进行特定次数的重试
- stage如果失败会自动进行特定次数的重试,而且只会计算失败的分片
- checkpoint和persist,数据计算之后持久化缓存
- 数据调度弹性,DAG TASK调度和资源无关
- 数据分片的高度弹性,分片可以合并或重新分区
1.3 血缘关系¶
多个连续的RDD的依赖关系,称之为血缘关系;每个RDD会保存血缘关系。
1.4 宽窄依赖¶
| 依赖类型 | 定义 | 特点 |
|---|---|---|
| 窄依赖 | 父RDD的一个分区的数据只会被子RDD的一个分区依赖 | 可并行计算,无shuffle |
| 宽依赖 | 父RDD的一个分区的数据会被子RDD的多个分区依赖 | 涉及shuffle,需等待上一阶段完成 |
1.5 Application、Job、Stage、Task 关系¶
- 初始化一个SparkContext就会生成一个Application
- 一个行动算子就会生成一个Job
- Stage(阶段)等于宽依赖的个数+1
- 一个阶段中,最后一个RDD的分区个数就是Task的个数
- 总结:Application > Job > Stage > Task
二、Spark 核心原理¶
2.1 Spark 任务执行流程(Spark on YARN)¶
- Spark 客户端提交作业给 YARN 的 RM
- RM 分配 container,启动 ApplicationMaster
- AM 启动 Driver,紧接着向 RM 申请资源启动 Executor
- Executor 进程启动后会向 Driver 反向注册
- 全部注册完后 Driver 开始执行 main 函数
- 当执行到行动算子,触发一个 Job,并根据宽依赖开始划分 stage
- 每个 stage 生成对应的 TaskSet,之后将 task 分发到各个 Executor 上执行
2.2 Stage 划分¶
- 从后往前,遇到窄依赖则加入该Stage,遇到宽依赖则断开,划分新的Stage
- 窄依赖尽量放在同一个 stage 中,实现流水线并行计算
- 宽依赖由于有 shuffle,只能在父 RDD 处理完成后,才能开始接下来的计算
2.3 Spark 有哪几种部署模式¶
| 模式 | 说明 |
|---|---|
| 本地模式 | 常用于本地开发程序与测试 |
| Standalone模式 | 只使用Spark自身节点的运行模式 |
| YARN模式 | YARN作为资源调度框架的运行模式 |
| Mesos模式 | Mesos作为资源调度管理系统 |
2.4 Client 和 Cluster 模式的区别¶
- client模式:driver运行在客户端
- cluster模式:driver运行在YARN集群
2.5 Spark job 和谁保持通信¶
和 Driver 保持通信,Driver 端的 HeartbeatReceiver 负责接收 Executor 心跳报文,监控 Executor 存活状态。
三、Spark Shuffle 机制¶
3.1 Spark Shuffle 实现¶
Spark的shuffle分为两种实现: - HashShuffle(Spark 1.2以前) - SortShuffle(Spark 1.2以后)
3.2 HashShuffle¶
分为普通机制和合并机制: - Write阶段:根据key进行分区,写入对应的 buffer 中,写满后溢写到磁盘 - Read阶段:reduce去拉取各个 maptask 产生的同一个分区的数据 - 合并机制:让多个 mapper 共享 Buffer,落盘文件数量等于 reduce数量 乘以 core的个数,减少磁盘IO
3.3 SortShuffle¶
分为普通机制和bypass机制:
| 机制 | 说明 | 触发条件 |
|---|---|---|
| 普通机制 | mapTask的结果先放到5M内存结构,溢写时先按key分区和排序,最后合并成大文件并生成索引 | 默认 |
| bypass机制 | 去掉排序过程,其他不变 | mapTask数量 < 200 且 不是聚合类Shuffle算子 |
3.4 能产生 Shuffle 的算子¶
reduceByKey、sortByKey、repartition、coalesce、join、cogroup
四、Spark 算子¶
4.1 Transform 和 Action 算子的区别¶
- 转换算子:将旧的RDD包装成新的RDD(lazy)
- 行动算子:触发任务的调度和作业的执行
4.2 常用转换算子¶
| 算子 | 说明 |
|---|---|
| map | 将数据逐条进行转换 |
| flatMap | 先map再扁平化 |
| filter | 根据指定规则筛选 |
| coalesce | 减少分区个数,默认不进行shuffle |
| repartition | 增加或减少分区个数,一定发生shuffle |
| union | 两个RDD求并集 |
| zip | 将两个RDD中的元素以键值对形式合并 |
| reduceByKey | 按照key对value进行聚合 |
| groupByKey | 按照key对value进行分组 |
| cogroup | 按照key对value进行合并 |
4.3 常用行动算子¶
| 算子 | 说明 |
|---|---|
| collect | 将数据采集到Driver端,形成数组 |
| take | 返回RDD的前n个元素组成的数组 |
| foreach | 遍历RDD中的每一个元素(executor端) |
4.4 groupByKey 和 reduceByKey 的区别¶
| 对比项 | reduceByKey | groupByKey |
|---|---|---|
| 功能 | 分组 + 聚合 | 只能分组,不能聚合 |
| 预聚合 | shuffle前对分区内相同key进行预聚合 | 无预聚合,直接shuffle |
| 性能 | 更高(减少落盘数据量) | 较低 |
五、Spark SQL¶
5.1 三种 Join 实现¶
Broadcast Hash Join¶
- 适用场景:小表 + 大表
- 流程:
- Broadcast阶段:将小表广播到所有Executor
- Hash Join阶段:在每个Executor上执行hash join,小表构建hash table,大表作为probe table
Shuffle Hash Join¶
- 适用场景:较大的小表 + 大表
- 流程:
- Shuffle阶段:对两张表分别按照join字段重分区
- Hash Join阶段:对每个分区执行hash join
Sort Merge Join¶
- 适用场景:两张大表
- 流程:
- Shuffle阶段:按join字段重分区
- Sort阶段:对每个分区内数据排序
- Merge阶段:对排好序的分区表进行join
5.2 执行计划优化¶
RBO(Rule-Based Optimization)¶
基于规则对逻辑计划进行优化: - 谓词下推:将过滤条件(where/on)尽可能提前执行 - 列裁剪:只读取查询相关的字段,减少IO - 常量替换/折叠:常量表达式提前计算
CBO(Cost-Based Optimization)¶
计算所有可能物理计划的代价,挑选代价最小的:
- 前提:需要收集统计信息
- 表级别:ANALYZE TABLE 表名 COMPUTE STATISTICS
- 列级别:ANALYZE TABLE 表名 COMPUTE STATISTICS FOR COLUMNS 列1,列2...
- 开启:spark.sql.cbo.enabled = true
SMB Join(Sort Merge Bucket Join)¶
- 原理:两表都按 join 列分桶并排序,相同 key 在同一个桶中,join 时只需桶内匹配
- 条件:两表都是分桶表,分桶个数相等,且
join列 = 排序列 = 分桶列
5.3 Spark SQL 去重方法¶
| 方法 | 说明 |
|---|---|
| DISTINCT | 简单去重 |
| 窗口函数 | 条件去重 |
| GROUP BY | 聚合去重 |
六、Spark 与 MapReduce 对比¶
6.1 主要区别¶
| 对比项 | Spark | MapReduce |
|---|---|---|
| 计算方式 | 内存计算,RDD+DAG | IO读写,中间结果落盘 |
| 任务调度 | Task为线程级别,复用线程池 | Task为进程级别 |
| Shuffle | 只有部分场景需要排序,支持Hash分布式聚合 | Shuffle前需大量时间排序 |
| 运行模型 | 多线程 | 多进程 |
| 通用性 | 提供transformation和action两大类API,还有Streaming、GraphX等模块 | 只提供map和reduce两种操作 |
6.2 Spark 为什么更快¶
- 内存计算:中间结果以RDD形式存放在内存中,减少磁盘IO
- 减少排序:Shuffle时如果选择基于hash的计算引擎,无需排序
- 多线程模型:每个task是运行在executor中的线程,减少启动开销
七、Spark Streaming¶
7.1 简介¶
Spark Streaming 是一种准实时、微批次的数据处理框架。 - 使用离散化流(DStream)作为抽象表示 - DStream 内部是一系列连续的 RDD,每个 RDD 含有一段时间间隔内的数据
7.2 基本工作原理¶
- 接受实时输入数据流
- 将数据封装成 batch(如设置1秒延迟)
- 将每个 batch 交给 Spark 计算引擎处理
- 输出结果数据流
7.3 窗口函数原理¶
在原来定义的批次大小基础上再次封装,每次计算多个批次的数据,同时传递滑动步长参数设置下次计算起始位置。
八、Spark AQE 新特性¶
8.1 什么是 AQE¶
AQE(Adaptive Query Execution,自适应查询执行)是 Spark SQL 的动态优化机制。
- 运行时,每当 Shuffle Map 阶段执行完毕,AQE 结合统计信息动态调整、修正尚未执行的计划
- 开启:spark.sql.adaptive.enabled = true(Spark 3.0+ 默认开启)
8.2 AQE 三大特性¶
| 特性 | 作用 |
|---|---|
| 自动分区合并 | Shuffle 后自动合并过小的数据分区 |
| Join 策略调整 | 如果过滤后表变小,从 Sort Merge Join 降级为 Broadcast Hash Join |
| 自动倾斜处理 | 自动拆分 Reduce 阶段过大的数据分区 |
九、Spark 持久化¶
9.1 为什么需要持久化¶
因为 RDD 不存储数据,要重用的话需要从头执行,持久化可提高重用性。
9.2 持久化方式对比¶
| 方式 | 说明 | 是否切断血缘 |
|---|---|---|
| cache | 默认调用 persist(MEMORY_ONLY),临时存储在内存 | 否 |
| persist | 可指定存储级别(内存/磁盘),程序结束自动删除 | 否 |
| checkpoint | 长久保存在磁盘 | 是 |
十、Spark 调优¶
10.1 资源参数调优¶
| 参数 | 作用 | 建议 |
|---|---|---|
--num-executors |
Executor数量 | 根据集群规模设置 |
--executor-memory |
每个Executor内存 | 4G~16G,避免OOM |
--executor-cores |
每个Executor核数 | 官方建议2~5,企业常用4 |
--driver-memory |
Driver内存 | collect大量数据或广播大变量时调大 |
spark.sql.shuffle.partitions |
Shuffle后分区数 | 默认200,数据量大可调到500~1000 |
spark.default.parallelism |
RDD Shuffle后分区数 | CPU总核数的2~3倍 |
10.2 内存参数调优¶
| 参数 | 作用 |
|---|---|
| spark.memory.fraction | 堆内内存用于执行、shuffle、缓存的比例(默认0.6) |
| spark.memory.storageFraction | 存储内存不会被逐出的比例(默认0.5) |
| spark.kryoserializer.buffer.max | Kryo序列化缓存大小 |
10.3 Shuffle 调优¶
- 调大
spark.shuffle.file.buffer:默认32k,建议64k,减少溢写次数 - 增加重试次数:
spark.shuffle.io.maxRetries = 10,spark.shuffle.io.retryWait = 20s - 合并小文件:开启
hive.merge.mapfiles、hive.merge.mapredfiles
10.4 并行度优化¶
- 并行度设置为集群 CPU 总核数的 2~3倍
- Spark SQL 调
spark.sql.shuffle.partitions(默认200) - RDD 调
spark.default.parallelism
10.5 数据本地化优化¶
Spark 数据本地化级别从优到劣: | 级别 | 含义 | |------|------| | PROCESSLOCAL | 数据和Task在同一个Executor JVM中(最优) | | NODELOCAL | 数据和Task在同一个节点不同Executor | | NOPREF | 数据在哪里访问都一样(如MySQL) | | RACKLOCAL | 数据和Task在同机架不同节点 | | ANY | 数据在任意节点(最差) |
优化:调大 spark.locality.wait(默认3s),让Task更耐心等待数据本地化。
10.6 数据倾斜调优¶
| 方案 | 说明 |
|---|---|
| 过滤空值/异常key | 无意义脏数据直接过滤 |
| 提高并行度 | 增加分区数,减少每个Task处理量 |
| 两阶段聚合 | 先加随机前缀做局部聚合,再去掉前缀做全局聚合 |
| Broadcast Join | 大表join小表,广播小表避免Shuffle |
| 加盐打散 | 倾斜key加随机前缀分散到不同分区 |
| AQE自动处理 | 开启AQE后自动检测并拆分倾斜分区 |
10.7 存储与序列化调优¶
- 序列化:推荐 Kryo 序列化,比Java序列化更小更快
- 持久化级别:内存足够用
MEMORY_ONLY,否则MEMORY_ONLY_SER或MEMORY_AND_DISK - 文件格式:列式存储(Parquet、ORC)比行存储更优,配合压缩减少IO
十一、Spark SQL vs RDD¶
什么情况用 RDD 而不用 Spark SQL: 1. 自定义分区器:需自定义分区逻辑时 2. 复杂数据处理:如图算法、迭代计算 3. 非结构化数据:处理文本、日志等 4. 低级别 API 操作:需要直接操作数据分区时 5. 自定义序列化:需要更多序列化控制时