
只要你碰过大数据处理大概率绕不开Spark。不管是做数仓ETL、实时日志分析还是跑复杂机器学习模型Spark几乎成了默认的分布式计算框架。很多人一到生产环境就出事换了内存参数、加了Executor数量还是跑不动问题往往不是代码逻辑不对而是没把Spark的内存计算这四个字真正吃透。我从Spark 1.x一直用到3.x踩过的坑比看过的文档多得多今天这篇就把内存计算原理、集群搭建、资源调度、Spark SQL实战和常见故障排查完整串一遍从入门到精通的半条捷径在这里了。这篇文章的适合人群是两类刚接触Spark、想搞清楚它为什么快的新手以及已经跑到生产环境、被OOM和CPU空转搞得头大的工程师。读完你会明白Executor内存怎么分、YARN上为什么Executor总是只分到1个vCore、Spark SQL里的Join为什么这么慢也能拿到一套可以直接复制的排查套路。1. 先搞明白Spark的内存计算到底快在哪1.1 从MapReduce的落盘说起很多人接触Spark之前用的都是MapReduce它慢真不是因为分布式计算这个模型不对而是每一轮Map和Reduce之间都会把中间结果写到磁盘。更难受的是一到迭代算法比如PageRank、Kmeans每迭代一次就落盘一次几十轮迭代下来I/O开销直接吞掉大部分计算时间。Spark的策略完全不同它把中间结果尽量放在内存里让计算在内存中反复流转只有内存不够了才把数据溢写到磁盘。这个转变听起来很朴素但实际效果是量级层面的同一个逻辑MapReduce可能跑40分钟Spark Spark跑三五分钟就结束了。我第一次在生产环境验证这个差距的时候整个团队都惊了后来我们干脆把原来的MR任务全部迁到了Spark。不过这里要纠正一个误区内存计算不是说把所有数据都塞进内存而是尽可能减少不必要的磁盘读写和序列化开销。数据分布在一堆节点上每个节点把自己的那部分数据在内存里搞定跨节点的数据交换能省则省这才是Spark快的内核。1.2 DAG、缓存和血缘一套组合拳Spark把一次作业组织成DAG全称是Directed Acyclic Graph。你用代码写出一堆transformation算子Spark不会傻乎乎一个一个执行而是先构建出一张计算图统一优化之后划分成一个个Stage。Stage里尽量把能合并的算子合在一起把不必要的Shuffle干掉。这种先看图、再动手的做法跟MapReduce那种跑一步算一步完全是两个时代的思维。RDD弹性分布式数据集还有一个杀手级特性可缓存。你可以把中间结果用cache或persist存在内存里后续其他计算链路要用直接取不用从头再算。这一点在做多轮迭代或者同一个数据源被多个报表复用时尤其香。另外Spark的容错也很有讲究它不复制数据副本靠的是血缘也就是每个RDD都记录了自己是从哪里怎么算出来的。某个分区挂了只需要从父RDD重新计算这一个分区而不像MapReduce整个作业重新跑。这套少落盘、少传输、按需重算的组合拳才是内存计算真正的完整形态。2. 内存模型拆解Executor内存到底怎么分配2.1 Executor堆内内存由哪几块组成很多新手调Spark参数只知道调spark.executor.memory调了半天还是OOM就是因为不知道这个内存进去之后还要被切分成几块。Spark把Executor的堆内内存分成四块Reserved Memory、User Memory、Execution Memory、Storage Memory。Reserved Memory是系统保留的默认300MB用来存Spark内部对象你没法动它。剩下的用户内存和计算缓存内存怎么分呢默认规则是全部可分配内存减去Reserved后用户内存占40%统一内存池占60%。统一内存池里Execution和Storage默认五五开各占一半。我举个具体例子假设spark.executor.memory4g也就是4096MB。扣除300MB保留内存剩3796MB统一内存池占60%也就是2277MB左右其中Execution和Storage各占1138MB左右剩下来的1518MB归User Memory。很多人把spark.memory.fraction调到0.9觉得给引擎多一点内存总是好的结果用户对象没地方放GC疯狂回收反而不如默认值。2.2 ExecutionMemory和StorageMemory的动态博弈这俩内存不是完全隔离的Spark用Unified Memory Manager统一管理它们之间可以互相借用。Execution内存是给Shuffle、Join、Aggregation这类计算用的Storage内存是给缓存RDD/DataFrame用的。比如一个任务正在做大规模Shuffle计算内存不够了它就可以去抢Storage那部分空闲内存反过来如果之后你要缓存的数据特别大也可以抢回Execution已经释放的内存。这个借用机制让内存利用率大幅提升。但有个关键的例外Storage被借用之后如果缓存数据要写回来是需要等Execution主动释放的。极端情况下你缓存了大量数据同时又在跑重型Shuffle缓存可能被一点点挤压掉甚至发生强制落盘。这也是为什么保留一个spark.memory.storageFraction很重要它像一道防线保证Storage至少有一半内存不会被抢走。生产环境里遇到缓存频繁丢、任务反复重算的先检查storageFraction和Execution内存的抢占情况。2.3 堆外内存什么时候该碰除了堆内内存Spark还有一个Off-Heap堆外内存。默认是关闭的需要同时打开spark.memory.offHeap.enabled并把spark.memory.offHeap.size设成具体大小比如1g。堆外内存的好处是绕过JVM不走GC稳定性和吞吐量在某些场景会好一些。代价是你要自己管理序列化和生命周期而且配置不当很容易出现堆外OOM报错信息还特别抽象排查起来很痛苦。我的建议是默认场景先别用堆外内存。如果你有20个以上的ExecutorGC时间超过总运行时间的15%数据序列化之后存储再考虑开。但别忘了还有一层系统层面的隐式堆外就是spark.executor.memoryOverhead它默认是executor内存的10%至少也要384MB主要给JVM的线程栈、Direct Buffer、NIO这些用。YARN模式下物理内存超了Container被NodeManager杀掉先看看是不是memoryOverhead太小了。3. 从零部署Spark的安装与集群搭建要点3.1 用Standalone模式20分钟跑通第一个任务先不说生产怎么搭新手入门最快的就是Standalone模式。去官网下载对应Hadoop版本的Spark包比如spark-3.5.1-bin-hadoop3解压、配环境变量基本就完成了90%。wget https://archive.apache.org/dist/spark/spark-3.5.1/spark-3.5.1-bin-hadoop3.tgz tar -zxvf spark-3.5.1-bin-hadoop3.tgz mv spark-3.5.1-bin-hadoop3 /opt/spark export SPARK_HOME/opt/spark export PATH$SPARK_HOME/bin:$SPARK_HOME/sbin:$PATH然后复制模板配置主节点和Worker资源cp conf/spark-env.sh.template conf/spark-env.sh在spark-env.sh里写入SPARK_MASTER_HOSTnode01 SPARK_WORKER_CORES8 SPARK_WORKER_MEMORY16g启动主节点和工作节点跑一个经典测试$SPARK_HOME/sbin/start-master.sh $SPARK_HOME/sbin/start-worker.sh spark://node01:7077 spark-submit --master spark://node01:7077 \ --class org.apache.spark.examples.SparkPi \ $SPARK_HOME/examples/jars/spark-examples_2.12-3.5.1.jar 100看到输出一个接近3.14的数整个环境就算通了。第二步才是加节点把其他机器也装一遍Spark改一下SPARK_MASTER_HOST指向主节点启动worker注册进来集群就扩起来了。3.2 Standalone、YARN、Kubernetes怎么选很多团队在部署方式上纠结我直接给结论如果你公司的集群已经跑着HDFS和Hive那YARN模式是首选因为YARN能统一调度内存和CPU跟Hive共用资源池还自带队列隔离如果你们是纯Spark场景不想ZooKeeper、HDFS这些组件Standalone足够简单如果已经在容器化平台上面Kubernetes模式更自然Spark的Pod可以自动伸缩。部署方式不能只看好不好装还要看团队运维能力和周边生态。部署方式优点缺点适合场景Standalone部署最简单无额外依赖没有多租户资源隔离测试环境、纯Spark集群YARN与Hadoop生态统一资源管控成熟需要维护YARN集群已有HDFS/Hive的生产环境Kubernetes弹性好容器化标准运维门槛高网络配置复杂云原生基础设施成熟的公司4. 生产环境必修YARN上Executor总是只分到1个vCore怎么办4.1 问题现象与根因分析很多人在集群上用YARN模式提交Spark作业打开YARN的ResourceManager界面一看发现每个Executor的Vcores只有1明明节点有32核资源完全没利用起来。这个问题特别典型根因也很简单你在spark-submit里只设了--executor-memory没设--executor-coresSpark在YARN模式下默认executor.cores1。也就是说一个Executor一个任务并行度直接拉不起来。如果业务SQL里分区数量大、Task数量多Task会在一个核上排队整个job看起来像是卡住了日志里全是Waiting for resources。另一个隐藏因素是YARN的调度上限。yarn.scheduler.maximum-allocation-vcores和yarn.nodemanager.resource.cpu-vcores如果配得偏小也会限制你申请到的核数。我见过一个集群NodeManager明明有64核但yarn-site.xml里只配了16后台上千个任务挤在16个vCore里跑怎么调Spark参数都没用。4.2 正确配置Executor核数与内存的方法这里我直接给一套通用配置公式。假设一台Worker节点是32 vCore、128GB内存生产上建议留出20%左右给系统、HDFS、YARN自身实际可用大概25个vCore和100GB内存。那每个Executor给8个vCore、32GB内存一个节点放3个Executor最合适。提交命令长这样spark-submit \ --master yarn \ --deploy-mode cluster \ --num-executors 6 \ --executor-cores 8 \ --executor-memory 30g \ --driver-memory 4g \ --conf spark.yarn.executor.memoryOverhead4g \ --class com.example.DataProcessor \ app.jar注意executor-memory我写的是30g而不是32g因为堆外还需要memoryOverhead比如4g加起来才是实际占用的物理内存。如果申请32g堆内加4g堆外Container总内存36g已经超过节点可用的约33g了任务会一直申请不到资源。核数和内存的合理搭配是并行度和GC开销之间的平衡。Executor核数太多比如16核并发Task多但内存Shared得厉害GC压力很大核数太少并行度上不去资源浪费。我踩坑踩出来的经验是单Executor控制在4到8核单Executor总内存控制在30到60G一段任务跑起来之后看GC时间低于10%算是健康状态。4.3 关于CPU分配最常见的三个误解第一个误解vCore等于物理线程数。实际上YARN里的vCore只是一个逻辑概念默认情况下一个物理核可以映射为多个vCore具体看yarn.nodemanager.resource.cpu-vcores怎么配的。如果集群是开启超线程的机器还要考虑是不是该按物理核的一半来配。第二个误解spark.task.cpus没什么用。它默认是1表示每个Task占用的核数。如果Task内部又调用了多线程库比如某些图像处理、模型推理需要把这个参数调大否则多个Task并行时CPU会互相打架。第三个误解开了动态资源分配就不用管核数。SPark的动态分配会按负载伸缩Executor数量但最大上限由spark.dynamicAllocation.maxExecutors决定通常默认值可能不是你要的。我建议配置好之后再观察几个任务确认Executor数量的伸缩范围合理而不是全交给默认值。5. Spark SQL与数据分析实战内存计算在业务里长什么样5.1 一段能直接用的Spark SQL示例Spark SQL是生产里用得最多的模块没有之一。它最大的价值是可以像写传统SQL一样处理分布式数据引擎底层自动把SQL翻译成执行计划再跑在内存计算框架上。我给你一个最常见的场景读入一份用户行为日志按天统计Top用户。from pyspark.sql import SparkSession spark SparkSession.builder \ .appName(spark-sql-demo) \ .enableHiveSupport() \ .getOrCreate() df spark.read.option(header, True).csv(/data/user_log.csv) df.createOrReplaceTempView(logs) result spark.sql( SELECT user_id, count(*) AS cnt FROM logs WHERE dt 2024-06-01 GROUP BY user_id ORDER BY cnt DESC LIMIT 10 ) result.show() result.write.mode(overwrite).parquet(/data/top_users)这段代码背后发生的事你看到的是SQL引擎里其实经历了一遍Catalyst优化器把过滤条件下推到数据源能少读的数据就少读再做聚合把中间结果尽量留在内存里。这也是Spark SQL比直接写RDD代码跑起来通常更快的原因之一。5.2 Join在内存里是怎么算的为什么说数据倾斜最要命Spark SQL里最影响性能的操作就是Join。小表与大表关联时可以用Broadcast Hash Join把小表广播到所有Executor完全不做Shuffle但如果你没把spark.sql.autoBroadcastJoinThreshold调好或者表实际大小超过了默认10MB阈值引擎就会老老实实走Sort Merge Join过程中会跑一次全量Shuffle。Shuffle一多内存压力就上来了。我调过最多的生产事故就是数据倾斜某个user_id占了99%的数据所有记录都被Hash到同一个Executor那个Executor内存直接被打爆其他Executor闲着没事做整个任务一直报OOM。遇到这种情况单靠加内存没用应该先做数据探查确认倾斜key的情况再用加随机前缀、两阶段聚合这样的手段打散数据。5.3 Catalyst优化器如何帮你省内存Catalyst是Spark SQL的优化引擎它的工作方式很像传统数据库的查询优化器。我挑三个对我们最实用的优化点谓词下推、列裁剪、动态分区裁剪。谓词下推是指SQL里的where条件尽量在读取数据时就过滤掉不把无用数据读进内存列裁剪是指只读SQL用到的列比如你有100个字段SQL里只用3个Parquet这种列式存储可以只扫3列的数据动态分区裁剪更进一步在运行时根据前面的过滤结果自动裁剪后面要扫描的分区。这三个优化组合起来实际读入内存的数据量能减少一个量级。这也提醒我们一件事写Spark SQL不要无条件select *不要觉得数据反正放内存无所谓。内存计算的前提是能被优化器裁剪到足够小数据规模控制住了执行效率自然高。6. 高频问题排查实录日志提示、OOM与Spill处理6.1 启动时的log4j提示并不是报错很多人在启动Spark任务时看到一行提示Using Sparks default log4j profile: org/apache/spark/log4j-defaults.properties以为出了什么大问题。别慌这只是一个提示意思是Spark没找到用户自定义的日志配置所以用了内置的默认配置。如果你想去掉或自定义它很简单在Spark的conf目录下放一份日志配置文件2.x版本用log4j.properties3.x用log4j2.properties。复制模板改一改就行cp $SPARK_HOME/conf/log4j2.properties.template $SPARK_HOME/conf/log4j2.properties然后把日志级别从INFO改成WARN或者按需调整重启Spark作业即可。这行提示就不会再以默认路径出现了。如果日志文件里还有大量INFO信息刷屏也建议先把级别调高再排查问题不然关键报错会被淹没。6.2 Executor失联、OOM和Spill过高怎么排查生产环境最大的三类问题Container被Kill、ExecutorOOM、Shuffle溢出。排查手法其实高度统一先打开Spark UI看Executors页签重点看三块GC时间、Shuffle Read/Write大小、以及Storage里的缓存占用。如果看到GC时间占比很高说明堆内内存不足调大executor-memory或者减少单Executor核数都有帮助如果看到Shuffle Write很大但Read很小重点看是不是Map端溢写太多可能是spark.shuffle.sort.bypassMergeThreshold之类的参数需要调如果是数据倾斜导致的OOM先去Spark UI的SQL页签看哪个Stage读入的记录数严重不均衡再针对倾斜键做处理。我遇到过一个诡异案例任务跑到一半Executor一个个被NodeManager杀掉日志里没有任何Java异常。后来去排查系统dmesg才发现是物理内存超了Container直接被杀。这就是spark.executor.memoryOverhead给太小了申请Executor内存时预留不够把物理内存算满导致的。6.3 常见问题速查表现象可能原因解决办法YARN上Executor每个只有1个vCore未设置executor-cores默认值为1spark-submit加--executor-cores任务一直卡在WaitingYARN队列资源不足或cores配置过大清理队列、降低executor规格Executor被NodeManager杀掉堆外内存不足或memoryOverhead过小调大memoryOverhead并保留系统余量SQL任务执行慢且GC高Execution内存不足调大executor-memory或减少单Executor并行任务数缓存数据频繁丢失Storage内存被Execution抢占调大storageFraction或减少Shuffle并发某个Executor OOM而其他正常数据倾斜定位倾斜Key加随机前缀或两阶段聚合每次处理完问题我都会把现象、配置、解决方式记下来这比自己翻书有效得多。有时候排查几天的问题最后原因就是一行参数写错了速查表能帮你少走很多弯路。7. 进阶方向GPU集群上的Spark部署与国产数据库适配7.1 DGX Spark这类GPU一体机上的部署心得最近NVIDIA的DGX Spark这类桌面级AI超算很火有些人直接在它上头部署Spark做数据预处理。我在这类GPU一体机上的实测经验是Spark的安装和传统x86环境差异并不大照样是下载二进制包、配Spark环境、启动Worker。但有一个点特别关键就是CPU资源的分配。GPU一体机往往同时跑着模型训练任务CPU核数如果被Spark占满训练这边的数据加载和预处理会非常慢。我建议给Spark只分配物理核的50%左右并设置好spark.task.cpus防止Spark和GPU任务互相抢资源。Spark作为一个数据处理引擎在GPU环境里更多承担的是干净高效的ETL工作把数据清洗好、组织好喂给后端的深度学习框架。如果你想更进一步可以关注RAPIDS加速器这样的插件它能在Spark里调用GPU做算子加速但不要一上来就全量开启先挑几个重度SQL跑benchmark确认收益之后再铺开。7.2 达梦数据库与Spark的JDBC适配集成国内很多政企项目用达梦数据库遇到Spark要和达梦做数据集成是常有的事。达梦兼容Oracle部分语法但驱动、URL和方言跟MySQL、Oracle都有区别。达梦JDBC驱动是dm.jdbc.driver.DmDriverURL格式是jdbc:dm://host:5236这一点在第一次配置时就很容易踩坑。使用时用Spark的标准JDBC方式接入即可df spark.read \ .format(jdbc) \ .option(url, jdbc:dm://10.0.0.5:5236) \ .option(dbtable, T_USER) \ .option(user, SYSDBA) \ .option(password, your_password) \ .option(driver, dm.jdbc.driver.DmDriver) \ .load()要注意三个实操点第一达梦表名列名会默认转成大写SQL和配置文件里要注意大小写匹配否则会报找不到表第二如果要做并行读取必须提供一个均匀分布的分区字段否则Spark默认单分区读一张大表速度慢到让人绝望第三写入时达梦对批量提交和事务的控制跟MySQL不同建议先小批量测试再切到大表避免中途事务冲突导致作业失败。这类适配集成没有多难但特别考验细心程度因为错误信息往往不会直接告诉你是驱动或方言的问题而要一层层往下查。8. 面试题方向与经验沉淀8.1 几个高频Spark面试题怎么答如果在准备面试下面这几个问题出现频率很高而且都能和内存计算扯上关系。第一RDD的弹性体现在哪里至少从三方面说血缘容错、分区可重算、资源自动伸缩。第二Shuffle为什么慢因为要落盘、要序列化、还有网络传输这三条每一条都跟内存和I/O相关。第三Spark和MapReduce的本质区别是什么别只说一句基于内存要往DAG调度、缓存机制、统一的执行引擎方向答。第四Executor内存溢出怎么排查从Spark UI看GC、Shuffle、Storage再结合代码和数据分布定位。更进阶一点面试官可能问你spark.memory.fraction和spark.memory.storageFraction的区别是什么前者决定了统一内存池在总堆内中的占比后者决定了统一内存池内Storage和Execution的初始比例。答出初始比例但可以互相借用才算到位。8.2 给新人的三个建议最后分享几条我在实际项目里用血泪换来的经验。第一遇到问题先看Spark UI和日志不要一上来就加内存很多性能瓶颈根本不是内存不够而是参数配错或数据倾斜。第二不要照抄网上的性能调优参数每个集群的CPU、内存、数据量都不一样拿一套配置打天下最终只会害了自己。第三大缓存要谨慎缓存虽然快但缓存多了会侵占Execution内存一旦跑Shuffle就会出现Spill和GC所以用完之后要及时unpersist。Spark的内功就是把内存分配、任务调度、资源申请这一套底层逻辑摸透。你把这些搞明白了看大多数报错都不会慌因为你知道它在哪个环节出了问题也知道该往哪个方向使劲。