一、先搞懂两个核心工具的本质(大白话版)
很多做大数据处理的朋友都有过这种崩溃经历:跑一个要十几个小时的Spark作业,跑了大半突然断了,要么是集群节点挂了,要么是中间数据丢了,只能从头再来,既费时间又占资源。今天就给大家拆解两个能解决这个问题的核心工具——Checkpoint和Cache,先不说术语,用大白话讲清楚它们到底是干嘛的。
1.1 Cache是什么?
你可以把Cache理解成电脑里的临时缓存文件夹。比如你用PS修图,修到一半会把当前的编辑进度存在内存里,下次打开直接接着弄,不用重新加载原图。Spark的Cache就是干这个的:它会把你计算过程中产生的中间数据(比如你从数据库拉出来的原始数据、经过过滤后的数据集)存在内存或者本地磁盘里,下次再用到这些数据的时候,直接从缓存里拿,不用重新计算,这就是它能加速计算的原因。 但Cache有个大问题:它是临时的,作业跑完或者集群重启,缓存的数据就没了;而且如果集群的内存不够用,Spark还会把部分缓存数据丢了,或者存在速度很慢的本地磁盘里,反而会拖慢速度。
1.2 Checkpoint是什么?
Checkpoint就像是你写论文时存在U盘里的备份文件。你写论文写到一半,怕电脑死机,就把当前的版本存在U盘里,哪怕电脑坏了,U盘里的文件还能接着用。Spark的Checkpoint就是把中间数据存在可靠的存储介质里(比如HDFS、S3这些专门存数据的地方),哪怕集群重启、节点挂了,这些数据也不会丢。 但Checkpoint也有缺点:它是把数据序列化后存到磁盘上的,读写速度比Cache慢很多,而且每次Checkpoint都会触发一次全量的计算和存储,会占用额外的集群资源,不能随便用。
二、什么时候需要用这两个工具?(高频场景盘点)
不是所有Spark作业都需要用Cache和Checkpoint,只有满足特定场景的作业,用了才会有效果,我给大家列几个最常见的需要用的场景:
2.1 长周期作业频繁失败
比如你要处理近一年的用户行为数据,作业要跑10个小时以上,而且经常因为集群波动、数据倾斜等问题失败,这种时候就需要用Checkpoint做容错,用Cache做加速。
2.2 中间数据被反复使用
比如你先对原始数据做了过滤、去重,得到了一个基础数据集,后面要基于这个数据集做统计、分类、关联等多个操作,这种时候把这个基础数据集Cache起来,就能避免每次都重新计算。
2.3 数据处理逻辑复杂
比如你要做多层嵌套的计算,先算A,再用A算B,再用B算C,这种时候中间数据的计算成本很高,一旦失败,重新计算的代价很大,就需要用Checkpoint做备份。
三、怎么组合用才合理?(完整示例演示)
光说理论没用,我给大家写一个完整的示例,把Cache和Checkpoint的组合用法讲清楚。示例用的技术栈是Spark 3.3.0(目前主流的稳定版本),语言用Scala,代码里加了详细的注释,大家可以直接参考。 首先,先说明示例的场景:我们要处理一个用户行为数据集,先做过滤、去重,得到基础数据集,然后基于这个基础数据集做两个统计:一个是用户的总访问次数,一个是用户的平均访问时长。这个作业跑一次大概需要8个小时,经常因为集群波动失败,所以我们需要用Cache和Checkpoint来优化。 首先,我们先写基础的Spark作业代码:
// 导入Spark相关的包
import org.apache.spark.sql.SparkSession
import org.apache.spark.storage.StorageLevel
object SparkCacheCheckpointDemo {
def main(args: Array[String]): Unit = {
// 1. 初始化SparkSession,设置应用名称
val spark = SparkSession.builder()
.appName("UserBehaviorAnalysis")
.getOrCreate()
// 2. 设置Checkpoint的存储路径,这里用HDFS的路径,也可以用S3等其他可靠存储
// 注意:这个路径必须是集群所有节点都能访问的路径
spark.sparkContext.setCheckpointDir("hdfs://hadoop-namenode:9000/spark/checkpoint/user_behavior")
// 3. 读取原始用户行为数据,假设数据存在HDFS上,格式是Parquet
val rawData = spark.read.parquet("hdfs://hadoop-namenode:9000/data/user_behavior/raw/")
// 4. 对原始数据做过滤、去重,得到基础数据集
// 过滤掉无效的用户ID、访问时长为负的数据,去重
val baseData = rawData
.filter("user_id is not null and duration > 0") // 过滤无效数据
.dropDuplicates("user_id", "visit_time", "page_id") // 按用户ID、访问时间、页面ID去重
// 5. 组合使用Cache和Checkpoint
// 第一步:先把基础数据集Cache起来,用MEMORY_ONLY_SER(内存存储,序列化),适合数据量较大的情况
// 如果数据量较小,可以用MEMORY_ONLY(不序列化,速度更快)
baseData.persist(StorageLevel.MEMORY_ONLY_SER)
// 第二步:触发一次Action操作,让Spark把数据缓存到内存里,比如用count()
baseData.count()
// 第三步:对基础数据集做Checkpoint,把数据存到HDFS上
baseData.checkpoint()
// 第四步:再触发一次Action操作,让Spark完成Checkpoint的存储,比如用count()
baseData.count()
// 6. 基于基础数据集做第一个统计:用户总访问次数
val userVisitCount = baseData
.groupBy("user_id")
.count()
.withColumnRenamed("count", "visit_count")
// 7. 基于基础数据集做第二个统计:用户平均访问时长
val userAvgDuration = baseData
.groupBy("user_id")
.avg("duration")
.withColumnRenamed("avg(duration)", "avg_duration")
// 8. 把统计结果写入HDFS
userVisitCount.write.parquet("hdfs://hadoop-namenode:9000/data/user_behavior/result/visit_count/")
userAvgDuration.write.parquet("hdfs://hadoop-namenode:9000/data/user_behavior/result/avg_duration/")
// 9. 关闭SparkSession
spark.stop()
}
}
3.1 示例里的组合逻辑解释
大家看上面的代码,我们的组合逻辑是:先Cache,再Checkpoint,中间用Action操作触发计算。为什么要这么做呢? 首先,Cache是存在内存里的,速度快,所以先把基础数据集Cache起来,后面做统计的时候,直接从Cache里拿数据,速度很快。然后,我们再做Checkpoint,把数据存到HDFS上,这样哪怕后面作业失败了,下次启动的时候,Spark会先去Checkpoint的路径里找数据,找到的话就直接用,不用重新计算原始数据的过滤、去重操作,节省了大量的时间。 这里要注意两个点:一是Checkpoint之前一定要触发Action操作,不然Spark不会真的计算出基础数据集;二是Checkpoint之后,Spark会自动把这个数据集的血统(也就是计算过程的记录)切断,所以后面的操作都直接基于Checkpoint的数据来做,这样哪怕后面的操作失败了,也不会影响前面的Checkpoint数据。
3.2 不同场景的组合调整
上面的示例是基础的组合方式,大家可以根据自己的场景调整: 如果你的基础数据集非常大,内存存不下,那Cache的StorageLevel可以改成MEMORY_AND_DISK_SER(内存和磁盘都存,序列化),这样内存不够的话,会把数据存在本地磁盘里,速度比纯内存慢,但比重新计算快。 如果你的作业失败频率特别高,那可以在Checkpoint之后,再把Checkpoint的数据Cache起来,这样后面的操作速度更快。 如果你的作业中间有多个需要重复使用的数据集,那可以给每个数据集都做Cache和Checkpoint,但要注意不要重复Cache,不然会占用过多的内存。
四、参数配置建议(让工具发挥最大作用)
光会用还不够,合适的参数配置能让Cache和Checkpoint发挥最大的作用,我给大家列几个最关键的参数:
4.1 Cache相关参数
- spark.storage.memoryFraction:这个参数设置Spark用来存储Cache数据的内存比例,默认是0.6,也就是60%的内存用来存Cache数据。如果你的作业中间数据很多,需要Cache的比例大,可以把这个参数调到0.7或者0.8,但不要超过0.8,不然会影响计算的内存使用。
- spark.storage.unrollFraction:这个参数设置Spark用来展开Cache数据的内存比例,默认是0.2,也就是20%的内存用来展开数据。如果你的数据结构比较复杂,可以把这个参数调到0.3。
- spark.shuffle.file.buffer:这个参数设置Shuffle操作的缓冲区大小,默认是32k。如果你的作业Shuffle操作很多,可以把这个参数调到64k或者128k,能减少磁盘的读写次数。
4.2 Checkpoint相关参数
- spark.checkpoint.dir:这个参数设置Checkpoint的存储路径,一定要设置成集群所有节点都能访问的路径,比如HDFS的路径、S3的路径,不要设置成本地路径,不然节点挂了,Checkpoint数据就丢了。
- spark.checkpoint.period:这个参数设置Checkpoint的周期,默认是-1,也就是不自动做Checkpoint,需要手动调用checkpoint()方法。如果你的作业很长,可以把这个参数调到1800(也就是30分钟),每30分钟自动做一次Checkpoint,这样哪怕作业失败,也只会损失最近30分钟的计算。
- spark.checkpoint.cleanup.enabled:这个参数设置是否自动清理旧的Checkpoint数据,默认是true。如果你的Checkpoint路径空间不够,可以把这个参数设置成true,自动清理旧的数据。
4.3 其他相关参数
- spark.executor.memory:这个参数设置每个Executor的内存大小,默认是1G。如果你的作业需要Cache大量的数据,需要把这个参数调到4G、8G甚至更高,根据你的数据量来定。
- spark.driver.memory:这个参数设置Driver的内存大小,默认是1G。如果你的作业需要处理大量的元数据,需要把这个参数调到2G或者4G。
五、优缺点分析和注意事项
5.1 组合使用的优缺点
优点:
- 容错性高:Checkpoint把数据存在可靠的存储里,哪怕集群重启、节点挂了,数据也不会丢,作业失败后可以快速恢复。
- 计算速度快:Cache把数据存在内存里,后面的操作直接从Cache里拿数据,不用重新计算,能节省大量的时间。
- 资源利用率高:组合使用能平衡容错和速度,不用为了容错牺牲速度,也不用为了速度牺牲容错。 缺点:
- 配置复杂:需要根据作业的场景调整参数,配置不当反而会降低速度,比如Cache的内存比例设置太高,会导致计算内存不足,反而变慢。
- 存储成本高:Checkpoint会占用大量的存储资源,比如HDFS的空间,需要定期清理旧的Checkpoint数据。
- 血统切断的影响:Checkpoint会切断数据集的血统,所以如果你的作业需要回溯计算过程,可能会有影响。
5.2 注意事项
- 不要对所有数据集都做Checkpoint:只有那些计算成本高、被反复使用的数据集才需要做Checkpoint,不然会占用大量的存储资源。
- Checkpoint之前一定要触发Action操作:不然Spark不会真的计算出数据集,Checkpoint就会失败。
- 定期清理旧的Checkpoint数据:Checkpoint数据会越来越多,占用大量的存储资源,所以要定期清理旧的、没用的Checkpoint数据。
- 测试参数配置:不同的作业场景,参数配置也不一样,所以要先在测试环境里测试参数,再用到生产环境里。
六、文章总结
对于频繁失败的长周期Spark作业,合理组合使用Cache和Checkpoint,是实现容错保护和计算加速的双赢策略。核心的逻辑是:用Cache来加速重复计算,用Checkpoint来做容错备份,中间通过Action操作触发计算,保证数据的正确性。 在实际使用中,大家要根据自己的作业场景调整组合方式和参数配置,比如数据量小的作业可以用纯内存的Cache,数据量大的作业可以用内存加磁盘的Cache,失败频率高的作业可以增加Checkpoint的频率。同时,要注意定期清理旧的Checkpoint数据,避免占用过多的存储资源。 最后,要记住,Cache和Checkpoint不是万能的,还要结合其他的优化手段,比如数据倾斜的优化、Shuffle的优化等,才能让Spark作业跑得更快、更稳定。
Comments