请阐述 Spark Streaming 在实时数据流处理中,如何应对故障恢复与任务重启?
考察说明
考查对 Spark Streaming 容错机制(如 checkpoint、RDD 血缘)及任务重启流程的理解。
回答思路
- 【回答框架 1】Spark Streaming 依赖 Spark 核心的 RDD 血缘(lineage)实现容错,每个批次数据对应一个 RDD 集合,若任务失败可基于血缘重算。
- 【回答框架 2】同时提供 checkpoint 机制,分为元数据 checkpoint(如配置、DStream 操作)和数据 checkpoint(如 RDD 数据),用于在 Driver 故障时恢复流计算状态。
- 【回答框架 3】具体恢复流程:Driver 重启后从 checkpoint 目录加载元数据和上次保存的偏移量,然后重建 DStream 的依赖关系,并重新计算未完成的批次。
- 【回答框架 4】对于 Receiver 方式,如果 Worker 故障,存储在 executor 中的已接收数据可能丢失,需配合 WAL(预写日志)或 Kafka 等外部存储来保证数据不丢,实现 at-least-once 语义。
- 【关键点 1】Spark Streaming 核心容错是 RDD 血缘与 checkpoint 结合。
- 【关键点 2】Driver 恢复需要 checkpoint 保存元数据和偏移量,任务重算基于血缘。
- 【关键点 3】Receiver 方式需 WAL 或外部存储保证数据不丢,但可能造成重复处理。
- 【易错点 1】不要忽略 checkpoint 的写开销和频繁 checkpoint 对性能的影响。
- 【易错点 2】不要误以为所有场景都能做到 exactly-once,默认是 at-least-once,需要幂等或事务性写入配合。
- 【易错点 3】不要只依赖 RDD 血缘,因为依赖链过长时重算代价很高。