全文共 5,268 字 预计阅读 16 分钟
bg

bg5.mapreduce->spark

一旦上分布式,就需要自己处理一堆脏活:数据怎么切分、任务怎么分发到各机器、某台机器挂了怎么办、中间结果怎么在机器间传输、如何知道所有机器都算完了。这些和业务逻辑毫无关系,却占了 90% 的代码。 计算引擎要解决的,就是把这些脏活抽象掉。

Google 在 2004 年(Dean & Ghemawat, MapReduce: Simplified Data Processing on Large Clusters, OSDI 2004)提出 MapReduce。它的核心洞察是:海量数据处理的绝大多数场景,都能拆成两个纯函数

map(key, value)     → 产出一批中间 (中间key, 中间value)
reduce(中间key, [中间value...]) → 产出最终结果

中间夹着一步框架自动完成的 shuffle(混洗):把所有 map 产出的、相同中间 key 的数据,搬到同一个 reduce 上。

用 SQL 类比(先给定义再注明近似):一个 SELECT city, SUM(amount) FROM t GROUP BY city 大致对应——

  • map:读每行,产出 (city, amount)
  • shuffle:把相同 city 的记录聚到一起;
  • reduce:对每个 city 的 amount 列表求和。

MapReduce 解决了什么:程序员只写 mapreduce 两个函数,分发、容错、shuffle 全由框架兜底。容错方式很朴素——某个 map/reduce 任务失败,框架在别的机器上重跑这一个任务(输入来自 HDFS,可重读)。这是"重算式容错"的雏形,Spark 后面把它发扬光大。

自问:既然抽象这么优雅,为什么后来被 Spark 取代?

答:贵在 IO,烦在编程模型。

  1. 每个作业的结果都要落 HDFS,多副本写。 一个 MapReduce 作业结束,输出写回 HDFS(默认 3 副本 + 跨网络复制)。而真实业务很少一步算完——一个 SQL 里的 join + 聚合,可能被编译成 3~5 个串行的 MapReduce 作业。每个作业之间都要"写 HDFS → 下个作业再读 HDFS",磁盘 IO 和网络复制成了主瓶颈。
  2. 迭代型作业尤其惨。 机器学习(逻辑回归跑 N 轮)、图算法(PageRank 迭代),每一轮都是一个 MapReduce 作业,每轮都完整地 "读 HDFS → 算 → 写 HDFS"。数据本可以留在内存里复用,MapReduce 却每轮都往返磁盘 N 次。
  3. 编程模型太原始。 什么都得掰成 map/reduce 两个函数,一个稍复杂的逻辑要手工串接多个作业、自己管中间目录,代码冗长、可读性差。

Spark 出自 UC Berkeley AMPLab,2010 年首篇论文 Spark: Cluster Computing with Working Sets(HotCloud 2010)提出思路,2012 年 Resilient Distributed Datasets(RDD 论文,NSDI 2012)正式形式化了核心抽象 RDD血缘(lineage)容错

Spark 针对上面三个痛点分别下药:

痛点 MapReduce Spark 的做法
多阶段落盘 每个作业输出写 HDFS,串行 DAG 把整条流水线一次规划,只在必要处(shuffle)落盘
迭代往返磁盘 每轮读写 HDFS 中间结果 cache/persistExecutor 内存,后续轮直接读内存
编程繁琐 只有 map/reduce RDD/DataFrame 提供 map/filter/join/groupBy... 几十个算子,链式组合
容错靠重跑 + HDFS 输入 任务级重跑 血缘(lineage):记录"这份数据由谁经什么转换得来",丢了按血缘只重算丢失的分区
代价(Spark 不是银弹)
  • 内存贵且有限。中间结果放内存,放不下就 spill(溢写)到磁盘,性能退化回接近 MapReduce。
  • shuffle 依然昂贵——这一步 Spark 也绕不开落盘 + 网络。
  • 所以"Spark 全程在内存所以快 100x"是营销话术,不是普适事实。RDD 论文报告的是:在迭代式机器学习负载上比 Hadoop 快约 20×(NSDI 2012 基准,数字随负载变化)。

