☰
Spark3.x核心概念详解:架构角色、API选择与执行模型
2026/9/26 23:06:11 网站建设 项目流程

1. 为什么Spark让人越学越乱:先从整体心智模型说起

先说个现象。我接触过不少Spark学习者,也包括团队里进来的新人,大多数人的状态是:能照着网上demo把代码跑起来,但改了参数就懵,报了错就慌,问“你这作业到底怎么分配到各个节点上的”,回答含含糊糊。这不是学习者不够聪明,而是Spark本身就是一个复合体——它既是一个编程模型,又是一套执行引擎,还叠加了资源调度和多种运行模式。如果一开始眼里只有API和代码,没有在脑子里建立一张“这东西到底由哪些部分组成、各部分各干什么”的图景,后面所有细节都会变成一堆散沙。

这篇是Spark3.x指北系列的第一篇,先把Spark3.x的基础概念按我的理解彻底捋一遍。我尽量不堆砌教科书定义,而是往“这东西到底是什么、为什么要这么设计、实际使用时需要注意什么”这三个方向去讲。适合三类人看:刚接触Spark想系统入门的新人;用过Hive或MapReduce、想搞清Spark和它们本质差异的迁移者;以及已经在用Spark但一直靠试错来解决问题、想补一补底层认知的工程师。

先把最重要的一句话放前面:Spark本质上是一个统一的分布式数据分析引擎。所谓“统一”,意思是它不只做某一种计算——批处理、交互式查询、流式处理、机器学习、图计算,都能在同一套框架里完成;所谓“分布式”,意思是数据分散在多个节点的内存和磁盘上,计算也被切分成很多小任务并行执行;所谓“引擎”,意思是它本身不存储数据,只管算——你的数据可以放在HDFS、S3、本地文件系统,也可以来自数据库或消息队列。

很多人对Spark的误解,是把Spark和Hadoop MapReduce放在同一条赛道里比较,觉得“Spark就是比MapReduce快十倍的替代品”。这个说法不算错,但视野太窄了。Spark真正改变的是计算范式:MapReduce把每一步计算都强制落盘,Spark则尽量把中间结果留在内存,并且用DAG(有向无环图)把一系列计算组织成一个整体来调度。再加上Spark SQL、Structured Streaming、MLlib、GraphX这一整层生态,它早就超出了“替代品”的定位。

所以我的建议是:学习Spark3.x,第一步不是在IDE里写代码,而是先搞懂它是怎么组织的。这一篇就是干这件事。我会把Spark的架构角色、运行模式、核心编程抽象、执行模型这几个关键支柱依次拆开。等这几块概念通了,再去看具体API和调优参数,你会发现一切都顺理成章。

2. 集群运行时的核心角色:Driver、Executor、Master、Worker各管什么事

Spark运行时的角色划分,是初学者最容易产生混乱的地方。混乱的原因在于,这些角色的名字一部分来自Spark自身,一部分来自它依赖的外部资源管理器。得先把这两套体系分开。

2.1 资源管理器与Spark运行时是两套并行体系

先说外部资源管理器。Spark本身不分配CPU和内存,它需要向一个资源调度器申请资源。常见的选择有三种:Spark自带的Standalone、Hadoop生态里的YARN、云原生时代的Kubernetes。如果你用的是本地模式,那其实连资源管理器都不需要,所有东西跑在一个JVM进程里。

真正属于Spark自身的核心进程角色有两个:Driver和Executor。另外还有一个Master角色,但Master只在Standalone模式下才由Spark自己承担,在YARN或Kubernetes模式下,Master的职责被对应集群的ResourceManager等组件接管了。

这里有个很重要的概念边界:Driver和Executor是计算层面的角色,它们关心的是“怎么把一个计算任务拆开并执行”;Master/Worker是资源层面的角色,它们关心的是“哪些机器有空闲资源”。Spark的任务提交过程,简单来说就是:用户写好代码,通过spark-submit提交给集群,Driver启动后向资源管理器申请Executor资源,资源管理器分配可用节点,Executor在节点上启动并等待接收任务。

