操作手册 · 大数据运维

待环境验证

Spark 流任务 checkpoint 与恢复边界

联查输入位点、checkpoint、状态与输出提交,设计有限恢复演练,识别 Kafka 起点、重复写入与状态兼容限制。

Apache Spark大数据
阅读导引 · 理解后再操作

这篇知识解决什么问题

本篇把源位点、checkpoint、状态和输出提交作为一个恢复单元,解释新查询与原查询重启的区别。有限文件示例只演示基础链路,更强的处理保证必须以源保留、状态兼容与目标端重试语义逐项验证。

恢复点必须与输出语义共同解释

checkpoint 记录进度和状态,目标系统仍可能在失败窗口接受重复写入。自定义 sink 需要把稳定批次身份与目标端原子提交或幂等机制配合,不能只靠进程成功退出证明端到端正确。

兼容性是复用 checkpoint 的前提

输入源、状态键和 schema 的变化可能不允许从原 checkpoint 恢复。变更前核对对应 Spark 版本的限制,再通过隔离演练验证;不能删除目录来绕过不兼容错误。

进入文章正文
关联架构图解6 个组件 · 点击展开

沿 Spark 作业图观察独立 checkpoint 存储、Driver 进度写入与 Executor 状态写入,再联系业务输入和输出提交路径,区分进程重启与查询状态恢复。参考图不证明 sink 事务、源保留期或 checkpoint 备份一致,实际恢复演练必须记录这些额外依赖。

查看场景架构与实施步骤
参考架构 · 非实时拓扑

Spark 作业诊断与流式恢复运维架构

从 Driver、Executor、Stage 与 SQL 计划定位 Spark 作业瓶颈,结合事件日志和 checkpoint 核验批处理与流式恢复边界。

  • 数据 / 请求
  • 控制 / 管理
  • 观测 / 查询

点击组件,在图下方查看职责;连线编号对应流向解读。小屏可横向滚动,或直接展开文字说明。

Spark 作业诊断与流式恢复运维架构:组件关系图从 Driver、Executor、Stage 与 SQL 计划定位 Spark 作业瓶颈,结合事件日志和 checkpoint 核验批处理与流式恢复边界。 Driver → 资源管理器:资源申请;资源管理器 → Executor:资源分配;Driver → Executor:Task 调度;Executor → 输入与输出存储:读取与提交;Driver → 事件日志与历史服务:事件持久化与读取;Driver → 流式持久恢复点:流式进度协调;Executor → 流式持久恢复点:有状态计算持久化。箭头说明见下方流向解读。
逻辑参考图,待环境验证。采用 Spark 3.5.7 的批处理与 Structured Streaming 微批模式。YARN、Kubernetes 和 Standalone 的资源管理不同,部署平台状态与 Spark 作业状态需分别检查。

Driver

控制 / 治理

构造作业与阶段并调度任务,需要独立观察 Driver 内存。

查看关联工具
全部组件职责 6 个组件
资源管理器
按所选部署平台分配 Executor 资源,三者不是同时必需。
Driver
构造作业与阶段并调度任务,需要独立观察 Driver 内存。工具介绍 Driver
Executor
执行计算任务,任务重试和丢失 Executor 应与 Stage 时间线对应。工具介绍 Executor
输入与输出存储
保存业务输入输出;外部提交语义与算子重算独立核验。
事件日志与历史服务
event log 提供已结束作业的可观测来源,History Server 不调度任务。工具介绍 事件日志与历史服务
流式持久恢复点
仅 Structured Streaming 分支使用,每个查询采用独立持久目录;包含进度及有状态查询所需状态,不能与业务输出或 event log 混用。
流向解读 7 条连接
  1. 1

    Driver 资源管理器

    控制 / 管理 · 资源申请

    Driver 与部署层协作获得执行资源,具体生命周期取决于运行模式。

  2. 2

    资源管理器 Executor

    控制 / 管理 · 资源分配

    资源管理器提供运行 Executor 的资源。

  3. 3

    Driver Executor

    控制 / 管理 · Task 调度

    Driver 调度 task 并跟踪执行与重试。

  4. 4

    Executor 输入与输出存储

    数据 / 请求 · 读取与提交

    任务访问业务输入输出;重算安全性需看目标提交协议。

  5. 5

    Driver 事件日志与历史服务

    观测 / 查询 · 事件持久化与读取

    Driver 写 event log,History Server 读取;此节点合并表示日志与历史查看服务。

  6. 6

    Driver 流式持久恢复点

    控制 / 管理 · 流式进度协调

    Driver 协调微批进度和提交记录,恢复时需与同一查询及输入输出语义相容。

  7. 7

    Executor 流式持久恢复点

    数据 / 请求 · 有状态计算持久化

    执行侧保存有状态查询所需状态;状态后端与检查点存储插件须支持当前部署。

