Spark 硬核指南(上篇):为什么是 Spark,核心概念与第一个程序
三部曲结构:
- 上篇(本文):Spark 是什么、和竞品的多维度对比、整体架构、核心抽象(RDD/DataFrame/Dataset)、第一个程序
- 中篇:Catalyst 优化器、Tungsten 引擎、AQE、Shuffle 内幕、内存管理、调度系统、Structured Streaming、Lakehouse 生态全景
- 下篇:性能调优实战手册、生产环境踩坑实录、监控调试、成本优化、深度对比补充、前沿技术、面试与职业发展
0. 先讲个故事:为什么会有 Spark
2009 年,加州大学伯克利分校的 AMPLab(Algorithms, Machines, People Lab)里,一群人正被 Hadoop MapReduce 折磨。
他们在做机器学习研究,核心痛点是:逻辑回归这类算法需要对同一份数据反复迭代计算成百上千次,而 MapReduce 的每一轮 Map/Reduce 之间,中间结果都要老老实实写回 HDFS,再从 HDFS 读出来进入下一轮。磁盘 I/O 成了压垮性能的最后一根稻草——不是算法慢,是"读写磁盘"这件事本身慢得离谱。
于是 Matei Zaharia(后来的 Databricks 创始人之一)和团队做了一个朴素到近乎"废话"的决定:把中间结果尽量留在内存里,别老往磁盘上写。这就是 Spark 最初的种子。2010 年论文发表,2013 年捐给 Apache 基金会,2014 年成为 Apache 顶级项目,同年 Databricks 成立,开始商业化运作。
这个"废话"级别的想法,十几年后长成了整个大数据生态里最重要的计算引擎之一——不只是"跑得快的 MapReduce",而是统一了批处理、流处理、SQL、机器学习、图计算的一整套分析引擎,并且成为 Delta Lake / Iceberg / Hudi 这些现代 Lakehouse 架构事实上的默认计算层。
理解这段历史,你会更容易理解 Spark 设计里几乎所有的取舍:为什么它对内存这么执着,为什么它的流处理是"微批"而不是"真流",为什么它的核心抽象叫 RDD(弹性分布式数据集)而不是别的名字。这些不是随意的设计,而是从"别再把中间结果写盘"这一个念头里,一路长出来的必然结果。
1. 多维度对比:Spark 放在整个技术版图里看
在深入 Spark 内部之前,先把它放到坐标系里,搞清楚它到底解决什么问题、不擅长什么问题。这是很多"学了 Spark 却不知道什么时候该用它"的人最缺的一课。
1.1 Spark vs Hadoop MapReduce
| 维度 | Hadoop MapReduce | Spark |
|---|---|---|
| 中间结果存储 | 写回 HDFS(磁盘) | 优先驻留内存,可溢写磁盘 |
| 编程模型 | 只有 Map / Reduce 两阶段 | DAG(有向无环图),任意多阶段 |
| 迭代计算 | 每轮迭代都要落盘,代价高 | 内存复用,迭代计算快 100 倍量级 |
| API 抽象层级 | 底层、命令式 | RDD(底层)→ DataFrame/Dataset(声明式)→ SQL |
| 适用场景 | 超大规模、对延迟不敏感的离线批处理 | 批处理 + 类流处理 + 交互式查询 + ML |
| 现状(2026) | 仍在部分存量系统中作为 HDFS/YARN 基础设施存在,但计算引擎地位已被 Spark/Flink 取代 | 事实上的批处理/ETL 默认引擎 |
结论:今天几乎没有人会为新项目选 MapReduce 作为计算引擎,但 HDFS 和 YARN 作为存储层/资源调度层依然广泛存在——这也是为什么配套的《搭建Spark集群》笔记里,Spark 集群搭建总是绕不开 Hadoop 环境。理解这一点很重要:Spark 取代的是 MapReduce 这个"计算模型",而不是整个 Hadoop 生态。
1.2 Spark vs Flink:批处理王者 vs 流处理王者
这是目前工程面试里问得最多的对比题,也是最容易被讲浅的一个。核心分歧只有一句话:
Spark Structured Streaming 用"微批(micro-batch)"模拟流处理;Flink 是真正的逐事件(event-at-a-time)流处理引擎。
这个架构选择直接决定了下游一切差异:
| 维度 | Spark(Structured Streaming) | Flink |
|---|---|---|
| 处理模型 | 微批:把一小段时间窗口内的数据攒成一个批次处理 | 真流:每个事件到达即被处理 |
| 典型延迟 | 秒级(严格模式下可到亚秒级,但仍非事件级) | 毫秒级,p99 可做到 100ms 以内 |
| 吞吐 vs 延迟权衡 | 提高吞吐容易,但会推高延迟 | 兼顾高吞吐与低延迟的设计目标 |
| 状态管理 | 内置 State Store,支持 RocksDB 后端 | 原生为流式状态设计,Checkpoint/Savepoint 机制成熟度更高 |
| 批流统一 | "流是批的特例"(微批模型) | "批是流的特例"(有界流) |
| 生态与易用性 | 生态更大、SQL/DataFrame API 统一、学习曲线平缓 | 专精流场景,配置和调优复杂度更高 |
| 典型选型场景 | ETL、准实时报表、离线+近实时统一平台 | 风控、实时计费、超低延迟事件驱动系统 |
行业里目前的共识大致是:Flink 通常比 Spark Streaming 在同等负载下低 10-100 倍延迟;在近期使用 Flink 1.19 与 Spark 3.5 的基准测试中,Flink 的 p99 延迟能稳定做到 100ms 以内,而 Spark Streaming 在类似数据量下通常在 2-5 秒区间。但原始吞吐基准测试常常会误导人,实际选型要看具体业务场景。对于追求快速开发、最小化运维成本、且需要与既有生态深度集成的团队,Spark Structured Streaming 仍然是首选;Flink 更适合需要精细控制、对延迟极致敏感的场景。
给你的建议:如果你的团队已经在用 Spark 做批处理/ETL,业务对"实时"的定义是"分钟级、秒级可接受",选 Structured Streaming 几乎总是对的——不用引入第二套技术栈、不用招 Flink 专才、状态管理和监控都能复用批处理的经验。只有当你的 SLA 真的卡在毫秒级(风控拦截、实时竞价、金融交易撮合),才值得为 Flink 付出额外的运维复杂度。
1.3 Spark vs Trino/Presto:交互式查询引擎
| 维度 | Spark SQL | Trino/Presto |
|---|---|---|
| 设计目标 | ETL + 批处理 + 交互式查询 + ML 一体化 | 专精 MPP 交互式即席查询 |
| 启动开销 | 有 JVM/Driver 启动开销,适合长任务 | 常驻集群,查询响应更快 |
| 容错模型 | Stage 级重试,牺牲部分延迟换稳定性 | 传统上容错较弱(新版本有改进),追求极致查询速度 |
| 联邦查询 | 需要额外配置数据源连接器 | 原生擅长跨多数据源联邦查询(MySQL、Kafka、Iceberg 等一次 JOIN) |
| 适用场景 | 复杂 ETL 管道、机器学习特征工程、大规模数据转换 | BI 报表、Ad-hoc 分析、跨源联邦查询 |
一句话总结:Spark 是"干活的引擎",Trino 是"查数的引擎"。很多成熟的数据平台会两者并存——Spark 跑夜间批处理把数据整理进 Lakehouse,Trino 给分析师和 BI 工具提供秒级的交互式查询。
1.4 Spark vs Ray / Dask:Python 原生分布式计算
如果你更多接触 Python/AI 方向,一定听说过 Ray 和 Dask:
| 维度 | Spark | Ray | Dask |
|---|---|---|---|
| 原生语言 | Scala(JVM),Python 是"客户端" | Python 原生 | Python 原生 |
| 强项 | 结构化数据 ETL、SQL、大规模数据处理 | 分布式训练、超参搜索、强化学习、Actor 模型 | 让 Pandas/NumPy 代码"无痛"扩展到集群 |
| 与深度学习框架集成 | 需要额外桥接(Horovod、TorchDistributor) | 原生集成 PyTorch/TensorFlow 分布式训练 | 较弱 |
| 学习曲线 | 需要理解 JVM、DataFrame/RDD 抽象 | 对 Python 工程师非常友好 | 几乎是 Pandas 的直接替代品 |
给你的建议:数据工程/ETL 场景优先 Spark;如果你的团队核心痛点是"模型训练需要分布式"、"超参搜索要并行跑几百个 trial",Ray 通常是更合适的工具。两者不是完全互斥——很多公司用 Spark 做特征工程,用 Ray 做模型训练,中间用 Parquet/Delta 文件或 Feature Store 衔接。
1.5 一张全景对比表
| 引擎 | 定位 | 一句话记忆点 |
|---|---|---|
| Hadoop MapReduce | 历史引擎 | 磁盘中转站,已被淘汰 |
| Spark | 统一分析引擎 | 内存优先的批流一体全能选手 |
| Flink | 真流处理引擎 | 毫秒级延迟的实时之王 |
| Trino/Presto | MPP 查询引擎 | 秒级响应的联邦查询之王 |
| Ray | 分布式 AI 计算框架 | Python 原生的训练/推理调度器 |
| Dask | Pandas 分布式扩展 | "别改代码就能上集群"的轻量选择 |
| DuckDB/Polars | 单机向量化引擎 | 数据量不大时,比启动一个 Spark 集群快得多 |
一个常被忽略的真相:不是所有"大数据"任务都需要分布式引擎。如果你的数据能塞进一台机器的内存(哪怕是几十 GB),DuckDB 或 Polars 单机跑往往比启动一个 Spark 集群(哪怕是 local 模式)更快,因为省掉了整套分布式调度、序列化、网络传输的开销。这也是"下篇"里会详细讨论的一个成本优化角度:先问"我真的需要分布式吗",再决定要不要上 Spark。
2. Spark 整体架构:先建立一张全局地图
在写任何代码之前,先把这几个核心概念钉死在脑子里,后面所有内容都是在这张图上做展开。
几个必须精确掌握、面试常考的概念区分:
| 术语 | 到底是什么 | 常见混淆点 |
|---|---|---|
| Application | 你提交的整个 Spark 程序,对应一个 SparkContext/SparkSession 生命周期 | 一个 Application 可以包含多个 Job |
| Job | 一次 Action(如 collect()、count())触发的完整计算 | 一个 Application 通常触发多个 Job |
| Stage | 一个 Job 按 Shuffle 边界切分出的若干阶段 | 常和 Job 混淆;Stage 内部是流水线执行的窄依赖操作 |
| Task | Stage 中作用在单个分区上的最小执行单元 | Task 数量 = 该 Stage 涉及的分区数 |
| Executor | Worker 节点上为某个 Application 启动的进程,内部用多线程并发跑多个 Task | 不是"一个 Executor 只能跑一个 Task" |
| Driver | 运行你 main() 函数、维护 DAG、做任务调度的进程 | 不直接参与数据计算,只做调度和协调 |
用一句话串起来记:一个 Application 里有多个 Job(每个 Action 一个),每个 Job 按 Shuffle 拆成多个 Stage,每个 Stage 按分区拆成多个 Task,Task 才是真正在 Executor 上跑的最小单元。
2.1 部署模式速览(详细版留到下篇的运维部分)
- Local 模式:单机、单 JVM 进程,Driver 和 Executor 都在同一进程内,适合开发调试;
- Standalone:Spark 自带的资源管理器,部署简单,适合中小规模自建集群;
- YARN:复用 Hadoop 生态的资源调度,企业里最常见的生产部署方式;
- Kubernetes:云原生时代的主流选择,容器化部署、弹性伸缩、与云厂商托管 K8s 服务无缝集成,是目前新建集群的默认推荐;
- Mesos:曾经的第三种资源管理器,已于 Spark 3.2 起被标记为废弃,2026 年的新项目不应再考虑它。
3. 环境搭建:三分钟跑起第一个程序
3.1 最快路径:本地模式
不需要任何集群,下载即用:
# 下载并解压(请以官网当前最新稳定版为准,4.x 系列已成为主流)
tar -zxf spark-4.0.0-bin-hadoop3.tgz -C /usr/local/
export SPARK_HOME=/usr/local/spark-4.0.0-bin-hadoop3
export PATH=$PATH:$SPARK_HOME/bin
# 交互式 Scala 环境
spark-shell
# 交互式 Python 环境
pyspark
# 交互式 SQL 环境
spark-sql3.2 更现代的路径:pip 一行装 PySpark
如果你主要用 Python,甚至不需要下载完整发行包:
pip install pyspark
python3 -c "
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName('quickstart').getOrCreate()
spark.range(10).show()
"3.3 Docker 路径:不想污染本机环境
docker run -it --rm apache/spark:4.0.0 /opt/spark/bin/spark-shell三种路径任选,重点不是"装成功",而是接下来要建立的核心心智模型。
4. 核心抽象:RDD → DataFrame → Dataset → SQL 的演进逻辑
很多教程会把这几个概念平铺着讲,但更有效的理解方式是:它们是同一套底层引擎上,逐渐"变笨(对用户更友好)、变聪明(对引擎更可优化)"的四层包装。
4.1 RDD:理解引擎,但生产代码里已很少直接手写
RDD(Resilient Distributed Dataset,弹性分布式数据集)是最原始的抽象:一个不可变、可分区、可并行操作的分布式对象集合。用经典的 WordCount 来体会它的编程风格:
// RDD 版本:命令式,你需要手动指定每一步怎么做
val counts = sc.textFile("hdfs:///data/input.txt")
.flatMap(line => line.split(" "))
.map(word => (word, 1))
.reduceByKey(_ + _)
counts.saveAsTextFile("hdfs:///data/output")RDD 的核心问题:引擎不知道你在处理什么"语义",它只看到一堆匿名函数(黑盒)。所以引擎没法做任何智能优化——比如它不知道能不能把 filter 提前、能不能只读某几列。这也是为什么 Spark 团队后来要发明 DataFrame。
RDD 什么时候还值得手写:处理非结构化数据(比如自定义二进制格式解析)、需要精细控制分区和物理执行方式的极端性能场景、或者你在读 Spark 源码/理解底层机制时。日常 ETL 业务代码,几乎不应该再直接写 RDD。
4.2 DataFrame:生产代码的绝对主力
DataFrame 本质上是"带 schema 的 RDD",可以理解为分布式版本的关系表(或者说,分布式版的 Pandas DataFrame,但语义更接近 SQL 表)。因为引擎知道每一列的类型和名字,Catalyst 优化器(中篇详细讲)就能像数据库查询优化器一样,对整个计算链路做重排、剪枝、下推。
import org.apache.spark.sql.functions._
val df = spark.read.textFile("hdfs:///data/input.txt")
val counts = df
.select(explode(split(col("value"), " ")).as("word"))
.groupBy("word")
.count()
counts.write.parquet("hdfs:///data/output")同样的逻辑,用 Python:
from pyspark.sql.functions import explode, split, col
df = spark.read.text("hdfs:///data/input.txt")
counts = (df
.select(explode(split(col("value"), " ")).alias("word"))
.groupBy("word")
.count())
counts.write.parquet("hdfs:///data/output")或者,最简洁的一种:直接写 SQL
df.createOrReplaceTempView("lines")
spark.sql("""
SELECT word, count(*) as cnt
FROM (SELECT explode(split(value, ' ')) as word FROM lines)
GROUP BY word
""").show()三种写法(DataFrame API / SQL / RDD)编译到同一个物理执行计划——这是理解 Spark SQL 引擎最重要的一句话。选哪种写法纯粹是团队习惯和可读性的问题,不存在"SQL 比 DataFrame 慢"这种说法(中篇会用 explain() 实际验证这件事)。
4.3 Dataset:Scala/Java 的"两者兼得"方案
Dataset 试图同时拥有 RDD 的编译期类型安全和 DataFrame 的执行优化能力,只在 Scala 和 Java 中存在(Python 是动态类型语言,天然没有这个需求,PySpark 里 DataFrame 就是 Dataset[Row]):
case class WordCount(word: String, count: Long)
val ds: Dataset[WordCount] = counts.as[WordCount]
ds.filter(_.count > 10).show() // 这里的 _.count 是编译期类型安全的该不该用 Dataset:如果你的团队是 Scala 重度用户、代码库需要强类型保证(避免拼错列名这种低级错误在运行时才暴露),Dataset 值得投入。如果团队是 Python/SQL 为主,这一层可以直接跳过——PySpark 用户不会用到它。
4.4 该用哪个?一张决策表
| 场景 | 推荐 |
|---|---|
| 日常 ETL / 数据清洗 / 聚合分析 | DataFrame API 或 Spark SQL |
| 分析师 / BI 工具对接 | Spark SQL |
| Scala 重度团队,要求编译期类型安全 | Dataset |
| 非结构化数据、自定义分区逻辑、底层机制学习 | RDD |
| 机器学习特征工程 | DataFrame(配合 MLlib 的 Pipeline API) |
5. 转换与行动:惰性求值的再理解
如果你看过配套的《搭建Spark集群》笔记,会记得这个结论:转换操作(Transformation)是懒惰的,只有行动操作(Action)才会触发真正执行。这里再深一层:为什么要这样设计?
答案不是"为了省事",而是为了让 Catalyst 优化器有机会看到完整的计算链路,再决定怎么执行最划算。类比一下:如果你告诉一个装修工人"先刷墙,再铺地板,再装灯",工人如果傻乎乎地严格按顺序执行,可能要多次进出同一个房间;但如果工人先听完你的完整需求,可能会发现"其实可以先装灯再铺地板,省一趟"。Spark 的惰性求值,就是把"听完完整需求再规划"这件事,在计算领域实现了出来。
# 这一整段代码,此刻什么都没有真正执行
df1 = spark.read.parquet("large_table.parquet")
df2 = df1.filter(col("status") == "active")
df3 = df2.select("user_id", "amount")
df4 = df3.groupBy("user_id").sum("amount")
# 直到这一行,Catalyst 才会看到上面完整的4步操作,
# 一起分析、优化、生成最终的物理执行计划
df4.show()Catalyst 拿到这个完整链路后,可能会做的优化包括:谓词下推(把 filter 尽量提前到读文件那一刻,减少读取的数据量)、列裁剪(Parquet 是列式存储,只读 user_id、amount、status 三列,其余列完全不读)。这些优化,如果你手写循环去实现,需要引擎知道文件格式、需要手动重排语句顺序——而惰性求值 + 声明式 API,让引擎自动帮你做了。这也是为什么"DataFrame 天生比 RDD 快":不是执行引擎本身不同,而是引擎能看懂你想干什么,从而做更聪明的规划。
6. 分区(Partition)初步:并行度的物理基础
RDD/DataFrame 在物理上被切成若干分区,分布存储在集群不同节点的内存/磁盘中——这是 Spark 能并行计算的物理基础。理解这一点是理解一切性能问题的起点,也是下篇"数据倾斜"章节的前置知识。
df = spark.read.parquet("large_table.parquet")
print(df.rdd.getNumPartitions()) # 查看当前分区数
# 增加分区(会触发 shuffle,代价较高)
df_more = df.repartition(200)
# 减少分区(不会触发全量 shuffle,代价较低,常用于写文件前合并小文件)
df_fewer = df.coalesce(10)一条经验法则先记下来(下篇会展开到具体计算公式):分区数太少 → 并行度不够,浪费集群资源;分区数太多 → 调度开销压过计算收益,小文件问题随之而来。理想的分区大小通常落在 128MB~256MB 区间(正好贴近 HDFS/对象存储的默认 block size),这不是巧合,而是让"一个分区对应大约一次高效的顺序读"。
本篇小结
到这里,你应该已经建立起这些心智模型:
- Spark 为什么存在:用内存换掉 MapReduce 的磁盘中转站,从"跑得快的批处理引擎"逐渐长成统一分析引擎;
- 它在技术版图里的位置:批处理选 Spark,超低延迟流处理选 Flink,交互式联邦查询选 Trino,Python 原生分布式训练选 Ray;
- 整体架构:Driver 负责调度,Executor 负责执行,Application → Job → Stage → Task 是理解一切性能问题的坐标系;
- 四层核心抽象:RDD(底层可控)→ DataFrame(生产主力)→ Dataset(Scala 类型安全)→ SQL(最低门槛),四者共享同一套执行引擎;
- 惰性求值的真正意义:不是偷懒,是给优化器"看全局再规划"的机会;
- 分区是并行度的物理基础,也是下篇性能调优一切讨论的起点。
中篇预告:我们会打开 Spark 的引擎盖,看看 Catalyst 优化器到底怎么把你的 DataFrame 代码变成高效的物理执行计划、Tungsten 引擎如何用堆外内存和整段代码生成压榨出接近手写代码的性能、AQE 如何在运行时动态调整执行计划、Shuffle 内部到底发生了什么、以及 Iceberg/Delta/Hudi 这些 Lakehouse 格式如何和 Spark 结合,撑起现代数据平台的地基。