04Spark

约 27 分钟 · 7338 字评论
请注意,本文编写于 1345 天前,最后修改于 43 天前,其中某些信息可能已经过时。

Spark

Spark 的几种部署模式比较

image-20250802144726609

  1. local 模式:单台计算机部署,用于测试。

  2. standalone 模式(单集群模式):需要单独配置 Spark 的集群,使用 Spark 自带的资源调度器,即资源和任务分配都由 Spark 自己负责。

  3. 集群模式

    • 存算分离,比如使用 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 的基本概念

  1. 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 上的位置,并遵循“移动计算而非移动数据”的原则进行任务调度。
  2. RDD 的弹性体现

    • _分区的弹性_:每个 RDD 的分区数量是由数据切片的数量决定的,处理不同的数据时,分区会动态变化。
    • _计算的弹性_:每个 RDD 都记录了需要处理数据切片的位置(例如在 HDFS 上的位置)。在任务分配时,会尽量遵循本地化原则(移动计算而非移动数据)来分配计算任务。
    • _容错的弹性_:每个 RDD 都记录了对其他 RDD 的依赖关系(血缘)。当某个分区的数据丢失时,可以根据血缘关系自动重试任务,恢复分区数据。
    • _存储的弹性_:RDD 在计算过程中,会优先使用内存来存储中间数据。当内存资源不足时,会自动将数据溢写到磁盘。数据存储位置是弹性的。
  3. 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 的区别

  1. Spark 需要做数据持久化的原因

    • 一个 RDD 在多个 Job 中被重复使用。
    • 一个 Job 的数据处理链路特别长,如果不进行持久化,一旦发生错误,重算成本很高,容错率会降低。
  2. 缓存和 Checkpoint 的比较

持久化方式使用方式数据保存位置RDD 依赖关系是否保留
缓存rdd.cache()rdd.persist()内存或者磁盘
Checkpointrdd.checkpoint()HDFS 等共享存储否(会切断血缘关系)
  • cachepersist 的区别cache()persist() 的一种特殊情况,等价于 persist(StorageLevel.MEMORY_ONLY)
  • 常用的存储级别:
    • StorageLevel.MEMORY_ONLY:一般用于小数据量场景。
    • StorageLevel.MEMORY_AND_DISK:一般用于大数据量场景。

任务划分:Application、Job、Stage、Task 关系

  1. 切分 Stage 的流程: 根据最后一个调用 Action 算子的 RDD,从后往前沿着其依赖关系(血缘)进行回溯,一直到 Job 的第一个 RDD。在回溯过程中,每遇到一个宽依赖,就切分出一个新的 Stage。
  2. _Application数量_ = SparkContext 的个数。一个 main 方法中通常只有一个 SparkContext,即一个 Application。
  3. _Job的个数_ 由 Action 算子(行动算子)的个数决定。每当代码执行到一个 Action 算子,就会触发一个 Job 的提交。常见的 Action 算子有 collect()foreach()count()take() 等。
  4. _Stage的个数_ = 宽依赖的个数 + 1。
  5. _Task的个数_(在一个 TaskSet 中)= 一个 Stage 中,最后一个 RDD 的分区数。

常用的转换算子和行动算子,及其比较

GitHub 示例代码地址:GitHub 相关的例子

  1. Transform 算子(转换算子) 转换算子是惰性执行的,它们不会立即触发 Job 的执行,只负责定义 RDD 之间的转换关系。只有当遇到 Action 算子时,这些转换才会真正执行。

    plaintext
    1. map 2. flatMap 3. reduceByKey(预聚合算子) 4. groupByKey 5. distinct
    • reduceByKeygroupByKey 的比较reduceByKey 相比 groupByKey,在 Shuffle 之前会先在 Map 端进行预聚合(Combiner),从而减少了 Shuffle 传输的数据量,性能更高。应优先使用 reduceByKey
    • mapflatMap 的比较flatMapmap 更灵活,对于输入的每个元素,map 只会产生一个输出元素,而 flatMap 可以产生 0 个、1 个或多个输出元素。
  2. Action 算子(行动算子) 行动算子会触发 Spark Job 的实际计算和执行。

    plaintext
    1. collect 2. count 3. foreach 4. foreachPartition 5. **reduce**
    • foreachforeachPartition 的比较foreach 会对 RDD 中每个分区里的每一条数据都执行一次指定的函数,函数的调用次数等于数据总数。而 foreachPartition 则是对 RDD 的每个分区执行一次指定的函数,函数的调用次数等于分区数。在需要进行数据库连接等昂贵操作时,使用 foreachPartition 可以显著减少连接的创建和销毁次数,提高执行效率。

