资讯动态

Spark缓存机制深度解析:Cache Engine与BlockManager的协同原理

发布时间:2026/9/14 23:57:32 来源:尧图企业网站定制
做大数据开发的兄弟应该都遇到过这种情景同一个DataFrame反复复用每跑一次都要重算好几分钟于是果断cache()一下性能立马上来。但你可能没细想过cache()按下之后数据到底被谁接管了背后真正干活的就是Cache Engine和BlockManager这两大核心组件。Cache Engine负责判断数据块要不要缓存、按什么级别缓存、缓存在哪个粒度BlockManager则在更底层管理数据块的存储位置、序列化格式、磁盘落点以及跨节点复制。两者协同工作构成了Spark等分布式计算框架存储体系的核心骨架。这篇文章我就结合自己的排障和调优经验把这个协同机制彻底拆开讲清楚适合正在用Spark做离线数仓、实时特征或者想深入理解分布式存储底层逻辑的读者。1. 先分开看Cache Engine 与 BlockManager 各自的职责边界1.1 Cache Engine 是缓存调度的大脑很多刚接触Spark源码的朋友会被名字绕晕特别是Spark 3.3之后代码里突然冒出来一个CacheEngine接口原来的CacheManager被改名成LegacyCacheManager。这里先把概念校准一下无论叫CacheManager还是CacheEngine它在整个存储体系里都是“语义层”的角色负责处理RDD和DataFrame缓存请求以及缓存块的新增、读取、失效和清理。Cache Engine的核心逻辑在getOrCompute方法里它可以理解成一个带缓存的迭代器入口。每次RDD的iterator被调用时都会走这个方法根据rdd.id和partition.index拼出RDDBlockId通过BlockManager尝试读取这个块如果命中缓存直接返回缓存数据如果没命中则调用RDD的compute方法重新计算并把计算结果按指定的StorageLevel写入存储体系写入成功后把它登记到一个待失效集合中用于处理依赖RDD缓存失效的场景。这个流程看起来简单但它是整个缓存机制的核心枢纽。它不关心数据到底存在堆内还是堆外也不关心块在网络里怎么复制这些脏活累活全都甩给BlockManager。Cache Engine只关心“这个数据块算出来没有”“要不要存”“从哪个存储级别取”也就是数据块的语义生命周期。1.2 BlockManager 是真正动手存储的底层管家BlockManager的概念比Cache Engine更底层它几乎是整个Spark存储体系的“门面”。每一个Executor启动时都会创建一个BlockManager实例并且向Driver端的BlockManagerMaster注册自己的BlockManagerId告诉Driver“我这里能存数据块来了放我这”。BlockManager管理的最小单位叫Block通过BlockId来唯一标识。BlockId是一个抽象类实际场景里常见的子类包括RDDBlockId(rddId, partitionIndex)缓存RDD分区数据ShuffleBlockId(shuffleId, mapId, reduceId)Shuffle过程中map端产出的中间数据BroadcastBlockId(broadcastId, field)广播变量数据TaskResultBlockId(taskId)任务结果数据块。BlockManager内部还维护着两个关键存储组件MemoryStore负责内存中的数据块存储和淘汰DiskStore负责磁盘上的数据块读写和落盘。块的位置信息会周期性通过心跳上报给BlockManagerMasterMaster端的BlockManagerMasterEndpoint会维护一张块位置表当某个Executor需要读取一个远程块时就是靠这张表找到数据在哪里的。一句话总结BlockManager不关心这个块是缓存块还是Shuffle块它只按BlockId做统一的写入、读取、复制和删除是一个纯正的“存储执行层”。1.3 为什么中间要隔一层解耦让机制可以复用很多人会问既然BlockManager已经这么强了为什么还要单独搞一个Cache Engine直接让BlockManager来管缓存不香吗答案是解耦。BlockManager要服务的对象太多了RDD缓存只是它众多客户之一。Shuffle数据、广播变量、任务结果这些都需要往BlockManager里塞。如果让BlockManager直接对接RDD缓存语义那它就会变成一个既管“存什么”又管“怎么存”的巨无霸每增加一种缓存策略都要改底层存储逻辑非常痛苦。Cache Engine这一层把高层语义和底层执行切开了Cache Engine知道RDD的依赖关系、知道存储级别、知道如何判断一个块是否失效BlockManager只需要实现“给我一个BlockId和一段字节你让我存哪我就存哪”这样的通用能力。这样的分层设计让BlockManager可以复用到Shuffle和Broadcast里而Cache Engine可以针对RDD缓存做更精细的调度互不污染。2. 协同链路拆解一次缓存写入到读取的完整旅程2.1 写入路径计算迭代器如何变成可复用的缓存块我把一次rdd.cache()之后发生的事完整走一遍。当你调用cache()的时候本质上只是设置了一个StorageLevel.MEMORY_ONLY的标记真正干活发生在后续某个Job的Action执行阶段。某个分区的数据需要被计算时会进入Cache Engine的getOrCompute发现块不存在于是执行rdd.compute得到数据迭代器。接下来Cache Engine会把迭代器交给BlockManager的putIterator方法。putIterator收到数据后会根据存储级别决定怎么写入MemoryStore。这里有个很关键的分流如果存储级别不带_SER后缀走putIteratorAsValues数据以反序列化对象的集合形式存进内存。对象数组的内部表示是一个PartiallySerializedBlock或SerializedMemoryEntry之类的结构具体看版本。如果带_SER后缀走putIteratorAsBytes数据会被序列化成字节数组再决定是存堆内还是堆外。如果存储级别里带了_DISK并且内存写入失败或内存放不下BlockManager会把数据溢写到磁盘。磁盘上的块在逻辑上仍然属于同一个BlockId但物理位置已经落到由spark.local.dir指定的一块本地目录里。如果存储级别里指定了副本数大于1比如MEMORY_AND_DISK_2BlockManager还会在写入完成后把块的字节复制到另一个Executor节点上并通过BlockManagerMaster更新块位置信息。这一步是异步的复制失败不会影响主块写入但会在日志里打出replicate失败的警告。2.2 读取路径本地命中、远程拉取与副本选择缓存的目的就是为了读得快。当同一个RDD被再次使用时Cache Engine会直接来BlockManager要数据。BlockManager的get方法首先查本地的MemoryStore或DiskStore命中就直接返回。如果本地没有BlockManager会向Master查询这个BlockId的locations也就是块位置列表。这里有一个细节BlockManager在读取远程块时并不像很多初学者想象的那样把块整个下载回本地再反序列化。它会先通过BlockTransferService发起一个fetchBlockSync请求把数据块以字节流形式拉回当前Executor然后在本地完成反序列化。这就意味着一旦决定走远程读取反序列化工作永远发生在读取方而不是存储方。远程拉取的块可能来自同一个进程的另一个BlockManager实例也可能来自另一台机器上的Executor。Spark在发起远程读时会按块大小做聚合避免大量小请求打爆网络具体由spark.reducer.maxSizeInFlight控制单次拉取的最大字节数。2.3 清理与失效缓存驱逐、Executor退出时的块处理数据能存进去自然也要能清理掉。协同机制里最容易被忽视的就是清理路径。内存紧张的时候MemoryStore会根据自己的淘汰策略驱逐一些块。驱逐是有选择性的优先驱逐那些没有磁盘副本、且当前任务不再使用的块如果块带_DISK驱逐时会先尝试把数据刷到磁盘再释放内存。这里有个很隐蔽的坑如果存储级别是MEMORY_ONLY又采用了DISK_ONLY的混合策略驱逐逻辑会严格遵循存储级别的约束该内存淘汰就内存淘汰不会自作主张落盘。Executor进程异常退出时Driver端的BlockManagerMaster会通过心跳超时检测到这个Executor失联然后清理这个Executor对应的BlockManagerId以及它持有的所有块。此时如果有其他任务傻傻地来读这块数据会直接抛MissingBlockException。这也是为什么我在生产环境一直建议重要RDD的持久化级别至少要带磁盘副本最好还能有副本节点。RDD的unpersist()则是主动失效的入口它会向Master发送RemoveRdd消息要求所有持有该RDD块的Executor把自己本地的对应块删掉。这个动作是全局的调用完再访问这个RDD就会重新计算。3. Shuffle 场景里的另一层协作容易被忽视的 Block 传输3.1 Shuffle 写阶段ShuffleBlockId 的生成与落盘提到BlockManager很多人只想到缓存其实Shuffle才是它服务量最大的场景。Shuffle写阶段每个map任务会把输出结果按分区写到本地磁盘并生成ShuffleBlockId(shuffleId, mapId, reduceId)这样的块标识。这个过程和Cache Engine没有直接关系但一样发生在BlockManager里。map端通过DiskBlockObjectWriter把数据流式写入磁盘文件写完以后调用blockManager的putBlockData方法把文件的逻辑处理交给DiskStore管理。随后块位置信息会被上报给Masterreduce端拿着这个位置信息才能找到数据。Shuffle数据块在内存里也会有一小段缓冲区但和缓存块不一样的是Shuffle块的整个生命周期通常很短拉取完就删。所以BlockManager在Shuffle场景里更多扮演的是“临时中转仓”的角色而不是“长期货架”。3.2 Shuffle 读阶段BlockTransferService 的并发拉取reduce阶段拉取Shuffle数据走的是ShuffleBlockFetcherIterator。它从Master拿到一批ShuffleBlockId的位置后会按Executor分组对每个远程BlockManager发起并发请求。这里也要通过BlockManager的BlockTransferService。这个传输服务是Netty实现的支持流式传输和零拷贝。如果拉取数据量大还会配合spark.shuffle.io.retryWait做重试。最经典的问题就是网络抖动导致的“connection closed”报错这时候去看对应Executor的GC日志和网络监控往往能发现内存频繁Full GC或者网卡打满。3.3 Shuffle 与缓存共享 BlockManager 时的相互影响Shuffle和缓存都共用BlockManager意味着它们会互相挤占资源。Shuffle写阶段的内存占用属于Execution池缓存占用属于Storage池。在统一内存管理器下两者可以互相借用空闲内存但方向有限制。这里要小心一个场景某台Executor上存了大量RDD缓存块同时又在跑一个大Shuffle。当Shuffle需要更多Execution内存时会尝试回收Storage池里被借出去的部分如果缓存块占用了Execution借来的内存缓存就会被强制驱逐。我见过不少案例明明设置了MEMORY_AND_DISK业务方却反馈“缓存了也没变快”十有八九就是缓存块反复被驱逐每次访问都在重新计算或重新反序列化。所以设计缓存策略时一定要把同节点的并发Shuffle负载考虑进去不能只看Storage空闲大小。4. 存储格式是协同机制里的隐藏变量4.1 内存中展开对象还是序列化字节数据块进入MemoryStore后到底以什么格式存在直接影响缓存命中率、GC压力以及反序列化开销。这就是“数据存储格式”在缓存机制里的真实分量。putIteratorAsValues存储的是Java对象集合每个分区一个ArrayBuffer或PartiallyUnrolledIterator对象完整保留在堆内。优点是读取时零反序列化拿到就能直接遍历。缺点是对象引用和头信息极其占内存100万条测试数据用对象数组存比字节数组至少多出0.5到2倍空间。如果在堆内放大量的展开对象GC会变得异常频繁老年代涨得飞快。putIteratorAsBytes则先把数据序列化成字节数组再入堆。Spark默认使用Java序列化也可以换成Kryo。字节数组紧凑、GC友好但每次读取都要反序列化CPU开销上来了。我自己的经验是如果数据量小、读取频率极高、存储空间充足用MEMORY_ONLY如果数据量大、对象庞大优先用MEMORY_ONLY_SER或MEMORY_AND_DISK_SER配合Kryo注册类内存能省一半还多。4.2 磁盘上的落盘格式与压缩选项当数据从内存溢写到磁盘或者直接采用DISK_ONLY级别时存储格式同样关键。磁盘块本质上是序列化后的字节流文件文件由BlockManager的DiskStore统一管理文件名是一段包含BlockId编码的十六进制字符串。磁盘缓存块默认是不压缩的spark.rdd.compress默认是false。如果你存的本来就是文本型数据或稀疏特征序列化后体积很大建议打开压缩。但要注意压缩会增加CPU开销读取时解压也耗时所以压缩只适合“写多读少”“磁盘容量紧张”的场景。Spark SQL里还有一套独立的列式存储格式DataFrame缓存在内存中默认采用InMemoryRelation底层是ColumnarBatch整列压缩存储。你可以通过spark.sql.inMemoryColumnarStorage.compressed控制是否压缩列默认是true对重复值高的列压缩效果非常明显。这是很多人忽略的“数据存储格式”优化点。4.3 参数配置如何影响协同行为存储格式不是自动决定的而是由一系列参数和StorageLevel共同决定的。我列几个直接影响协同行为的参数建议收藏参数默认值影响spark.rdd.compressfalse是否压缩RDD缓存块压缩格式走spark.io.compression.codecspark.io.compression.codeclz4Snappy、Lz4、Zstd影响磁盘和网络传输压缩效率spark.kryo.registrationRequiredfalseKryo是否强制注册类开启后未注册类会直接报错spark.sql.inMemoryColumnarStorage.compressedtrueDataFrame列式缓存是否压缩spark.memory.storageFraction0.5Storage池占统一内存的比例影响缓存最小保障空间spark.memory.fraction0.6Execution和Storage合计占堆内存比例其余留给用户代码和元数据spark.memory.offHeap.enabledfalse是否启用堆外内存存储需配合offHeap.size这里面最容易被踩的是storageFraction和offHeap.enabled。把storageFraction调得太小缓存块很容易被Execution挤掉调得太大Shuffle内存又不够。我一般在交互式分析场景把它设为0.4到0.5在纯ETL离线计算场景调到0.3左右让更多内存给Shuffle和聚合。5. 内存模型与调优给两者分配空间的依据5.1 统一内存管理器中的 Execution 与 Storage 池理解了存储格式之后再看内存分配逻辑就会清晰很多。Spark默认使用UnifiedMemoryManager它把堆内存中的spark.memory.fraction部分划成总预算再在这个预算里按storageFraction切分Storage池和Execution池。描述机制时我尽量简洁Storage池和Execution池不是物理隔离的而是软边界。Execution可以借用Storage空闲的内存Storage在Execution空闲时也能借用Execution的内存。区别在于归还方式Execution借用Storage后如果Storage需要内存Execution在完成任务后会归还Storage借用Execution后如果Execution需要内存缓存块会被强制驱逐来归还。这也解释了为什么缓存块放在内存里并不是“绝对安全”的。你在Storage UI里看到缓存占用150GB实际上里面可能有一部分是从Execution池借来的一旦并发来一个大Shuffle这些借来的缓存会被清出去。5.2 一次缓存任务的内存估算实例给你一个具体的估算例子。假设每个Executor堆内存是10GBspark.memory.fraction0.6统一内存预算就是6GBstorageFraction0.5Storage最小保障就是3GB。你想要缓存一份序列化后大约5GB的中间结果。如果存储级别是MEMORY_ONLY_SER5GB根本放不进3GB的保障空间。但实际运行时只要Execution的3GB内存没有被使用Storage可以借过来所以一开始可能成功存入5GB。但紧接着另一个Stage开始跑ShuffleExecution需要内存就会强制驱逐4GB的缓存来归还。结果就是缓存表面上存在过实际作用约等于零。这也是为什么我不建议用MEMORY_ONLY缓存超大结果。生产里要么把spark.memory.storageFraction调高到0.6要么直接改成MEMORY_AND_DISK_SER让溢出部分自动落盘才不至于缓存了个寂寞。5.3 调优建议与避坑清单基于上面的内存模型我把缓存调优的几条实战经验整理成清单给你只缓存真正复用的数据不是所有中间结果都值得cache临时变量缓存了反而是负担。设置持久化级别时默认优先考虑MEMORY_AND_DISK_SER内存不够自动落盘读取性能依然比重新计算高一两个数量级。优先用Kryo替代Java序列化Kryo序列化后的块体积通常只有Java的一半GC压力也小。记得提前注册类避免运行时反复的全量序列化。Spark SQL场景下优先用DataFrame缓存而不是RDD缓存列式存储的压缩比和Scan效率比行式对象高很多。如果任务读写比高、反复Scan打开spark.sql.inMemoryColumnarStorage.compressed压缩慢一点但省内存效果拔群。不要盲目调大storageFraction要结合Executor并发度和Shuffle量级做权衡。我见过有人调到0.8结果缓存是稳了Shuffle疯狂溢写磁盘整体反而更慢。堆外内存开启前先确认机器物理内存spark.memory.offHeap.size设置不当会直接打爆进程。6. 常见问题排查与实战记录6.1 问题速查表Cache Engine和BlockManager协调过程中我真实遇到过的坑不少整理成一张速查表方便你对症下药现象可能原因排查方向缓存命中率极低任务反复重算Storage内存不足缓存被Executive借用后驱逐看Storage UI的Cache Miss查spark.memory.storageFraction报MissingBlockExceptionExecutor异常退出持有的块全部丢失看Executor日志和心跳记录换带副本的存储级别内存充足但缓存写入失败堆外内存未开启或offHeap.size设置过小检查spark.memory.offHeap.enabled反序列化报ClassNotFoundExceptionKryo未注册类或注册类改动打开spark.kryo.registrationRequired验证磁盘缓存目录写满spark.local.dir空间不足溢写量过大扩容磁盘目录或缩小storageFractionFull GC频繁内存中存储了大量展开对象未序列化换成MEMORY_ONLY_SER配合Kryo远程块拉取超时Netty线程阻塞Executor GC停顿过长查GC日志减小单次拉取块大小sapark.reducer.maxSizeInFlight缓存RDD在UI里显示“Not Serialized”设置了_SER但实际存储用了putIteratorAsValues检查数据源是否已预先序列化确认StorageLevel拼写6.2 两个典型故障复盘第一个是内存驱逐导致的“假缓存”。有个数仓任务每天凌晨固定时间段特别慢日志里一堆EvictionWarningStorage UI显示的缓存占用从早上的几百GB一路掉到几十GB。查到最后是同一时间点上调度系统跑了一个大Shuffle的ETL任务把Storage借出去的内存全挤回来了RDD缓存在几分钟内被驱逐干净。后续任务再访问RDD只能重新计算整体比不缓存还慢。排查后我把这个缓存任务单独放到一个专用队列并把存储级别从MEMORY_ONLY改成MEMORY_AND_DISK_2问题才算根治。第二个是Executor节点异常退出导致的MissingBlockException。有一次某个计算节点因为硬件故障宕机而该节点上恰好存了一批MEMORY_ONLY级别的核心维度数据。当时下游一堆任务在等这批数据结果全部报块缺失。因为MEMORY_ONLY没有磁盘副本宕机就彻底丢了。后来我把这批维度数据换成了MEMORY_AND_DISK_2副本分散到两个节点稳定性好了很多。6.3 排查工具箱与定位思路如果你遇到存储相关的诡异问题我给你一套比较顺手的排查顺序打开Web UI的Storage标签页先看每个RDD的Storage Level、Cached Partitions和Fraction Cached一眼就能看出缓存是否完整。看Executor标签页对比每个Executor的Storage Memory和Disk Used判断是否存在负载不均。在Spark Shell或代码里执行sc.getPersistentRDDs拿到所有持久化RDD的Id和存储级别可以精确判断哪些RDD还没被清理。用BlockManagerMaster的相关API比如sc.env.blockManager.master.getStorageStatus拿到每个BlockManager的块分布统计。看Driver日志里的UpdateBlockInfo和BlockManagerHeartbeat延迟判断上报是否堵塞。这套流程走下来绝大多数缓存和存储问题都能定位到一个比较小的范围。我个人在实际操作中还有一个习惯遇到缓存相关的性能问题先不急着改代码先看一遍Storage UI里缓存块有没有被频繁驱逐。只要驱逐日志消失命中率上去了性能基本就回来了。Cache Engine和BlockManager这套协同机制说复杂也复杂说简单也简单核心就一句话高层调度负责“该不该存”底层存储负责“能不能存”两者配合得当你的数据管道就能安安静静地跑出好性能。最后再分享一个小技巧排查Shuffle相关故障时用blockManagerId配合blockId过滤日志比直接搜异常栈要高效得多这个习惯帮我省了不少时间。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价