2.2 Driver进程:整个应用的“总指挥”

Driver是跑你main函数的那个进程,里面包含SparkContext(3.x习惯上通过SparkSession间接创建)。它承担了以下几件事:

  • 把用户代码转换成计算逻辑DAG图;
  • 把DAG进一步划分成不同的Stage;
  • 把Stage里的任务集分发给不同的Executor节点;
  • 汇总任务执行状态、接受计算结果。

一句话总结:Driver负责“思考该算什么、怎么分、谁去算”,不实际承担大量数据的计算。它更像一个项目经理,而Executor是干活的一线员工。

在实际生产环境中,Driver如果运行在集群之外(比如你本地电脑或调度机),它和集群之间的网络通信带宽、自身的内存大小都要预留足够,否则它在处理大量结果回传或心跳信息时可能先崩掉。

2.3 Executor进程:真正干活的“一线员工”

Executor是运行在工作节点上的JVM进程,会在应用运行期间常驻。它做两件事:

  • 执行分配给它的Task(一个Executor里可以跑多个Task,取决于CPU核心数);
  • 把计算结果、累积器等数据存储或回传给Driver。

每个Executor能并发执行的Task数量,由分配给它的CPU核心数决定。比如--executor-cores 2表示分配两个虚拟核心,那同一时间最多并行跑两个Task。Executor的内存则分为几个区域,其中最重要的是给RDD缓存用的Storage Memory和给执行运算用的Execution Memory,从Spark 1.6开始这俩是统一管理的,后面细说。

2.4 一个经典问题:spark-submit参数到底怎么设置

因为Executor是整个计算真正的执行单元,生产上最常被问的就是参数怎么定。先把常用参数意思对照列一下:

参数含义影响
--master运行模式local、yarn、k8s、standalone 等
--num-executors启动多少个Executor进程决定横向并行度
--executor-cores每个Executor占用的CPU核心数决定每个进程内并发Task数
--executor-memory每个Executor占用的内存JVM堆大小
--driver-memoryDriver占用的内存Driver堆大小

实际的配置并不是越大越好。我见过太多人把--executor-memory直接堆到几十G,以为内存越多越快。但实际上,当内存超过一定阈值后,JVM的垃圾回收暂停时间会明显增长,尤其G1GC模式下如果分配不当,一次Full GC可能卡上几秒。还有个大坑是内存总量不能超过节点剩余资源,否则资源管理器直接把你的作业挂起等待,看起来像是任务不跑,实际是资源不够。

我自己的经验值:生产环境单个Executor内存设置不超过8G-12G为宜,CPU核心数通常2-4个,然后根据数据集大小和作业并发度调整Executor数量。这样即使单个Executor挂掉,影响范围也有限。Executor天然有容错机制——一个失效后会由资源管理器重新拉起,但如果单Executor内存被撑爆导致OOM,重新拉起的Executor依然会挂,所以别指望靠执行力度的容错来处理OOM。

另外,无论集群多大多小,有个不能忘记的角色叫BlockManager。它其实是Executor内部的一个组件,负责数据块的管理和传输,包括RDD缓存数据、shuffle过程中的数据读取。因为shuffle数据往往需要跨节点传输,BlockManager是执行链路里的关键节点,它在Spark架构中确实承担了最底层的数据流转职能。

3. RDD、DataFrame和Dataset:三个编程API为什么共存,怎么选择

Spark3.x的编程模型在API层面有三套面孔:底层的RDD、主流的DataFrame,以及Scala/Java里才有的Dataset。很多人刚开始完全分不清它们之间的关系,看到的资料又多,越看越晕。我从设计缘由讲起。

3.1 RDD:最底层的分布式数据集抽象

RDD(Resilient Distributed Dataset)是Spark最早的核心抽象,名字里有三个关键词:

  • Resilient(弹性):节点故障时,可以基于数据血缘关系重新计算丢失的分区;
  • Distributed(分布式):数据以分区(Partition)的形式分布在不同节点上;
  • Dataset(数据集):它代表一个只读的、可分区的数据集合,而不是具体的数据存储格式。