共享变量

分布式只读共享变量

广播变量:默认情况下,Spark 会将 Driver 端的数据给每个 Task 都发送一份,所以该数据占用的内存大小 = Task 个数 * Driver 数据大小。但是如果将该 Driver 数据广播,此时数据会广播给每个 Executor,该 Executor 中的所有 Task 共用这一份数据,所以此时该数据占用的内存大小 = Executor 个数 * Driver 数据大小,相比之前节省了内存空间。

分布式只写共享变量

累加器:先在 RDD 的每个分区中进行累加,然后将各个分区的累加结果返回给 Driver 进行汇总。

SparkSQL

hive on sparkspark 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 中常见的数据抽象及其区别

  1. RDD, DataFrame, DataSet

    • RDD(弹性分布式数据集)是 Spark 的基础数据结构,它关注于对数据的细粒度处理,提供了丰富的算子,但缺乏结构化信息。
    • DataFrame 为数据引入了 Schema(模式),类似于关系型数据库中的表。它将每条数据封装为 Row 对象,使 Spark 能更高效地处理结构化数据,并进行查询优化。但 DataFrame 是弱类型的,在编译时不做类型检查。
    • DataSet 结合了 RDD 的类型安全和 DataFrame 的性能优势。它既有强类型检查(面向对象),也保留了 DataFrame 的非类型化操作接口,是 Spark 推荐使用的数据结构。

SparkSession 的概念

SparkSession 是 Spark 2.0 以后进行 Spark 编程的统一入口点,它整合了 SparkContextSQLContextHiveContext 的功能,是与 Spark 进行交互的起点。

SparkCore

RDD 的序列化

  • RDD 需要序列化的原因:Spark 算子(如 map)的 call 方法中的代码是在 Executor 的 Task 中执行的,而 call 方法外部的代码是在 Driver 中执行的。如果 Task 中需要引用在 Driver 端创建的对象,Spark 会将该对象序列化后传输给 Task 使用。因此,这个对象及其所有字段都必须是可序列化的。
  • 如何序列化 RDD:Spark 默认使用 Java 序列化,但为了获得更好的性能,推荐使用 Kryo 序列化库。

Spark 的 Shuffle

数据的重新分配和收集可以称为 shuffle

plaintext
MR计算过程: 数据 -> InputFormat -> map方法 -> 环形缓冲区[分区、排序, 80%溢写] -> [Combiner] -> 磁盘[小文件] -> 合并小文件 -> Reducer拉取数据[归并排序] -> reduce方法 -> outputFormat -> 磁盘 MR shuffle阶段: -> 环形缓冲区[分区、排序, 80%溢写] -> [Combiner] -> 磁盘[小文件] -> 合并小文件 -> Reducer拉取数据[归并排序] -> Spark shuffle: -> 缓冲区[分区、[排序]] -> [Combiner] -> 磁盘[小文件] -> 合并小文件 -> 分区拉取数据[归并、[排序]] ->

MR 的 Shuffle

spark 几种 shuffle 机制

  1. hashshuffle(未优化)
  2. hashshuffle(Consolidated 优化)
  3. sortshuffle
    • bypassh shuffle
    • unsafe shuffle(Tungsten shuffle)
  4. AQE 中的动态 shuffle
  5. 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 瓶颈,节点宕机容错等问题

task-ms4evm7h14iiz

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 数) 的场景基准表示如下:

plaintext
Map 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 相关层级关系如下:

plaintext
YARN 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情况

plaintext
Cores & 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 个物理文件

中间文件数 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文件

plaintext
Map 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 个文件

中间文件数 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中的排序操作

plaintext
Map 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 个最终物理文件

中间文件数 M *= 2×1000 = 2000 个物理文件

GC 情况: 相比 Sort Shuffle 省去了排序数据结构的内存占用,Young GC 压力有所缓解。

unsafe shuffle

对于不满足使用bypass shuffle的场景,sort shuffle中的排序操作是必要的,为了减轻堆内 Java 对象序列化开销大(需要反序列化为java对象,才可排序),GC压力(排序操作带来的),采用直接向操作系统请求操作序列化的数据,对数据的指针进行排序操作,绕过了jvm堆,gc问题得以减轻

