Skip to content

Spark 完全入门指南:从零到硬核的大数据之旅 ​

你好,欢迎来到 Spark 的世界。

这不是一篇"Hello World"式的水文,而是一篇真正能让你理解 Spark 为什么快、怎么用、跟其他技术比怎么样、未来往哪走的硬核入门指南。

读完这篇,你不仅能写 Spark 程序,还能在面试时把面试官聊懵。


目录 ​


一、为什么是 Spark?—— 大数据的"速度与激情" ​

1.1 大数据时代的痛点 ​

想象一下,你有 1TB 的日志数据要处理。

用单机 Python 脚本? 等一天可能都跑不完,内存还会爆。

用 Hadoop MapReduce? 能跑完,但是慢。为什么慢?因为 MapReduce 把中间结果都写磁盘了,每一步都要落盘,磁盘 IO 成了瓶颈。一个复杂的计算可能要跑好几个小时甚至几天。

💡 打个比方: MapReduce 就像一个超级谨慎的厨师,每切完一道菜就把刀放下、把菜放冰箱、洗手、然后再打开冰箱拿菜、拿刀、继续切。安全是安全,但是慢死了。

Spark 呢?Spark 就像一个熟练的大厨,菜放案板上,刀不离手,一口气切完所有菜,中间不放下。能放内存的绝不放磁盘。

1.2 Spark 快在哪里? ​

官方说 Spark 比 Hadoop MapReduce 快 100 倍(内存计算时),磁盘上也快 10 倍。

凭什么这么快?三个核心原因:

原因一:内存计算(最核心) ​

  • MapReduce:中间结果写磁盘 → 读磁盘 → 下一步计算
  • Spark:中间结果放内存 → 直接下一步计算
  • 内存比磁盘快多少?大约 100 万倍(纳秒 vs 毫秒级别)

原因二:DAG 执行引擎 ​

  • MapReduce:只有 Map 和 Reduce 两个阶段,复杂任务要串很多个 MR,每个之间都要落盘
  • Spark:有 DAG(有向无环图)调度引擎,能把多个操作优化成一个流水线,中间不落盘

原因三:丰富的算子 ​

  • MapReduce:你得自己实现 Map 和 Reduce 逻辑,很多操作要自己写
  • Spark:内置几十种算子(map、filter、reduceByKey、join...),一行代码搞定复杂操作

🎯 一句话总结:Spark = 内存计算 + DAG 优化 + 丰富 API = 快 + 好用

1.3 Spark 的应用场景 ​

场景例子用什么组件
批处理日志分析、数据清洗、ETLSpark Core / Spark SQL
交互式查询即席查询、数据探索Spark SQL / Spark Shell
实时计算实时监控、实时推荐Spark Streaming / Structured Streaming
机器学习推荐系统、分类预测MLlib
图计算社交网络、路径分析GraphX

是的,你没看错,一个框架全搞定。这就是 Spark 被称为"大数据瑞士军刀"的原因。


二、Spark 是什么?—— 一个全能选手的诞生 ​

2.1 Spark 的身世 ​

  • 2009 年:诞生于 UC Berkeley AMPLab(加州大学伯克利分校)
  • 2010 年:开源
  • 2013 年:捐赠给 Apache 基金会
  • 2014 年:成为 Apache 顶级项目
  • 2016 年:Spark 2.0 发布,统一 DataFrame API
  • 2020 年:Spark 3.0 发布,重大升级
  • 2024 年:Spark 3.5+,持续进化中

创始人 Matei Zaharia 是个狠人,读博期间搞出了 Spark,后来又搞出了 MLflow 和 Delta Lake,现在是 Databricks 的 CTO。

2.2 Spark 的设计哲学 ​

1. 速度(Speed) ​

  • 内存计算
  • DAG 优化
  • 钨丝计划(Project Tungsten):直接操作二进制数据,跳过 JVM 对象开销

2. 易用(Ease of Use) ​

  • 支持 Scala、Java、Python、R、SQL 多种语言
  • 丰富的 API,几十种算子
  • 写起来像写单机程序一样自然

3. 通用(Generality) ​

  • 批处理、流处理、SQL、机器学习、图计算,一个框架全搞定
  • 不用学好几个框架,一套技术栈通吃

4. 兼容(Compatibility) ​

  • 跑在 Hadoop YARN、Mesos、K8s、独立集群上
  • 能读 HDFS、HBase、Hive、Kafka、Cassandra... 各种数据源
  • 不跟现有技术栈冲突,直接接入