你可以把它理解成分布在多台机器上的数组集合,支持两类操作:Transformation(转换,懒执行)和Action(行动,触发真正计算)。

RDD的好处是灵活可控,几乎所有数据形态都能通过RDD处理,你可以自定义分区器、定义复杂的键值操作。缺点也明显:没有内建的结构化信息,Spark拿到的只是一个元素集合,无法针对“字段类型”“列名”做深度优化。同时,用RDD写代码通常比DataFrame冗长,处处都是map、flatMap、reduceByKey这些细节。

3.2 DataFrame:一张分布在集群上的“表”

DataFrame在概念上可以类比成关系型数据库里的表,它增加了Schema(模式),也就是每个字段的名称和类型。因为有了Schema,Spark就能做很多聪明事:

  • 知道某个字段是数值类型,就可以用更高效的二进制格式存储和序列化;
  • 可以把用户写的高层操作翻译成执行计划,再做优化;
  • 可以和Spark SQL无缝互通,因为DataFrame本身就是带Schema的分布式行集合。

更大的价值在于,DataFrame的API背后是Catalyst优化器和Tungsten执行引擎。Catalyst会把你的代码转换成一个逻辑计划,再用规则进行优化(比如谓词下推、列剪枝、常量折叠),最终变成最优的物理计划;Tungsten则通过直接操作二进制内存、减少Java对象开销的方式大幅提升执行效率。这些优化对RDD是不存在的,所以相同逻辑用DataFrame写通常比用RDD快数倍。

3.3 Dataset:类型安全的DataFrame

Dataset是DataFrame在类型层面的增强。DataFrame叫Dataset[Row],而Dataset则在编译时就检查字段类型——你在Scala里定义了一个Person类,那么ds.filter(_.age > 20)这类操作会在编译期就发现类型错误,而DataFrame里filter("age > 20")这类字符串表达式要等运行时才知道对错。

这里有个很多人踩过的坑:在Python/PySpark里,其实没有独立的Dataset API。你在PySpark里写的都是DataFrame,或者通过RDD转DataFrame。Dataset的强类型优势只在Scala和Java中才有。所以做技术选型之前,先看团队的开发语言——如果团队用Python,就别纠结Dataset,DataFrame就是最优解;如果团队主力是Scala,那复杂业务里用Dataset确实更好。

给一个选择建议:新项目没有特殊需求,一律以DataFrame为主。只有当你要处理的数据是极端非结构化、或者需要自定义分区和操作底层RDD功能时,再把数据转回RDD处理,处理完之后再转回DataFrame。

4. 执行模型的核心:宽窄依赖、Shuffle与Stage切分

Spark能跑得这么快,不只是“把中间结果放内存”这么简单。真正理解它的执行模型,需要抓住三个概念:依赖关系、Shuffle和Stage。这三个概念串起来,你就能读懂Spark UI上的DAG图,也能定位大部分性能问题。

4.1 窄依赖与宽依赖:数据要不要跨节点搬移

RDD之间的操作会构成血缘关系。根据父RDD和子RDD分区的对应关系,依赖被分为两类:

窄依赖:父RDD的每个分区最多被一个子RDD分区使用。典型操作有map、filter、union。这种依赖下,每个分区只和数据在本分区内发生关联,不需要数据跨节点搬运。

宽依赖:父RDD的每个分区可能被多个子RDD分区使用。典型操作有groupByKey、reduceByKey、join。宽依赖会触发Shuffle——父分区里的数据必须按照某种规则重新排列,划分给不同的下游分区,而这些下游分区往往在别的节点上。

用一个生活中的例子:窄依赖就像每个班级自己改本班作业,改完直接上报;宽依赖像全校重新分班,先要把所有学生名单汇总,再按新规则分发到各个班级,中间省不掉一个全局调配的动作。

4.2 Shuffle:Spark性能的第一瓶颈