plaintext
Off-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 个文件

中间文件数 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 执行)。

    plaintext
    Map 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 的数据预先合并为大块顺序文件

    plaintext
    Map 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 区别

  1. Hadoop 的 MapReduce Shuffle 无论任务是否需要排序,都必须进行排序。而 Spark 的 BypassShuffle 在满足一定条件时可以跳过排序,提高效率。
  2. 在 MapReduce 中,MapTask 的环形缓冲区大小是固定的。而 Spark 中 Map 阶段的缓冲区是动态的,可以根据内存使用情况进行扩容。
  3. 在 Hadoop 中,ReduceTask 在部分 MapTask 完成后就可以开始启动并拉取数据。而在 Spark 中,必须等待前一个 Stage 的所有 Task 全部执行完成后,下一个 Stage 才能开始执行。

Spark 比 MR 快的原因

  1. 内存计算:Spark 优先使用内存进行数据处理和存储(默认存储级别为 MEMORY_ONLY),减少了磁盘 I/O,而 MapReduce 在任务传递过程中必须依赖磁盘进行 Shuffle。
  2. 执行模型:Spark 的任务执行是基于线程的(一个 Executor 中可以运行多个 Task 线程),而 MapReduce 的任务是基于进程的(每个 MapTask 或 ReduceTask 都是一个独立的 JVM 进程)。线程的创建和调度开销通常远小于进程。
  3. DAG 优化:Spark 使用有向无环图(DAG)来表示计算流程,可以进行更复杂的优化,例如将多个操作流水线化(Pipelining)在一个 Stage 中执行,避免了不必要的中间数据落地。

Repartition 和 Coalesce 区别

  1. repartitioncoalesce 都有改变分区数的功能。repartitioncoalesce 的一种特殊情况,即 repartition(n) 等价于 coalesce(n, shuffle = true)
  2. repartition 通常用于 增加 分区数。它会进行一次完整的 Shuffle,将数据随机打散并重新分布到所有分区中。
  3. coalesce 通常用于 减少 分区数。默认情况下,它会尽量避免 Shuffle,通过合并相邻的分区来实现,效率更高。如果需要将数据均匀分布,可以设置 shuffle = true

Spark 的内存管理

Spark 1.6 版本后引入了统一内存管理(Unified Memory Management),Spark 可以同时使用堆内内存(On-heap)和堆外内存(Off-heap)。

执行内存(Execution Memory)和存储内存(Storage Memory)的动态占用机制 这两块内存区域可以相互借用。当一方空闲而另一方资源不足时,空闲方可以将部分内存借给对方使用。但是,执行内存有优先权,如果存储内存借用了执行内存,当执行内存需要时,存储内存必须归还,可能会导致缓存的数据被移出内存。

Spark 的持久化保存

见“缓存(cache,persist),checkpoint 的区别”一节

  1. 数据持久化原因
  2. 使用方式
  3. 比较

Spark-YARN 的 Cluster 模式的任务提交流程

sequenceDiagram autonumber actor User as 用户 (Client/Script) participant Submit as SparkSubmit 进程 participant HDFS as HDFS 存储 participant RM as YARN ResourceManager participant NM as YARN NodeManager participant AM as AM (Spark Driver) participant Exec as Spark Executor %% 1. 应用提交 User->>Submit: 执行 spark-submit 命令 activate Submit Submit->>HDFS: 上传 应用Jar包, 配置文件, Task切片 Submit->>RM: 发送 submitApplication 请求 deactivate Submit %% 2. AM (Driver) 启动与初始化 RM->>NM: 指令: 分配 Container 并启动 AM activate NM NM->>AM: 创建容器并启动 ApplicationMaster deactivate NM activate AM Note over AM: AM 进程启动内部 Driver 线程<br/>初始化 SparkContext, DAGScheduler, TaskScheduler %% 3. Executor 资源申请与启动 AM->>RM: registerApplicationMaster (申请 Executors 资源) RM-->>AM: 返回可用的 Containers 节点列表 AM->>NM: 请求在指定 Container 中启动 Executor 进程 activate NM NM->>Exec: 创建 Container 并启动 Executor 进程 deactivate NM activate Exec %% 4. Executor 反向注册与心跳 Exec->>AM: 通过 RPC 主动向 Driver 进行反向注册 (RegisterExecutor) AM-->>Exec: 注册成功响应 (RegisteredExecutor) loop 持续心跳与度量上报 Exec-->>AM: 发送 Heartbeat & Task Metrics end %% 5. 任务划分与分发执行 Note over AM: 遇到 Action 算子,触发 Job<br/>1. DAGScheduler 按宽依赖划分为 Stages<br/>2. TaskScheduler 组合 TaskSet par 并行发送 Tasks 到各个 Executor AM->>Exec: 发送 Task 描述及依赖 (Launch Task) Note over Exec: 提取线程池中的 Thread<br/>在内存中执行数据处理流水线 (Pipelining) Exec-->>AM: 异步返回 Task 执行结果 (StatusUpdate) end %% 6. 结束与资源释放 Note over AM: 所有 Tasks 完成,主函数 run() 结束 AM->>RM: unregisterApplicationMaster (注销应用) deactivate AM RM->>NM: 回收所有 Container 资源 NM->>Exec: 强行终止 (Kill) Executor 进程 deactivate Exec

