适用范围与恢复目标
本文以 Apache Spark 3.5.7 Structured Streaming 的 micro-batch 模式为基线,讨论 checkpoint、输入位点、状态和输出提交之间的恢复关系。旧 DStream、Continuous Processing 及异步进度追踪不在本文示例范围;升级至其他版本时应重新核对恢复兼容性。本文没有启动流任务或修改 checkpoint,验证状态为 PENDING。
恢复目标要同时说明“从哪里继续读”“历史状态是否完整”“重试会不会重复写”和“多久恢复业务可见性”。checkpoint 存在只证明找到了一份恢复输入,不能单独证明端到端 exactly-once。先记录 Spark 和连接器版本、应用修订、查询标识、源与目标、checkpoint 路径及保留策略。
把位点、状态与输出看成一个恢复单元
Structured Streaming 使用 checkpoint 保存恢复所需的进度和状态信息;有状态计算还依赖与查询结构相容的状态数据。生产 checkpoint 应位于适合引擎访问且具有持久保障的存储中,每个独立查询使用独立路径,权限和容量也需纳入运维。
输出是否可重复执行取决于 sink。foreachBatch 默认提供至少一次写入,需要使用 batchId 等机制配合目标端原子提交或去重,才能实现相应的更强保证。不能把“先查批次是否存在,再写数据,再标记完成”的三个非原子步骤当成已解决崩溃窗口。基础机制见 Structured Streaming 编程指南。
首次启动与 checkpoint 恢复不同
以 Kafka 为源时,startingOffsets 用于新查询起点;已有 checkpoint 的恢复会继续此前的进度,修改该参数不能被当作重置历史消费的办法。Kafka 日志保留期还必须覆盖停机与重放窗口:checkpoint 引用的位点若已被源删除,恢复无法凭空找回记录。
诊断时比较 checkpoint 对应查询、Kafka topic/partition 和保留边界,保留缺失区间;不要直接关闭数据丢失检查来制造绿色运行状态。Kafka sink 可能产生重复输出,业务去重还需设计稳定事件键。相关语义依据 Spark 3.5.7 Kafka 集成指南。
用有限文件演练创建与恢复
以下仅供隔离学习环境:输入是一份预先准备好的不可变 JSON 样本目录,总计不超过五个文件、一千条记录,字段与 schema 一致;输出和 checkpoint 是两个独立的新测试目录。示例会写出文件,本文没有执行它。
source = "hdfs://REPLACE_WITH_NAMESERVICE/REPLACE_WITH_FIXED_JSON_SAMPLE"
output = "hdfs://REPLACE_WITH_NAMESERVICE/REPLACE_WITH_NEW_TEST_OUTPUT"
checkpoint = "hdfs://REPLACE_WITH_NAMESERVICE/REPLACE_WITH_NEW_TEST_CHECKPOINT"
events = (spark.readStream
.schema("event_id STRING, event_time TIMESTAMP, amount DECIMAL(18,2)")
.option("maxFilesPerTrigger", 1)
.json(source))
query = (events.writeStream.format("parquet").outputMode("append")
.option("path", output)
.option("checkpointLocation", checkpoint)
.trigger(availableNow=True).start())
query.awaitTermination()
print(query.lastProgress)availableNow 处理本轮可用数据后停止,maxFilesPerTrigger 限制每批读取的文件数量,并不是全次查询的数据上限;真正的总量边界来自预先固定的样本。lastProgress 是最近一个触发批次的进度,不等于所有批次记录数之和;尚无进度时也可能为空。样本仅演示无状态文件输入输出,不证明 Kafka、状态恢复或自定义 sink 已通过验证。
失败时先保留现场再分层诊断
- 输入:需要保存的证据:源分区、位点、权限与保留窗口;候选问题:数据已过期、订阅变化、读取失败。
- checkpoint:需要保存的证据:查询标识、路径、存储错误;候选问题:路径错误、权限、容量或存储故障。
- 状态计算:需要保存的证据:状态行数、状态算子、版本与 schema;候选问题:状态增长、恢复不兼容、资源压力。
- 输出:需要保存的证据:batchId、目标提交与重复事件;候选问题:部分成功、重试重复、提交阻塞。
先解释首个失败批次,再分析后续重试。不能只因最新日志出现恢复错误就删除 checkpoint;那会改变查询的恢复起点并丢失重要证据。处理速度、输入速度与批次耗时需要在相同窗口比较,输入速率降低也可能来自限流或源侧问题。
水位线与状态增长的判断边界
对有状态聚合和 Join,明确事件时间字段、窗口、允许迟到的业务约定以及输出模式,再查看 watermark 和状态指标。水位线用于状态与迟到数据处理,不是可靠的墙上时钟倒计时;源缺乏新事件或多输入推进差异都会影响解释。不能看到状态很大就任意缩短迟到容忍窗口,否则可能改变业务结果。
诊断先区分业务键持续增加、事件时间异常、算子本来需要长历史,以及水位线未按预期推进。把状态总行数、更新行数和输入业务分布结合起来,选取已知迟到与边界时间样本解释结果。不要用“内存降下来了”作为有状态计算正确性的唯一证据。
变更后复用 checkpoint 的兼容性
同一 checkpoint 不能随意改变输入源数量或类型,也不能随意改变状态算子的分组键、去重列和状态 schema。过滤或投影等变更虽然有部分允许情形,仍需检查输出兼容与业务含义。对 file sink 更换输出路径也不能当作无条件兼容的变更。详细限制应逐项对照对应版本的恢复语义章节。
升级或逻辑变更之前保留应用修订、连接器包与依赖配置,在隔离环境使用协调一致的恢复输入演练。不要随意复制一个正在变化的 checkpoint 目录并声称它是已验证备份,也不要让两个活跃查询共用同一 checkpoint。存在不兼容变更时,应设计新查询、明确起点与回补区间,并处理与旧输出的重叠。
恢复演练如何证明结果正确
在前述有限样本中,先记录输入事件 ID 清单与输出行数,再以相同逻辑、相同源目标和同一 checkpoint 重启,检查已经处理的输入是否被错误重复写出。进一步的故障注入应在独立测试运行中覆盖“读取后、目标提交前”和“目标已提交、进度确认前”等窗口,不能仅测试正常结束后的重启。
对有状态任务另外检查窗口汇总、迟到事件、重复 ID 和边界时间;对自定义 sink 检查数据写入与批次去重登记能否共同提交。重建新 checkpoint 后 batchId 的编号可能重新开始,因此去重键需要包含稳定的查询或业务身份,不能假设 batchId 在所有查询间全局唯一。
验收、停止与回退边界
验收应包含输入覆盖区间、checkpoint 身份、状态恢复结果、目标端无意外缺失或重复、积压收敛与实际恢复耗时。若源保留期不足、状态不兼容或目标提交结果不明,停止反复重启并先确认数据范围。成功启动只说明进程开始工作,不等于历史结果已恢复。
旧应用能否回退取决于 checkpoint 与状态是否仍兼容;输出系统已写入的数据不会因换回旧代码而自动撤销。新查询已经接管时,回切前需处理新增输出与重叠位点,避免两个查询同时写同一业务结果。checkpoint 在 HDFS 的保护可结合 Hadoop 恢复规划,计算执行基础见 Spark 架构。