纠偏:Spark 快,不是因为"所有数据都在内存",而是两点——① 用 DAG 把多阶段流水线化,省掉了 MapReduce 作业之间反复写/读 HDFS 的落盘;② 允许把可复用的中间数据显式缓存在内存,让迭代不必每轮回磁盘。不缓存、又全是 shuffle 的作业,Spark 也快不到哪去。

Spark 是一个分布式计算引擎:你描述"要对数据做什么",它负责把计算拆分、调度到集群、容错、汇总。它不管数据存哪,也不天生管资源怎么申请,只管

RDD(Resilient Distributed Dataset,弹性分布式数据集) 是 Spark 最底层的核心抽象。一句话定义:RDD 是一个被切成多个分区、分布在集群内存/磁盘上的只读数据集,并记得自己是怎么算出来的。

RDD 的三个关键属性:

  • 分布式 + 分区:一个 RDD 逻辑上是一份数据,物理上是散在各 Executor 上的一批分区。一个分区 = 一个 Task 的处理单位(记住这句,后面反复用)。
  • 只读 + 转换:RDD 不可变。你对它做 mapfilter,不是原地改,而是生成一个新的 RDD。这一串"由谁生成谁"的关系,就是血缘。
  • 弹性(Resilient)= 靠血缘重算容错:见下。

自问:HDFS 靠 3 副本容错(第 03 篇),Spark 的中间数据大多在内存里,一台 Executor 挂了、内存数据没了,怎么办?难道也存 3 份内存副本?

答:不。Spark 存的不是数据副本,是"配方"(lineage 血缘)。 每个 RDD 都记得:我是由哪个(些)父 RDD、经过什么转换得到的。某个分区因为 Executor 崩溃丢了,Driver 就顺着血缘,只重算这一个丢失的分区,不影响其它分区。

因果链与代价:

  • 省了什么:不用为海量中间数据额外存副本(内存/磁盘/网络全省)。
  • 代价:重算不总是便宜。窄依赖重算只需重跑对应的父分区(便宜);宽依赖重算意味着要重跑上游一整个 shuffle,甚至级联到更上游(贵)。血缘链很长、或宽依赖很多时,一次故障可能触发大面积重算。
  • 对策与取舍:用 checkpoint 把某个 RDD 落到 HDFS(可靠存储),截断血缘——之后故障从 checkpoint 恢复,不再回溯。代价是 checkpoint 本身要写 HDFS(多副本、慢)。这是"重算成本 vs 落盘成本"的权衡。(persist/cache 只是把数据留在内存/本地盘加速复用,节点挂了仍要靠血缘重算;checkpoint 才是真正的可靠断点。)

直接写 RDD 算子太底层。DataFrame 是"带列名和类型(Schema)的分布式表",API 长得像 SQL(df.select().join().groupBy())。Dataset 是 DataFrame 的强类型版本(主要在 Scala/Java 里用,Python 只有 DataFrame)。

为什么它俩比裸 RDD 重要?因为引擎看得懂结构,就能自动优化:RDD 里的 map(lambda x: ...) 对引擎是黑盒,只能照跑;而 DataFrame 的 groupBy("city") 引擎知道语义,可以谓词下推、列裁剪、选 join 策略——这就是 Catalyst 优化器。你写的 Spark SQL,底层就是 DataFrame。

一条 SELECT ... JOIN ... GROUP BY ... 的 SQL,经过 SQL 解析 → 逻辑计划 → Catalyst 优化 → 物理计划,最终变成一张 RDD 的 DAG;这张 DAG 被切成若干 Stage,每个 Stage 展开成一批 Task,分发到 Executor 上执行。

Spark 运行架构:Driver / Executor / Cluster Manager