2.3 Spark 架构总览 ​

                    ┌──────────────────┐
                    │   Driver Program │  ← 你的代码跑在这里
                    │  (SparkContext)  │
                    └────────┬─────────┘
                             │
              ┌──────────────┼──────────────┐
              ▼              ▼              ▼
        ┌──────────┐   ┌──────────┐   ┌──────────┐
        │ Worker 1 │   │ Worker 2 │   │ Worker N │
        │          │   │          │   │          │
        │ Executor │   │ Executor │   │ Executor │
        │  (JVM)   │   │  (JVM)   │   │  (JVM)   │
        │          │   │          │   │          │
        │  Task 1  │   │  Task 3  │   │  Task 5  │
        │  Task 2  │   │  Task 4  │   │  Task 6  │
        └──────────┘   └──────────┘   └──────────┘

角色说明 ​

角色职责类比
Driver负责任务调度、分配任务、汇总结果项目经理
Worker管理本节点的计算资源包工头
Executor实际执行 Task 的进程(JVM)工人
Task最小执行单元,一个分区对应一个 Task具体的活

💡 理解要点:

  • Driver 只有一个,Executor 有很多个
  • 数据是分布式存储的,每个 Executor 处理自己本地的数据
  • Driver 负责"指挥",Executor 负责"干活"

三、核心概念深度拆解 —— 理解了这些才算真懂 Spark ​

3.1 RDD:弹性分布式数据集 ​

什么是 RDD? ​

RDD(Resilient Distributed Dataset)是 Spark 最核心的抽象,翻译过来叫"弹性分布式数据集"。

名字很唬人,拆开理解:

  • 分布式(Distributed):数据分布在多台机器上
  • 数据集(Dataset):就是一堆数据
  • 弹性(Resilient):丢了能自动恢复,容错性强

💡 通俗理解: RDD 就像一个分布式的购物清单。

  • 清单很长,一页写不下,拆成好几页(分区),分给几个人拿着(分布式)
  • 如果某个人的那页丢了,可以根据原始数据重新算出来(弹性/容错)
  • 你可以对清单做各种操作:过滤、映射、排序...

RDD 的五大特性 ​

  1. 一系列分区(Partitions):数据被分成很多份,每份在不同节点上
  2. 每个分区有一个计算函数:对分区内的数据做什么操作
  3. 一系列依赖关系:RDD 之间的血缘关系(Lineage)
  4. (可选)分区器(Partitioner):键值对 RDD 才有,决定数据怎么分区
  5. (可选)首选位置(Preferred Locations):数据在哪台机器上,计算就调度到哪台(移动计算不移动数据)

RDD 的容错机制:Lineage(血缘) ​

Spark 不需要像传统数据库那样频繁 checkpoint(检查点),因为 RDD 记录了血缘关系——我是从哪个 RDD 怎么变过来的。

如果某个分区的数据丢了,根据血缘关系重新算一遍就行。

原始数据 → map → filter → reduceByKey → 结果
            ↑       ↑          ↑
         血缘关系记录了每一步

🎯 为什么这很牛?

  • 不需要频繁写磁盘做 checkpoint,省了大量 IO
  • 容错成本低,只需要重算丢失的分区,不用全量重算
  • 这就是"弹性"的含义

3.2 转换算子 vs 行动算子 ​

Spark 的操作分为两类,这是面试必考题。

转换算子(Transformation) ​

  • 作用:从一个 RDD 生成另一个 RDD
  • 特点:懒执行(Lazy)——调用的时候不会立刻执行,只是记录下来
  • 例子:map、filter、flatMap、reduceByKey、join...

行动算子(Action) ​

  • 作用:触发计算,返回结果或者写出数据
  • 特点:立刻执行——调用的时候才真正开始跑
  • 例子:collect、count、first、take、saveAsTextFile...

💡 打个比方: 转换算子就像你在购物 APP 里往购物车里加东西,加多少都不会扣钱(懒执行)。 行动算子就是"结算"按钮,一点击才真正开始计算价格、扣钱、发货。

好处是什么?Spark 可以在你点"结算"的时候,看看你购物车里的东西能不能优化一下打包方式,省运费。

为什么要懒执行? ​

因为懒执行才能做优化!

Spark 会把所有转换算子攒起来,等到行动算子触发的时候,整体优化一下执行计划,比如:

  • 把多个 map 合并成一个
  • 调整操作顺序减少 shuffle
  • 生成最优的 DAG

如果每一步都立刻执行,就没法优化了。

3.3 宽窄依赖:Spark 最核心的概念之一 ​

什么是依赖? ​

RDD 之间的关系叫依赖。子 RDD 依赖父 RDD。

宽依赖 vs 窄依赖 ​

类型定义类比例子
窄依赖父 RDD 的每个分区最多被子 RDD 的一个分区使用独生子女:一个爸只有一个娃map、filter、union
宽依赖父 RDD 的每个分区可能被子 RDD 的多个分区使用多胞胎:一个爸有好几个娃reduceByKey、groupByKey、join、distinct

💡 怎么快速判断? 看数据需不需要"重新洗牌"(Shuffle)。

  • 不需要重新分布数据 → 窄依赖
  • 需要重新分布数据(比如按 key 重新分组) → 宽依赖

为什么宽窄依赖很重要? ​