说到Shuffle,我建议所有Spark学习者都把“Shuffle慢”这件事刻在脑子里。从MapReduce时代开始,Shuffle就一直是分布式计算里最昂贵的一环,Spark也没能消除它,只是比MapReduce更优化了一些。

在Shuffle过程中,每个Mapped任务(上游Task)会把结果按照分区器规则写到本地磁盘,然后下游的Reduce任务需要通过网络从上游节点拉取属于自己那份的数据。这中间涉及磁盘I/O、网络传输、数据排序/聚合,任何一个环节都可能成为瓶颈。

因此Spark里有一个经典优化法则:减少Shuffle,或避免Shuffle。比如:

  • 用reduceByKey代替groupByKey,因为前者会在上游节点先做一次本地聚合,让Shuffle的数据量显著减少;
  • join前先做布隆过滤或提前裁剪数据集;
  • 对于一些可预测的重复计算,使用分区缓存甚至用repartition后多次复用同一分区规则。

4.3 DAG与Stage:Spark怎么“切”计算图

用户写的一系列转换操作会形成一张计算图,也就是DAG。DAG中从某个地方开始到下一个宽依赖的边界,会被切成一个Stage。换句话说,遇到宽依赖就切Stage,窄依赖不会切。

Stage又被划分为一组可以并行执行的Task。每个分区对应一个Task,Task在Executor上执行同一个逻辑但处理不同分区数据。

具体到WordCount例子,它被划分成几个Stage?从代码来看:

rdd = sc.textFile("hdfs:///input.txt") # 加载 words = rdd.flatMap(lambda line: line.split()) # 切词 pairs = words.map(lambda w: (w, 1)) # 映射 counts = pairs.reduceByKey(lambda a, b: a + b) # 聚合 counts.saveAsTextFile("hdfs:///output")

这里reduceByKey是宽依赖操作,所以DAG会在它之前切一刀:Stage 0 执行textFile到map,并把结果按key分区写入shuffle文件;Stage 1 拉取shuffle数据后执行reduceByKey聚合,最后写结果。如果这里用的是groupByKey,逻辑也是两段,但shuffle数据量会大得多。

4.4 血缘与容错:为什么不轻易Checkpoint

RDD还有一个特色机制——血缘图谱(Lineage)。每个RDD都记得自己是怎么从父RDD计算出来的。一旦某个分区的数据丢失,Spark可以根据血缘关系重新计算该分区,而不是整个数据集重建。这种容错方式比MapReduce直接重放整个任务要轻量得多。

血缘机制的基础知识很容易理解,但生产经验是:当血缘链条过长(比如几百个RDD连续转换),或者上游基线数据重建成本极高时,重建一个分区的代价也可能很大。这时候需要做Checkpoint(检查点),把某个阶段的RDD数据直接持久化到可靠存储(如HDFS)上,相当于砍断血缘关系,以存储换时间。但Checkpoint有明确代价,它需要重新计算一次并落盘,所以应该用在血缘链长且中间数据经过高代价处理后仍有复用价值的位置。

5. 四种运行模式的适用场景与选择建议

Spark3.x支持多种运行模式,每个模式都有自己的边界条件。选错了,要么开发验证效率极低,要么生产环境维护困难。我按实际使用场景来讲。

5.1 Local模式:学习、验证、开发调式首选

--master local[4]这种参数开启的就是本地模式,其中数字代表本地启动的线程数。它不需要任何集群,Driver和Executor都在同一个进程里(某些实现中会为Executor单独开线程或进程)。Local模式下数据不会跨网络传输,非常适合跑通逻辑、单机调试、单元测试。

这里有个细节值得注意:local模式的并发能力受限于本机CPU核心数,所以local[*]表示用所有可用核心。如果你在本地测试了一个map消耗很大的作业,从local模式得出的性能结论几乎完全不能反映集群表现——集群里shuffle和网络传输会变成主导因素。

5.2 Standalone模式:Spark自己管理资源

Standalone是Spark自带的简单资源管理方式。它有一个Master进程和多个Worker进程,Master负责决策哪个Worker进程给当前作业分配Executor资源,Worker则负责启动Executor供Driver调用。

