
你有没有遇到过这种情况同一个Spark任务在测试环境跑得飞快一上生产就慢到让人怀疑人生。排查了半天发现某个stage的shuffle read反复出现同一个RDD被从头算了一遍又一遍。问题大概率不在代码逻辑而在持久化机制没用好。Spark的persist()和存储级别StorageLevel是性能调优里最常见的切入点但也是被误解最多的。很多人拿到RDD就顺手cache()一下然后该慢还是慢甚至更慢有些人分不清MEMORY_ONLY和MEMORY_AND_DISK到底差在哪还有些人在循环里反复持久化同一个对象把内存撑爆。这篇文章我打算从原理到实战把Spark持久化机制完整拆一遍为什么要持久化、persist()和cache()的真实关系、每种存储级别的取舍逻辑、以及我踩过的坑。适合刚接触Spark的初学者也适合已经写了不少Spark任务但对性能调优还停留在缓存一下就好阶段的同学。1. 先从根上说Spark为什么默认偷懒不缓存数据很多初学者第一次接触Spark时会觉得RDD就是一个装了数据的分布式集合。但实际不是这样RDD内部并没有真正保存那份数据它保存的是一个菜谱——这个数据集怎么从源头一步步算出来的。这个设计是理解持久化的前提。1.1 血统Lineage机制计算的后悔药RDD会记录自己从父RDD经过哪些转换算子map、filter、join等得到这一串依赖关系叫血统Lineage。Spark官方设计逻辑是既然我记录了完整的血缘关系那么在任意一个分区数据丢失时我只需要顺着血缘重新计算那一个分区就行不需要像Hadoop那样维护复杂的副本机制。这个设计非常巧妙但也埋了一个隐患RDD本身并不持有数据数据是按需计算出来的。如果你反复使用同一个RDD它不会自动把结果保存下来而是每次都从头顺着血缘把计算再执行一遍。你可以想象成每次要点外卖不是从冰箱拿现成菜出来热一下而是从买菜、洗菜、切菜开始全部重做一遍。1.2 延迟计算Action之前一切皆虚Spark的计算分成两类转换Transformation和行动Action。map、filter、flatMap、reduceByKey、join这些都是转换它们只是记录操作不会真正执行只有count、collect、saveAsTextFile这些行动才会触发真正的计算。这个机制叫延迟计算Lazy Evaluation。延迟计算带来的直接后果是同一个RDD可以被多个Action重复消费而每次Action都会重新触发一次从头到尾的计算。比如代码里写了cleanData.count()统计总数后面又写了cleanData.collect()取样查看这两次行动如果中间没有任何缓存cleanData的血缘链就会完整执行两遍。数据规模大的时候两遍意味着两倍的磁盘读取、两倍的网络传输、两倍的shuffle。1.3 重复计算的代价到底有多大要量化这个代价最简单是看一个宽依赖的例子。假设你有两个大表A和B做了一次join得到RDD J然后后续两个Action都用到了J。如果J没有持久化第一次Action执行时A和B需要从HDFS读出来、shuffle到对应分区、完成join第二次Action执行时这些过程全部重来一遍。其中shuffle是最大的痛点。shuffle意味着数据要跨网络传输、落磁盘、再重新拉取这部分耗时往往占整个stage的大头。重算一个带shuffle的RDD代价绝不是简单的两倍CPU而是两倍的网络IO加上两倍的磁盘IO。实际生产里一个TB级数据的join结果被重复计算两次任务时间可能从20分钟直接飙升到40分钟以上。这也是为什么持久化是Spark性能调优里最优先考虑的手段之一。2. persist()与cache()最容易混淆的一组API在代码层面persist()和cache()是Spark使用者最常碰到的两个持久化入口。很多人以为它们是两个不同的功能实际上它们的关系简单得惊人。2.1 cache()就是persist()的默认参数版如果要我用一句话说明cache()等价于persist(StorageLevel.MEMORY_ONLY)。在Spark源码里RDD.cache()方法体内部就是一行persist(StorageLevel.MEMORY_ONLY)没有做任何额外的事。所以本质上它们调用的是同一套持久化机制区别只是persist()可以通过参数指定存储级别cache()则固定使用默认级别。这里有一个特别容易踩的坑针对RDD和DataFrame/Datasetcache()的默认行为并不一样。RDD.cache()默认是MEMORY_ONLY也就是只放内存而Spark 2.x之后的DataFrame/Dataset因为底层是列式存储cache()的默认行为通常是内存加磁盘MEMORY_AND_DISK。如果你以前写的是RDD代码后来改用DataFrame API还默认cache就是只放内存那就会对资源消耗判断失误。2.2 persist()的延迟生效机制这里有个高频误解以为调用persist()的那一行代码执行后数据就已经被放进内存了。不是的。persist()只是给这个RDD打上一个需要被缓存的标记真正把数据写入存储介质是等到第一个Action执行并且计算完成之后才发生的。举个例子val rdd sc.textFile(hdfs://data/logs).map(parseLog) rdd.persist(StorageLevel.MEMORY_AND_DISK) // 只是打标记此刻内存里没有数据 rdd.count() // 第一次触发计算计算完成后数据才真正写入缓存 rdd.collect() // 第二次Action直接读缓存不再重算明白这一点很重要。如果你在persist()之后立刻想通过读取缓存来验证效果是看不到任何东西的另外如果第一个Action因为某种原因失败重试缓存写入也会跟着重来。2.3 为什么Persist了还是慢缓存失效的几种可能我见过不少同学在代码里加了cache()但任务依然慢于是得出结论Spark缓存没用。其实缓存没生效的原因通常出在下面几个地方内存不足触发LRU淘汰。MEMORY_ONLY级别下数据放不进内存的部分会被直接丢弃后续Action访问到丢失的块时会重新计算。Executor重启或动态分配导致缓存丢失。持久化数据存在Executor内存里一旦Executor被杀掉缓存跟着消失TaskScheduler会重新调度计算。缓存的位置在血缘链的末端但血缘链中间那段昂贵计算没有被缓存。比如你在join后的结果上做了persist但每次重算还是要从源头读两个大表开销仍然很大。存储级别选择不当。用了MEMORY_ONLY_SER反序列化的CPU开销反而比重算还高这种情况在数据量小但每个对象很大的场景里特别明显。所以加缓存之前先想清楚你缓存的这个节点是否真的覆盖了最昂贵的计算路径。3. 存储级别全景拆解选错级别等于白缓存StorageLevel是Spark持久化机制的核心参数它决定了缓存数据放在哪里、以什么形式放、存几份。选错级别轻则内存浪费重则任务直接OOM。这一节我把所有内置存储级别拆开讲。3.1 一张表看完全部StorageLevelSpark内置了12个标准存储级别我先整体列出来后面逐个详解存储级别使用磁盘使用内存使用堆外内存序列化副本数NONE否否否否1DISK_ONLY是否否是1DISK_ONLY_2是否否是2MEMORY_ONLY否是否否1MEMORY_ONLY_2否是否否2MEMORY_ONLY_SER否是否是1MEMORY_ONLY_SER_2否是否是2MEMORY_AND_DISK是是否否1MEMORY_AND_DISK_2是是否否2MEMORY_AND_DISK_SER是是否是1MEMORY_AND_DISK_SER_2是是否是2OFF_HEAP是否是是1表格里使用磁盘和使用内存同时为是时表示内存优先、放不下的溢写到磁盘而不是两份都存。3.2 MEMORY_ONLY与MEMORY_AND_DISK内存不够时的两种命运MEMORY_ONLY是RDDcache()的默认级别数据只放内存。它的优点是省去了序列化和反序列化的CPU开销数据以Java对象形式直接存在堆内访问速度最快。缺点是只要内存放不下多余的分区直接丢弃不写磁盘后续一旦访问到丢失的分区就得从头重算。MEMORY_AND_DISK则提供了一层兜底内存放不下的分区会溢写到磁盘后续访问时从磁盘读回来。它不会出现某个分区突然消失的情况但代价是要么占内存要么占磁盘IO。这两者的选择逻辑很直接如果数据量远小于Executor可用内存并且重算成本高选MEMORY_ONLY如果数据量接近或超过内存上限宁愿多花一点磁盘IO也坚决不要触发重算那就选MEMORY_AND_DISK。我的经验是生产环境里数据量估算往往不准MEMORY_AND_DISK的兜底能力比那点磁盘IO开销更值钱。3.3 序列化SER不是洪水猛兽MEMORY_ONLY_SER与MEMORY_AND_DISK_SER这两个级别的关键差异在于是否对缓存对象做序列化。MEMORY_ONLY_SER和MEMORY_AND_DISK_SER会把数据先序列化成字节数组再存储这样内存占用大幅下降但每次读取缓存时都需要反序列化多了一道CPU开销。举一个直观的例子一个包含20个字段的日志对象在堆内存里可能占几百字节对象头、指针、padding都在消耗空间序列化后可能只有不到一半的大小。数据量一大这个差距可能直接决定任务能不能跑完。实际生产中MEMORY_AND_DISK_SER是我最常用的级别。因为它兼顾了内存占用可控和永远不重算这两个优点。前提是配合Kryo序列化框架使用否则默认的Java序列化性能会让你怀疑人生。开启Kryo的方式后面专门讲。3.4 带副本的级别_2后缀背后的容错账所有标准级别都有一个后缀带_2的版本例如MEMORY_ONLY_2、MEMORY_AND_DISK_SER_2。这个_2表示每个分区缓存两份副本分散存储在不同节点上。多副本的价值在于当某个Executor宕机或者缓存块丢失时Spark可以直接从另一份副本读取不需要触发重新计算。这在节点不稳定、任务运行时间很长的场景里非常有用。缺点也很明显存储开销翻倍写入缓存的时间也变长。我的建议是除非你的集群节点频繁故障、并且RDD重算代价极高否则不要轻易用_2级别。大部分场景下顺着血统重算一个分区并不是灾难多花的那点时间比双倍存储成本划算得多。3.5 OFF_HEAP堆外存储与Alluxio的适用场景OFF_HEAP是内置级别里的另类它把数据放到JVM堆之外需要配合Tachyon/Alluxio这样的外部系统使用。堆外内存的好处是减少JVM GC压力因为大对象不会被Full GC反复扫描坏处是需要额外部署依赖而且访问路径更长。在实际项目里OFF_HEAP的使用率远低于前面几个级别。除非你的Spark任务已经明确遇到GC瓶颈、并且有专门的Alluxio集群否则我不建议一上来就折腾它。先熟悉MEMORY_ONLY_SER和MEMORY_AND_DISK_SER性价比高得多。4. 看一眼源码StorageLevel到底存了什么前面都是从使用角度讲级别要真正理解存储级别的设计逻辑还是得看一眼源码。不用深挖把核心字段看懂就足够帮你做决策。4.1 StorageLevel的四个布尔开关加一个副本数StorageLevel在源码里就是一个普普通通的case class核心参数只有五个useDisk、useMemory、useOffHeap、deserialized、replication。四个布尔开关加一个副本数所有存储级别都是这五个参数的组合。比如DISK_ONLY就是useDisktrue、useMemoryfalse、useOffHeapfalse、deserializedtrue、replication1。由于deserializedtrue存在磁盘上的数据不需要反序列化直接以对象字节流存储读取时靠内部机制自动处理。MEMORY_AND_DISK_SER则是useDisktrue、useMemorytrue、deserializedfalse多了一个序列化存储的语义。4.2 deserialized参数是怎么伪装成序列化选项的既然参数叫deserialized为什么我们平时说的是是否序列化这里有个容易绕晕的点deserialized的含义是是否以反序列化后的Java对象形式存储。当它为true时意味着不序列化直接存对象为false时意味着要序列化成字节再存。所以MEMORY_ONLY的deserializedtrue而MEMORY_ONLY_SER的deserializedfalse。理解了这层关系你再看源码里那些标准级别的定义就不会觉得它们长得像玄学了。也可以基于这五个参数自定义存储级别比如想要三副本、内存加磁盘、序列化写一个new StorageLevel(true, true, false, false, 3)就行。4.3 标准级别为什么是单例序列化传输与相等比较你可能注意到我们使用时都写StorageLevel.MEMORY_ONLY_SER这种大写常量而不是new StorageLevel。原因是标准级别在Spark内部以单例形式定义这样在Driver和Executor之间传输时它们能保持同一引用方便做相等判断。源码里StorageLevel实现了Externalizable接口来定制序列化行为并且在反序列化时通过readResolve()方法把对象还原成对应的标准单例。这个设计让storageLevel StorageLevel.MEMORY_ONLY这样的比较在分布式环境下依然成立。对我们使用者的启发是能直接用标准级别就用标准级别除非有非常特殊的场景否则不要自己new存储级别否则可能踩到引用相等和序列化匹配的坑。5. 实战按场景决定该不该持久化、用哪个级别理论讲完回到最实际的问题代码里到底该怎么写这一节我按三类高频场景给出具体方案。5.1 迭代计算KMeans这类反复登场的训练数据机器学习的迭代算法是持久化的经典场景。拿KMeans举例训练样本的RDD在整个迭代过程中要被反复扫描每一轮都要重新计算每个点到簇中心的距离。如果不持久化每一轮迭代都会从HDFS重新读一遍全部训练数据代价是灾难性的。正确写法是先对训练数据做一次持久化再进入循环val data sc.textFile(hdfs://data/samples) .map(parseSample) .persist(StorageLevel.MEMORY_AND_DISK) var centroids initialCenters for (i - 0 until maxIterations) { val newCentroids data .map(p (nearestCentroid(p, centroids), (p, 1))) .reduceByKey { case ((sumP, count), (p, one)) (sumP p, count one) } centroids updateCentroids(newCentroids) } data.unpersist()这里我选MEMORY_AND_DISK而不是MEMORY_ONLY是因为训练数据量大我不敢赌它一定能全部塞进内存。迭代场景一旦中途某块缓存被淘汰后续每一轮都会反复重算这块任务基本就废了所以宁可让溢写磁盘兜底。5.2 一个数据多个Action复用清洗后结果的正确缓存姿势另一类高频场景是一份数据经过清洗后既要统计又要抽样还要写入外部存储。下面这段代码里cleanDF被三个Action使用如果不持久化清洗逻辑会执行三遍。val raw spark.read.parquet(hdfs://data/raw) val cleanDF raw.filter(...).withColumn(...) cleanDF.persist(StorageLevel.MEMORY_AND_DISK_SER) cleanDF.count() // 触发计算并写缓存 val sample cleanDF.limit(100).collect() // 读缓存 cleanDF.write.mode(overwrite).saveAsTable(ods.clean_data) // 读缓存 cleanDF.unpersist()这个场景我推荐MEMORY_AND_DISK_SER。原因很现实count()的shuffle结果本来就大后面还要collect和save每一步都在消耗内存如果数据不序列化三份引用同时活跃在Executor里很容易把内存挤爆。用SER版本内存压力小一个量级代价只是多出的反序列化CPU相比任务稳定性来说非常值得。5.3 什么样的数据不值得持久化不是所有RDD都值得持久化。下面这几种情况我建议连cache()都不要写只被一个Action使用的RDD比如只调用一次saveAsTextFile持久化纯粹是额外开销。数据量极小、重算只耗时几十毫秒的RDD持久化节省的时间微乎其微。血缘链很短的RDD比如从一个已经常驻内存的小集合parallelize()出来的重算非常廉价。每一步宽依赖被shuffle落盘过、且后续没有复用需求的结果shuffle本身已经写了一遍磁盘再持久化是重复劳动。尤其注意最后一点。很多新手以为凡是有reduceByKey就一定要缓存结果其实shuffle过程中数据已经在磁盘和内存之间流转过一轮如果后面只有一次Action缓存的意义并不大。判断标准就一条这个RDD会被多个Action消费吗会再考虑持久化不会别动。6. 我在生产环境踩过的坑与总结的检查清单持久化机制看着简单用起来全是细节。这一节我把这些年踩过的坑集中说一遍每一条都是真金白银换来的教训。6.1 cache()之后一定要unpersist()持久化的数据不会随着Action结束自动释放它会一直占着Executor内存直到SparkContext停止或者你显式调用unpersist()。在一个长任务里如果先后缓存了十几个RDD都不清理后面缓存新数据时就可能把前面的挤出内存形成缓存打架。我的习惯是每个persist()都配一个unpersist()放在数据不再被使用的位置如果害怕中间有异常导致跳过就用try/finally包一层。用SQL做缓存时也要记得执行UNCACHE TABLE否则查完就忘内存迟早被撑爆。6.2 Kryo序列化与类注册使用任何带SER的存储级别时如果你不配置KryoSpark默认用Java序列化性能和占用空间都很糟糕。正确做法是在提交任务时加上spark.serializer org.apache.spark.serializer.KryoSerializer spark.kryo.registrator com.example.MyRegistrator spark.kryoserializer.buffer.max 128mKryo性能好但有些自定义类需要注册不注册时只能用全类名反射序列化效率大打折扣甚至可能出现无法序列化的报错。如果对注册类的列表管理比较头疼也可以先注册最核心的领域对象其余让Kryo自动处理。总之用了_SER级别却不配Kryo等于白用。6.3 缓存块丢失时任务不会失败只会悄悄变慢这是最阴间的坑。当持久化的数据因为Executor退出或内存淘汰而丢失时Spark不会报错而是安静地顺着血统重算丢失的分区。表现在Spark UI上就是Storage标签页里出现Lost blocks同时任务时间莫名变长。所以当你发现一个任务比平时慢很多、且代码没有任何改动时第一件事是去Storage页看缓存状态确认有没有丢块。如果有通常意味着Executor不够稳定或者内存分配不足需要考虑调大spark.executor.memory、去掉动态分配或者换个存储级别。6.4 持久化对象不是数据快照还有一类误解以为持久化之后数据就固定不变了。实际上持久化保存的是第一次Action计算完成时的数据状态。如果这个RDD的源头是外部数据源比如每次读取都会变化的最新日志那么缓存里保存的永远是第一次读到的那份不会自动更新。想保证每次读到最新数据不要缓存或者要主动unpersist()后重新持久化。这也是checkpoint()和persist()的核心区别之一。checkpoint会把数据连同血缘写到可靠存储中切断过长的血缘链而persist只是缓存当前计算结果血缘链依然保留。需要长期保存中间结果、防止血缘链过长导致恢复代价大时用checkpoint而不是persist。6.5 上线前的持久化检查清单我在处理线上任务时习惯在提交前过一遍下面这些检查项你也可以直接拿去用梳理代码里所有被多个Action复用的RDD确认每条都用persist()覆盖了。根据数据量和Executor内存确定每个缓存的存储级别不确定时默认MEMORY_AND_DISK_SER。确认Kryo已配置并注册了必要的自定义类。检查每个persist()都有对应的unpersist()或者至少写在任务结束前。在Spark UI的任务运行中期和结束时分别看一眼Storage页确认没有Lost blocks。以我自己的实际经验来说判断要不要持久化就一句话这个RDD会不会被多个Action重复消费超过一次。会就持久化不会就什么都别加。存储级别上我默认选MEMORY_AND_DISK_SER只有当数据量小到能确定全部装进内存时才换成MEMORY_ONLY带副本的级别只放在核心链路且节点不稳定的集群上。这套判断逻辑帮助我处理过不少线上任务变慢和OOM的问题。如果你现在正被某个诡异的Spark性能问题困扰不妨先从Storage页和这五个字段开始查起答案往往就在那里。