因为它决定了 Stage 划分。

Spark 会把一个 DAG 拆成多个 Stage(阶段),从后往前推,遇到宽依赖就切一刀。

Stage 0                    Stage 1                    Stage 2
┌──────────────┐        ┌──────────────┐        ┌──────────────┐
│  map         │ 窄依赖  │  reduceByKey │ 窄依赖  │  collect     │
│  filter      │───────▶│  map         │───────▶│  (Action)    │
└──────────────┘        └──────────────┘        └──────────────┘
         ↑ 宽依赖的地方切一刀 ↑
  • 每个 Stage 内部是窄依赖,可以流水线执行
  • Stage 之间是宽依赖,必须等上一个 Stage 全部跑完才能开始下一个
  • 宽依赖 = Shuffle = 性能瓶颈

🎯 面试必问:Spark 怎么划分 Stage? 答:从后往前(从 Action 往前推),遇到宽依赖就切分,每个 Stage 包含一系列窄依赖的转换操作。

3.4 DAG:有向无环图 ​

DAG(Directed Acyclic Graph)就是把 RDD 之间的依赖关系画成一张图。

  • 有向:有方向,数据从前往后流
  • 无环:不能循环,不能绕回来

DAGScheduler 负责把 DAG 拆成 Stage,然后把 Task 分配给 Executor 执行。

💡 为什么 DAG 比 MapReduce 强? MapReduce 只有 Map 和 Reduce 两个阶段,再复杂的逻辑也要硬塞到这两个阶段里,中间还得落盘。 Spark 的 DAG 可以有任意多个 Stage,而且 Stage 内部是流水线执行的,中间不落盘。 这就是 Spark 快的另一个重要原因。


四、Spark 技术栈全景 —— 五大组件,一个框架搞定 ​

Spark 不是一个单一的工具,而是一个技术栈。核心是 Spark Core,上面构建了四大组件。

┌─────────────────────────────────────────────────────────────┐
│                    Spark 技术栈全景                           │
├─────────┬───────────┬─────────────┬───────────┬─────────────┤
│ Spark   │ Spark     │ Spark       │ Spark     │ GraphX      │
│ SQL     │ Streaming │ MLlib       │ GraphX    │ (图计算)    │
│ (SQL)   │ (流计算)   │ (机器学习)   │           │             │
├─────────┴───────────┴─────────────┴───────────┴─────────────┤
│                     Spark Core (RDD)                         │
│                   (核心:RDD、调度、内存管理)                  │
├─────────────────────────────────────────────────────────────┤
│  集群管理器:Standalone / YARN / Mesos / Kubernetes          │
├─────────────────────────────────────────────────────────────┤
│  数据源:HDFS / HBase / Hive / Kafka / Cassandra / ...      │
└─────────────────────────────────────────────────────────────┘

4.1 Spark SQL:用 SQL 处理大数据 ​

是什么? ​

让你能用 SQL 语句查询分布式数据,就像查 MySQL 一样,但能处理 TB 级数据。

核心抽象:DataFrame / Dataset ​

  • DataFrame:带 Schema 的分布式数据集合,类似关系型数据库的表
  • Dataset:DataFrame 的强类型版本(Scala/Java 才有)

为什么重要? ​

  • 写 SQL 比写 RDD 代码简单太多了
  • 非技术人员(数据分析师)也能用
  • Catalyst 优化器会自动优化 SQL,比手写 RDD 可能还快

适用场景 ​

  • 数据仓库
  • 即席查询(Ad-hoc)
  • 数据 ETL

💡 Spark SQL 有多快? Spark 2.0 以后,Spark SQL 的性能已经非常接近专业的 MPP 数据库(比如 Impala、Presto), 而且更通用、生态更好。

4.2 Spark Streaming:微批处理的流计算 ​

是什么? ​

把实时数据流切成很小的批次(比如 1 秒一批),然后用批处理的方式处理每一批。

核心抽象:DStream ​

DStream(Discretized Stream)就是一系列连续的 RDD,每个 RDD 代表一个时间窗口的数据。

特点 ​

  • 微批处理:不是真正的实时,是"准实时"
  • 延迟:秒级(最低 0.5 秒左右)
  • 容错:基于 RDD 血缘,天然容错

适用场景 ​

  • 实时日志分析
  • 实时监控
  • 实时统计

⚠️ 注意: Spark Streaming 现在已经处于维护模式了,官方推荐用 Structured Streaming。 但很多老项目还在用,所以还是得了解。

4.3 Structured Streaming:新一代流计算 ​

是什么? ​

Spark 2.0 推出的新一代流计算引擎,基于 Spark SQL,用 DataFrame/Dataset API。

核心思想 ​

把流数据当成一张"无限增长的表",你对这张表写查询,就像批处理一样。

特点 ​

  • 延迟更低:可以做到 100ms 级
  • API 更统一:批处理和流处理用同一套 API
  • 事件时间支持:支持基于事件时间的窗口、Watermark
  • Exactly-Once 语义:精确一次处理保证