提交重点:

  1. _应用的提交_:客户端解析任务的提交参数。
  2. _执行环境的准备_:在 ApplicationMaster 中建立 Driver 线程,并初始化 SparkContext 的执行环境,包括任务的切分和调度。
  3. _任务的调度和执行_:Driver 向 Executor 发送启动任务的命令,Executor 执行具体的 Task。

具体流程:

  1. 用户通过 spark-submit 脚本提交应用。Submit 进程解析提交参数,并封装成一个 YARN 应用。
  2. Submit 进程与 YARN 的 ResourceManager 通信,申请启动一个 ApplicationMaster(AM)。
  3. ResourceManager 在一个 NodeManager 上启动 ApplicationMaster。AM 内部会启动 Driver 线程,并初始化 Spark 的执行环境(如 DAGScheduler、TaskScheduler)。
  4. ApplicationMaster(作为 Driver 的一部分)向 ResourceManager 申请计算资源(Container)。
  5. ResourceManager 分配 Container 后,AM 在这些 Container 中启动 Executor 进程。
  6. Executor 启动后,会反向注册到 Driver,并与 Driver 建立心跳连接。至此,资源申请和执行环境准备完成。
  7. Driver 中运行用户的 main 方法。当遇到 Action 算子时,触发一个 Job 的提交。DAGScheduler 根据 RDD 的宽依赖关系,将 Job 划分为一个或多个 Stage。每个 Stage 被划分为一组 Task(TaskSet)。
  8. TaskScheduler 将 TaskSet 中的 Task 分配给可用的 Executor 执行,并监控任务的执行情况。
  9. 当所有 Task 执行完成,ApplicationMaster 向 ResourceManager 取消注册,YARN 回收所有资源,应用结束。

spark 的 join 选择策略

spark 中提供非常丰富的 spark join 策略来处理不同表数据量情况下的选择

  1. broadcast hash join
  2. shuffle hash join
  3. sort merge join
  4. cartesian join
  5. broadcast nested loop join

broadcast hash join

涉及广播,hash 表的形成,数据间的 join;使用场景为 小表 join 大表等值连接 场景,spark 默认小表标准为 10MB;driver 将小表数据 广播到各个存储大表数据的 executor,executor 将小表数据利用哈希函数形成一张《joinkey,小表行数据》的哈希表,遍历大表每一行,从哈希表取得小表数据,哈希读取效率高

image-20250802160226763

出发BHJ的条件:

  • join的连接必须是等值连接(因为hash table的形成依赖形成的hash值精确)
  • 满足小表阈值(可以使用hint语法跳过小表阈值检查,但是也要满足等值连接才可)

shuffle hash join

当 join 的表数据都超过小表存储的定义(10MB),数据数据量适中等值连接 情况,spark join 将会采用 shuffle hash join;涉及过程:

  • shuffle:将两个表都按照 joinkey 进行 shuffle,将相同关联 key 的数据发往同一节点的相同分区
  • hash:在相同分区内,取数据量较小的表制为 hash 表,遍历较大表的数据,完成 join
