尧图精选

Spark3.x核心概念详解:架构角色、API选择与执行模型

🕒 发布时间:2026/9/26 23:06:14 📁 来源:尧图网络
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函数的那个进程里面包含SparkContext3.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最底层的分布式数据集抽象RDDResilient Distributed Dataset是Spark最早的核心抽象名字里有三个关键词Resilient弹性节点故障时可以基于数据血缘关系重新计算丢失的分区Distributed分布式数据以分区Partition的形式分布在不同节点上Dataset数据集它代表一个只读的、可分区的数据集合而不是具体的数据存储格式。你可以把它理解成分布在多台机器上的数组集合支持两类操作Transformation转换懒执行和Action行动触发真正计算。RDD的好处是灵活可控几乎所有数据形态都能通过RDD处理你可以自定义分区器、定义复杂的键值操作。缺点也明显没有内建的结构化信息Spark拿到的只是一个元素集合无法针对“字段类型”“列名”做深度优化。同时用RDD写代码通常比DataFrame冗长处处都是map、flatMap、reduceByKey这些细节。3.2 DataFrame一张分布在集群上的“表”DataFrame在概念上可以类比成关系型数据库里的表它增加了Schema模式也就是每个字段的名称和类型。因为有了SchemaSpark就能做很多聪明事知道某个字段是数值类型就可以用更高效的二进制格式存储和序列化可以把用户写的高层操作翻译成执行计划再做优化可以和Spark SQL无缝互通因为DataFrame本身就是带Schema的分布式行集合。更大的价值在于DataFrame的API背后是Catalyst优化器和Tungsten执行引擎。Catalyst会把你的代码转换成一个逻辑计划再用规则进行优化比如谓词下推、列剪枝、常量折叠最终变成最优的物理计划Tungsten则通过直接操作二进制内存、减少Java对象开销的方式大幅提升执行效率。这些优化对RDD是不存在的所以相同逻辑用DataFrame写通常比用RDD快数倍。3.3 Dataset类型安全的DataFrameDataset是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就别纠结DatasetDataFrame就是最优解如果团队主力是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 ShuffleSpark性能的第一瓶颈说到Shuffle我建议所有Spark学习者都把“Shuffle慢”这件事刻在脑子里。从MapReduce时代开始Shuffle就一直是分布式计算里最昂贵的一环Spark也没能消除它只是比MapReduce更优化了一些。在Shuffle过程中每个Mapped任务上游Task会把结果按照分区器规则写到本地磁盘然后下游的Reduce任务需要通过网络从上游节点拉取属于自己那份的数据。这中间涉及磁盘I/O、网络传输、数据排序/聚合任何一个环节都可能成为瓶颈。因此Spark里有一个经典优化法则减少Shuffle或避免Shuffle。比如用reduceByKey代替groupByKey因为前者会在上游节点先做一次本地聚合让Shuffle的数据量显著减少join前先做布隆过滤或提前裁剪数据集对于一些可预测的重复计算使用分区缓存甚至用repartition后多次复用同一分区规则。4.3 DAG与StageSpark怎么“切”计算图用户写的一系列转换操作会形成一张计算图也就是DAG。DAG中从某个地方开始到下一个宽依赖的边界会被切成一个Stage。换句话说遇到宽依赖就切Stage窄依赖不会切。Stage又被划分为一组可以并行执行的Task。每个分区对应一个TaskTask在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 血缘与容错为什么不轻易CheckpointRDD还有一个特色机制——血缘图谱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-clientDriver跑在客户端机器上适合交互式使用yarn-clusterDriver跑在集群的一个AMApplicationMaster进程中适合生产批处理任务。因为Driver不在客户端任务提交后客户端即使断开连接也不会影响运行。YARN的优势是与Hadoop生态融合好可以利用YARN的队列机制做多租户资源隔离和调度策略。如果你的集群同时跑Hive、Flink等任务YARN模式通常是最通用的选择。5.4 Kubernetes模式云原生的未来方向Kubernetes模式是Spark 3.x重点发展的方向。它的基本思路是把Driver和Executor都做成Pod由Kubernetes动态调度。相比YARNK8s模式在弹性伸缩、资源利用率上更有优势也更适合混合部署。但代价是运维复杂度上升对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://主机名:7077Master默认在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对应一次ActionStage对应宽依赖之间的一个计算片段Task数量对应分区数Executor页里的每个进程对应一个JVM实例。这些概念不是死的名词而是UI上真实跳动的数字。以后你写调优参数的时候内心对照的就是这一张张页面。顺便说一句Spark UI默认端口是4040但如果你同时跑多个应用后启动的应用会自动递增端口。在实际的开发过程中我最常做的一件事就是盯着Stage上每个Task的耗时分布如果绝大多数Task几秒跑完、唯独几个Task要花几分钟那么大概率是数据倾斜——这种情况在后面的系列里会专门展开讲。8. 写在最后每个概念背后都是一个“为什么要这样设计”的问题作为系列的第一篇基础概念就到这里。总结起来学习Spark的路线其实是在几组关系之间反复对照第一组是“角色关系”Driver决定算什么Executor负责执行资源管理器决定在哪跑。第二组是“API关系”RDD灵活但费事DataFrame结构化且可优化Dataset在Scala里才是完整形态。第三组是“执行关系”窄依赖在本地完成流水线宽依赖引发ShuffleStage划分的标准就是宽依赖的边界。任何一道Spark性能调优题目最终都能回归到这三组关系上。我个人的体会是初学者最不应该做的事情就是拿着工具书一个API一个API去啃。因为Spark的API数量庞大但绝大多数都是RDD的mapfilterflatMapreduceByKey以及DataFrame的selectfiltergroupByjoin的组合使用。API只是表皮真正决定你能不能看懂Spark UI、能不能定位性能瓶颈、能不能设计出可扩展的数据管道的恰恰是这一篇里讲的那些看似枯燥的基础概念。下一期我计划写Spark SQL的Catalyst优化原理和Spark 3.x新版本的执行优化特性把“DataFrame为什么更快”这件事从头到尾讲透。先把这些底层模型装进脑子再去调优你会发现一切都顺理成章了。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →