ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SparkStreaming Driver HA:Checkpoint恢复与生产落地

SparkStreaming Driver HA:Checkpoint恢复与生产落地 你有没有经历过凌晨三点被电话叫醒打开监控一看Kafka Lag 以肉眼可见的速度往上飙SparkStreaming 作业的 Driver 已经不知道什么时候挂掉了如果你还没有给 Driver 做 HAHigh Availability所有所谓的“自动恢复”就都是空谈。等人工重启之后你会发现数据断了几小时下游报表全部错位这一晚上的数据链路基本等于瘫痪。这篇文章把 SparkStreaming 的 Driver HA 掰开揉碎讲清楚为什么 Driver 是流任务的命门、Checkpoint 恢复机制底层到底存了什么、ZooKeeper 与资源框架分别起到什么作用以及最终如何基于 getOrCreate 和一组参数在生产环境把 Driver HA 真正落起来。目标读者是已经在用 SparkStreaming 做实时计算但还没有给作业做高可用保护、或做了但恢复链路不清晰的工程师。1. 深夜断流与单点命门先搞懂Driver在Streaming里扛着什么很多新手把 SparkStreaming 的 HA 简单理解为“给任务配个重启脚本”这个认知是危险的。不搞清楚 Driver 在 Streaming 作业里具体管着什么你搭出来的所谓 HA 很可能只是在重启一个不断丢失状态的空壳。Driver 进程在 SparkStreaming 作业里承载的职责比普通 Spark 批作业更重。它不仅仅是 SparkContext 的宿主编排 Spark 任务还要管理实时计算特有的运行时状态。以最常见的 Kafka Direct 模式为例Driver 至少要扛住这么几件事DStreamGraph 的构建与维护这个数据流图记录了每个 DStream 之间的依赖关系、窗口操作逻辑、状态操作的算子链路是整个流计算的处理蓝图。JobScheduler 的运行每个 batch interval 会触发一个或多个 Spark Job这些 Job 的生成、排队、提交和状态跟踪都运行在 Driver 端。InputDStream 的 offset 管理Kafka Direct 模式下每次读取哪些分区的哪些 offset是由 Driver 端当前保存的 offset 决定的。Driver 挂了这个消费位置信息也就暂时丢了除非落到了 Checkpoint 或外部系统中。ReceivedBlockTracker 的状态如果你用的是 Receiver 方式Receiver 接收到的数据块分配给了哪个 batch、每个 batch 对应了哪些 Block这些信息同样记录在 Driver 端。一句话总结Streaming 作业的“大脑”就在 Driver 里。大脑死掉即使 Executor 全活着、Kafka 的数据还在堆积作业也无法继续调度。更麻烦的是没有 HA 保护的作业在 Driver 挂掉后整个作业进程往往直接消失不会自动拉起来。数据从挂在那一刻开始断流直到有人发现并重新提交作业。理解了这一点你就能明白为什么 Driver HA 是整个 SparkStreaming 高可用方案里优先级最高的一环。它要解决的其实是两个问题容器/进程层面Driver 进程挂了有没有外部机制把它重新拉起来状态恢复层面进程重新拉起之后能不能从挂掉的瞬间继续跑而不是从头消费数据或直接丢掉 Checkpoint 前的处理结果这两个问题缺一不可。只解决进程拉起不解决状态恢复拉起来之后数据可能从 Kafka earliest 重新消费造成海量重复计算只解决状态恢复不解决进程拉起那恢复逻辑写得再漂亮也没人执行。2. Driver恢复的地基Checkpoint元数据备份到底存了什么掉旧坑之前先把地基打牢。SparkStreaming 的 Driver HA 核心依赖是 Checkpoint它的本质是定期把 Driver 端的元数据序列化后写到可靠的共享文件系统生产环境一般就是 HDFS。恢复的时候新的 Driver 进程从这个目录反序列化出刚才提到的 DStreamGraph、未完成 batch、block 分配信息等从而重建出完整的运行现场。很多人以为 Checkpoint 就是“存了一份 RDD 数据”其实不对。SparkStreaming 的 Checkpoint 分为两种用途完全不同Checkpoint 类型存储内容主要用途Metadata CheckpointSparkConf 配置、DStreamGraph 逻辑、未完成 Batch 的元数据、Receiver 接收块的分配信息恢复 Driver 运行框架Data Checkpoint带状态算子如 updateStateByKey、reduceByKeyAndWindow产生的中间 RDD 数据恢复跨 batch 的计算状态Metadata Checkpoint 里最核心的是 DStreamGraph 的序列化结果。它保存的不是“代码逻辑”而是运行过程中已经构建好的 DStream 对象图——包括每个 InputDStream 当前消费到了哪个 offset、每次操作对应的函数对象、每个状态算子的配置。所以恢复时不依赖重新执行一遍代码来构建图而是直接加载这个对象图。Data Checkpoint 则专门服务于跨 batch 的状态计算。比如你维护了一个按天累计的 key 计数状态必须存在某个位置才能在下一个 batch 继续累加。实时状态默认在内存中Driver 一挂内存就没了因此 SparkStreaming 会按一定周期把状态 RDD 也 Checkpoint 到 HDFS。这个周期称为 checkpoint interval默认等于 batch interval但生产上一般建议设置为 batch interval 的 5 到 10 倍避免频繁写 HDFS 拖垮处理性能。恢复流程拆开看是这样的新的 Driver 进程启动通过StreamingContext.getOrCreate检测到 Checkpoint 目录下存在有效的元数据文件。框架从 HDFS 反序列化 Metadata Checkpoint恢复 SparkConf 和 DStreamGraph。基于恢复出的图重新创建 StreamingContext 的内部组件比如 JobScheduler、ReceiverTracker。找到上次尚未处理完成的 batch从那个时间点重新生成 Job 并提交执行。Data Checkpoint 中的状态 RDD 会作为状态算子的初始状态加载回去保证跨 batch 累计不中断。所以你看到的“自动续跑”不是魔法而是这套元数据恢复机制在背后兜底。但这里也有一个极其容易踩的坑Checkpoint 里存的是 DStreamGraph 对象不是代码。如果你修改了业务代码里的 DStream 转换逻辑从旧 Checkpoint 恢复时加载出来的仍然是旧的图结构新的代码逻辑根本不会生效。反之如果删掉 Checkpoint 目录再启动新代码才会通过 creatingFunc 重新构建。3. 谁来把Driver拉起来ZooKeeper与资源框架的角色分工Checkpoint 负责的是“恢复状态”但“把 Driver 拉起来”这件事本身需要外部机制介入。这就引出了两种常见部署模式下的不同实现方式Standalone 集群模式和 YARN 模式。3.1 Standalone 模式ZooKeeper 管集群supervise 管 Driver在 Spark 自带的 Standalone 集群里如果只启动一个 Master那么 Master 本身也是单点。所以要让 Driver 能被重新拉起先要让 Master 先具备 HA 能力常见做法是把 Master 的元数据放到 ZooKeeper 里让多个 Master 节点通过 ZK 选主active Master 挂掉后 standby Master 自动接管。这时候 Spark 的 Driver HA 实际上是两层配合spark.deploy.recoveryModeZOOKEEPER让 Master 把自己的状态包括已注册的 Driver 信息持久化到 ZooKeeper。提交作业时加上--superviseMaster 才会在 Driver 进程异常退出后考虑重新调度一个 DriverWrapper 到其他存活 Worker 上。拉起后的 Driver 进程仍然是一个全新的 JVM 进程。它启动后执行的还是你提交的那个 main 类中的代码所以代码里必须有StreamingContext.getOrCreate逻辑否则新进程只会重新创建一个全新的 StreamingContext不会去读 Checkpoint更不会自动续跑。Standalone 模式下恢复链路可以这样描述Worker 上的 Driver 进程挂掉机器宕机或 OOM→ Master 通过 ZK 中保存的信息感知 Driver 丢失 → 检查该 Driver 是否设置了 supervise → 在可用 Worker 上重新调度容器启动新 Driver → 新 Driver 从 HDFS Checkpoint 恢复 StreamingContext → 作业续跑。3.2 YARN 模式利用 ApplicationMaster 重试机制大多数公司的生产集群都在 YARN 上这种模式下 Driver 实际上运行在 ApplicationMasterAM内部。Spark 提交到 YARN 的任务本身就具备 AM 重启的能力只要 AM 异常退出ResourceManager 会根据spark.yarn.max.app.attempts配置判断是否重新调度一个 AM 容器。这个设计天然适合做 Driver HA因为重启后的新 AM 容器里会重新启动 Driver 进程。你只需要保证spark.yarn.max.app.attempts大于 1否则 AM 挂掉直接就是任务失败不会重试。Checkpoint 目录在共享存储上新 AM 里的 Driver 能正常读取。代码里正确使用getOrCreate。YARN 模式比 Standalone 省心的地方在于你不需要单独维护一套 Spark Master 的 HAYARN ResourceManager 本身就是高可用的它会妥善处理调度和重试。需要额外注意的一点是YARN 上的 AM 重试如果是由于“内存超限”“OOM Kill”等被 NodeManager 杀死的情况新的 AM 会被调度到其他节点但旧的容器可能残留一段时间。因此 Spark 内部会有防止多个 AM 同时处理同一个作业的机制这也是官方建议不要在 YARN 模式下手动频繁 kill AM 来测试的原因之一。4. 落地实操getOrCreate与关键参数一步步搭出Driver HA理论部分讲完下面直接进入可复制的落地步骤。我以 Kafka Direct Scala 为例实际项目里改动最大的就是这个入口类和一堆 spark-submit 参数。4.1 代码侧必须用 getOrCreate 包装 StreamingContext直接用new StreamingContext的作业没有办法享受 Driver 恢复能力。正确姿势是让启动逻辑被StreamingContext.getOrCreate接管import org.apache.spark.SparkConf import org.apache.spark.streaming.{Seconds, StreamingContext} import org.apache.spark.streaming.kafka010._ def createStreamingContext(): StreamingContext { val conf new SparkConf().setAppName(DriverHAExample) val ssc new StreamingContext(conf, Seconds(5)) // 这里根据你的数据源构造 DStreamGraph val kafkaParams Map[String, Object]( bootstrap.servers - kafka1:9092,kafka2:9092, key.deserializer - org.apache.kafka.common.serialization.StringDeserializer, value.deserializer - org.apache.kafka.common.serialization.StringDeserializer, group.id - driver-ha-demo, auto.offset.reset - earliest, enable.auto.commit - false ) val topics Array(demo_topic) val stream KafkaUtils.createDirectStream[String, String]( ssc, LocationStrategies.PreferConsistent, ConsumerStrategies.Subscribe[String, String](topics, kafkaParams) ) stream.map(_.value()) .map(record (record.split(,)(0), 1L)) .reduceByKeyAndWindow(_ _, _ - _, Seconds(60), Seconds(10)) .foreachRDD { rdd // 你的业务处理逻辑建议做幂等写入 rdd.foreachPartition { iter // 写外部存储 } } ssc.checkpoint(hdfs://nameservice/user/spark/checkpoint/driver-ha-demo) ssc } def main(args: Array[String]): Unit { val checkpointPath hdfs://nameservice/user/spark/checkpoint/driver-ha-demo val ssc StreamingContext.getOrCreate(checkpointPath, createStreamingContext _) ssc.sparkContext.setLogLevel(WARN) ssc.start() ssc.awaitTermination() }这段代码里有几个细节值得抠一下ssc.checkpoint的路径必须放在共享文件系统上本地盘路径在 Driver 换节点后读不到。StreamingContext.getOrCreate的行为是Checkpoint 目录存在且恢复成功 → 直接加载元数据不存在 → 调用传入的createStreamingContext函数创建新作业。所以检查点目录一旦存在创建函数里的 DStream 构建逻辑不会重新执行。恢复成功后不要再次调用ssc.checkpoint来覆盖路径路径以 Checkpoint 文件里的记录为准。enable.auto.commitfalse是为了不依赖 Kafka 消费者组管理 offset否则恢复时的 offset 管理会和 Checkpoint 冲突。4.2 Standalone 模式提交--supervise 是关键如果跑在 Standalone 集群spark-submit 要带着--supervise参数同时给 Spark 配好 ZKspark-submit \ --class com.example.DriverHAExample \ --master spark://master1:7077,master2:7077 \ --deploy-mode cluster \ --supervise \ --conf spark.deploy.recoveryModeZOOKEEPER \ --conf spark.deploy.zookeeper.urlzk1:2181,zk2:2181,zk3:2181 \ --conf spark.deploy.zookeeper.dir/spark-ha \ --conf spark.streaming.kafka.maxRatePerPartition1000 \ driver-ha-demo.jar这里--supervise是最容易被漏掉的一项。不加它Spark Master 也会感知到 Driver 退出但不会重新调度。加上之后Master 才会在你当前作业的 submitted driver 列表里维护一个需要重启的状态。4.3 YARN 模式提交调大 AM 重试次数YARN 模式不需要--supervise核心是下面几个参数spark-submit \ --class com.example.DriverHAExample \ --master yarn \ --deploy-mode cluster \ --conf spark.yarn.max.app.attempts5 \ --conf spark.yarn.am.attemptFailuresValidityInterval1h \ --conf spark.streaming.kafka.maxRatePerPartition1000 \ driver-ha-demo.jarspark.yarn.max.app.attempts默认值是 1不显式调大AM 挂掉直接任务失败。我见过不少团队在这里吃过亏辛辛苦苦搭好了 Checkpoint结果 AM 重试参数没改Driver 一挂作业就进入 FAILED 状态恢复逻辑完全没机会执行。另外要注意 YARN 的硬性限制spark.yarn.max.app.attempts不能超过 ResourceManager 侧的yarn.resourcemanager.am.max-attempts后者通常默认是 2。如果 RM 侧没调大你这里写 5 也会被强制压回 2。4.4 常用参数速查与说明参数推荐值说明spark.deploy.recoveryModeZOOKEEPER仅 Standalone 模式开启 Spark Master 元数据持久化spark.deploy.zookeeper.urlzk1:2181,zk2:2181Master 元数据存储的 ZK 地址spark.deploy.zookeeper.dir/spark-haZK 中存储 Spark 元数据的 znode 路径spark.yarn.max.app.attempts3 到 5YARN 模式 AM 最大重试次数决定 Driver 能自动拉起几次spark.yarn.am.attemptFailuresValidityInterval1h统计 AM 失败时间窗口避免历史失败累积导致不再重试spark.streaming.receiver.writeAheadLog.enabletrueReceiver 方式下开启 WAL减少数据丢失概率spark.streaming.receiver.writeAheadLog.rollingFile.maxSize512MB控制 WAL 文件滚动大小减少 HDFS 小文件spark.streaming.stopSparkContextByDefaultfalse恢复后 Driver 退出时是否连带关闭 SparkContext这里额外解释一下spark.streaming.stopSparkContextByDefault。它默认是 true含义是 StreamingContext 停止时自动联动停止 SparkContext。但在 Driver HA 场景下如果新 Driver 从 Checkpoint 恢复失败、或某个批处理抛出致命异常导致 ssc.stop()这个默认行为会让 SparkContext 一起退出最终把 AppMaster 也带崩溃。生产上建议显式设为 false至少保住 SparkContext方便排查和进一步处理。4.5 验证 Driver HA 是否真的生效搭完之后一定要做一次故障演练不要等到线上真的挂了才发现配置不对。实践中我一般按这个顺序验证先把作业正常提交确认 Checkpoint 目录在 HDFS 上生成能看到metadata子目录里有文件写入。登录 YARN Web UI 找到当前 ApplicationMaster 的运行节点和容器 ID。在对应节点上执行kill -9杀掉 AM 进程模拟最极端的崩溃场景。观察 YARN 是否在若干秒后自动调度新的 AM。查看新 AM 的日志正常情况下会看到Recovered from checkpoint之类的日志说明恢复流程已经触发。再确认 Kafka 消费位点不是从 earliest 重新开始而是从挂掉前的位置继续。如果发现恢复后作业确实从 Checkpoint 续跑但处理 Lag 明显高于常态这是正常现象说明恢复过程中积压的 batch 正在被补算。不要慌等它追到实时进度即可。5. 恢复不是保险箱WAL、幂等消费与那些易踩的坑Driver HA 能自动重启进程并恢复元数据但这绝不等于数据处理“既不丢也不重”。实际操作中很多从外表看已经开了 HA 的作业在数据语义上仍然有隐患。5.1 WAL 到底在防什么如果你是老派 Receiver 方式数据先由 Receiver 接收并存放在 Executor 内存中Driver 端只记录块元数据。这个过程有个天然弱点Receiver 收到了数据但还没等 Driver 记录“这个 batch 包含哪些 Block”Driver 挂了恢复后这些 Block 对应关系丢失数据也就丢了。WALWrite Ahead Log要解决的就是这个时间窗问题Receiver 收到数据后先写入 HDFS 上的 WAL再更新内存中的块信息。恢复时可以从 WAL 重放数据。Kafka Direct 模式其实不太依赖 WAL因为 offset 就保存在 Checkpoint 或外部存储里数据源本身就是可重放的。但如果你混合使用了 Receiver 和 Kafka Direct或任务里有其他第三方数据源开启 WAL 仍然是一个重要的安全垫。5.2 重复消费与幂等是必须面对的现实这里必须泼一盆冷水即使配置全部到位Driver HA 提供的也只是 At-Least-Once 语义不是 Exactly-Once。原因很简单恢复时那个“未完成 batch”到底执行到哪一步Spark 并没有精确记录。有可能 Job 已经跑完了大部分任务只是结果还没来得及写入外部存储Driver 就挂了。恢复后这个 batch 会被重新调度执行下游就会收到重复数据。所以生产落地的铁律是下游写入必须做幂等或幂等控制。比如写 MySQL用唯一键插入或 update 而不是无条件 insert写 ES用带 id 的 upsert写 Kafka给消息加全局唯一 ID消费端自己做去重。没有幂等保护HA 恢复的往往是“半截数据”比没有 HA 更让人头疼。5.3 Checkpoint 与代码变更的冲突前面提过一句这里值得展开。Checkpoint 恢复的是序列化后的 DStreamGraph 对象不是当前运行代码。这意味着如果你改了 DStream 的转换逻辑比如从 map 改成 flatMap或者改了窗口大小从旧 Checkpoint 恢复时运行的是旧逻辑代码变更完全不会生效。如果你改了类名、包名、算子里引用的类结构反序列化时还可能直接报 ClassNotFoundException 或序列化异常恢复失败。新逻辑要上线的标准做法是要么清空旧的 Checkpoint 目录重新开始消费要么另开一个全新的 Checkpoint 路径。无论如何都要接受重新消费已有数据或位点重置带来的业务影响。我踩过的坑是有一次只换了一个工具类的内部实现没改类名结果反序列化时旧对象带出来的字段结构和新类对不上恢复直接报错。后来吸取的教训是凡是涉及代码逻辑变更先主动验证 Checkpoint 能恢复不能恢复就果断清理检查点。5.4 Checkpoint 间隔与 HDFS 小文件Checkpoint 写得越频繁Driver 恢复时丢失的状态越少但代价是 HDFS 上会产生大量小文件同时对 NameNode 造成压力。Metadata Checkpoint 每个 batch 都会写一次默认情况下Data Checkpoint 则按 checkpoint interval 触发。生产调优建议Metadata Checkpoint 频率保持默认即可Data Checkpoint 间隔按状态数据量和恢复容忍度调整一般设为 batch interval 的 5 到 10 倍。如果你的 batch 是 5 秒那 checkpoint interval 可以设 30 秒左右如果 batch 是 1 分钟那 5 到 10 分钟都可以。写 HDFS 前也可以考虑用压缩ssc.checkpoint(hdfs://nameservice/...) ssc.conf.set(spark.rdd.compress, true)spark.rdd.compress能在一定程度上减少 Checkpoint 和 shuffle 中间数据的体积但对 CPU 有额外开销要结合集群资源情况权衡。5.5 Kafka Direct 模式 offset 恢复的两个细节第一个细节恢复后auto.offset.resetearliest不会因为 Checkpoint 里已经保存了 offset 而重新从最早消费。这个配置只在没有 Checkpoint 且 Kafka 里没有已提交 offset 时才会生效。所以不要在恢复失败后天真地以为调小 earliest 就能把数据捡回来。第二个细节如果你同时启用了 Kafka 自身的 offset 提交enable.auto.committrue那么 Checkpoint 和 Kafka 记录的 offset 可能不一致。SparkStreaming 官方建议在 streaming 场景下禁用 Kafka 自动提交以 Checkpoint 为准。否则恢复时两边 offset 打架会出现跳过数据或重复消费的奇怪现象。5.6 同机恢复与资源不足虽然 YARN 和 Standalone 都会尽量把新 Driver 调度到其他节点但如果集群资源紧张新的 AM 容器有可能被分配到同一台故障率较高的机器上。更麻烦的是Driver 恢复时需要重新创建 SparkContext如果之前的 Executor 动态分配策略没配好恢复后可能长期只有少量 Executor 在处理积压 batch。为此可以配合spark.dynamicAllocation.enabledtrue和spark.dynamicAllocation.maxExecutors来给恢复过程留一些弹性。6. 我这几年的实操心得与验证建议最后聊一点个人经验不一定能写进官方文档但都是真实趟过后觉得值得记下来的点。第一能上 YARN 尽量上 YARN别在 Standalone 上硬磕。Standalone 的 Driver HA 虽然原理上很清晰但它要求你额外维护 Spark Master 的 ZK HA而且--supervise之后新 Driver 经常被调度到任意一个 Worker 上客户端日志采集、监控 Agent 部署都得跟着适配。YARN 模式下 AM 重启是资源管理器原生能力配合 Checkpoint 就能把恢复链路走通运维负担小很多。第二把故障演练当成上线流程的一部分。我每搭好一个新的 Streaming 任务都会在测试环境做一次“杀 Driver 进程”演练确认日志里出现从 Checkpoint 恢复的记录确认 Kafka Lag 能在几分钟内回落。很多参数问题只有真刀真枪杀进程才暴露得出来。演练脚本很简单核心就是根据 YARN applicationId 找到 AM 所在节点然后 kill 对应进程。跑通一次之后心里才真正有底。第三监控不能只看 Kafka Lag还要看 Driver 重启事件。很多 HA 配置能自动恢复但恢复过程通常有几分钟空窗数据会积压。如果监控维度太粗可能业务方已经发现报表延迟了你还在因为 Lag 还没突破阈值而没收到告警。建议对 YARN AM 重启次数、SparkStreaming 的 job 提交时间间隔做单独告警这两个指标能第一时间反映 Driver 是否发生了重建。第四永远为幂等做好准备。即使你觉得当前下游系统不会产生重复数据也要在设计阶段把幂等逻辑加上。因为 Driver HA 恢复的“未完成 batch”在极端情况下可能被重复执行多次如果下游是无脑 append数据质量会直接崩掉。这个成本在开发阶段看似多余但在真实故障发生时是最能救命的一道防线。SparkStreaming 的 Driver HA 不是一个开关就能解决的问题它需要代码侧、参数侧、部署侧三条线同时配合。把 getOrCreate 用对、把 AM 重试次数调大、把 Checkpoint 放对地方、把下游幂等做好这一套走下来你才算是真正给流任务上了一道保险。
返回列表