Spark Streaming vs Structured Streaming ​

对比项Spark StreamingStructured Streaming
基于RDDDataFrame/Dataset
延迟秒级100ms 级
APIDStreamDataFrame
事件时间不支持(只能用处理时间)原生支持
状态管理需自己实现内置支持
推荐度维护模式官方推荐

4.4 MLlib:分布式机器学习 ​

是什么? ​

Spark 自带的机器学习库,能在分布式环境下跑机器学习算法。

两大 API ​

  • spark.mllib:基于 RDD 的老 API,维护模式
  • spark.ml:基于 DataFrame 的新 API,官方推荐

包含什么? ​

  • 特征工程:TF-IDF、Word2Vec、标准化、归一化、特征选择...
  • 分类:逻辑回归、决策树、随机森林、朴素贝叶斯、SVM...
  • 回归:线性回归、随机森林回归、GBRT...
  • 聚类:K-Means、高斯混合、LDA...
  • 推荐:ALS 协同过滤
  • 评估:各种评估指标
  • Pipeline:机器学习工作流

适用场景 ​

  • 数据量大到单机跑不动的机器学习任务
  • 数据已经在 Spark 生态里,不想搬来搬去

💡 MLlib vs TensorFlow/PyTorch? 不是一个赛道的。

  • MLlib:传统机器学习(非深度学习),分布式,大数据量
  • TensorFlow/PyTorch:深度学习,主要是单机/小集群,GPU 加速

深度学习要分布式的话,一般用 TensorFlow 分布式或者 PyTorch DDP,不是 MLlib。

4.5 GraphX:分布式图计算 ​

是什么? ​

处理图数据(顶点+边)的分布式计算框架。

核心抽象 ​

  • VertexRDD:顶点 RDD
  • EdgeRDD:边 RDD
  • Graph:图,由顶点和边组成

内置算法 ​

  • PageRank(网页排名)
  • 三角形计数
  • 连通分量
  • 最短路径
  • 社区发现

适用场景 ​

  • 社交网络分析
  • 推荐系统
  • 路径规划
  • 欺诈检测

⚠️ 现状: GraphX 现在也不是特别活跃了,图计算领域现在更火的是 Neo4j(图数据库)、JanusGraph、DGL(图神经网络)等。 但作为 Spark 的组件,了解一下还是有必要的。


五、硬核对比 —— Spark vs 世界 ​

5.1 Spark vs Hadoop MapReduce ​

对比项Hadoop MapReduceSpark
计算模型Map + Reduce 两阶段DAG 多阶段
中间结果写磁盘放内存
速度慢(磁盘 IO 瓶颈)快 10~100 倍
API简单,只有 Map/Reduce丰富,几十种算子
易用性难,代码量大简单,代码量少
场景简单的批处理批处理、SQL、流、ML、图
延迟高(分钟~小时级)低(秒~分钟级)
成本省内存,费磁盘费内存,省时间

🎯 结论: Spark 在几乎所有方面都碾压 MapReduce。 现在新项目基本不会用 MapReduce 了,都是用 Spark 或者 Flink。 MapReduce 只在一些老项目、或者对内存要求极低的场景还在用。

这是大数据领域最经典的"世纪之争"。

对比项SparkFlink
计算模型微批处理(Micro-batching)真流处理(Native Streaming)
延迟秒级(Structured Streaming 可到 100ms)毫秒级
吞吐量高很高
状态管理一般(Structured Streaming 有改进)强大(State Backend 很成熟)
事件时间Structured Streaming 支持原生支持,非常完善
容错基于血缘Checkpoint + Savepoint
SQL 能力强(Catalyst 优化器)强(持续改进中)
批处理强(出身就是批处理)也支持,但流是核心
机器学习MLlib(成熟)Flink ML(较弱)
图计算GraphXGelly(较弱)
生态非常成熟,生态庞大快速发展中,生态相对小
学习曲线较低较高
国内使用非常广泛越来越多,互联网公司用得多

怎么选? ​

场景推荐
主要是批处理,偶尔流处理Spark
主要是流处理,要求低延迟Flink
需要机器学习、图计算Spark
需要复杂的状态管理、事件时间Flink
公司技术栈是 Hadoop 生态Spark
公司是互联网、实时业务多Flink

💡 我的建议:

  • 两个都学,先学 Spark(生态好、工作多、入门容易),再学 Flink(流处理更强)
  • 批处理选 Spark,流处理选 Flink,这是目前的行业共识
  • 但 Spark 也在流处理上持续进步(Structured Streaming),Flink 也在批处理上发力(批流一体)
  • 未来的趋势是"批流一体",两个框架都在往这个方向走

5.3 Spark vs Storm ​

对比项Spark StreamingStorm
模型微批处理逐条处理
延迟秒级毫秒级
吞吐量高较低
容错基于血缘,自动恢复Ack 机制
状态管理有弱(需自己实现)
API丰富较底层
易用性好一般
现状主流逐渐被 Flink 替代

🎯 结论: Storm 已经是"上一代"流计算框架了,新项目基本不用。 要低延迟用 Flink,要吞吐量和易用性用 Spark Streaming。

5.4 Spark SQL vs Hive ​

对比项Spark SQLHive (MapReduce)Hive on Spark / Tez
底层引擎SparkMapReduceSpark / Tez
速度快慢较快
延迟秒~分钟分钟~小时分钟级
SQL 兼容性兼容 Hive SQLHive SQLHive SQL
并发能力一般好好
适用场景即席查询、ETL离线数仓、大批量离线数仓

💡 现状:

  • Hive 现在更多是作为"数据仓库"的元数据管理和 SQL 接口
  • 底层执行引擎很多都换成了 Spark 或者 Tez
  • Spark SQL 可以直接读 Hive 表,两者配合使用很常见

六、前沿方向 —— Spark 3.x 之后的新世界 ​

6.1 Spark 3.x 重大升级 ​

Spark 3.0 是一个里程碑版本,带来了很多重大改进:

1. Adaptive Query Execution(AQE,自适应查询执行) ​

  • 是什么:运行时动态调整查询计划
  • 解决什么问题:静态优化器估计不准,比如数据倾斜、join 策略选错
  • 怎么工作:
    • 动态调整 shuffle 分区数
    • 动态切换 join 策略(SortMergeJoin → BroadcastJoin)
    • 动态优化数据倾斜
  • 效果:性能提升 20~30%,很多场景不用手动调优了

2. Dynamic Partition Pruning(动态分区裁剪) ​

  • 是什么:运行时根据 join 的结果裁剪不需要读取的分区
  • 效果:事实表和维度表 join 时,性能大幅提升

3. Pandas API on Spark(原 Koalas) ​

  • 是什么:让你用 Pandas 的 API 写 Spark 程序
  • 意义:数据分析师不用学 Spark API,直接用熟悉的 Pandas 就能跑分布式计算
  • 一句话:Pandas 的语法,Spark 的性能

4. 其他重要改进 ​

  • ANSI SQL 兼容:更标准的 SQL 语法
  • DataSource V2:更灵活的数据源接口
  • Kubernetes 调度:更好的 K8s 支持
  • Python UDF 性能提升:用 Apache Arrow 加速,快了很多

6.2 Spark Connect:客户端-服务端架构 ​

Spark 3.4 引入的新特性:

  • 是什么:把 Spark 拆成客户端和服务端,通过 gRPC 通信
  • 解决什么问题:
    • 以前 Spark 客户端和 Driver 绑在一起,启动慢、资源占用大
    • 现在客户端很轻,连接到远程的 Spark 服务就行
  • 好处:
    • IDE 里写代码不用本地装 Spark
    • 多客户端共享一个 Spark 集群
    • 支持更多语言(不仅仅是 Scala/Python/Java/R)
    • 更像数据库的使用方式

6.3 Photon 引擎:向量化执行引擎 ​

Databricks 公司(Spark 创始团队的公司)搞的:

  • 是什么:用 C++ 写的向量化执行引擎,替代 JVM 执行
  • 有多快:官方说比 JVM 快 2~10 倍
  • 特点:
    • 向量化执行(一次处理一批数据,不是一条一条)
    • 直接操作二进制数据,没有 JVM 对象开销
    • 充分利用 CPU 缓存和 SIMD 指令
  • 现状:目前是 Databricks 商业版的功能,开源版 Spark 还没有完全跟进

💡 趋势: 向量化执行是数据库/大数据引擎的大趋势。 ClickHouse、DuckDB、Doris 都是向量化执行,性能非常猛。 Spark 也在往这个方向走,只是步子慢一点。

6.4 GPU 加速 ​

  • 是什么:用 GPU 来加速 Spark 计算
  • 怎么做:NVIDIA 的 RAPIDS Accelerator for Apache Spark
  • 效果:某些场景(比如机器学习、数据清洗)能快好几倍
  • 现状:还在发展中,不是所有算子都支持 GPU

6.5 湖仓一体(Lakehouse) ​

这是当前大数据最火的方向之一。

什么是湖仓一体? ​

  • 数据湖:什么数据都能存(结构化、半结构化、非结构化),便宜,但管理难、性能一般
  • 数据仓库:结构化数据,性能好、管理好,但贵、不灵活
  • 湖仓一体:把两者的优点结合起来——数据存数据湖上,但能提供数仓的性能和管理能力

Spark 在其中的角色 ​

Spark 是湖仓一体的核心计算引擎之一,配合:

  • Delta Lake:Databricks 搞的,给数据湖加 ACID 事务
  • Iceberg:Apache 的,类似的东西
  • Hudi:Uber 搞的,增量数据处理强

🎯 为什么湖仓一体这么火? 因为以前数据湖和数据仓库是两套系统,数据搬来搬去,又麻烦又贵。 湖仓一体想做到"一套存储、多种计算、ACID 保证",统一架构。 Spark + Delta/Iceberg/Hudi 是目前最主流的湖仓一体方案。

6.6 大模型时代的 Spark ​

AI 大模型火了之后,Spark 也在跟进:

  • Spark ML 支持大模型:虽然不是 Spark 的强项,但也在探索
  • 数据处理:大模型训练的数据预处理,Spark 是主力之一
  • 向量检索:结合向量数据库,做 RAG 之类的应用
  • MLflow:Spark 团队搞的 MLflow,现在是 MLOps 的事实标准

💡 个人观点: 大模型时代,Spark 的定位更多是"数据基础设施"—— 处理大模型训练需要的海量数据,而不是自己做大模型训练。 大模型训练那是 GPU 和深度学习框架的事。


七、实战指南 —— 从环境搭建到第一个程序 ​

7.1 环境搭建(三种方式) ​

方式一:本地模式(最简单,适合学习) ​

bash
# 1. 安装 JDK 8 或 11
# 2. 下载 Spark
wget https://archive.apache.org/dist/spark/spark-3.5.1/spark-3.5.1-bin-hadoop3.tgz

# 3. 解压
tar -zxvf spark-3.5.1-bin-hadoop3.tgz
cd spark-3.5.1-bin-hadoop3

# 4. 启动 Spark Shell
./bin/spark-shell

搞定,就这么简单。本地模式不需要 Hadoop,不需要集群。

方式二:Docker(推荐,环境隔离好) ​

bash
# 拉取镜像
docker pull apache/spark:3.5.1

# 启动 Spark Shell
docker run -it apache/spark:3.5.1 /opt/spark/bin/spark-shell

方式三:集群模式(生产环境) ​

需要先搭 Hadoop 集群(HDFS + YARN),然后部署 Spark on YARN。 步骤比较多,建议参考官方文档或者专门的搭建教程。

7.2 第一个程序:WordCount ​

WordCount 是大数据界的"Hello World"。

方式一:Spark Shell(Scala) ​

scala
// 启动 spark-shell 后输入

// 1. 读取文件,创建 RDD
val lines = sc.textFile("README.md")

// 2. 拆分单词
val words = lines.flatMap(line => line.split(" "))

// 3. 每个单词记一次数
val wordOne = words.map(word => (word, 1))

// 4. 按 key 累加
val wordCount = wordOne.reduceByKey(_ + _)

// 5. 触发计算,收集结果
wordCount.collect().foreach(println)

方式二:一行代码版 ​

scala
sc.textFile("README.md")
  .flatMap(_.split(" "))
  .map((_, 1))
  .reduceByKey(_ + _)
  .collect()
  .foreach(println)

方式三:Spark SQL 版 ​

scala
// 用 DataFrame + SQL
val df = spark.read.text("README.md")
df.createOrReplaceTempView("words")

spark.sql("""
  SELECT word, count(*) as cnt
  FROM (
    SELECT explode(split(value, ' ')) as word
    FROM words
  )
  GROUP BY word
  ORDER BY cnt DESC
""").show()

💡 三种方式对比:

  • RDD 版:最底层,最灵活,代码最多
  • DataFrame 版:有优化器,性能可能更好,代码简洁
  • SQL 版:最简单,非技术人员也能写

7.3 核心 API 速览 ​

RDD 常用算子 ​

转换算子(Transformation):

算子作用例子
map一对一转换rdd.map(x => x * 2)
flatMap一对多,然后打平rdd.flatMap(x => x.split(" "))
filter过滤rdd.filter(x => x > 10)
distinct去重rdd.distinct()
reduceByKey按 key 聚合rdd.reduceByKey(_ + _)
groupByKey按 key 分组rdd.groupByKey()
sortBy排序rdd.sortBy(x => x)
join连接两个 RDDrdd1.join(rdd2)
union合并两个 RDDrdd1.union(rdd2)
intersection交集rdd1.intersection(rdd2)

行动算子(Action):

算子作用
collect()收集所有数据到 Driver
count()统计数量
first()取第一个
take(n)取前 n 个
reduce()聚合
foreach()遍历执行
saveAsTextFile()保存到文件

⚠️ 重要提醒: collect() 会把所有数据拉到 Driver 端,数据量大的时候会 OOM! 生产环境慎用,一般用 take() 看几条就行,或者直接 saveAsTextFile。


八、性能调优圣经 —— 让你的 Spark 飞起来 ​

Spark 程序跑慢了?别慌,按这个清单一步步来。

8.1 调优总览 ​