胜在部署简单、环境依赖少,适合学习集群概念,也适合小规模内部环境。但缺点明显:没有YARN那样成熟的多租户队列和资源共享能力,出现节点故障时的资源管理和恢复能力相对竞品弱。现实中大公司不会用Standalone跑核心生产任务,它更多见于中小团队或测试集群。

5.3 YARN模式:生产环境的老牌选择

YARN模式下,Spark Master角色让位给YARN的ResourceManager。Driver提交到集群后,由ResourceManager协调启动Executor。YARN具体有两种提交方式:

  • yarn-client:Driver跑在客户端机器上,适合交互式使用;
  • yarn-cluster:Driver跑在集群的一个AM(ApplicationMaster)进程中,适合生产批处理任务。因为Driver不在客户端,任务提交后客户端即使断开连接也不会影响运行。

YARN的优势是与Hadoop生态融合好,可以利用YARN的队列机制做多租户资源隔离和调度策略。如果你的集群同时跑Hive、Flink等任务,YARN模式通常是最通用的选择。

5.4 Kubernetes模式:云原生的未来方向

Kubernetes模式是Spark 3.x重点发展的方向。它的基本思路是把Driver和Executor都做成Pod,由Kubernetes动态调度。相比YARN,K8s模式在弹性伸缩、资源利用率上更有优势,也更适合混合部署。但代价是运维复杂度上升,对Kubernetes集群本身的稳定性、网络插件(如CNI)的性能都有更苛刻的要求。如果你的团队已经有成熟的K8s平台,且任务以定时批处理为主,Spark on K8s是一个值得考虑的方案。

把四种模式放到一张表里直观对照:

模式资源管理适用阶段优缺点
Local无开发调试方便、快捷,但不能验证分布式问题
StandaloneMaster/Worker小集群/学习部署简单,功能少
YARNResourceManager生产环境(Hadoop生态)成熟稳定,多租户好
Kubernetes容器编排云原生环境弹性好,运维成本高

无论选哪种模式,有一个概念要贯穿始终:部署模式隔离了“资源从哪来”的差异,但不影响RDD/DAG的执行模型。换句话说,代码在local模式和集群模式跑起来执行计划是一致的,只是在哪个进程、哪些节点上执行不同。这也是为什么在本地“跑通了”的生产代码到了集群上跑不对时,首先要去查资源参数、数据分区数量,而不是反复看代码逻辑。

6. Spark3.x初体验:从下载到跑起第一个作业的完整过程

基础概念讲再多,不如亲手跑一次。这里给一个最快路径,让你在一台机器上完成Spark3.x的安装和验证。我以Spark 3.5.x版本和本机部署为例,操作系统以Linux/Mac为主。

6.1 下载与安装

不依赖Hadoop集群时,可以选择官方提供的预编译版本。到官网下载页里找spark-3.5.x-bin-hadoop3这类包,下载后解压,放到/opt/spark或用户目录下。不需要修改任何配置就能启动local模式。

两个环境变量建议配好:

export SPARK_HOME=/opt/spark export PATH=$SPARK_HOME/bin:$PATH

如果你想跑Standalone模式,启动命令也很简单:

$SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-worker.sh spark://<主机名>:7077

Master默认在8080端口提供Web UI(新版Spark默认5678端口的说法并不准确,需要确认实际配置文件)。启动后访问控制台可以看到Worker状态和资源使用情况。

6.2 用spark-shell快速体验

spark-shell是Spark自带的交互式命令行,Scala接口,local模式下启动:

$SPARK_HOME/bin/spark-shell --master local[2]

进去后,试着执行最简单的例子:

val textFile = sc.textFile("README.md") textFile.count()

这里count()是Action,会触发真实计算。你会在控制台看到一堆日志,关键是Info日志里的DAGScheduler相关输出,它会告诉你Stage如何划分、任务如何调度。

PySpark用户则可以启动pyspark,体验基本一致的交互流程:

$SPARK_HOME/bin/pyspark --master local[2]

