大家做SparkStreaming实时任务,问得最多、踩坑最深的一个话题,就是“到底怎么保证数据不丢”。作为跑了几年实时数仓的老兵,我在这块算是交了不少学费,今天把SparkStreaming容错机制这件事掰开揉碎讲清楚。文章不会只停留在“什么是checkpoint”这种概念层,而是从实际故障场景出发,把从数据源到输出端的完整链路拆开看,哪些地方会丢数据、Spark内部做了什么、我们自己还要补什么,一次性讲透。
先记住一个核心结论:SparkStreaming的容错不是某一个机制单独完成的,而是一套由数据源重放、RDD血缘、checkpoint、WAL、输出幂等性共同组成的组合拳。很多资料把它们混在一起讲,读完你还是不知道出事时该查哪块,这篇文章按我的实际排查经验重新组织。
1. 容错问题全景:故障发生之后,谁在负责什么
1.1 实时任务常见的三类故障
实时流任务不像离线任务,跑挂了重试一次就行。它一边消费上游数据,一边产结果,一旦中途崩溃,恢复的时候要面对三个问题:
- 数据丢没丢:正在处理、但还没来得及输出的数据,在故障后还能不能找回来。
- 数据重没重:Kafka里的offset可能已经提交,但结果没写出去,恢复后消费重复的数据,会不会把结果写两遍。
- 状态还在不在:实时任务往往有累计逻辑(比如按用户汇总、计数、session统计),这些状态是存放在内存里的,TaskManager一挂就全没了,没有持久化的话,恢复之后状态归零,结果完全对不上。
这三类问题的处理思路完全不同,你可以在心里先分个组:数据不丢靠数据源重放和WAL;数据不重靠输出幂等;状态不丢靠checkpoint。后面我按这条线展开。
1.2 SparkStreaming容错的整体分层
网上很多文章画架构图,把我绕晕过很久。后来我自己做项目,把容错按职责拆成了四层,排查起来很清楚:
| 层级 | 负责机制 | 解决什么问题 |
|---|---|---|
| 输入层 | 数据源重放机制(Receiver/Direct) | 故障后未处理的数据能否重新获取 |
| 计算层 | RDD血缘(lineage) | 已生成的分区数据丢失后能否重新计算 |
| 元数层 | checkpoint持久化 | 应用配置、DAG逻辑、状态数据能否恢复 |
| 输出层 | 输出幂等性(事务/唯一键) | 重复计算导致的结果重复写入 |
这个分层特别重要。面试的时候,如果被人问“SparkStreaming是怎么做容错的”,你按这四层回答,面试官会觉得你脑子里有一张清晰的地图,而不是背了两个术语。
1.3 什么是真正意义上的“数据不丢”
大家在讨论容错时,经常把“数据不丢”挂在嘴边,但很少有人先定义清楚:到底怎么算丢?
我们做实时任务的最终目的是把结果写到目标系统(数据库、HDFS、消息队列)。所以“数据不丢”的精确含义是:每条进入Kafka的数据,在故障之后,最终都能体现在目标系统里一次以上(至少一次语义)。
这里不动声色地引入了一个真实理解:SparkStreaming默认给出的是at-least-once(至少一次)的保证。如果你想做到精确一次,那还得看输出端的处理方式和数据源的消费方式,这是整个容错话题里最微妙的地方,后面单独展开。
2. 核心机制深度剖析:checkpoint绝不是只存个进度
2.1 checkpoint到底存了什么
很多人的理解只停留在“checkpoint就是存Kafka offset”,这个理解太浅了,而且在实际排查问题时会误导你。
SparkStreaming的checkpoint分为两类,安全起见我们用表格看清楚:
| 类型 | 存储内容 | 恢复目的 |
|---|---|---|
| Metadata Checkpoint | 应用配置、DAG逻辑、未完成批次的边界 | Driver故障后重新构建执行计划 |
| Data Checkpoint | RDD的血缘信息、State状态数据(配合WAL) | Executor故障后重新计算数据、恢复状态 |
Metadata Checkpoint存在的必要性,在于Driver挂掉之后,整个任务的拓扑信息需要能重新拉起来。如果你只是重启一个SparkSubmit,配置和代码写死了还好说,但那些正在运行的批次、zipkin式的调度状态,全部需要从checkpoint里恢复。
Data Checkpoint则是把stateful计算的中间状态(比如updateStateByKey里的状态)持久化。这里有个最容易踩的坑——如果你用了状态算子,但没在StreamingContext里开启checkpoint,任务重启后状态直接归零,你看着程序没报错,实际上统计结果全错了,而且这种错是悄无声息的,往往等对账才发现,代价极高。
2.2 WAL(Write Ahead Log)为什么常常被误解
WAL和checkpoint的关系很多资料说不清楚。我用一句话帮你彻底区分:
checkpoint是“记一笔账”,WAL是“为这笔账留底稿”。
在Receiver模式下,Executor收到上游数据后,数据先存在内存里,如果这个Executor突然挂掉,内存里的数据就没了,这时候靠什么恢复?靠Kafka或上游数据源的重放。但问题来了,如果数据早就从Kafka消费走了,Kafka那边在短时间内可能已经清理掉这段数据(默认保留7天,但高吞吐场景会调短),数据源也没法重放了。
WAL解决的就是这个窗口问题:Receiver收到数据后,先把它同步落盘到HDFS上一个预写日志里,再做下一层的计算。等到数据源确认不再需要负责这段数据时,WAL日志再做清理。
但我必须说一个残酷的现实:WAL不是免费的午餐,开WAL的代价是吞吐量大幅下降。每条数据先写一次HDFS再到内存,磁盘IO的瓶颈立刻卡上来。我们生产环境实测过,开了WAL之后吞吐下降约40%~60%。所以如果你用的是Direct模式(推荐),压根不需要WAL,因为Kafka本身就在扮演日志的角色。
2.3 checkpoint存储的选型与生命周期
checkpoint最好放在HDFS上,这不是随口一说,而是有具体原因的:
- HDFS提供跨节点的高可用,Task换一台机器还能读到checkpoint数据。
- 本地文件系统在机器崩溃时,数据跟着机器一起没了,checkpoint白存。
- HDFS的写入和读取是分布式并行的,状态大的任务恢复起来快。
还有数据版本的问题。很多人以为checkpoint是“最新的一份”,其实checkpoint目录下会保留多个版本。SparkStreaming默认会在每次batch之后更新checkpoint,恢复时选择最近一个已完成且一致的版本。
这里要提醒一个常见事故:checkpoint目录残留了旧版本数据时,重启会把历史状态加载回来。如果你改了业务逻辑,特别是改了状态存储结构(比如从单字段改成case class),老checkpoint反序列化会直接报错,而且错误信息往往很绕,不是一眼能看出来的。我们做过一次大版本升级,就因为没清理checkpoint目录,任务起来就挂,折腾了两个小时才定位到是旧状态文件类型不匹配。
3. 数据源容错选型:Receiver与Direct模式的生死抉择
3.1 Receiver模式:WAL救了命,但也埋了雷
Receiver模式是SparkStreaming最早期的方式,它有一个专门的Receiver线程常驻在Executor上,负责去拉数据到内存。这个模式天然地存在两个问题:一是数据先进入内存再处理,内存压力大;二是Receiver把数据拉进内存后,如果任务还没处理就挂掉,数据就直接丢失。
所以这个模式设计出来的时候,就配套了WAL机制。开启WAL之后,数据先写日志再处理,看起来解决了丢数据问题,但代价是带来了重复消费问题。你想想这个场景:数据写进了WAL,业务也处理了,但输出结果还没写入目标库,此时任务崩溃。重启之后WAL里的数据会重新被读取和处理,而这部分业务数据之前已经算过一版,只是没来得及写库,于是恢复之后它会再算一遍,导致输出结果重复。
想处理重复问题?那你就得去输出端做幂等保障。但要命的是,这个模式下你连判断哪些数据已经被处理过都很困难,因为WAL没有记录完整的事务边界。所以在真实项目里,我几乎不推荐用Receiver模式。除非你的上游数据源不是Kafka,而是某种普通socket、文件目录这种没有offset概念的输入,否则没有任何理由放弃Direct模式。
3.2 Direct模式:没有WAL,也一样不丢数据
Direct模式,官方名字叫Direct Approach,是Spark 1.3引入的模式,之后成了绝对主流。它和Receiver模式最大的区别是:
Receiver是“推”数据进Spark,Direct是“拉”数据进Spark。
Direct模式下,没有Receiver这个组件,Driver自己负责计算Kafka每个分区这次要消费哪些offset,生成对应的RDD分区。这里有一个关键的机制变化:
- 数据不会先落到Spark的内存中,再去依赖WAL做保护。
- Kafka本身在扮演“日志”的角色——如果Executor处理过程中挂了,只要offset没有提交,下次调度的时候Driver会重新创建RDD,并指定相同的offset范围,重新拉取这些数据。
这个模式的好处是不需要WAL那一层写盘,吞吐直线提升;坏处是数据“可能”会被重复消费。为什么说可能?因为如果一批数据处理到一半挂掉了,这一批的offset已经写入本次调度上下文,但结果可能已经部分输出。重启后Driver重新指定这批次全量offset范围,所以这批次的数据会被再消费一次。于是计算结果重复。想要精确一次,还是在输出端想办法。
3.3 为什么推荐Direct模式:三个决定性理由
- 语义清晰:batch和Kafka partition一一映射,相当于把Kafka消息在批次内的顺序关系天然保留下来,后续排查问题时可以精准定位到某个offset区间。
- 资源占用低:不再需要常驻Receiver线程,不需要WAL写盘,能省下不少Executor资源,瓶颈在Kafka吞吐而非Sparke处理极限。
- 背压机制更自然:Direct模式消费速度由Spark的处理速度直接控制,天然有反压能力。Receiver模式下数据会先积压在Receiver内存里,即使有背压还是可能积压。
从我们线上经验来看,所有要求稳定的生产任务全部采用Direct模式。除非你用的数据源是一个没有提交offset概念的系统,否则我都建议你把任务改成Direct模式。
4. 计算层的容错:RDD血缘是如何自动兜底的
4.1 lineage的本质
实时任务的每个batch,会生成一个RDD。RDD不是凭空产生的,它记录了自己的“祖宗”:父RDD是谁、经过了哪些算子、输入数据从哪来。这套依赖链就是lineage(血缘)。
如果某个Executor在计算过程中挂掉,Driver会根据血缘信息重新调度,重新拉取输入数据,重新执行这段计算。而且因为RDD本身是只读的、每个分区只计算一次的设计,这个过程对上层透明,不需要开发人员写任何恢复代码。
这就是SparkStreaming与其他流处理框架不一样的地方——很多框架需要自己做状态备份和恢复,而RDD的设计天然自带“重建”的能力。
4.2 血缘重建存在的边界
RDD血缘不是万能的,它有一个前提:上游数据还在。
如果输入源的offset已经被清理了(比如Kafka Retention到期),或者读取的是文件、但文件已经被删除了,那RDD血缘就算再完整,也拉不回原始数据。这个边界在数据量大的场景特别要命。我们遇到过一次Kafka集群磁盘告急,运维同事把retention调短到4小时,结果早上实时任务处理故障恢复时,需要的offset早就被清了,错误码是OffsetOutOfRangeException,直接导致那个时段的数据永久性丢失。
所以这里要强调一个原则:计算层的容错机制再好,也依赖一个可靠的数据源保留策略。Kafka的retention建议设置为你允许的“最长故障恢复时间 + 2小时”以上,别为了省磁盘空间把自己逼到绝境。
4.3 一个容易误会的点:checkpoint和lineage在恢复时的分工
恢复时,checkpoint和lineage到底谁先起作用?这个很多人容易搞混,画过不少错误的图。
实际的恢复逻辑是这样的:如果Driver挂了,应用会从最近的checkpoint恢复元数据信息和执行计划,然后通过lineage来重建那些还未完成的batch。checkpoint里存了当前“进度”和自己想要的“地图”,但具体脚印还是要靠lineage一步步重走。如果只是Executor挂了但Driver还活着,是不需要重启整个应用的,Driver会直接根据lineage重新调度丢失的分区计算,走的是局部恢复流程。
这个差异在运维上的含义是:Executor故障几乎对用户无感,而Driver故障会导致几秒到几十秒的停机恢复,期间任务无法产出新数据。调度上要通过监控及时发现并重新拉起。
5. 输出端幂等:整个容错体系中最容易被忽视的玩家
5.1 为什么说输出端决定最终一致性
到这里你会发现,无论我们选Receiver还是Direct模式,无论checkpoint和lineage多完善,最终都可能面临同一个问题:一批数据被重新计算了,而计算结果会再次写入目标系统。
如果你输出的目标是支持事务的关系型数据库,可以考虑用事务来控制:先delete当批次的数据范围,再insert。但如果你输出到HBase、ES、Redis这种非事务性系统,那事务完全不好使,你必须依靠“幂等写入”。
什么是幂等写入?简单说就是:同一份数据,不管写多少遍,最终在目标系统里的结果都和只写一遍相同。
5.2 实现幂等写出的三种常见方案
| 方案 | 适用场景 | 操作方式 | 局限 |
|---|---|---|---|
| 唯一键覆盖 | HBase、带主键的SQL表 | 以业务主键做upsert,重复写时直接覆盖旧值 | 主键设计要合理,否则会有脏数据 |
| 预批次删除 | SQL数据库、Hive分区表 | 按batch批次先删除该批次的输出区间,再重新写入 | 删除易误伤,事务边界要严格 |
| 幂等写入器 | ES(version/id)、HBase、Redis | 写入时使用稳定的批次唯一ID作为冲突判断 | ES有版本冲突,需开启乐观锁 |
这里重点说一下业内用得很普遍的唯一键覆盖方案。在做实时用户画像标签表时,我们用(userId, tagId)作为HBase的rowkey,每次写入都是put。就算一个用户的标签在恢复阶段被重复计算了三遍,每次put的最终效果是一样的——最后写入的值就是最新计算出的值,不会有多份数据。这就是典型幂等设计。
5.3 foreachRDD输出时的几个大坑
很多人在foreachRDD里写输出代码时,会犯一些想当然的错误,多亏我在生产环境里吃过亏,现在总结三条:
- 不要在foreachRDD内部创建连接:这会导致每条记录都建立一个连接,性能直接崩掉。正确做法是用
ForeachWriter,它在每个partition上创建并复用连接,同时对事务提交做了封装。 - 不要只往外部系统提交就完事:要确保输出动作和批次偏移提交具备联动。如果你在foreachRDD里输出成功后,下游系统收到了结果,但Spark这边还没把offset交出去,那么故障发生后这批数据会重跑,靠输出端幂等兜住;反之如果Spark把offset先交出去,下游系统还没写入成功,故障恢复后这批数据就真的丢了。所以绝不能让offset提交先行。
- 输出幂等不能只靠SQL唯一键:ES的写入默认是按版本号覆盖,如果恢复时重复计算的批次版本号比已存在的版本旧,会被丢弃,反而丢数据。要在写入时明确指定
version_type=external,强制使用外部传入的版本。
6. 状态管理与窗口运算:最容易出错的高危地带
6.1 状态算子与checkpoint的深度绑定
凡是用了updateStateByKey、mapWithState这类操作,任务就有了跨batch的状态。状态默认维护在Executor的内存里,可靠性的唯一来源就是checkpoint。
这里要区分一个细节:checkpoint状态数据和WAL是两码事。mapWithState虽然也会写state的快照,但它不是靠WAL,而是靠每个batch结束后的定时快照机制。如果你的任务用了状态算子,必须设置checkpoint目录,否则状态根本不会持久化。
很多人在测试环境跑任务一切正常,上了生产重启一次,状态清零,数据从0开始重新累加。原因很简单,测试环境跑得时间短,状态没有积累出影响力,生产环境一重启就原形毕露。
6.2 窗口操作的容错陷阱
窗口操作(比如滑动窗口计算近5分钟的量)天然涉及多个batch的数据聚合。故障恢复时,这些数据必须能回到窗口起点重新算。但这里有个你绝对要小心的地方:窗口长度越长,涉及的历史batch越多,checkpoint里需要持久化的元数据量越大,恢复时的计算量也就越大。
假设你做一个24小时的窗口,每分钟一个batch,一个窗口就要覆盖1440个batch。如果其中某几个batch的数据在Kafka中已经过期,窗口计算最终的结果就会少数据。而且窗口汇总结果不会主动告警,只能靠统计口径校验发现,是个隐蔽事故。
我们线上做过一次把窗口从1小时改为6小时的尝试,专门评估过这个风险,后面为了数据完整性还是放弃了。不是所有场景都适合用长窗口,如果业务上需要长时间窗口的准确汇总,建议用离线批处理做兜底校验,不要在流式窗口上铤而走险。
6.3 mapWithState vs updateStateByKey:容错视角的差异
有些人会问这两个状态算子选哪个。从功能上讲,mapWithState引入了超时机制(timeout),只输出变化的部分;updateStateByKey每次都会输出全量状态。但站在容错角度,真正的差异在于状态manageable的方式:
updateStateByKey的状态在executor端是一份粗粒度的整个RDD,发生变更时要把全量状态WAL到checkpoint,状态越大写盘开销越大。mapWithState内部维护了增量状态更新,checkpoint压力相对更好控制,且支持给状态设定TTL,避免状态无限膨胀。
生产场景,只要状态规模可能比较大,都建议选mapWithState。我见过有人在状态十几亿条时用updateStateByKey,每次checkpoint都卡几分钟,任务跟死了一样。后来换成mapWithState加合理超时,磁盘和CPU压力瞬间下来了。
7. 背压机制:容错体系中的“隐形稳定器”
7.1 背压到底在防什么
流处理任务最怕的就是消费速度大于处理速度。如果Kafka的写入劲头很大,Spark的处理追不上,数据会在内存里越堆越多,最后要么OOM,要么GC把整个Task卡死。不管哪种情况发生,结果都会导致批量数据来不及处理,随之而来的就是任务失败、检查点过期、甚至是不可恢复的丢失。
背压机制(Backpressure)就是干这个的——它通过动态调整Spark的接收速率,保证处理速度和消费速度匹配,从源头上避免OOM这种不可控故障的发生。
7.2 SparkStreaming背压的配置与经验值
SparkStreaming从1.5版本开始支持背压,核心参数是spark.streaming.backpressure.enabled=true。在这个开关下面,有两个更细的参数经常被忽略:
spark.streaming.backpressure.initialRate:任务启动时接收端允许达到的最大速率,如果Kafka里积压了大量数据,这个值设小了会拖慢接入速度,设大了又起不到保护作用。spark.streaming.kafka.maxRatePerPartition:每个Kafka Partition在每批次中最多能消费多少条记录。这个参数非常关键,它是你控制处理速度的“总闸”,背压机制调的就是它。
一个比较稳妥的配置思路:先根据Kafka分区数、数据平均大小、Executor核心数,估算一个理论吞吐上限,把maxRatePerPartition设为这个上限的70%左右,再开启背压让它自动微调。真实环境我见过不少同学不理解这个参数,把它设得极大,等于没有限速,然后一有大流量波动任务就挂。
7.3 背压和容错的关系
背压看起来和“故障恢复”没关系,但它其实是容错体系的第一道防线。OOM导致的任务崩溃,是所有故障中最难恢复的一种,因为它往往伴随着JVM进程死掉,Executor直接失联,状态数据可能来不及checkpoint。与其事后花大量时间恢复,不如提前用背压控制住摄入速度。
我建议每一个生产SparkStreaming任务都检查一下spark.streaming.backpressure.enabled是不是true,很多本地Demo默认是不开的。真正跑业务的时候必须开,你不开就是拿任务的稳定性去赌。
8. 生产环境常见故障与排查实录
8.1 问题快查表
我在线上踩坑多年,整理了这几个最常见的问题及排查思路,先呈现快速定位表,再逐一解释:
| 故障现象 | 可能原因 | 排查方向 |
|---|---|---|
| 重启后状态丢失 | 未开启checkpoint或checkpoint目录错误 | 检查StreamingContext是否设置checkpoint路径,确认恢复时用的路径一致 |
| 消费offset越界 | Kafka retention被调短,数据被清理 | 比对Kafka最早的offset与Spark记录的offset |
| 任务反复重启但恢复不了 | checkpoint中的数据损坏或版本不兼容 | 查看日志里的反序列化异常,必要时清理checkpoint并接受状态丢失 |
| 重复数据很多 | 输出未做幂等 | 检查目标系统的唯一键设计,跟进重复写入的具体主键 |
| 吞吐下降一半 | 打开了WAL或者背压参数设置太保守 | 检查SparkUI的输入速率与处理速率,动态调整参数 |
| 长窗口结果不准确 | 窗口跨度过大、中间有数据过期或状态过期 | 缩短窗口,或用离线计算兜底 |
这个表你可以直接截图保存,排查线上问题时对着看,能省不少时间。
8.2 一次真实的事故复盘:offset提交顺序引发的重复风暴
有一次我们做了一个到ES的实时推荐标签任务,用的是Direct模式,输出逻辑是:foreachRDD里先写ES,然后用stream.awaitTermination的逻辑定期提交offset。某次ES集群抖动,导致一批数据ES写入超时,部分成功部分失败,但Spark侧把offset已经提交了。恢复之后,那段offset范围被Kafka标记为已消费,那部分失败的数据就再也不会被读取,推荐标签缺了很多数据;同时,已经写入ES成功的数据,在任务重启后又被重新拉了一次,出现了覆盖旧版本的情况,用旧版本把新版本的标签顶掉了。
这个事故的原因就是输出与offset没有做到事务性联动。解决方式是我前面提到过的:确保输出成功后才提交offset,或者用ES的version参数配合外部批次ID来强制按批次覆盖。把这个机制做完之后,类似的重复和数据丢失就再没出现过。
排查过程中我们还发现一个问题:如果想要精确判断某批数据“输出成功”,就必须拿到一个明确的事务边界。但在ES里,批量写入成功和部分失败混在一起,连返回结果都很难判断。后来我们换了一种策略,先写到HBase(支持单行原子写入),再异步同步到ES,这样容错逻辑就大大简化了。
8.3 关于checkpoint目录的一个冷门坑
某个时间点我们为了做业务升级,把任务代码从Spark 2.3升到了Spark 2.4,然后重启任务,发现恢复失败。日志里报错是ClassCastException或NoSuchMethodError,指向的类都是我们自己的业务类。
排查了很久才发现,checkpoint目录里持久化了旧的DAG数据和业务对象序列化结果。新代码里改了一个字段名,旧的持久化数据没办法反序列化到新结构,恢复就失败了。
这里有个经验,如果你觉得代码里改了会影响到状态结构,重启时最好的做法是:做一个新的checkpoint目录并清空旧的状态,让任务结合Kafka里的offset从当前位点继续消费,宁可放弃一点历史状态,也别卡在恢复流程里动不了。当然,如果你改的只是计算逻辑,没有改状态schema,一般可以安全复用旧checkpoint。
8.4 监控告警配置建议
容错机制做得再好,如果故障发生时你完全不知道,那一切白搭。完整监控应该包括三块:
- 处理延迟监控:实时查看
Processing Time和Scheduling Delay,后者一旦持续走高,说明任务快要撑不住了。 - offset堆积监控:Kafka消费组的
Lag指标这个不用多说,一旦持续增长,不是加并行度就是调背压。 - checkpoint成功监控:发布一个新任务,关键看每个batch是否正常执行
checkpoint操作,如果interval设置不当导致checkpoint一直失败,状态就存在隐患。
告警群的命名也要规范,实时任务出问题的时候,第一反应谁都不想再去找这找那。
9. 实操配置参考:一套可直接落地的生产参数
最后给出一个我在生产环境使用过、被验证过比较稳妥的核心配置模板,你直接抄作业,但要根据自己集群规模调参:
val ssc = new StreamingContext(sparkConf, Seconds(5)) ssc.checkpoint("hdfs://nameservice/flink/checkpoint/your_app") // 关键:状态保留与幂等输出相关 sparkConf.set("spark.streaming.backpressure.enabled", "true") sparkConf.set("spark.streaming.backpressure.initialRate", "20000") sparkConf.set("spark.streaming.kafka.maxRatePerPartition", "2000") sparkConf.set("spark.streaming.stopGracefullyOnShutdown", "true") sparkConf.set("spark.task.maxFailures", "4") // 窗口和状态最大保留时间(避免状态膨胀) ssc.remember(Minutes(10))这里几个参数的意图我解释一下:
stopGracefullyOnShutdown=true保证优雅停机,停机时会处理完当前批次再退出,减少人为变更时的数据损失。task.maxFailures=4给短时的Executor故障留出重试空间,但也要配合Kafka保留时间来看。remember要大于最大窗口长度或状态依赖的历史批次范围,避免状态被清理导致窗口结果不完整。
关于checkpoint间隔,官方默认是batch间隔的5倍左右。设得太短会频繁写HDFS,产生大量小文件,白白消耗集群磁盘IO;设得太长,故障恢复时的状态丢失窗口变大。建议生产环境手动调整spark.streaming.checkpoint.interval,找一个你们能接受的平衡值。
写在最后的一点个人体会
容错这个东西,说到底是“你愿意在哪个环节付出代价”的取舍问题。WAL给Receiver模式兜底,但牺牲吞吐;checkpoint给状态兜底,但牺牲存储和恢复时间;幂等输出给重复数据兜底,但要求你在设计之初就考虑好目标系统的主键和事务模型。没有一种方案是零成本搞定所有故障的,你只能结合自己的业务需求,选一个主导方案,再把其他边角补上。
我个人跑了这么多年实时任务,最大的感受是:大多数容错事故都不是Spark本身不强大,而是我们对它的机制理解得不够精确。比如offset到底谁在提交、状态数据存在哪、重复写入由谁负责去重,这些问题想透了,容错就不是什么玄学。
如果你正准备上线一个SparkStreaming任务,我的建议是先做全链路故障演练:模拟一次Executor宕机、一次Driver重启、一次Kafka被清空,分别看看任务能不能恢复、数据最终对不对得上。把这些场景过一遍,比你读十篇原理文章都有用。