graph TD subgraph "初始状态: 两个大表分布在不同分区" A_P1["<b>Table A - Partition 1</b><br/>(k1, valA1)<br/>(k2, valA2)"] A_P2["<b>Table A - Partition 2</b><br/>(k1, valA3)<br/>(k3, valA4)"] B_P1["<b>Table B - Partition 1</b><br/>(k1, valB1)<br/>(k3, valB2)"] B_P2["<b>Table B - Partition 2</b><br/>(k2, valB3)<br/>(k3, valB4)"] end subgraph "第一阶段: Shuffle (按 Join Key 重新分区)" direction LR Shuffle_Process["<font size=5><b>Shuffle</b></font><br/>按 Key (k1, k2, k3...)<br/>重新分发数据"] end A_P1 --> Shuffle_Process A_P2 --> Shuffle_Process B_P1 --> Shuffle_Process B_P2 --> Shuffle_Process subgraph "Shuffle 后: 相同 Key 的数据汇集到同一分区" New_P1["<b>Partition X (Executor 1)</b><br/><i>--- Data for Key k1 ---</i><br/>from A: (k1, valA1), (k1, valA3)<br/>from B: (k1, valB1)"] New_P2["<b>Partition Y (Executor 2)</b><br/><i>--- Data for Key k2 ---</i><br/>from A: (k2, valA2)<br/>from B: (k2, valB3)"] New_P3["<b>Partition Z (Executor 3)</b><br/><i>--- Data for Key k3 ---</i><br/>from A: (k3, valA4)<br/>from B: (k3, valB2), (k3, valB4)"] end Shuffle_Process --> New_P1 Shuffle_Process --> New_P2 Shuffle_Process --> New_P3 subgraph "第二阶段: Hash Join (在单个分区内执行)" subgraph "在 Partition Z 内部" Build["<b>1. 构建哈希表</b><br/>(用较小的表A的数据)"] HashTable["<b>Hash Table</b><br/>{ k3 -> (k3, valA4) }"] Build --> HashTable Probe["<b>2. 遍历大表B的数据 (Probe)</b><br/>(k3, valB2)<br/>(k3, valB4)"] Probe -- "用 Key 'k3' 查找" --> HashTable Result["<b>3. 输出Join结果</b><br/>(k3, valA4, valB2)<br/>(k3, valA4, valB4)"] HashTable --> Result end end New_P3 --> Build style A_P1 fill:#cde4ff style A_P2 fill:#cde4ff style B_P1 fill:#d5e8d4 style B_P2 fill:#d5e8d4 style New_P1 fill:#f5f5f5,stroke:#333,stroke-width:2px style New_P2 fill:#f5f5f5,stroke:#333,stroke-width:2px style New_P3 fill:#fff2cc,stroke:#d6b656,stroke-width:2px,color:#000 style Shuffle_Process fill:#dae8fc,stroke:#6c8ebf,stroke-width:2px

Spark Sort Merge Join

Sort Merge Join (SMJ) 是 Spark 中处理 大表等值 Join 的核心策略。当参与 Join 的两张表都过大,无法利用广播机制时,SMJ 提供了一种稳健且内存友好的解决方案。虽然它通常被视为 Spark 的默认 Join 策略(由 spark.sql.join.preferSortMergeJoin 控制),但在现代 Spark 版本中,其最终应用会受到自适应查询执行(AQE)的动态调整。

SMJ 的整个过程可以分为三个核心阶段:ShuffleSortMerge 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 对内存的要求极低,非常适合处理海量数据。