flowchart TB
    subgraph 提交端
      D["Driver 进程<br/>持有 SparkContext<br/>内含 DAGScheduler / TaskScheduler<br/>≈ 你的 main + 一个调度中心"]
    end
    CM["Cluster Manager 集群管理器<br/>(YARN / K8s / Standalone)<br/>只负责分配资源, 详见第09篇"]
    subgraph W1[Worker 节点 1]
      E1["Executor 进程<br/><b>本身就是一个 JVM 进程</b><br/>内部用线程池跑 Task"]
    end
    subgraph W2[Worker 节点 2]
      E2["Executor 进程<br/><b>一个 JVM 进程</b><br/>线程池跑 Task"]
    end
    D -->|1 申请资源| CM
    CM -->|2 在节点上拉起 Executor| E1
    CM -->|2 拉起 Executor| E2
    D -->|3 分发 Task / 收结果| E1
    D -->|3 分发 Task / 收结果| E2
    E1 <-->|4 shuffle: 相互拉取数据| E2
  • Driver(驱动器):运行 main 逻辑的进程,持有 SparkContext。它是大脑:把你的代码翻成 DAG、切 Stage、生成 Task、把 Task 派给 Executor、汇总结果、处理容错。你在 PySpark 里 df.collect() 拉回的数据,就是拉回到 Driver 的内存——所以 Driver 内存也会被撑爆(后面"坑"会讲)。
  • Executor(执行器)它就是一个 JVM 进程,跑在 Worker 节点上。一个 Executor 分到若干 CPU 核(spark.executor.cores)和一块堆内存(spark.executor.memory),内部用线程池并发跑多个 Task——一个核跑一个 Task。这和你在 Java 里 ExecutorService 提交 Runnable同一个模型Executor 这个词的来历也在此)。
  • Cluster Manager(集群管理器):只干一件事——给这个 Spark 应用分配"要几个 Executor、每个多大"。它可以是 YARN、K8s 或 Spark 自带的 Standalone。资源怎么申请、容器从哪来,是 YARN/编排 的主题,这里你只需知道:Executor 不是 Spark 自己变出来的,是向 Cluster Manager 要来的资源里拉起的 JVM。

一句话记忆:Task 是提交到 Executor(一个 JVM)线程池里的一个"处理单个分区"的任务;Driver 是那个决定有哪些 Task、派给谁的调度中心。

Spark 的执行是惰性(lazy)的:map/filter/join 这些转换(transformation)只是在搭 DAG,不真跑;直到遇到一个行动(action)——collectcountwrite.save 等——才真正触发计算。

一次 action 触发一个 Job(作业)。Job 被切成多个 Stage(阶段),每个 Stage 展开成一批并行的 Task(任务)。切分的依据,是 RDD 之间是窄依赖还是宽依赖

flowchart LR
    subgraph 窄依赖 narrow
      A1[父分区1] --> B1[子分区1]
      A2[父分区2] --> B2[子分区2]
      A3[父分区3] --> B3[子分区3]
    end
    subgraph 宽依赖 wide / shuffle
      C1[父分区1] --> D1[子分区1]
      C1 --> D2[子分区2]
      C2[父分区2] --> D1
      C2 --> D2
      C3[父分区3] --> D1
      C3 --> D2
    end
  • 窄依赖(narrow dependency):父 RDD 的每个分区,最多被子 RDD 的一个分区使用。典型算子:mapfilterselect。这类操作不需要跨机器搬数据——子分区要的输入就在本地父分区里,可以在同一个 Task 里流水线(pipeline)连着算完(读一条 → filter → map → 下一条)。
  • 宽依赖(wide dependency,又称 shuffle dependency):子 RDD 的一个分区,要依赖父 RDD 多个(乃至全部)分区的数据。典型算子:groupByreduceByKeyjoindistinctrepartition。这类操作必须 shuffle——因为同一个 key 的数据分散在所有父分区里,得先搬到一起。

Stage 边界,就切在宽依赖(shuffle)处。 为什么?因为 shuffle 是一道"栅栏(barrier)":下游分区必须等上游所有分区都把数据吐出来、重分发完成,才能开始。一连串窄依赖可以塞进一个 Stage 里流水线执行;一遇到宽依赖,就必须断开,前面是一个 Stage(负责 shuffle write),后面是下一个 Stage(负责 shuffle read + 继续算)。

一个 Stage 里有多少 Task?= 这个 Stage 处理的分区数。 Stage 之间串行(下游等上游),Stage 内部的 Task 并行。

Shuffle 是理解 Spark 性能的核心。自问:不就是把数据按 key 重新分个组吗,能有多贵?

答:贵在它同时踩了磁盘、网络、CPU 三个慢环节。以 GROUP BY city 为例,假设上游有 M 个 map 分区、下游要 R 个 reduce 分区:

Map 端 (Stage N, shuffle write)                 Reduce 端 (Stage N+1, shuffle read)
┌───────────────────────────┐
│ 每个 map Task:            │
│  1. 算出 (city, amount)  │
│  2. 按 city 的 hash 分成 │   跨网络拉取           ┌────────────────────────┐
│     R 个桶 (partition)   │  ───────────────────▶ │ 每个 reduce Task:      │
│  3. 排序 + 序列化        │   M×R 个数据块        │  从 M 个 map 各拉自己  │
│  4. 写本地磁盘           │  ───────────────────▶ │  那一份, 归并, 聚合    │
└───────────────────────────┘                       └────────────────────────┘

三笔代价,逐条给因果:

  1. 落盘(disk IO):map 端把输出按目标分区分桶后,写到本地磁盘(shuffle write),不是直接走网络发出去。为什么要落盘?① reduce 端是异步来拉的,map 不能一直把结果攥在内存里等;② 数据量常常超过内存;③ 容错——落盘后,某个 reduce Task 失败重跑时,能重新来拉这份 shuffle 数据,不必把上游 map 整个重算。代价就是一次磁盘写 + 一次磁盘读。
  2. 网络传输(network IO):reduce 端要从所有 M 个 map Executor 上,各拉取属于自己的那一份。逻辑上是 M × R 组数据传输,跨机器、跨网卡。数据量大时,网络带宽直接成为瓶颈。
  3. 排序 + 序列化(CPU):Spark 默认用 sort-based shuffle,map 端要对输出排序;数据跨进程/跨网络必须序列化再反序列化,这是纯 CPU 开销,还会产生大量临时对象、加重 GC。

对比一下代价量级:一串窄依赖是"本地内存里流水线跑完",只有 CPU + 少量内存;一次 shuffle 则叠加了本地磁盘写 + 磁盘读 + 跨网络传输 + 序列化 CPU + 排序。这就是为什么"减少 shuffle 次数 / 减小 shuffle 数据量"是 Spark 调优的第一要义

一个直接的优化例子:reduceByKeygroupByKey 快,因为前者在 map 端先做局部聚合(partial aggregation / combine),比如各分区先把本地的 amount 求个局部和,再 shuffle——跨网络传的数据量小得多。SQL 的 SUM/COUNT 走的正是这条 map 端预聚合路径。

SQL/DataFrame,不是照着字面执行的,中间有两层"编译器"。

Catalyst(优化器)——负责"算什么"的优化,全程操作的是计划树:

SQL/DataFrame
   → 未解析逻辑计划 (Unresolved Logical Plan)
   → 解析: 绑定表名/列名/类型 (查 Catalog)  → 逻辑计划
   → 逻辑优化: 谓词下推、列裁剪、常量折叠、算子合并 ...
   → 物理计划: 为每个逻辑算子选具体执行策略 (如 join 选 SortMergeJoin 还是 BroadcastHashJoin)
   → (基于统计信息 CBO 选代价最低的物理计划)
   → 生成可执行的 RDD DAG

几个在执行计划里能直接看到的优化:

  • 谓词下推(predicate pushdown)WHERE dt='2026-07-28' 被推到 Scan 里,配合 Parquet和分区,在读文件时就跳过无关数据,而不是全读进来再过滤。
  • 列裁剪(column pruning):只 SELECT city, amount,就只从 Parquet 读这两列,其余列根本不读。

Tungsten(钨丝计划)——负责"怎么算得省"的执行层优化,两把利器,正好都能用你的 JVM 知识理解:

  • 堆外内存管理(off-heap):Spark 处理的是海量结构化数据,如果都用 Java 对象表示(每个 Row 一个对象、每个字段一个 Integer/String),对象头开销 + GC 压力会非常大(你懂 JVM 的对象内存布局和 Full GC 停顿)。Tungsten 用 sun.misc.Unsafe 直接申请并操作二进制内存,把一行数据紧凑地摆成字节,绕开 JVM 对象开销、减少 GC。这是"更省内存 + 更少 GC 停顿"的因果,代价是实现复杂、绕过了 JVM 的安全检查。
  • 全阶段代码生成(whole-stage code generation):一个 Stage 里连续的多个算子(Filter → Project → 部分聚合),传统做法是每个算子一个对象、逐条调用(大量虚函数调用、中间对象)。Tungsten 在运行时把这些算子融合、动态生成一段紧凑的 Java 代码并编译成字节码——相当于把一串算子手写成一个 for 循环。减少虚函数调用和中间数据,逼近手写代码的性能。在执行计划里,被 codegen 融合的算子前面带 * 号。

