本节摘要:流作业 7×24 常驻,容错与性能是生存问题。本节讲清检查点与预写日志的恢复分工、exactly-once 语义的三个成立条件、背压机制对积压的自适应节流,以及流侧独有的一批调优旋钮。
流作业的故障恢复依赖两份档案,作用不同:
ssc = StreamingContext(sc, 5) ssc.checkpoint("hdfs-stream-ckpt") # 元数据 + 状态归档位置 ssc.sparkContext.setLogLevel("WARN") # 预写日志开关:可靠源已自带副本时可关闭,省一轮写放大 ssc.conf.set("spark.streaming.receiver.writeAheadLog.enable", "true")
巡检要点:从可靠且可重放的源(如消息队列带位点提交)读取时,数据恢复靠源重放,WAL 可以关;从不可重放源(自定义 TCP)读取时,WAL 是唯一的补课记录。开销换安全,按源的性质定。
"恰好处理一次"不是引擎单方面承诺,而是三方合力:
def dump(batch_rdd): if batch_rdd.isEmpty(): return ts = batch_rdd.name() # 用批次时间戳当事务批次号 sink.write(batch_rdd.collect(), batch_id=ts) # 幂等 upsert 或事务提交 counts.foreachRDD(dump)
第 3 条最容易漏。foreachRDD 里的普通 append 写入,重跑一次就多一份——恰好一次在输出关口退化成至少一次,账面上看不出来,对账时才爆炸。
生产事故的标准剧本:上游流量洪峰,批次处理时间超过批次间隔,积压滚雪球,作业被拖死。两道防线:
ssc.conf.set("spark.streaming.backpressure.enabled", "true") # 按处理能力动态限流接收速率 ssc.conf.set("spark.streaming.backpressure.initialRate", "1000") # 冷启动限速 ssc.conf.set("spark.streaming.kafka.maxRatePerPartition", "2000") # 每分区每秒最大拉取条数
背压让接收速率跟随处理能力自适应收缩,本质是"宁可上游积压在消息队列,不让引擎内部积压"——队列的存储比引擎的内存便宜得多。
| 旋钮 | 方向 | 巡检依据 |
|---|---|---|
| 批次间隔 | 处理时间稳定在间隔的三分之一内 | UI 批次耗时折线 |
| 并行度 | 接收块数或 Shuffle 后分区数覆盖核心数 | 空 Task 多则减,排队多则加 |
| 状态超时 | 键规模收敛 | 状态存储体积曲线 |
| 序列化 | 状态与 Shuffle 用紧凑序列化 | GC 时间占比 |
| 数据块复制 | 默认两副本防 Executor 单点 | 丢失块的恢复日志 |
💡 关键直觉:流作业的性能问题九成能在 UI 的"批次耗时 vs 批次间隔"这一张图上确诊。处理时间贴着间隔跑就是红灯,先减计算(增量窗口、状态瘦身),再动参数。
背景:风控实时规则引擎跑在 DStream 上,每周例行"杀 Driver"演练验证恢复能力。第一次演练:kill 后自动拉起,作业 40 秒恢复,但风控同学反馈恢复后的十分钟内命中量翻倍。操作:对账恢复窗口内的输出记录,与上游重放的消息 id 比对。结果:重复全部来自输出侧——恢复后检查点回退了一个批次,消息重放重算,而下游写入是纯 append,三方条件里的"输出幂等"缺失,恰好一次在关口退化。处置:给写入加批次事务号,恢复后先按事务号探测已提交批次再决定跳过,重复清零。第二次演练又暴露第二个缺口:恢复耗时 40 秒里有 25 秒花在重放 WAL——源本就是可重放的 Kafka,WAL 属于双保险白付。关闭 WAL 后恢复降到 12 秒。第三个缺口是冷启动:重启后第一批次没有限流,瞬间拉满积压又抖了一轮,配上 initialRate 后平稳。解读:容错演练的价值在于按"源可重放性、计算确定性、输出幂等"三条逐一验伤,三条各对应一次真实故障模式;变式:若源不可重放(自建 TCP),WAL 必须保留,恢复时长与数据安全的交换方向相反。
# 输出幂等的落地核心:事务批次号先探后写 def dump(batch_rdd): if batch_rdd.isEmpty(): return bid = batch_rdd.name() # 批次时间戳即事务号 if sink.has_committed(bid): # 恢复场景:已提交批次直接跳过 return sink.write_upsert(batch_rdd.collect(), batch_id=bid) sink.mark_committed(bid) # 提交记录与数据同事务落库