跳转至

Spark

一、Spark 基础概念

1.1 RDD 是什么

RDD(Resilient Distributed Dataset,弹性分布式数据集)是Spark的核心数据结构。 - 特点:不保存实际数据,仅封装计算逻辑 - 代码层面:是一个抽象类

1.2 RDD 的弹性体现在哪

  1. 自动的进行内存和磁盘的存储切换
  2. 基于血缘的高效容错
  3. task如果失败会自动进行特定次数的重试
  4. stage如果失败会自动进行特定次数的重试,而且只会计算失败的分片
  5. checkpoint和persist,数据计算之后持久化缓存
  6. 数据调度弹性,DAG TASK调度和资源无关
  7. 数据分片的高度弹性,分片可以合并或重新分区

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)

  1. Spark 客户端提交作业给 YARN 的 RM
  2. RM 分配 container,启动 ApplicationMaster
  3. AM 启动 Driver,紧接着向 RM 申请资源启动 Executor
  4. Executor 进程启动后会向 Driver 反向注册
  5. 全部注册完后 Driver 开始执行 main 函数
  6. 当执行到行动算子,触发一个 Job,并根据宽依赖开始划分 stage
  7. 每个 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 为什么更快

  1. 内存计算:中间结果以RDD形式存放在内存中,减少磁盘IO
  2. 减少排序:Shuffle时如果选择基于hash的计算引擎,无需排序
  3. 多线程模型:每个task是运行在executor中的线程,减少启动开销

七、Spark Streaming

7.1 简介

Spark Streaming 是一种准实时、微批次的数据处理框架。 - 使用离散化流(DStream)作为抽象表示 - DStream 内部是一系列连续的 RDD,每个 RDD 含有一段时间间隔内的数据

7.2 基本工作原理

  1. 接受实时输入数据流
  2. 将数据封装成 batch(如设置1秒延迟)
  3. 将每个 batch 交给 Spark 计算引擎处理
  4. 输出结果数据流

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 = 10spark.shuffle.io.retryWait = 20s
  • 合并小文件:开启 hive.merge.mapfileshive.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_SERMEMORY_AND_DISK
  • 文件格式:列式存储(Parquet、ORC)比行存储更优,配合压缩减少IO

十一、Spark SQL vs RDD

什么情况用 RDD 而不用 Spark SQL: 1. 自定义分区器:需自定义分区逻辑时 2. 复杂数据处理:如图算法、迭代计算 3. 非结构化数据:处理文本、日志等 4. 低级别 API 操作:需要直接操作数据分区时 5. 自定义序列化:需要更多序列化控制时

评论