6.3 用Python代码实现第一个作业

如果要提交一个完整的独立脚本,用Python最直观。下面这段代码,覆盖了从初始化SparkSession到执行查询再到停止会话的完整生命周期:

from pyspark.sql import SparkSession spark = SparkSession.builder \ .appName("spark-basics-demo") \ .master("local[2]") \ .getOrCreate() # 从本地文件系统创建一个DataFrame df = spark.read.text("README.md") # 切词统计(体现flatMap + groupBy的转换过程) from pyspark.sql import functions as F word_counts = ( df.select(F.explode(F.split(F.col("value"), "\\s+")).alias("word")) .groupBy("word") .count() .orderBy(F.col("count").desc()) ) word_counts.show(10) spark.stop()

这个例子很短,但你可以从Spark UI(默认4040端口)看到这两个Stage是怎么划分的:第一个Stage做文件扫描和切词,遇到groupBy触发的宽依赖之后,第二个Stage做聚合。用UI里的DAG Viz功能看执行计划,能比自己读代码理解得更直观。

一个小提醒:SparkSession是Spark 2.0之后统一入口,它内部包含了SQLContext和HiveContext的功能。所以在3.x里,不需要再显式创建SparkContext,直接继承自SparkSession即可。

7. 从这页纸开始去读你的第一份Spark UI

很多人拿到Spark UI不知道看什么。我建议所有刚学完基础概念的人,把跑完一个作业后的Spark UI当作第一份读图练习来对待。

打开UI,分成几个关键部分:

  • Jobs页:列出当前应用触发的所有Action任务,可以看作一个个完整计算请求;
  • Stages页:展示每个Job被切分出的Stage列表,看一眼就知道你的作业里到底有多少次宽依赖;
  • Storage页:显示RDD/DataFrame被缓存到内存中的情况,能帮你确认缓存策略是否生效;
  • Executors页:列出集群中所有Executor的资源使用、GC消耗和运行时状态。

看完之后再回去想我们这一篇的基础概念:Job对应一次Action;Stage对应宽依赖之间的一个计算片段;Task数量对应分区数;Executor页里的每个进程对应一个JVM实例。这些概念不是死的名词,而是UI上真实跳动的数字。以后你写调优参数的时候,内心对照的就是这一张张页面。

顺便说一句,Spark UI默认端口是4040,但如果你同时跑多个应用,后启动的应用会自动递增端口。在实际的开发过程中,我最常做的一件事就是盯着Stage上每个Task的耗时分布,如果绝大多数Task几秒跑完、唯独几个Task要花几分钟,那么大概率是数据倾斜——这种情况在后面的系列里会专门展开讲。

8. 写在最后:每个概念背后都是一个“为什么要这样设计”的问题

作为系列的第一篇,基础概念就到这里。总结起来,学习Spark的路线其实是在几组关系之间反复对照:

第一组是“角色关系”:Driver决定算什么,Executor负责执行,资源管理器决定在哪跑。第二组是“API关系”:RDD灵活但费事,DataFrame结构化且可优化,Dataset在Scala里才是完整形态。第三组是“执行关系”:窄依赖在本地完成流水线,宽依赖引发Shuffle,Stage划分的标准就是宽依赖的边界。任何一道Spark性能调优题目,最终都能回归到这三组关系上。

我个人的体会是,初学者最不应该做的事情,就是拿着工具书一个API一个API去啃。因为Spark的API数量庞大,但绝大多数都是RDD的mapfilterflatMapreduceByKey以及DataFrame的selectfiltergroupByjoin的组合使用。API只是表皮,真正决定你能不能看懂Spark UI、能不能定位性能瓶颈、能不能设计出可扩展的数据管道的,恰恰是这一篇里讲的那些看似枯燥的基础概念。

下一期我计划写Spark SQL的Catalyst优化原理和Spark 3.x新版本的执行优化特性,把“DataFrame为什么更快”这件事从头到尾讲透。先把这些底层模型装进脑子,再去调优,你会发现一切都顺理成章了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询