你不需要背 Catalyst/Tungsten 的实现细节,只需记住:Catalyst 决定"执行计划长什么样"(含 shuffle 在哪、join 用什么策略),Tungsten 决定"每个算子在 Executor 里跑得多省"。 具体规则和默认阈值以对应 Spark 版本官方文档为准。

数据倾斜(data skew)定义:shuffle 后按 key 分区时,某个(些)key 的数据量远大于其它 key,导致对应的那个 Task 处理的数据量畸大。

成因很实在:null/默认值特别多(比如没登录用户 user_id 全是 -1)、热点大客户(某个 city='上海' 占了 30%)、爬虫/机器刷的脏数据集中在少数 key。

自问:那不就是一个 Task 慢点吗,为什么说拖垮整个 Stage、甚至整个作业?

答:木桶效应 + Stage 栅栏。

  • Stage 完成时间 = 该 Stage 里最慢那个 Task 的完成时间。 假设 200 个 Task,199 个处理均匀分区、10 秒跑完,剩下 1 个 Task 分到了倾斜 key、要处理 50 倍数据、跑 8 分钟——整个 Stage 就得等这 8 分钟。这和你在 Java 里 CountDownLatch 等所有子任务 countDown、或 CompletableFuture.allOf(...) 等全部完成是同一个道理(近似类比:都是"最慢的决定整体")。CPU 和内存资源在这 8 分钟里大量闲置。
  • 下游 Stage 被卡死:shuffle 是栅栏,下游 Stage 必须等上游整个 Stage 完成才能启动。所以一个倾斜 Task 不只慢了自己这个 Stage,后面所有依赖它的 Stage 都跟着顺延——这就是"拖慢整个作业"。
  • 雪上加霜:那个超大分区还容易把单个 Executor 的内存撑爆(OOM),一旦挂掉又触发血缘重算,恶性循环。

缓解思路(概念级,各有代价)

  • 加盐打散(salting):给倾斜 key 拼上随机后缀(上海_0~`上海_9`),先分散聚合,再去掉后缀二次聚合。把一个大分区拆成多个。代价:逻辑变复杂、要两阶段聚合。
  • 广播小表避免 shuffle(Broadcast Join):如果 join 的一侧足够小,把它广播到每个 Executor,大表不 shuffle、本地完成 join,从根上消灭这次 shuffle 及其倾斜。Spark 在小表小于 spark.sql.autoBroadcastJoinThreshold(默认 10 MB)时会自动这么做。代价:广播表太大就会撑爆 Executor/Driver 内存。
  • AQE 自适应处理倾斜:见下。
  • 单独处理倾斜 key:把已知倾斜的 key 过滤出来单独算,其余正常算,最后合并。
  • 注意反直觉:单纯调大 spark.sql.shuffle.partitions(增加并行度)对单个超大 key 无效——因为同一个 key 的数据仍然只会进同一个分区。它只能缓解"key 多而每个不大"型的不均。

AQE(Adaptive Query Execution,自适应查询执行):Spark 3.0 引入、据官方发布说明自 Spark 3.2.0 起默认开启spark.sql.adaptive.enabled=true;版本行为请以对应版本 SQL 调优文档核实)。它在运行时根据每个 shuffle 真实产出的数据大小动态调整计划,主要能力:① 动态合并过小的分区(避免小任务过多);② 自动拆分倾斜分区(把那个 50 倍大的分区切成几块并行处理);③ 运行时把符合条件的 SortMergeJoin 降级为 BroadcastJoin。这也是为什么现代 Spark 里数据倾斜比早年好治了不少——但它不是万能,严重倾斜仍需在建模/SQL 层动手。

下面用 PySpark 写一个"订单表 join 用户维表,再按城市汇总金额"的典型作业,然后打印物理计划。

# PySpark 3.x —— 逐行注释(不默认你熟 Python 语法)

orders = spark.table("dwd_orders")   # 订单明细表, 按 dt 分区(第02篇); 底层是 Parquet(第04篇)
users  = spark.table("dim_users")    # 用户维表