故障域与操作边界

适用版本与部署模式

采用 Spark 3.5.7 的批处理与 Structured Streaming 微批模式。YARN、Kubernetes 和 Standalone 的资源管理不同,部署平台状态与 Spark 作业状态需分别检查。

数据与变更边界

不通过反复重提作业、清空 checkpoint、全量 collect 或扩大 Driver 内存代替诊断。批任务重跑与流式恢复必须考虑已提交输出和外部副作用。

架构依据与版本核对 2 篇官方资料

图解是本站基于官方资料整理的逻辑参考;实施前仍需核对实际部署版本、组件支持范围与变更审批。

适用范围与恢复目标

本文以 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 架构

参考资料

从现象到判断

先收集证据,再缩小范围。以下是判读路径,不代表已经确认根因或获准变更。

  1. 修改 startingOffsets 后重启 Kafka 查询,实际读取位置仍从原进度继续。

    只读核对
    核对是否复用原 checkpoint、查询身份与源分区,区分首次启动和已有进度恢复。
    如何判读
    该参数不是已有查询的重置开关;需要重建查询时必须另行定义起点与重叠输出。
  2. 任务成功启动但持续积压,状态规模上升,事件时间水位推进不符合预期。

    只读核对
    联查源输入、批次耗时、watermark、状态行数与业务事件时间分布。
    如何判读
    需要区分限流、事件时间异常和状态逻辑;不能只缩短迟到窗口换取更小状态。
  3. 失败重试后目标端出现重复记录,checkpoint 却继续向前推进。

    只读核对
    按事件键和 batchId 对齐目标提交与重试窗口,核对去重登记是否与数据共同提交。
    如何判读
    至少一次 sink 或非原子登记可能解释重复;进度推进不能自动证明输出幂等。
常见误区与判断边界 2 项

发现恢复错误就删除 checkpoint

这会丢失状态与原始进度证据,并可能形成不同起点的新查询。先定位源、存储、状态和输出哪个层次失败,再决定明确的重建与回补方案。

把每批文件上限当作整个演练上限

maxFilesPerTrigger 只控制单批,availableNow 仍会处理本轮可用输入。有限演练必须预先固定总文件与记录数量,输出和 checkpoint 也应是独立测试路径。

交接时应留下的证据

作为记录提纲使用,不是自动检查结果;未取得的证据应标记缺口,并注明负责人。

  • 记录精确版本、应用修订、查询身份、源分区、目标和 checkpoint 路径及其权限、容量与保留条件。
  • 保存首个失败批次的输入范围、状态算子、存储错误和目标提交记录,保留结果不明的窗口。
  • 归档有限样本的事件覆盖、重复、窗口与迟到结果,以及正常恢复和故障窗口的分别验收。
  • 交接源已过期区间、状态兼容限制、实际恢复耗时和新旧查询输出归属,保持未完成项可追踪。

记录需包含环境、版本、时间与时区;分享前脱敏,不附访问令牌、密码或完整业务敏感数据。

继续阅读与资料核对

补充相关主题,再结合当前环境的实施记录形成结论。

返回原理导读

DOUYA OPS ECOSYSTEM

贡献你的经验,帮助更多运维人

把故障复盘、标准流程和最佳实践沉淀为可检索、可复用的知识内容。