04Spark
Spark
Spark 的几种部署模式比较
-
local 模式:单台计算机部署,用于测试。
-
standalone 模式(单集群模式):需要单独配置 Spark 的集群,使用 Spark 自带的资源调度器,即资源和任务分配都由 Spark 自己负责。
-
集群模式
- 存算分离,比如使用 yarn,k8s ,mesos 等资源调度平台结合 spark 使用,经典的比如 spark 和 yarn 的结合
YARN 模式:
-
client 模式:driver 执行在提交任务的本地客户端,任务提交失败,需要手动干预
shell# 提交spark作业时指定yarn模式下的部署方式 --deploy-mode client bin/spark-submit \ --class org.apache.spark.examples.SparkPi \ --master yarn \ --deploy-mode client \ ./examples/jars/spark-examples_2.12-3.3.0.jar -
cluster 模式:driver 运行在 spark 集群中的任一节点,任务提交失败,spark 会自动重试
shell# 提交spark作业时指定yarn模式下的部署方式 --deploy-mode cluster bin/spark-submit \ --class org.apache.spark.examples.SparkPi \ --master yarn \ --deploy-mode cluster \ ./examples/jars/spark-examples_2.12-3.3.0.jar
RDD
RDD 的基本概念
-
RDD 的五大属性
java/* * Internally, each RDD is characterized by five main properties: * * - A list of partitions * - A function for computing each split * - A list of dependencies on other RDDs * - Optionally, a Partitioner for key-value RDDs (e.g. to say that the RDD is hash-partitioned) * - Optionally, a list of preferred locations to compute each split on (e.g. block locations for * an HDFS file) */_一些分区的集合_:每个 RDD 读取数据后,会对文件进行分片,每个文件分片对应一个 RDD 分区。_切片的计算函数_:每个分区的计算逻辑是相同的,计算逻辑即call()方法里的数据处理逻辑。java// 按照制表符分割,并且取第五个元素 JavaRDD<String> productProvincesRdd = filterProductRdd.map(new Function<String, String>() { @Override public String call(String line) throws Exception { // 计算函数 return line.split("\t")[4]; } });_其他RDD的血缘依赖_:RDD 记录了对其他 RDD 的依赖关系。当 RDD 的分区数据丢失时,RDD 会根据血缘关系重新计算以恢复数据。_KV型数据的分区器_:对于处理 K-V 型数据的 RDD,拥有自己的分区器,比如HashPartitioner(根据 Key 的哈希值计算数据应发送的分区)和RangePartitioner(对数据进行抽样,选出几个特定的 Key 作为划分分区的边界,从而确定数据应发往的分区)。_记录切片位置的list_:一个记录分区内数据切片的优先位置列表。例如,如果数据存储在 HDFS 上,则每个 RDD 会记录要处理的数据块在 HDFS 上的位置,并遵循“移动计算而非移动数据”的原则进行任务调度。
-
RDD 的弹性体现
_分区的弹性_:每个 RDD 的分区数量是由数据切片的数量决定的,处理不同的数据时,分区会动态变化。_计算的弹性_:每个 RDD 都记录了需要处理数据切片的位置(例如在 HDFS 上的位置)。在任务分配时,会尽量遵循本地化原则(移动计算而非移动数据)来分配计算任务。_容错的弹性_:每个 RDD 都记录了对其他 RDD 的依赖关系(血缘)。当某个分区的数据丢失时,可以根据血缘关系自动重试任务,恢复分区数据。_存储的弹性_:RDD 在计算过程中,会优先使用内存来存储中间数据。当内存资源不足时,会自动将数据溢写到磁盘。数据存储位置是弹性的。
-
Spark 常见的端口号
18080:Spark 历史服务器的 WebUI 端口。
RDD 的血缘关系
- 宽依赖(Wide Dependency):同一个父 RDD 的 Partition 被多个子 RDD 的一个 Partition 使用(一对多关系)。
- 窄依赖(Narrow Dependency):每一个父 RDD 的 Partition 最多被子 RDD 的一个 Partition 使用(一对一或多对一关系)。
如果父 RDD 的每个 Partition 都只被一个子 RDD 的 Partition 所使用,就是窄依赖;否则就是宽依赖。
划分宽窄依赖的作用是用于判断数据流向。窄依赖必定不会产生 Shuffle,而宽依赖则相反,通常意味着 Shuffle 的发生(父 RDD 的数据将被发送到不同子 RDD 的分区,这些分区可能分布在不同的节点上)。
缓存(cache,persist),checkpoint 的区别
-
Spark 需要做数据持久化的原因
- 一个 RDD 在多个 Job 中被重复使用。
- 一个 Job 的数据处理链路特别长,如果不进行持久化,一旦发生错误,重算成本很高,容错率会降低。
-
缓存和 Checkpoint 的比较
| 持久化方式 | 使用方式 | 数据保存位置 | RDD 依赖关系是否保留 |
|---|---|---|---|
| 缓存 | rdd.cache() 或 rdd.persist() | 内存或者磁盘 | 是 |
| Checkpoint | rdd.checkpoint() | HDFS 等共享存储 | 否(会切断血缘关系) |
cache和persist的区别:cache()是persist()的一种特殊情况,等价于persist(StorageLevel.MEMORY_ONLY)。- 常用的存储级别:
StorageLevel.MEMORY_ONLY:一般用于小数据量场景。StorageLevel.MEMORY_AND_DISK:一般用于大数据量场景。
任务划分:Application、Job、Stage、Task 关系
- 切分 Stage 的流程: 根据最后一个调用 Action 算子的 RDD,从后往前沿着其依赖关系(血缘)进行回溯,一直到 Job 的第一个 RDD。在回溯过程中,每遇到一个宽依赖,就切分出一个新的 Stage。
_Application数量_=SparkContext的个数。一个main方法中通常只有一个SparkContext,即一个 Application。_Job的个数_由 Action 算子(行动算子)的个数决定。每当代码执行到一个 Action 算子,就会触发一个 Job 的提交。常见的 Action 算子有collect()、foreach()、count()、take()等。_Stage的个数_= 宽依赖的个数 + 1。_Task的个数_(在一个 TaskSet 中)= 一个 Stage 中,最后一个 RDD 的分区数。
常用的转换算子和行动算子,及其比较
GitHub 示例代码地址:GitHub 相关的例子
-
Transform 算子(转换算子) 转换算子是惰性执行的,它们不会立即触发 Job 的执行,只负责定义 RDD 之间的转换关系。只有当遇到 Action 算子时,这些转换才会真正执行。
plaintext1. map 2. flatMap 3. reduceByKey(预聚合算子) 4. groupByKey 5. distinctreduceByKey和groupByKey的比较:reduceByKey相比groupByKey,在 Shuffle 之前会先在 Map 端进行预聚合(Combiner),从而减少了 Shuffle 传输的数据量,性能更高。应优先使用reduceByKey。map和flatMap的比较:flatMap比map更灵活,对于输入的每个元素,map只会产生一个输出元素,而flatMap可以产生 0 个、1 个或多个输出元素。
-
Action 算子(行动算子) 行动算子会触发 Spark Job 的实际计算和执行。
plaintext1. collect 2. count 3. foreach 4. foreachPartition 5. **reduce**foreach和foreachPartition的比较:foreach会对 RDD 中每个分区里的每一条数据都执行一次指定的函数,函数的调用次数等于数据总数。而foreachPartition则是对 RDD 的每个分区执行一次指定的函数,函数的调用次数等于分区数。在需要进行数据库连接等昂贵操作时,使用foreachPartition可以显著减少连接的创建和销毁次数,提高执行效率。
共享变量
分布式只读共享变量
广播变量:默认情况下,Spark 会将 Driver 端的数据给每个 Task 都发送一份,所以该数据占用的内存大小 = Task 个数 * Driver 数据大小。但是如果将该 Driver 数据广播,此时数据会广播给每个 Executor,该 Executor 中的所有 Task 共用这一份数据,所以此时该数据占用的内存大小 = Executor 个数 * Driver 数据大小,相比之前节省了内存空间。
分布式只写共享变量
累加器:先在 RDD 的每个分区中进行累加,然后将各个分区的累加结果返回给 Driver 进行汇总。
SparkSQL
hive on spark 和 spark on hive 的比较
- Hive on Spark:Hive 作为 SQL 引擎,负责 SQL 的解析、生成逻辑计划、优化、生成物理计划等,但将底层的计算引擎从 MapReduce 替换为 Spark。执行任务时,Hive 将物理计划转换为 Spark 作业来执行。(参考 03_HIVE.md 中的 "hive 的 Driver 组成")
- Spark on Hive:Spark SQL 作为 SQL 引擎,负责 SQL 解析和执行的所有环节。Hive 只负责提供元数据管理服务(Metastore),让 Spark 可以读写 Hive 表。
SparkSQL 中常见的数据抽象及其区别
-
RDD, DataFrame, DataSet
- RDD(弹性分布式数据集)是 Spark 的基础数据结构,它关注于对数据的细粒度处理,提供了丰富的算子,但缺乏结构化信息。
- DataFrame 为数据引入了 Schema(模式),类似于关系型数据库中的表。它将每条数据封装为
Row对象,使 Spark 能更高效地处理结构化数据,并进行查询优化。但 DataFrame 是弱类型的,在编译时不做类型检查。 - DataSet 结合了 RDD 的类型安全和 DataFrame 的性能优势。它既有强类型检查(面向对象),也保留了 DataFrame 的非类型化操作接口,是 Spark 推荐使用的数据结构。
SparkSession 的概念
SparkSession 是 Spark 2.0 以后进行 Spark 编程的统一入口点,它整合了 SparkContext、SQLContext 和 HiveContext 的功能,是与 Spark 进行交互的起点。
SparkCore
RDD 的序列化
- RDD 需要序列化的原因:Spark 算子(如
map)的call方法中的代码是在 Executor 的 Task 中执行的,而call方法外部的代码是在 Driver 中执行的。如果 Task 中需要引用在 Driver 端创建的对象,Spark 会将该对象序列化后传输给 Task 使用。因此,这个对象及其所有字段都必须是可序列化的。 - 如何序列化 RDD:Spark 默认使用 Java 序列化,但为了获得更好的性能,推荐使用 Kryo 序列化库。
Spark 的 Shuffle
数据的重新分配和收集可以称为 shuffle
plaintextMR计算过程: 数据 -> InputFormat -> map方法 -> 环形缓冲区[分区、排序, 80%溢写] -> [Combiner] -> 磁盘[小文件] -> 合并小文件 -> Reducer拉取数据[归并排序] -> reduce方法 -> outputFormat -> 磁盘 MR shuffle阶段: -> 环形缓冲区[分区、排序, 80%溢写] -> [Combiner] -> 磁盘[小文件] -> 合并小文件 -> Reducer拉取数据[归并排序] -> Spark shuffle: -> 缓冲区[分区、[排序]] -> [Combiner] -> 磁盘[小文件] -> 合并小文件 -> 分区拉取数据[归并、[排序]] ->
MR 的 Shuffle:
spark 几种 shuffle 机制
- hashshuffle(未优化)
- hashshuffle(Consolidated 优化)
- sortshuffle
- bypassh shuffle
- unsafe shuffle(Tungsten shuffle)
- AQE 中的动态 shuffle
- push-baseed shuffle
简要概述:
spark1.0 世代 的 hash shuffle,consolidated shuffle,sort shuffle 主要解决可用性问题
spark2.0 世代 的 unsafe shuffle 解决 sort shuffle 中的排序带来的 GC 性能问题
spark3.0 世代 的 AQE 的动态 shuffle, push-based shuffle 主要是由于存算分离的发展,解决海量节点的网络瓶颈,磁盘 io 瓶颈,节点宕机容错等问题
hash shuffle
spark1.0 时代,基于数据的 key 做 hash 处理定义数据所属的下游 reducer,此时每个 maptask 都需要创建对应下游 reducetask 数量的中间文件,此时 shuffle 的中间文件数据量为 M(maptask 数量) * R(reducertask 数量),在处理海量数据的场景下,中间的文件将会非常多,导致后续的 reducetask 拉取效率低下(磁盘的随机读写和文件句柄泄露)
基于 M = 1000*M*= 1000 (MapTask 数),R = 1000*R*= 1000 (ReduceTask 数),C = 50*C*= 50 (并行 CPU Cores 数) 的场景基准表示如下:
plaintextMap Tasks (M=1000) Disk Output Files (M x R) ┌──────────────┐ ┌─► [File P0] │ MapTask 1 │ ───────────────┼─► [File P1] └──────────────┘ │ ... (1000个文件/Task) ┌──────────────┐ └─► [File P999] │ MapTask 2 │ ───────────────┼─► [File P0] ... [File P999] └──────────────┘ │ ... │ ┌──────────────┐ │ │ MapTask 1000 │ ───────────────┴─► [File P0] ... [File P999] └──────────────┘ Total: 1,000,000 个物理小文件
中间文件数:M × R = 1000×1000 = 1, 000, 000 个物理文件。
GC 情况: 频繁创建/销毁缓冲 Stream 对象及大量 Key-Value Java 对象,引发频繁 Young GC。
hash shuffle (Consolidated 优化)
hash shuffle 在海量数据场景下,任务的 map,reduce 任务数量是非常大的,导致中间文件数据爆多,文件句柄可能存在泄露,通过引入 ShuffleFileGroup 机制,让运行在同一个 cpu core 上的 maptask 复用同一批输出文件,数据采用 append 模式追加;由于 cpu 的 core 数量决定了同时运行的 task 数量。以 yarn 为资源调度框架为例,spark 相关层级关系如下:
plaintextYARN Cluster └── NodeManager(物理节点上的代理进程,负责启动容器) └── Container(YARN 分配的资源单元,里面跑 Executor) └── Executor(Spark 的 JVM 进程,真正干活的地方) └── 多个 CPU Core(线程池大小 = spark.executor.cores) └── 同时只能跑 等于 Core 数量的 Task(MapTask 或 ReduceTask)
此时 hash shuffle 的每个 maptask 需要 shuffle write 下游 reduce task 数量的文件转变为了共享一个 cpu core 的 maptask 们复用一批输出文件(即为 1000 数量),此时shuffle write情况
plaintextCores & Task Exec (C=50) ShuffleFileGroup (C x R) ┌───────────────────────────┐ │ Core 1 (运行 MapTask 1..20)│ ───► [File P0] ◄──(Task 1~20 串行追加写入) └───────────────────────────┘ ───► [File P1] ... [File P999] (共1000个文件) ┌───────────────────────────┐ │ Core 2 (运行 MapTask 21..40)──► [File P0] ... [File P999] (共1000个文件) └───────────────────────────┘ ... ┌───────────────────────────┐ │ Core 50 │ ───► [File P0] ... [File P999] (共1000个文件) └───────────────────────────┘ Total: 50 x 1000 = 50,000 个物理文件
中间文件数:C× R *= 50×1000 = 50000 个物理文件。
GC 情况: 由于文件还是过多,也会频繁创建/销毁缓冲 Stream 对象及大量 Key-Value Java 对象,引发频繁 Young GC。
在shuffle read阶段,即使多个maptask共享同一个物理文件,但在中间文件中会标识 mapid + reduceid,后续reduce task拉取时会根据中间文件的这个标记拉取自己想要的数据
sort shuffle
sort shuffle的出现是为了进一步减轻hash shuffle场景下的中间文件过多的问题,在hash shuffle通过ShuffleFileGroup后中间文件数据一般和后续reducetask的数据挂钩,当数据量增大,文件数据问题还是会再次暴露。
sort shuffle机制中每个maptask只产生2个文件,.index和.data两个文件。maptask边读数据边写数据到内存缓冲区(特定的数据结构),缓冲区不断溢写并排序数据,记录数据的位置,数据形成一个个临时文件,最后归并排序形成一个大的.data文件和.index文件
plaintextMap Tasks (M=1000) In-Memory Sort & Merge Disk Output Files (2 x M) ┌──────────────┐ ┌──────────────────┐ ┌────────────────────────────┐ │ MapTask 1 │ ───────────────► │ Partitioned Map │ ──► │ Task1.data (P0~P999连续数据) │ └──────────────┘ └──────────────────┘ │ Task1.index (Offsets 索引) │ ┌──────────────┐ ├────────────────────────────┤ │ MapTask 2 │ ────────────────────────────────────────► │ Task2.data / Task2.index │ └──────────────┘ ├────────────────────────────┤ ... │ ... │ ┌──────────────┐ ├────────────────────────────┤ │ MapTask 1000 │ ────────────────────────────────────────► │ Task1000.data / .index │ └──────────────┘ └────────────────────────────┘ Total: 1000 x 2 = 2,000 个文件
中间文件数:2× M *= 2×1000 = 2000 个物理文件
GC 情况: 由于涉及到排序等耗cpu和内存的操作,GC情况也比较严重
shuffle read阶段:reducetask通过.index文件中的数据索引找到自己需要的数据,进行后续操作
bypass shuffle
上述sort shuffle对所有的场景都需要中间的排序,在无需预聚合(Combine)场景下,Sort Shuffle 强制在内存和磁盘做排序带来 CPU 开销,因此对于以下特定场景,可以改良为bypass shuffle:
- 下游分区数少于200
- maptask阶段无需预聚合
bypass shuffle省略了sortshuffle中的排序操作
plaintextMap Tasks (M=1000) No-Sort Direct Write Temp Part Files Final Merge (2 x M) ┌──────────────┐ ┌────────────────┐ ┌─────────────────┐ ┌────────────────────────────┐ │ MapTask 1 │ ─────► │ Bypass Writer │ ───► │ Temp_P0...P999 │ ────► │ Task1.data (P0~P999拼合) │ └──────────────┘ └────────────────┘ └─────────────────┘ │ Task1.index │ (删除临时小文件) ├────────────────────────────┤ │ Task2 ~ 1000 (.data/.index)│ └────────────────────────────┘ Total: 2,000 个最终物理文件
中间文件数:2× M *= 2×1000 = 2000 个物理文件
GC 情况: 相比 Sort Shuffle 省去了排序数据结构的内存占用,Young GC 压力有所缓解。
unsafe shuffle
对于不满足使用bypass shuffle的场景,sort shuffle中的排序操作是必要的,为了减轻堆内 Java 对象序列化开销大(需要反序列化为java对象,才可排序),GC压力(排序操作带来的),采用直接向操作系统请求操作序列化的数据,对数据的指针进行排序操作,绕过了jvm堆,gc问题得以减轻
plaintextOff-Heap Memory (堆外连续空间) Memory Pointer Sorting Disk Output Files (2 x M) ┌─────────────────────────────┐ ┌──────────────────┐ ┌────────────────────────────┐ │ Raw Binary Bytes (Serialized│ ────► │ 8-Byte Pointers │ ────► │ Task1.data (P0~P999二进制) │ └─────────────────────────────┘ │ [PID|Page|Offset]│ │ Task1.index │ (无 Java 对象头,0 堆内 GC 压力) └──────────────────┘ ├────────────────────────────┤ (Radix Sort 指针数组) │ Task2 ~ 1000 (.data/.index)│ └────────────────────────────┘ Total: 2,000 个文件
中间文件数:2× M *= 2×1000 = 2000 个物理文件
GC 情况: GC问题得以减轻
unsafe shuffle之所以称为unsafe,得名于使用了java 的unsafe接口,而并非是不安全的意思。
AQE 中的动态 shuffle
-
解决问题: 解决静态配置 Shuffle 分区数 R导致的“小 Task 过多”或“数据倾斜 Task 卡死”问题。
-
实现机制: 运行时根据 Map 阶段输出的统计信息(Statistics),在 Shuffle Read 侧动态合并连续的小分区,或动态切分倾斜的大分区。
-
中间文件数: Map 端落盘文件数不变,仍为 2×M=2,0002×M=2,000 个;下游 Reduce 侧读取的分区从静态 R=1000 动态压减合并(如自动合并为 100 个 Task 执行)。
plaintextMap Output Files (2 x M = 2000) AQE Runtime Engine Dynamic Reduce Execution ┌────────────────────────────┐ ┌──────────────────┐ │ Task1~1000.data (P0) │ ────► │ Preserve (1 Task)│ ───► ReduceTask 1 (只处理 P0) ├────────────────────────────┤ ├──────────────────┤ │ Task1~1000.data (P1) │ ┐ │ Dynamic Coalesce │ │ ... │ ├───► │ (数据量极小, │ ───► ReduceTask 2 (合并处理 P1~P999) │ Task1~1000.data (P999) │ ┘ │ 合并为 1 Task) │ (实际Reduce Task数从1000压减至2个) └────────────────────────────┘ └──────────────────┘ -
GC 情况: 显著减少因倾斜 Task 内存暴涨导致的 OOM 和 Full GC。
push-based shuffle
-
解决问题: 解决 Shuffle Read 阶段下游节点跨网“随机 Fetch”导致的大规模分布式集群磁盘 IO 及网络吞吐瓶颈。
-
实现机制: 架构转为存算分离/Push 模式。MapTask 边生成数据边主动 Push 到远端 Shuffle 服务,Shuffle 服务按 Partition 将来自 1000 个 MapTask 的数据预先合并为大块顺序文件。
plaintextMap Tasks (Push Mode) Remote Shuffle Services Merged Files on RSS (R=1000) ┌──────────────┐ Push Stream (P0) ┌─────────────────────┐ ┌────────────────────────────┐ │ MapTask 1 │ ────────────────► │ Service Node 1 │ ─►│ Partition_0.merged │ ├──────────────┤ │ (Merge P0 Streams) │ │ (含 Task1~1000 的 P0 数据) │ │ MapTask 2 │ ────────────────► └─────────────────────┘ ├────────────────────────────┤ │ ... │ Push Stream (P999)┌─────────────────────┐ │ Partition_1 ~ 999.merged │ │ MapTask 1000 │ ────────────────► │ Service Node N │ ─►│ (各包含对应分区的全量数据) │ └──────────────┘ │ (Merge P999 Stream) │ └────────────────────────────┘ (Map 端无本地 Sort/Index 文件) └─────────────────────┘ Total: 1,000 个远程连续大文件 -
中间文件数: 计算节点本地 0 物理文件;远程 Shuffle 服务端生成 R=1000个连续的大文件。
-
GC 情况: 减少上游 Task 维护连接和 Read Block 的内存开销,主计算节点的 GC 更加平稳。
Spark Shuffle 和 Hadoop Shuffle 区别
- Hadoop 的 MapReduce Shuffle 无论任务是否需要排序,都必须进行排序。而 Spark 的 BypassShuffle 在满足一定条件时可以跳过排序,提高效率。
- 在 MapReduce 中,MapTask 的环形缓冲区大小是固定的。而 Spark 中 Map 阶段的缓冲区是动态的,可以根据内存使用情况进行扩容。
- 在 Hadoop 中,ReduceTask 在部分 MapTask 完成后就可以开始启动并拉取数据。而在 Spark 中,必须等待前一个 Stage 的所有 Task 全部执行完成后,下一个 Stage 才能开始执行。
Spark 比 MR 快的原因
- 内存计算:Spark 优先使用内存进行数据处理和存储(默认存储级别为
MEMORY_ONLY),减少了磁盘 I/O,而 MapReduce 在任务传递过程中必须依赖磁盘进行 Shuffle。 - 执行模型:Spark 的任务执行是基于线程的(一个 Executor 中可以运行多个 Task 线程),而 MapReduce 的任务是基于进程的(每个 MapTask 或 ReduceTask 都是一个独立的 JVM 进程)。线程的创建和调度开销通常远小于进程。
- DAG 优化:Spark 使用有向无环图(DAG)来表示计算流程,可以进行更复杂的优化,例如将多个操作流水线化(Pipelining)在一个 Stage 中执行,避免了不必要的中间数据落地。
Repartition 和 Coalesce 区别
repartition和coalesce都有改变分区数的功能。repartition是coalesce的一种特殊情况,即repartition(n)等价于coalesce(n, shuffle = true)。repartition通常用于 增加 分区数。它会进行一次完整的 Shuffle,将数据随机打散并重新分布到所有分区中。coalesce通常用于 减少 分区数。默认情况下,它会尽量避免 Shuffle,通过合并相邻的分区来实现,效率更高。如果需要将数据均匀分布,可以设置shuffle = true。
Spark 的内存管理
Spark 1.6 版本后引入了统一内存管理(Unified Memory Management),Spark 可以同时使用堆内内存(On-heap)和堆外内存(Off-heap)。
执行内存(Execution Memory)和存储内存(Storage Memory)的动态占用机制 这两块内存区域可以相互借用。当一方空闲而另一方资源不足时,空闲方可以将部分内存借给对方使用。但是,执行内存有优先权,如果存储内存借用了执行内存,当执行内存需要时,存储内存必须归还,可能会导致缓存的数据被移出内存。
Spark 的持久化保存
见“缓存(cache,persist),checkpoint 的区别”一节
- 数据持久化原因
- 使用方式
- 比较
Spark-YARN 的 Cluster 模式的任务提交流程
提交重点:
_应用的提交_:客户端解析任务的提交参数。_执行环境的准备_:在 ApplicationMaster 中建立 Driver 线程,并初始化 SparkContext 的执行环境,包括任务的切分和调度。_任务的调度和执行_:Driver 向 Executor 发送启动任务的命令,Executor 执行具体的 Task。
具体流程:
- 用户通过
spark-submit脚本提交应用。Submit 进程解析提交参数,并封装成一个 YARN 应用。 - Submit 进程与 YARN 的 ResourceManager 通信,申请启动一个 ApplicationMaster(AM)。
- ResourceManager 在一个 NodeManager 上启动 ApplicationMaster。AM 内部会启动 Driver 线程,并初始化 Spark 的执行环境(如 DAGScheduler、TaskScheduler)。
- ApplicationMaster(作为 Driver 的一部分)向 ResourceManager 申请计算资源(Container)。
- ResourceManager 分配 Container 后,AM 在这些 Container 中启动 Executor 进程。
- Executor 启动后,会反向注册到 Driver,并与 Driver 建立心跳连接。至此,资源申请和执行环境准备完成。
- Driver 中运行用户的
main方法。当遇到 Action 算子时,触发一个 Job 的提交。DAGScheduler 根据 RDD 的宽依赖关系,将 Job 划分为一个或多个 Stage。每个 Stage 被划分为一组 Task(TaskSet)。 - TaskScheduler 将 TaskSet 中的 Task 分配给可用的 Executor 执行,并监控任务的执行情况。
- 当所有 Task 执行完成,ApplicationMaster 向 ResourceManager 取消注册,YARN 回收所有资源,应用结束。
spark 的 join 选择策略
spark 中提供非常丰富的 spark join 策略来处理不同表数据量情况下的选择
- broadcast hash join
- shuffle hash join
- sort merge join
- cartesian join
- broadcast nested loop join
broadcast hash join
涉及广播,hash 表的形成,数据间的 join;使用场景为 小表 join 大表 的 等值连接 场景,spark 默认小表标准为 10MB;driver 将小表数据 广播到各个存储大表数据的 executor,executor 将小表数据利用哈希函数形成一张《joinkey,小表行数据》的哈希表,遍历大表每一行,从哈希表取得小表数据,哈希读取效率高
出发BHJ的条件:
- join的连接必须是等值连接(因为hash table的形成依赖形成的hash值精确)
- 满足小表阈值(可以使用hint语法跳过小表阈值检查,但是也要满足等值连接才可)
shuffle hash join
当 join 的表数据都超过小表存储的定义(10MB),数据数据量适中 的 等值连接 情况,spark join 将会采用 shuffle hash join;涉及过程:
- shuffle:将两个表都按照 joinkey 进行 shuffle,将相同关联 key 的数据发往同一节点的相同分区
- hash:在相同分区内,取数据量较小的表制为 hash 表,遍历较大表的数据,完成 join
Spark Sort Merge Join
Sort Merge Join (SMJ) 是 Spark 中处理 大表等值 Join 的核心策略。当参与 Join 的两张表都过大,无法利用广播机制时,SMJ 提供了一种稳健且内存友好的解决方案。虽然它通常被视为 Spark 的默认 Join 策略(由 spark.sql.join.preferSortMergeJoin 控制),但在现代 Spark 版本中,其最终应用会受到自适应查询执行(AQE)的动态调整。
SMJ 的整个过程可以分为三个核心阶段:Shuffle、Sort 和 Merge Join。
0.基本流程
-
Shuffle (重分区): 对两张表的数据按照 Join Key 进行 Shuffle。此阶段的核心目标是确保所有具有相同 Key 的数据,无论它们最初位于哪个节点或分区,最终都会被发送到同一个 Executor 的同一个分区中,为后续的 Join 做好准备。这是整个流程中最昂贵的操作之一,涉及大量的网络 I/O。
-
Sort (排序): 在每个分区内,Spark 会 独立地 对来自两张表的数据块按照 Join Key 进行排序。排序是 SMJ 的基石,它将无序的数据转化为有序的流,使得后续的合并操作可以高效地进行。此阶段同样消耗大量的磁盘 I/O 和 CPU 资源。
-
Merge Join (合并连接): 这是 SMJ 的精髓所在。Spark 使用类似“拉链”的算法,通过两个指针分别遍历两个已排序的数据流。它逐个比较 Key,根据比较结果移动指针并生成匹配的 Join 结果。由于数据是有序的,已处理过的数据可以立即从内存中释放,这使得 SMJ 对内存的要求极低,非常适合处理海量数据。
1. SMJ 的触发条件:何时被选中?
SMJ 并非总是被选中,它的应用有明确的前提条件:
- key能够排序:由于smj需要按照join key来进行排序方便后续的merge join
- 等值 Join: Join 条件必须包含至少一个等值表达式(例如
ON a.id = b.id)。对于非等值 Join(如a.id > b.id),Spark 只能使用BroadcastNestedLoopJoin或CartesianProduct,效率极低。 - 无法广播: 这是关键。SMJ 是在 Broadcast Hash Join (BHJ) 无法应用时的备选方案。也就是说,当参与 Join 的 两张表,没有任何一张的大小小于
spark.sql.autoBroadcastJoinThreshold(默认 10MB)时,Spark 才会优先考虑 SMJ。如果一张表可广播,Spark 会毫不犹豫地选择效率更高的 BHJ。
2. 静态计划与动态调整
虽然 spark.sql.join.preferSortMergeJoin 默认为 true,但在现代 Spark (3.0+) 中,其行为更加智能:
- 静态计划阶段: 当 Spark 生成初始的物理执行计划时,如果判定无法使用 BHJ,它会遵循
preferSortMergeJoin的配置,将 SMJ 作为计划的 Join 策略。 - 动态执行阶段 (AQE): 如果开启了 Adaptive Query Execution (AQE) (
spark.sql.adaptive.enabled=true,默认开启),Spark 会在运行时根据实际的数据统计信息 动态调整 Join 策略。- 场景: 初始计划为 SMJ,因为统计信息显示表 A 和 B 都很大。
- 动态优化: 在 Shuffle 完成后,AQE 获取到了每个分区数据的精确大小。如果发现其中一张表(比如 B)在 Shuffle 后的实际数据大小小于广播阈值,AQE 会 将 SMJ 在运行时切换为 Broadcast Hash Join。
- 结论: 因此,更准确的说法是,SMJ 是 Spark 处理大表 Join 的 基础和首选策略,但并非不可更改的“最终”策略。AQE 的存在使得 Spark 能够做出更优的实时决策。
3. Merge Join 阶段的详细机制
“对分区内的两张表数据进行遍历”这句话背后是一个精巧的算法。
在 Merge 阶段,Spark 为每个分区内的两个已排序数据流各维护一个指针(可以想象成两个手指分别指着两份名单)。
- 初始化: 两个指针
ptrA和ptrB分别指向两个已排序数据流的起始位置。 - 循环比较: 进入一个循环,不断比较两个指针当前位置的 Join Key (
keyA和keyB):keyA == keyB(匹配成功):- 找到了匹配项!将两行数据合并,生成一条 Join 结果。
- 处理一对多: 此时,需要处理可能存在的“一对多”关系。Spark 会缓存
ptrB当前指向的行,然后 仅移动ptrA,检查A流中是否还有相同的keyA。只要有,就持续与缓存的B行进行匹配并输出结果,直到A流的key发生变化。然后,再移动ptrB。
keyA < keyB(A 表的 Key 较小):- 这意味着
A表中当前这个keyA在B表中不存在(因为B表已排序,后面的key只会更大)。 - 移动
ptrA,跳过A表中这条不匹配的记录。
- 这意味着
keyA > keyB(B 表的 Key 较小):- 同理,这意味着
B表中当前这个keyB在A表中不存在。 - 移动
ptrB,跳过B表中这条不匹配的记录。
- 同理,这意味着
- 终止: 当任一指针到达其数据流的末尾时,整个 Merge 过程结束。
4. SMJ 的核心优势:内存效率与稳定性
- 流式处理,内存友好: Merge 阶段是流式的,它不需要在内存中构建任何一张表的完整哈希表。它只需要在内存中保留当前正在处理的 Key 及其关联的少量行。一旦某个 Key 处理完毕,相关数据即可被垃圾回收。这使得 SMJ 对 Executor 内存的要求非常低,极大地降低了 Out Of Memory (OOM) 的风险。
- 对数据倾斜有更好的容忍度: 即使某个 Key 的数据量特别大(数据倾斜),SMJ 也能稳定处理,因为它不会因为单个 Key 过大而导致内存崩溃。它只会花费更多的时间来处理这个倾斜 Key 的所有匹配项。相比之下,Hash Join 在遇到严重数据倾斜时,构建的哈希表可能会撑爆内存。
5. 终极优化:跳过 Shuffle 和 Sort
SMJ 的性能瓶颈在于 Shuffle 和 Sort。如果数据在 物理存储上 就已经是按照 Join Key 分区和排序的,Spark 的优化器就能识别到这一点,并 完全跳过 Shuffle 和 Sort 阶段,直接执行 Merge 操作。这是 SMJ 性能的极限状态,通常通过使用 clustered by创建 Hive 表或 Delta Lake 表来实现。
cartesian join
使用场景为等值/非等值,仅支持内连接(inner join),没有明确指明关联键;
broadcast nested loop join
像 BHJ 或 Sort Merge Join (SMJ) 这种高效策略,都非常依赖等值条件(即必须有 a.id = b.id)来建立 Hash 表或进行排序比较。 但如果 SQL 条件是非等值连接(Non-Equi Join)(例如:a.date > b.date 或 a.val BETWEEN b.min_val AND b.max_val),无法建立 Hash 表,Spark 就必须使用 BNLJ 通过双重遍历来逐一校验条件。broadcast nested loop join的效率最低
与 Broadcast Hash Join 的核心区别:
- BNLJ: 对于大表的每一行,都要 遍历扫描 小表,效率最低
- BHJ: 对于大表的每一行,只需在小表构建的 哈希表中进行一次 O(1) 的快速查找
执行流程:
plaintext[ Driver 端 ] └─ 收集表 B (Small Table),并广播发送至所有 Executor 内存 [ 各 Executor 节点内部 ] 大表分区数据 (Table A Partition) 广播的表 B 副本 (Broadcasted Table B) Row A1 ─────────────┬──────────────────> [ Row B1, Row B2, Row B3, ... Row BN ] │ (遍历表 B 的所有 N 条记录,逐一校验条件) Row A2 ─────────────┼──────────────────> [ Row B1, Row B2, Row B3, ... Row BN ] │ (遍历表 B 的所有 N 条记录,逐一校验条件) Row A3 ─────────────┘ ...
伪代码:
plaintext# 每个 Executor 上的执行过程 for row_A in local_partition_A: # 外层循环:遍历本地大表分区的每行 for row_B in broadcasted_table_B: # 内层循环:遍历广播过来的整个小表 if evaluate_condition(row_A, row_B): # 校验任意复杂的 Join 条件 emit(combine(row_A, row_B))
影响 spark 选择 join 策略关键因素
Spark 在选择 Join 策略时,主要依据三个核心维度:
-
Join 条件类型:是等值连接(
a.id = b.id)还是非等值连接(a.time > b.time或无条件)。 -
数据量大小与 Hint:是否有一方足够小(默认 10MB)可以广播,或者是否写了 Hint。
-
数据可排序性:Join Key 是否支持排序。
spark 的 join 策略选择优先级:
源码注释说明:
plaintext// If it is an equi-join, we first look at the join hints w.r.t. the following order: // 1. broadcast hint: pick broadcast hash join if the join type is supported. If both sides // have the broadcast hints, choose the smaller side (based on stats) to broadcast. // 2. sort merge hint: pick sort merge join if join keys are sortable. // 3. shuffle hash hint: We pick shuffle hash join if the join type is supported. If both // sides have the shuffle hash hints, choose the smaller side (based on stats) as the // build side. // 4. shuffle replicate NL hint: pick cartesian product if join type is inner like. // // If there is no hint or the hints are not applicable, we follow these rules one by one: // 1. Pick broadcast hash join if one side is small enough to broadcast, and the join type // is supported. If both sides are small, choose the smaller side (based on stats) // to broadcast. // 2. Pick shuffle hash join if one side is small enough to build local hash map, and is // much smaller than the other side, and `spark.sql.join.preferSortMergeJoin` is false. // 3. Pick sort merge join if the join keys are sortable. // 4. Pick cartesian product if join type is inner like. // 5. Pick broadcast nested loop join as the final solution. It may OOM but we don't have // other choice. // ============================================================================================== // If it is not an equi-join, we first look at the join hints w.r.t. the following order: // 1. broadcast hint: pick broadcast nested loop join. If both sides have the broadcast // hints, choose the smaller side (based on stats) to broadcast for inner and full joins, // choose the left side for right join, and choose right side for left join. // 2. shuffle replicate NL hint: pick cartesian product if join type is inner like. // // If there is no hint or the hints are not applicable, we follow these rules one by one: // 1. Pick broadcast nested loop join if one side is small enough to broadcast. If only left // side is broadcast-able and it's left join, or only right side is broadcast-able and // it's right join, we skip this rule. If both sides are small, broadcasts the smaller // side for inner and full joins, broadcasts the left side for right join, and broadcasts // right side for left join. // 2. Pick cartesian product if join type is inner like. // 3. Pick broadcast nested loop join as the final solution. It may OOM but we don't have // other choice. It broadcasts the smaller side for inner and full joins, broadcasts the // left side for right join, and broadcasts right side for left join.
| Join 策略 | Join 条件限制 | 表大小要求 | 触发条件 / 参数 | 适用场景 |
|---|---|---|---|---|
| 1. Broadcast Hash Join (BHJ) | 必须是等值 | 任意表 + 极小表 | 小表大小 < 10MB 或 BROADCAST Hint | 小表 join 大表(最优先选) |
| 2. Shuffle Hash Join (SHJ) | 必须是等值 | 大表 + 中等表 | preferSortMergeJoin=false 且中等表能构建 Partition Hash Table | 中表 join 大表(或数据无序且不想 Sort) |
| 3. Sort Merge Join (SMJ) | 必须是等值 | 大表 + 大表 | 默认策略(preferSortMergeJoin=true),且无法广播 | 海量数据 / 大表 join 大表(最稳健) |
| 4. Cartesian Join (CJ) | 支持非等值 | 大表 + 大表 | 仅 Inner-like Join,无广播小表,显式开启 crossJoin | 大表间的非等值 Join / 交叉连接 |
| 5. Broadcast Nested Loop Join (BNLJ) | 支持任意条件 | 任意表 + 小表 | 兜底策略,或非等值 Join 下有小表可广播 | 小表非等值 Join / 兜底保底方案 |
-
等值连接
BHJ(有小表/Hint)--》 SMJ(默认首选)--》 SHJ(关闭SMJ偏好)--》 Cartesian Join --》 BNLJ(兜底)
-
非等值连接
BNLJ(有广播Hint)--》 Cartesian Join (仅Inner/大表)--》 BNLJ(最终兜底)
Spark 的通信原理
Driver 与 Executor 之间的通信依靠 RPC(远程过程调用)实现。早期版本使用 Akka 框架,后期版本(2.x 及以后)则完全迁移到了 Netty 框架。
Spark 的任务调度流程和 YARN 的资源调度
Task 调度要点:
- 任务的切分:将不同的 RDD 按照宽依赖划分成 Stage,每个 Stage 再按照分区数划分为等量的 Task。
- Task 的调度和执行:将 Task 分配到 Executor 上运行。
调度细节:
- 用户代码在执行时遇到行动算子,会触发一个 Job 的提交。
- DAGScheduler 接收到 Job 后,会根据 RDD 的血缘关系从后往前遍历,每遇到一个宽依赖就划分出一个 Stage。Stage 内部的任务(Task)数量由该 Stage 最后一个 RDD 的分区数决定。这些 Task 被打包成一个 TaskSet。
- TaskScheduler 接收到 TaskSet 后,通过 SchedulerBackend 向资源管理器(如 YARN)申请资源。
- 在资源充足的情况下,TaskScheduler 将 Task 发往不同的 Executor 执行。
- TaskScheduler 在调度 Task 时,有两种策略:FIFO(先进先出)和 FAIR(公平调度),默认使用 FIFO。 FIFO 是为了保证划分的 Stage 能按照依赖顺序执行 。
调度的原则:
遵循数据本地化(Data Locality)的原则,尽量将计算任务分配到数据所在的节点上,以减少网络传输开销。调度优先级为:PROCESS_LOCAL > NODE_LOCAL > RACK_LOCAL > ANY。
Spark 优化
Hive on Spark 的优化
针对 Hive 的优化策略((参考 03_HIVE.md 中的 "Hive 框架优化"))对 Hive on Spark 大部分都有效。 例如,当通过 Spark 历史服务器发现某个任务的 GC time 过长时,可以调整 Spark 内存模型中 "Other" 部分的内存大小。"Other" 部分内存用于存储 Spark 内部的元数据和用户自定义的一些数据结构。
单独对 Spark 的优化
通过合理设置 Spark Application 的运行参数,可以最大程度地提高资源利用效率。
参数配置示例
假设集群中每台节点配置为 128G 内存,64 个 CPU 核心。
- 资源分配给 YARN:通常将节点 80%的资源分配给 YARN。则每台节点可用资源为
128G * 0.8 ≈ 100G内存,64 * 0.8 ≈ 48个核心。- Executor 配置:每个 Executor 分配的 core 数不宜过多,一般为 3-6 个。这里我们设置为 6 个 core。
- 计算 Executor 数量:每台 NodeManager 可以启动
48 / 6 = 8个 Executor。- 计算 Executor 内存:每个 Executor 可分配的内存为
100G / 8 ≈ 12G。- 总 Executor 数:假设集群共有 7 台 NodeManager,则总共可启动
7 * 8 = 56个 Executor。最终配置参数:
bash--num-executors 56 --executor-cores 6 --executor-memory 12G --driver-cores 6 --driver-memory 12G
一些其他常用运行参数:
--master:指定将任务提交到哪个资源调度器(如yarn)。--class:指定待运行的带有main方法的全类名。--deploy-mode:指定 YARN 作为资源调度器时的部署模式(client/cluster)。--queue:指定任务提交到 YARN 的哪个资源队列中。
版权声明