result = (
    orders
    .join(users, "user_id")          # 按 user_id 关联 —— 宽依赖, 会触发 shuffle
    .groupBy("city")                 # 按城市分组 —— 又一个宽依赖, 再触发一次 shuffle
    .sum("amount")                   # 对每个城市的 amount 求和
)

result.explain()                     # 打印物理执行计划(不真跑, 只看计划)

模拟输出(简化后的物理计划,Spark 3.x,explain() 的树形要从下往上、从内往外读):

== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false                              -- AQE 已开启, 计划运行时还会调整
+- HashAggregate(keys=[city], functions=[sum(amount)])           -- 聚合第二阶段: 汇总各分区局部和
   +- Exchange hashpartitioning(city, 200)                       -- 【shuffle #2】按 city 重分区
      +- HashAggregate(keys=[city], functions=[partial_sum(amount)])  -- 聚合第一阶段: map 端预聚合
         +- Project [city, amount]                               -- 列裁剪: 只保留这两列
            +- SortMergeJoin [user_id], [user_id], Inner         -- 大表 join: 排序归并策略
               :- Sort [user_id ASC]
               :  +- Exchange hashpartitioning(user_id, 200)     -- 【shuffle #1a】orders 按 user_id
               :     +- Filter isnotnull(user_id)                -- 谓词下推到扫描附近
               :        +- Scan parquet dwd_orders[user_id,city,amount,dt]  -- 只读用到的列
               +- Sort [user_id ASC]
                  +- Exchange hashpartitioning(user_id, 200)      -- 【shuffle #1b】users 按 user_id
                     +- Filter isnotnull(user_id)
                        +- Scan parquet dim_users[user_id]

