你有没有遇到过这种情况:同一个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/Dataset,cache()的默认行为并不一样。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 一张表看完全部StorageLevel
Spark内置了12个标准存储级别,我先整体列出来,后面逐个详解:
| 存储级别 | 使用磁盘 | 使用内存 | 使用堆外内存 | 序列化 | 副本数 |
|---|---|---|---|---|---|
| NONE | 否 | 否 | 否 | 否 | 1 |
| DISK_ONLY | 是 | 否 | 否 | 是 | 1 |
| DISK_ONLY_2 | 是 | 否 | 否 | 是 | 2 |
| MEMORY_ONLY | 否 | 是 | 否 | 否 | 1 |
| MEMORY_ONLY_2 | 否 | 是 | 否 | 否 | 2 |
| MEMORY_ONLY_SER | 否 | 是 | 否 | 是 | 1 |
| MEMORY_ONLY_SER_2 | 否 | 是 | 否 | 是 | 2 |
| MEMORY_AND_DISK | 是 | 是 | 否 | 否 | 1 |
| MEMORY_AND_DISK_2 | 是 | 是 | 否 | 否 | 2 |
| MEMORY_AND_DISK_SER | 是 | 是 | 否 | 是 | 1 |
| MEMORY_AND_DISK_SER_2 | 是 | 是 | 否 | 是 | 2 |
| OFF_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就是useDisk=true、useMemory=false、useOffHeap=false、deserialized=true、replication=1。由于deserialized=true,存在磁盘上的数据不需要反序列化,直接以对象字节流存储,读取时靠内部机制自动处理。MEMORY_AND_DISK_SER则是useDisk=true、useMemory=true、deserialized=false,多了一个"序列化存储"的语义。
4.2 deserialized参数是怎么"伪装"成序列化选项的
既然参数叫deserialized,为什么我们平时说的是"是否序列化"?这里有个容易绕晕的点:deserialized的含义是"是否以反序列化后的Java对象形式存储"。当它为true时,意味着不序列化,直接存对象;为false时,意味着要序列化成字节再存。
所以MEMORY_ONLY的deserialized=true,而MEMORY_ONLY_SER的deserialized=false。理解了这层关系,你再看源码里那些标准级别的定义,就不会觉得它们长得像玄学了。也可以基于这五个参数自定义存储级别,比如想要三副本、内存加磁盘、序列化,写一个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的存储级别时,如果你不配置Kryo,Spark默认用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页和这五个字段开始查起,答案往往就在那里。