性能调优金字塔:

        ┌─────────────┐
        │  代码优化    │  ← 性价比最高,先做这个
        ├─────────────┤
        │  参数调优    │  ← 然后调参数
        ├─────────────┤
        │  资源调优    │  ← 加机器、加内存
        ├─────────────┤
        │  架构优化    │  ← 最后考虑,成本最高
        └─────────────┘

8.2 代码优化(最重要) ​

1. 避免 Shuffle ​

Shuffle 是 Spark 性能的第一杀手。因为要跨节点传输数据,还要写磁盘。

常见触发 Shuffle 的操作:

  • reduceByKey、groupByKey
  • join、distinct
  • sortBy、repartition
  • coalesce(shuffle=true)

怎么减少 Shuffle:

  • 能用 reduceByKey 就不用 groupByKey(reduceByKey 有本地预聚合)
  • 小表 join 大表用 broadcast join(把小表发到每个节点,不用 Shuffle)
  • 过滤尽早做,减少数据量
  • 能在 map 端做的事,不要等到 reduce 端

2. 用好缓存 ​

如果一个 RDD/DataFrame 要被多次使用,缓存起来!

scala
// 缓存
rdd.cache()  // 等价于 persist(MEMORY_ONLY)

// 或者指定存储级别
df.persist(StorageLevel.MEMORY_AND_DISK)

// 不用了释放
rdd.unpersist()

什么时候缓存:

  • 同一个 RDD 被多次使用
  • 计算这个 RDD 很耗时
  • 数据量不是特别大(放得下内存)

3. 避免数据倾斜 ​

数据倾斜就是:有的分区数据特别多,有的特别少。 结果就是快的很快,慢的很慢,整体被最慢的那个拖死。

怎么发现数据倾斜:

  • Spark UI 上看 Stage,有的 Task 很快,有的特别慢
  • 看 Shuffle Read 数据量,差异很大

怎么解决:

  1. 过滤异常 key:如果是某个 null 或者特殊值导致的,直接过滤掉
  2. 加盐(加盐打散):给倾斜的 key 加上随机前缀,分成多个分区处理,然后再合并
  3. Broadcast Join:如果是 join 倾斜,小表用 broadcast
  4. AQE 自动优化:Spark 3.x 开 AQE,自动处理部分倾斜

4. 用好 Spark SQL 优化器 ​

  • 能用 DataFrame/SQL 就不用 RDD(Catalyst 优化器会自动优化)
  • 开 AQE(自适应查询执行):spark.sql.adaptive.enabled=true
  • 开动态分区裁剪

8.3 参数调优 ​

常用重要参数 ​

参数作用建议值
spark.executor.memory每个 Executor 的内存根据机器配置,一般 4~16g
spark.executor.cores每个 Executor 的 CPU 核数2~8 核
spark.executor.instancesExecutor 数量根据集群规模
spark.driver.memoryDriver 内存一般 2~4g,collect 多的话调大
spark.sql.shuffle.partitionsShuffle 分区数默认 200,数据量大可以调大
spark.sql.adaptive.enabled自适应查询执行true(Spark 3.x 推荐开)
spark.sql.adaptive.coalescePartitions.enabled动态合并分区true
spark.sql.adaptive.skewJoin.enabled倾斜 join 优化true
spark.memory.fraction执行内存占比默认 0.6,一般不用改
spark.serializer序列化器org.apache.spark.serializer.KryoSerializer

Shuffle 分区数怎么设? ​

  • 原则:每个分区大约 128MB ~ 256MB
  • 公式:总数据量 / 128MB ≈ 分区数
  • 开了 AQE 的话可以设大一点,AQE 会自动合并小分区

8.4 资源调优 ​

资源分配原则 ​

  • CPU 和内存配比:一般 1 核配 2~4G 内存比较合理
  • 不要把机器资源占满:留点给系统和 HDFS 等其他进程
  • Executor 数量不是越多越好:太多了 Driver 调度不过来

常见配置示例(8 核 32G 机器) ​

spark.executor.memory    = 16g
spark.executor.cores     = 4
spark.executor.instances = 2
(每台机器 2 个 Executor,每个 4 核 16G)

8.5 调优步骤(建议顺序) ​

  1. 看 Spark UI:找到慢的 Stage,看是哪个操作慢
  2. 看数据量:输入多少,Shuffle 多少,是不是数据倾斜
  3. 优化代码:减少 Shuffle、加缓存、处理倾斜
  4. 调参数:分区数、内存、开 AQE
  5. 加资源:最后才考虑加机器加内存

🎯 调优心法: 调优不是瞎调参数,是先找到瓶颈,再针对性优化。 用 Spark UI 定位问题,这是最重要的技能。


九、排错手册 —— 常见问题与解决方案 ​

9.1 内存相关 ​

OOM:OutOfMemoryError ​

Driver OOM:

  • 原因:collect() 太多数据,或者 Driver 内存太小
  • 解决:
    • 不要 collect 全量数据,用 take 或者直接保存
    • 调大 spark.driver.memory

Executor OOM:

  • 原因:数据倾斜、单个分区太大、内存不够
  • 解决:
    • 增加分区数,让每个分区小一点
    • 处理数据倾斜
    • 调大 spark.executor.memory
    • 用 persist(MEMORY_AND_DISK),内存不够溢写磁盘

9.2 Shuffle 相关 ​

Shuffle 文件找不到 ​

  • 原因:Executor 挂了,Shuffle 文件丢了
  • 解决:
    • 检查 Executor 日志,看为什么挂(一般是 OOM)
    • 调大 Executor 内存
    • 开启 Shuffle 服务(External Shuffle Service)

Shuffle 数据量太大 ​

  • 原因:数据量大,或者 join 后数据膨胀
  • 解决:
    • 尽早过滤,减少数据量
    • 调大分区数
    • 用 broadcast join 如果是小表

9.3 序列化相关 ​

NotSerializableException ​

  • 原因:你在算子里面用了一个不能序列化的对象
  • 解决:
    • 把对象改成可以序列化的(实现 Serializable)
    • 或者在算子内部创建对象,不要从外面传
    • 用 Kryo 序列化器

9.4 数据倾斜 ​

怎么判断是不是数据倾斜? ​

  • Spark UI 上看 Task 执行时间,有的几秒,有的几小时
  • Shuffle Read 数据量差异很大

怎么解决? ​

  1. 过滤异常 key:先看看是不是某个 key 特别多,能不能过滤掉
  2. 加盐打散:给倾斜的 key 加随机前缀,分两阶段聚合
  3. Broadcast Join:小表 join 大表的倾斜
  4. AQE 倾斜优化:Spark 3.x 开了 AQE 自动处理

9.5 其他常见问题 ​

任务跑的特别慢 ​

  • 检查是不是数据倾斜
  • 检查是不是 Shuffle 太多
  • 检查是不是没加缓存(重复计算)
  • 看 Spark UI 定位瓶颈

结果不对 ​

  • 检查是不是有数据倾斜导致的问题
  • 检查数据类型转换有没有问题
  • 检查 null 值处理
  • 用小数据量验证逻辑

十、学习路径与资源 —— 从入门到精通 ​

10.1 学习路径建议 ​

第一阶段:入门(1~2 周) ​

  • 了解 Spark 是什么、能做什么
  • 搭环境,跑通 WordCount
  • 学习 RDD 基本算子:map、filter、reduceByKey...
  • 理解核心概念:RDD、宽窄依赖、Stage、DAG

第二阶段:进阶(2~4 周) ​

  • Spark SQL / DataFrame(重点,工作中用得最多)
  • Spark Streaming / Structured Streaming
  • 性能调优基础
  • 集群模式:YARN、spark-submit

第三阶段:深入(1~2 个月) ​

  • 源码阅读(可选,想进大厂建议看)
  • 高级调优
  • MLlib / GraphX(根据需要)
  • 项目实战

第四阶段:拓展 ​

  • 学习 Flink(流处理方向)
  • 学习湖仓一体:Delta Lake / Iceberg / Hudi
  • 学习云原生:Spark on K8s

10.2 推荐资源 ​

官方资源 ​

书籍 ​

  • 《Spark 快速大数据分析》 —— 入门经典,就是有点老
  • 《Spark 权威指南》 —— 比较全面,Bill Chambers 写的
  • 《深入理解 Spark 核心思想与源码分析》 —— 想深入源码可以看

实战项目 ​

  • 日志分析
  • 推荐系统
  • 实时数据仓库
  • 用户行为分析

10.3 面试常考知识点 ​

按频率排序:

  1. 宽窄依赖、Stage 划分 —— 几乎必问
  2. Spark 为什么快 —— 必问
  3. reduceByKey vs groupByKey —— 高频
  4. cache vs persist —— 高频
  5. coalesce vs repartition —— 高频
  6. Spark 运行架构(Driver/Executor) —— 高频
  7. 数据倾斜怎么解决 —— 高频
  8. 性能调优 —— 高频
  9. Spark vs MapReduce / Flink —— 中频
  10. Shuffle 过程 —— 中高频
  11. 容错机制(Lineage) —— 中频
  12. Spark SQL 优化(Catalyst、AQE) —— 中高级岗常问

写在最后 ​

Spark 是大数据领域的"通用语言",几乎所有大数据岗位都要求会。

它不是银弹,也有自己的局限(比如流处理不如 Flink,图计算不如专门的图数据库),但它是目前生态最完善、应用最广泛的大数据计算框架。

学 Spark,最重要的是动手。光看没用,一定要自己搭环境、写代码、调 Bug。 写得多了,自然就懂了。

祝你在大数据的世界里,玩得开心。

—— 2026 年 8 月


本文档版本:V1.0 最后更新:2026 年 8 月 适用 Spark 版本:3.5.x

基于 Vite 强力驱动 | 纯静态轻量托管