在 Spark Streaming 中,系统通过哪些具体机制来保证 Fault Tolerance(容错性)?请描述其实现原理。
考察说明
考查对 Spark Streaming 容错机制的深入理解,包括数据可靠性保障和计算恢复策略。
回答思路
- 【回答框架 1】Spark Streaming 通过基于 RDD 的 lineage(血统)机制实现容错,每个批次数据对应一个 RDD,若计算节点故障,可重放上游数据并重新执行操作以恢复状态。
- 【回答框架 2】输入数据默认从可靠数据源(如 HDFS)读取,或通过预写日志 WAL 将接收到的数据持久化,避免 Driver 或 Executor 故障导致数据丢失。
- 【回答框架 3】对有状态计算,使用 updateStateByKey 或 mapWithState 时,通过 checkpoint 机制定期保存状态和元数据到可靠存储,故障后可从最近 checkpoint 恢复并重放未完成数据。
- 【回答框架 4】在 exactly-once 语义实现上,结合数据源的可重放性(如 Kafka offset)和输出操作的事务性(如幂等写入或事务 API),确保结果一致性。
- 【回答框架 5】容错不仅依赖数据冗余,还需考虑 Driver 故障恢复,通过 checkpoint 元数据实现 Driver 重启后的应用恢复。
- 【关键点 1】lineage 机制实现基于 RDD 的容错重算。
- 【关键点 2】WAL 和可靠数据源保障输入数据不丢失。
- 【关键点 3】checkpoint 保存状态和元数据,支持故障恢复。
- 【关键点 4】实现 exactly-once 需结合可重放源和数据输出幂等。
- 【易错点 1】不要忽略 Driver 故障恢复,仅靠 Executor 重算不足。
- 【易错点 2】默认配置下某些输入源可能丢数据,需显式启用 WAL。
- 【易错点 3】有状态计算需要定期 checkpoint,否则状态恢复耗时过长。