graph TD subgraph "Phase 1: Initial State (Scattered Data)" A1["Table A<br/>Partition 1"] A2["Table A<br/>Partition 2"] B1["Table B<br/>Partition 1"] B2["Table B<br/>Partition 2"] end subgraph "Phase 2: Shuffle (Group by Join Key)" SHUFFLE["<font size=5><b>SHUFFLE</b></font><br/>Move data so same keys<br/>are in the same partition"] end A1 & A2 & B1 & B2 --> SHUFFLE subgraph "Phase 3: Sort (Inside Each Partition)" P1_unsorted["<b>Partition X (Unsorted)</b><br/>All data for Keys 1, 5, 2"] P2_unsorted["<b>Partition Y (Unsorted)</b><br/>All data for Keys 8, 4, 6"] SORT1["<font size=4><b>SORT</b></font><br/>by Join Key"] SORT2["<font size=4><b>SORT</b></font><br/>by Join Key"] P1_sorted["<b>Partition X (Sorted)</b><br/><i>Table A part:</i> [k1, k2, k5]<br/><i>Table B part:</i> [k1, k5]"] P2_sorted["<b>Partition Y (Sorted)</b><br/><i>Table A part:</i> [k4, k6]<br/><i>Table B part:</i> [k4, k6, k8]"] P1_unsorted --> SORT1 --> P1_sorted P2_unsorted --> SORT2 --> P2_sorted end SHUFFLE --> P1_unsorted SHUFFLE --> P2_unsorted subgraph "Phase 4: Merge Join (Pointer-based Iteration)" subgraph "Focus on Partition X" DataA["<b>Sorted A data</b><br/>[ (k1, vA1), (k2, vA2), (k5, vA3) ]<br><font color=blue><b>↑</b></font><br>Pointer A"] DataB["<b>Sorted B data</b><br/>[ (k1, vB1), (k5, vB2) ]<br><font color=green><b>↑</b></font><br>Pointer B"] MERGE["<b>3. Merge & Join</b><br/>Advance pointers and match keys"] DataA --> MERGE DataB --> MERGE Result["<font size=4><b>Joined Result</b></font><br/>(k1, vA1, vB1)<br/>(k5, vA3, vB2)"] MERGE --> Result end end P1_sorted --> DataA P1_sorted --> DataB %% Styling style A1 fill:#cde4ff,stroke:#6c8ebf style A2 fill:#cde4ff,stroke:#6c8ebf style B1 fill:#d5e8d4,stroke:#82b366 style B2 fill:#d5e8d4,stroke:#82b366 style SHUFFLE fill:#dae8fc,stroke:#6c8ebf,stroke-width:2px,stroke-dasharray: 5 5 style P1_unsorted fill:#f5f5f5,stroke:#333 style P2_unsorted fill:#f5f5f5,stroke:#333 style P1_sorted fill:#fff2cc,stroke:#d6b656,stroke-width:2px style P2_sorted fill:#fff2cc,stroke:#d6b656,stroke-width:2px style SORT1 fill:#e1d5e7,stroke:#9673a6 style SORT2 fill:#e1d5e7,stroke:#9673a6 style Result fill:#d5e8d4,stroke:#82b366,stroke-width:2px

1. SMJ 的触发条件:何时被选中?

SMJ 并非总是被选中,它的应用有明确的前提条件:

  • key能够排序:由于smj需要按照join key来进行排序方便后续的merge join
  • 等值 Join: Join 条件必须包含至少一个等值表达式(例如 ON a.id = b.id)。对于非等值 Join(如 a.id > b.id),Spark 只能使用 BroadcastNestedLoopJoinCartesianProduct,效率极低。
  • 无法广播: 这是关键。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 为每个分区内的两个已排序数据流各维护一个指针(可以想象成两个手指分别指着两份名单)。

  1. 初始化: 两个指针 ptrAptrB 分别指向两个已排序数据流的起始位置。
  2. 循环比较: 进入一个循环,不断比较两个指针当前位置的 Join Key (keyAkeyB):
    • keyA == keyB (匹配成功):
      • 找到了匹配项!将两行数据合并,生成一条 Join 结果。
      • 处理一对多: 此时,需要处理可能存在的“一对多”关系。Spark 会缓存 ptrB 当前指向的行,然后 仅移动 ptrA,检查 A 流中是否还有相同的 keyA。只要有,就持续与缓存的 B 行进行匹配并输出结果,直到 A 流的 key 发生变化。然后,再移动 ptrB
    • keyA < keyB (A 表的 Key 较小):
      • 这意味着 A 表中当前这个 keyAB 表中不存在(因为 B 表已排序,后面的 key 只会更大)。
      • 移动 ptrA,跳过 A 表中这条不匹配的记录。
    • keyA > keyB (B 表的 Key 较小):
      • 同理,这意味着 B 表中当前这个 keyBA 表中不存在。
      • 移动 ptrB,跳过 B 表中这条不匹配的记录。
  3. 终止: 当任一指针到达其数据流的末尾时,整个 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.datea.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))
