1299 字
约 4 分钟
1
Spark 执行与任务生命周期管理

Spark 执行与任务生命周期管理

Spark 执行与任务生命周期管理1. 要解决的问题Spark Job 是跨进程、长耗时、有外部 application id 的任务。一次 spark-submit 调用不能解决状态持久化、失败重试、服务重启恢复、取消和审计问题。本项目先建立统一 Job 领域模型,再把提交器和存储实现放到端...

字数: 1206 | 语雀原文


Spark 执行与任务生命周期管理

1. 要解决的问题

Spark Job 是跨进程、长耗时、有外部 application id 的任务。一次 spark-submit 调用不能解决状态持久化、失败重试、服务重启恢复、取消和审计问题。本项目先建立统一 Job 领域模型,再把提交器和存储实现放到端口后面。

2. 领域模型与端口

JobRecord

字段 作用
run_id 关联一次 Agent 运行
job_id 系统内部 Job 主键
job_type distributed query、pipeline、graph 或 vector
status 当前领域状态
payload 提交参数、SQL、数据集和 Spark 配置
attempts 每次提交尝试的历史
artifacts 与 Job 相关的产物路径

JobAttempt

每次尝试保存 attempt 序号、状态、外部 ID、错误信息和 metrics。重试不会覆盖首次失败证据。

两个端口

Diagram Python Protocol 在运行时不强制继承,图中的实现关系表达结构化契约,而不是显式基类继承。

3. 状态模型

领域枚举定义了 createdqueuedsubmittedrunningsucceededfailedretryingcancelled。当前 JobOrchestrator 实际主动产生的状态比枚举少:它创建 created,持久化 queued,然后直接采用 submitter 返回的状态;retrying 只在命中可重试失败时出现。

Diagram 虚线意义上的“外部状态同步”和取消在模型中已表达,但当前尚无完整 watcher/cancel API,面试中不能说已经实现。

4. 提交与重试链路

Diagram

重试判断

_can_retry() 同时检查:

  • attempt 数量未超过 max_retries
  • 存在 latest attempt;
  • message 中包含 oomshufflefetch failedtimeouttransient。 这是可解释的最小实现,但依赖错误字符串,不如结构化错误码稳定,也没有指数退避和 jitter。

5. SQLite 持久化

表结构只有一张 jobs

含义
job_id 主键
run_id Agent 运行关联键
status 便于直接查询的冗余状态
payload_json 完整 JobRecord.model_dump_json()
updated_at 默认更新时间字段

save() 使用 INSERT ... ON CONFLICT(job_id) DO UPDATE。这样同一个 Job 多次状态更新不会新增行。完整领域对象保存在 JSON 中,开发简单,但缺点是很难用 SQL 高效查询 attempt、external_id 和 metrics;生产数据库应把高频查询字段正规化。

注意:当前 upsert 更新语句没有显式更新 updated_at,所以默认值只在首次插入时产生。这是现有实现的一个细节缺口。

6. 执行后端选择

Diagram 这里有两个维度:JobSubmitterPort 管理 Job 领域提交记录;DockerSparkRunner/SparkExecution 负责真实计算。当前两者还不是完全统一的异步状态源,生产化时需要以 external_id 为中心收敛。

7. 一致性与恢复难点

幂等性

当前 job_id 使用随机 UUID 片段,同一请求重复调用会创建两个 Job。需要增加业务幂等键,例如 tenant + task_type + biz_date + request_hash,并在数据库建立唯一约束。

提交与落库双写

外部 Spark 提交成功但最终 save() 失败时,可能出现“外部作业在跑、内部无记录”。可采用 transactional outbox:先持久化待提交事件,再由 worker 提交外部系统并回写 external_id。

多实例并发

SQLite 适合单机,不适合多个 worker 抢任务。生产实现可用 PostgreSQL 的状态 version 做乐观锁,或 FOR UPDATE SKIP LOCKED 抢占 queued Job,确保同一 Job 只有一个执行者。

重启恢复

启动时扫描 submitted/running/retrying Job,根据 external_id 查询 Spark/Livy 状态。没有 external_id 的中间状态要结合 outbox 判断是否重提,不能直接把所有 running 标成 failed。

取消

领域状态已定义 cancelled,但取消必须贯穿 UI、Orchestrator、submitter 和 Spark 后端;还要处理“取消请求与成功回调同时到达”的竞态。

8. 源码索引

  • sparkos/domain/job.py:状态、JobRecord 和 JobAttempt。
  • sparkos/application/job_orchestrator.py:状态流转与重试。
  • sparkos/infrastructure/persistence/sqlite_job_store.py:表结构和 upsert。
  • sparkos/infrastructure/spark/docker_runner.py:容器执行。
  • sparkos/infrastructure/spark/livy_submitter.py:远端提交适配器。
  • sparkos/application/spark_execution.py:计算执行封装。
  • tests/test_workbench.py::test_job_orchestrator_persists_history:持久化测试。

9. 如何向面试官概括

我先定义 Job 和 attempt 的领域模型,再用 Store/Submitter Protocol 隔离 SQLite 与 Livy。Orchestrator 持久化 queued 和最终状态,并对 OOM、Shuffle、超时等失败追加重试 attempt。当前是单机原型,生产化的关键不是换数据库这么简单,而是幂等键、外部提交双写、状态并发控制和重启恢复。

Spark 执行与任务生命周期管理
http://www.clxhxhhr.top/posts/782/
作者
clxstart
发布于
2026-09-18
许可协议
CC BY-NC-SA 4.0
评论
0 条
还没有评论,先写一条吧。