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 解决了什么:程序员只写 map 和 reduce 两个函数,分发、容错、shuffle 全由框架兜底。容错方式很朴素——某个 map/reduce 任务失败,框架在别的机器上重跑这一个任务(输入来自 HDFS,可重读)。这是"重算式容错"的雏形,Spark 后面把它发扬光大。
自问:既然抽象这么优雅,为什么后来被 Spark 取代?
答:贵在 IO,烦在编程模型。
- 每个作业的结果都要落 HDFS,多副本写。 一个 MapReduce 作业结束,输出写回 HDFS(默认 3 副本 + 跨网络复制)。而真实业务很少一步算完——一个 SQL 里的 join + 聚合,可能被编译成 3~5 个串行的 MapReduce 作业。每个作业之间都要"写 HDFS → 下个作业再读 HDFS",磁盘 IO 和网络复制成了主瓶颈。
- 迭代型作业尤其惨。 机器学习(逻辑回归跑 N 轮)、图算法(PageRank 迭代),每一轮都是一个 MapReduce 作业,每轮都完整地 "读 HDFS → 算 → 写 HDFS"。数据本可以留在内存里复用,MapReduce 却每轮都往返磁盘 N 次。
- 编程模型太原始。 什么都得掰成 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/persist 到 Executor 内存,后续轮直接读内存 |
| 编程繁琐 | 只有 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 不可变。你对它做
map、filter,不是原地改,而是生成一个新的 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)——collect、count、write.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 的一个分区使用。典型算子:
map、filter、select。这类操作不需要跨机器搬数据——子分区要的输入就在本地父分区里,可以在同一个 Task 里流水线(pipeline)连着算完(读一条 → filter → map → 下一条)。 - 宽依赖(wide dependency,又称 shuffle dependency):子 RDD 的一个分区,要依赖父 RDD 多个(乃至全部)分区的数据。典型算子:
groupBy、reduceByKey、join、distinct、repartition。这类操作必须 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. 写本地磁盘 │ ───────────────────▶ │ 那一份, 归并, 聚合 │
└───────────────────────────┘ └────────────────────────┘
三笔代价,逐条给因果:
- 落盘(disk IO):map 端把输出按目标分区分桶后,写到本地磁盘(shuffle write),不是直接走网络发出去。为什么要落盘?① reduce 端是异步来拉的,map 不能一直把结果攥在内存里等;② 数据量常常超过内存;③ 容错——落盘后,某个 reduce Task 失败重跑时,能重新来拉这份 shuffle 数据,不必把上游 map 整个重算。代价就是一次磁盘写 + 一次磁盘读。
- 网络传输(network IO):reduce 端要从所有 M 个 map Executor 上,各拉取属于自己的那一份。逻辑上是 M × R 组数据传输,跨机器、跨网卡。数据量大时,网络带宽直接成为瓶颈。
- 排序 + 序列化(CPU):Spark 默认用 sort-based shuffle,map 端要对输出排序;数据跨进程/跨网络必须序列化再反序列化,这是纯 CPU 开销,还会产生大量临时对象、加重 GC。
对比一下代价量级:一串窄依赖是"本地内存里流水线跑完",只有 CPU + 少量内存;一次 shuffle 则叠加了本地磁盘写 + 磁盘读 + 跨网络传输 + 序列化 CPU + 排序。这就是为什么"减少 shuffle 次数 / 减小 shuffle 数据量"是 Spark 调优的第一要义
一个直接的优化例子:
reduceByKey比groupByKey快,因为前者在 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很小,这里会变成BroadcastHashJoin,users侧的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)、SortMergeJoin、HashAggregate。 - 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。