graph TD subgraph "Phase 1: Initial State & Broadcast" BigTable["<b>Big Table (Distributed)</b><br/>Partition 1<br/>Partition 2<br/>..."] SmallTable["<b>Small Table</b><br/>(on Driver)"] SmallTable -- "1. Collect to Driver" --> Driver["<font size=5><b>Driver</b></font>"] Driver -- "2. Broadcast to All Executors" --> Executor1 & Executor2 end subgraph "Phase 2: Nested Loop Join (Inside Each Executor)" subgraph "Focus on Executor 1" P1["<b>Partition 1 of Big Table</b><br/>(row A1)<br/>(row A2)<br/>(row A3)<br/>..."] Broadcasted_SmallTable["<b>Broadcasted Small Table (in Memory)</b><br/>(row B1)<br/>(row B2)<br/>..."] Loop["<font size=5><b>Nested Loop</b></font><br/>For each row in Big Table,<br/>iterate through <b>ALL</b> rows in Small Table"] P1 --> Loop Broadcasted_SmallTable --> Loop Check["<b>Join Condition Check</b><br/>e.g., big.date > small.date<br/>(N x M comparisons)"] Loop --> Check Result["<font size=4><b>Joined<br/>Result</b></font>"] Check -- "If True" --> Result end Executor2["<b>Executor 2</b><br/>(Performs the same<br/>nested loop on its partition)"] end BigTable --> P1 %% Styling style BigTable fill:#cde4ff,stroke:#6c8ebf style SmallTable fill:#d5e8d4,stroke:#82b366 style Driver fill:#f8cecc,stroke:#b85450,stroke-width:2px style Executor1 fill:#fff2cc,stroke:#d6b656,stroke-width:2px style Executor2 fill:#f5f5f5,stroke:#333 style Loop fill:#f8cecc,stroke:#b85450,stroke-width:2px,stroke-dasharray: 5 5 style Check fill:#e1d5e7,stroke:#9673a6 style Result fill:#d5e8d4,stroke:#82b366,stroke-width:2px

影响 spark 选择 join 策略关键因素

Spark 在选择 Join 策略时,主要依据三个核心维度:

  1. Join 条件类型:是等值连接a.id = b.id)还是非等值连接a.time > b.time 或无条件)。

  2. 数据量大小与 Hint:是否有一方足够小(默认 10MB)可以广播,或者是否写了 Hint。

  3. 数据可排序性: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 / 兜底保底方案
  1. 等值连接

    BHJ(有小表/Hint)--》 SMJ(默认首选)--》 SHJ(关闭SMJ偏好)--》 Cartesian Join --》 BNLJ(兜底)

  2. 非等值连接

    BNLJ(有广播Hint)--》 Cartesian Join (仅Inner/大表)--》 BNLJ(最终兜底)

Spark 的通信原理

Driver 与 Executor 之间的通信依靠 RPC(远程过程调用)实现。早期版本使用 Akka 框架,后期版本(2.x 及以后)则完全迁移到了 Netty 框架。

Spark 的任务调度流程和 YARN 的资源调度

Task 调度要点:

  • 任务的切分:将不同的 RDD 按照宽依赖划分成 Stage,每个 Stage 再按照分区数划分为等量的 Task。
  • Task 的调度和执行:将 Task 分配到 Executor 上运行。

调度细节:

  1. 用户代码在执行时遇到行动算子,会触发一个 Job 的提交。
  2. DAGScheduler 接收到 Job 后,会根据 RDD 的血缘关系从后往前遍历,每遇到一个宽依赖就划分出一个 Stage。Stage 内部的任务(Task)数量由该 Stage 最后一个 RDD 的分区数决定。这些 Task 被打包成一个 TaskSet。
  3. TaskScheduler 接收到 TaskSet 后,通过 SchedulerBackend 向资源管理器(如 YARN)申请资源。
  4. 在资源充足的情况下,TaskScheduler 将 Task 发往不同的 Executor 执行。
  5. 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 核心。

  1. 资源分配给 YARN:通常将节点 80%的资源分配给 YARN。则每台节点可用资源为 128G * 0.8 ≈ 100G 内存,64 * 0.8 ≈ 48 个核心。
  2. Executor 配置:每个 Executor 分配的 core 数不宜过多,一般为 3-6 个。这里我们设置为 6 个 core。
  3. 计算 Executor 数量:每台 NodeManager 可以启动 48 / 6 = 8 个 Executor。
  4. 计算 Executor 内存:每个 Executor 可分配的内存为 100G / 8 ≈ 12G
  5. 总 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 的哪个资源队列中。

版权声明

本文作者:hedeoer

本文链接:/post/__gh__69

许可说明:本博客所有文章除特别声明外,均采用 许可协议。转载请注明出处!

04Spark

0

匿名访客仅支持文本评论。

暂无评论