一句话解释这份计划:全计划里出现了 3 个 Exchange 节点,就是 3 处 shuffle——join 两侧各按 user_id 重分区各一次(#1a、#1b),聚合再按 city 重分区一次(#2)。每个 Exchange 就是一道 Stage 边界,所以这个作业大致被切成 4 个 Stage:扫 orders、扫 users(两者并行)→ join + 局部聚合 → 最终聚合。

把它画成带 Stage 划分的 DAG(实线=窄依赖,可在同一 Stage 流水线;虚线=shuffle,Stage 之间的栅栏):

flowchart LR
    subgraph St1["Stage 1"]
      A1[Scan orders] --> A2["Filter user_id 非空"] --> A3["shuffle write<br/>按 user_id 分桶"]
    end
    subgraph St2["Stage 2"]
      B1[Scan users] --> B2["Filter user_id 非空"] --> B3["shuffle write<br/>按 user_id 分桶"]
    end
    subgraph St3["Stage 3"]
      C1[SortMergeJoin] --> C2[Project city amount] --> C3["partial_sum<br/>map端预聚合"] --> C4["shuffle write<br/>按 city 分桶"]
    end
    subgraph St4["Stage 4"]
      D1["sum by city<br/>最终聚合"] --> D2[输出结果]
    end
    A3 -. "shuffle #1a" .-> C1
    B3 -. "shuffle #1b" .-> C1
    C4 -. "shuffle #2" .-> D1

Stage 1 和 Stage 2 之间没有依赖,可并行跑;Stage 3 必须等 1、2 的 shuffle 数据都到齐才能启动;Stage 4 又等 Stage 3。这就是"下游等上游、Stage 内并行、Stage 间串行"的具体形态——也是数据倾斜能拖垮整条链的结构原因。

再点几个和前文的呼应:

  • Scan parquet dwd_orders[user_id,city,amount,dt] 只列出 4 列 → 列裁剪生效(Parquet 列存的价值)。
  • Filter isnotnull(user_id) 贴着 Scan谓词下推
  • HashAggregate ... partial_sum 出现在 shuffle 之前map 端预聚合,减少 shuffle 数据量(reduceByKey 优化)。
  • SortMergeJoin 说明两张表都不够小(超过 10 MB 广播阈值),走的是"两侧各 shuffle + 排序 + 归并"。如果 dim_users 很小,这里会变成 BroadcastHashJoinusers 侧的 Exchange 消失。
  • 顶上的 AdaptiveSparkPlan isFinalPlan=false → AQE 开着,真实运行时若发现某个 city/user_id 分区倾斜,会自动拆分。

想看得更细,用 result.explain(mode="formatted")(Spark 3.0+)会给每个节点编号并列出详细属性;explain(True) 会连同解析/逻辑/优化后逻辑/物理四份计划一起打。输出格式随版本略有差异,以你所用 Spark 版本为准。

spark 的边界 / 坑:

  • Executor / Driver OOM(内存溢出):Executor 就是 JVM,超过"堆 + 堆外"就 OOM 或陷入频繁 Full GC 而假死。常见诱因:单分区过大(往往就是数据倾斜)、df.collect() 把全量数据拉回 Driver 内存、cache 了太多用不上的数据、broadcast 了一张其实并不小的表。排查从 Spark UI 的 Stage/Task 数据量分布看起。
  • Shuffle 分区数不当spark.sql.shuffle.partitions 默认 200,是个"一刀切"值。数据量大时 200 太少 → 单分区过大、易 OOM/倾斜;数据量小时 200 太多 → 一堆几乎空的小 Task,调度开销和小文件反而拖累。现代靠 AQE 动态合并缓解,但大作业仍建议按数据量显式设置。
  • 小文件问题:如果最终有 N 个 reduce 分区,write 就产出 N 个文件。分区数偏大或长期增量追加,会在 HDFS/对象存储上堆出海量小文件,拖垮元数据服务(HDFS NameNode 内存、对象存储的 list 操作)——这正是存储时的痛点在计算侧的映射。缓解:写出前 coalesce/repartition 合并,或靠 AQE 合并小分区。
  • 不适合低延迟 / 流处理:Spark 批作业从申请 Executor、构建 DAG 到调度启动就有秒级开销,本质是"攒一批、算一批"。它的 Structured Streaming 用**微批(micro-batch)**模拟流,延迟通常在亚秒到秒级,仍不是逐条、毫秒级的连续处理。真正的低延迟流,见 Flink。
  • 何时不该用 Spark
    • 数据量不大(GB 级以内):单机 pandas / DuckDB 往往更快,Spark 的分布式调度开销和序列化成本不划算。
    • 高并发点查 / 交互式低延迟查询:那是 OLAP 引擎的活,Spark 是跑批的、不是给前端页面做毫秒响应的。
    • 毫秒级流式处理:交给 Flink。

火山引擎 LAS(湖仓一体分析服务) 上写的 Spark SQL,走的正是本篇这套流程。把心智图对应过去:

  • Spark SQL → 平台侧的 Catalyst 生成执行计划 → 切 Stage/Task → 分发到 Executor(JVM 进程) 执行;你能在作业的执行计划/UI 里看到本文讲的 Exchange(shuffle)、SortMergeJoinHashAggregate
  • Executor 容器由平台的资源层分配——对应本文的 Cluster Manager 角色(云上通常是 K8s/YARN 化的托管资源)。"资源怎么申请、容器从哪来"是第 09 篇的主题,这里你只需知道 Executor 是向资源层要来的。
  • LAS 上表常以 Iceberg(第 08 篇)组织,底层文件是 Parquet(第 04 篇),Scan 阶段的分区裁剪对应第 02 篇的分区——执行计划里谓词下推 + 分区裁剪能帮你少读大量数据。
你的 SQL
  │  Catalyst 优化(谓词下推/列裁剪/选 join 策略)
  ▼
物理计划 = RDD DAG
  │  按【宽依赖 / shuffle】切 Stage
  ▼
Job ── Stage₁ ─(shuffle: 落盘+网络+排序, 贵)─ Stage₂ ── ...
        │每个 Stage 展开成 N 个 Task(N=分区数)
        ▼
     Task 跑在 Executor(=JVM 进程)的线程池里, 一核一 Task
        │  Tungsten: 堆外内存 + 代码生成, 榨性能
        ▼
     窄依赖→Task 内流水线; 宽依赖→跨 Stage 的栅栏
        │  某分区丢失→按 lineage 血缘重算
        ▼
  数据倾斜: 一个大 key → 一个巨型 Task → 木桶效应拖垮整个 Stage/作业

记住三句话就够用了:shuffle 在宽依赖处发生、也在那里切 Stage;shuffle 贵是因为它同时踩了磁盘+网络+排序;倾斜之所以致命,是因为 Stage 只能等最慢的 Task。

Back to Blog