
摘要普通 Spark 作业 Driver 挂了重跑就行但 Spark Streaming 的 Driver 挂着的是流处理的进度——挂了不自动恢复整个流就断在那一刻。这篇讲清 SparkStreaming Driver HA 靠什么实现checkpoint 持久化 YARN 自动重启checkpoint 到底存了哪些东西以及一套能直接用的 getOrCreate supervise 搭建代码和三个必须注意的坑。关键词Spark Streaming, Driver HA, checkpoint, getOrCreate, supervise, 故障恢复一、SparkStreaming 的 Driver 为什么特殊普通 Spark 批处理作业Driver 挂了任务重算就行——RDD 血缘容错天然支持。但 Spark Streaming 不一样Driver 挂着的是整个流处理的调度中枢DStream 的 lineage、批处理的节奏。有状态算子updateStateByKey 之类的累计状态也只在 Driver 侧内存里。所以 Driver 一挂如果不做特殊处理流就永久停摆。Driver HA 就是让 Driver 挂了之后能自动重启、从断点接着处理。二、checkpoint 存了什么恢复的根基SparkStreaming 的 Driver HA 核心是checkpoint。它是恢复的根基持久化了四类东西DStream Lineage整个 DStream 的依赖链也就是计算逻辑本身。配置信息SparkConf、StreamingContext、batchInterval 等。有状态算子的状态updateStateByKey 等跨 batch 累计的状态。未处理元数据还没消费完的 block/offset 元数据——这是断点续传的关键决定了从哪继续。在 HDFS 上checkpoint 目录里最核心的是receivedBlockMetadata它记录了每个 batch 的 block 元数据。Driver 恢复时就从这里重建处理链。三、故障恢复的完整流程一次完整的 Driver HA 恢复是这样走的Driver 崩溃进程挂了或所在节点宕机。YARN 自动重启 Driver这要求部署在cluster 模式下Driver 跑在 ApplicationMaster 里并且开了–supervise。读 checkpoint新 Driver 从 checkpoint 目录恢复 lineage 和状态。从断点继续getOrCreate重建 StreamingContext接着处理没处理完的数据。三个要素缺一不可checkpoint存状态 cluster 模式Driver 在 AM 里 supervise挂了自动重启。client 模式下 Driver 跑在提交机上挂了没人拉起来HA 无从谈起。四、搭建实操第一步代码侧checkpoint getOrCreatedefcreateContext():StreamingContext{valsscnewStreamingContext(conf,Seconds(5))ssc.checkpoint(hdfs://namenode:8020/checkpoint/app)// ... 构建 DStream 处理逻辑 ...ssc}// 关键getOrCreate —— checkpoint 存在就恢复不存在就新建valsscStreamingContext.getOrCreate(checkpointPath,createContext _)ssc.start()ssc.awaitTermination()getOrCreate是整套机制的核心入口它先检查 checkpoint 目录是否存在——存在说明是故障恢复场景直接从中重建 StreamingContext不存在说明是首次启动调用createContext新建。注意createContext里要重新设置 checkpoint 目录ssc.checkpoint(...)这行因为新建场景也需要把 checkpoint 路径告诉 StreamingContext后续才会持续写 checkpoint。第二步提交侧cluster supervisespark-submit\--masteryarn\--deploy-mode cluster\--supervise\--classcom.example.StreamingApp\app.jar--deploy-mode clusterDriver 跑在 ApplicationMaster 里随集群管理。--superviseDriver 挂了YARN 自动重启它。五、三个必须注意的坑坑一checkpoint 目录不能变恢复时是从固定目录读的改了路径就等于找不回原来的状态恢复会失败。checkpoint 目录要稳定、用可靠的 HDFS 路径。坑二代码变更不兼容checkpoint 里存的是 DStream 的 lineage序列化的计算逻辑。你改了处理逻辑之后旧 checkpoint 里的 lineage 和新代码对不上恢复时会抛反序列化/兼容性异常。所以上线新逻辑时要换新的 checkpoint 目录或先删掉旧 checkpoint 再启动旧状态作废从零开始。坑三getOrCreate 的正确用法createContext里必须重新执行ssc.checkpoint(...)很多人漏了这一步导致首次启动后根本没有 checkpoint 目录HA 形同虚设。同时createContext里的逻辑要和恢复后的逻辑一致否则新旧 lineage 对不上。六、总结SparkStreaming 的 Driver 挂着流处理进度和状态HA 依赖 checkpoint 持久化 YARN 自动重启缺一不可。checkpoint 存四类东西lineage、配置、有状态算子状态、未处理元数据最后一项决定从哪续传。搭建三要素代码侧checkpointgetOrCreate提交侧--deploy-mode cluster--supervise。三个坑checkpoint 目录不能变、代码变更要换新目录、createContext 里要重设 checkpoint。作者大数据技术实践者博客blog.starzy.cnGitHubstarzy1990.github.io专注 AI Agent · LangGraph · RAG · 大数据架构 · 数据工程实践