切换主题
Worker 持久化与恢复
问题:收到命令后进程马上重启
执行器收到启动命令便回 ACK,然后在内存里启动协程。ACK 成功后协调者可能不再重发,但执行器在持久化前崩溃,工作丢失。反过来,执行已经开始但 ACK 丢失,重发又可能启动两份进程。
正确的顺序要围绕可恢复事实设计:先保存命令及身份,再确认接收;重复命令识别为同一动作。执行结果也先保存,再交给独立发送器可靠上报。
原理:本地库是恢复依据
命令 Inbox 记录接收与处理状态,Attempt 记录执行生命周期,Checkpoint 记录能力可继续的位置,事件 Outbox 保存尚未上报的事实。内存中的 goroutine、channel 和 map 只是当前进程的工作状态。
重启扫描时不能默认重跑所有 RUNNING。不同能力的恢复语义不同:可以从文件偏移续传、可以查询远端任务、可以安全重算、或者必须进入待人工确认。
最小例子:两段可靠边界
text
收到 command_id
→ 本地事务保存 Inbox 和待执行 Attempt
→ ACK 接收
→ 调用能力(可能失败或崩溃)
→ 本地事务保存结果/Checkpoint 和 event_id
→ 发送 event_id
→ Server 持久接收后,本地标记已上报文件系统与 SQLite 也不是同一个事务。如果文件已写、Checkpoint 未更新,要能通过临时路径、校验和或原子 rename 识别有效产物,不能默认全部已经一致。
方案比较:恢复粒度按能力决定
| 能力 | 恢复策略 | 必须保存的证据 |
|---|---|---|
| 纯计算 | 用稳定输入重新计算 | 输入版本与输出身份 |
| 大文件传输 | 续传或安全重新传输 | 偏移、校验与传输身份 |
| 远端长期任务 | 先查询再决定 | 远端 job/session ID |
| 不可查询外部副作用 | 幂等键或人工对账 | 请求身份、未知状态 |
Checkpoint 只说明记录到哪一步,不能证明外部系统没有多做一步。恢复策略必须解释这个窗口。
真实案例:单发送器与重试分类
Worker Runtime 持有 LocalStore、ControlPlane、CommandStream、CapabilityExecutor 等端口。执行协程保存事件并通过 channel 唤醒发送器,可靠数据留在本地库中。
sendEvents 独占当前会话的发送,网络失败退避,延迟上限 30 秒。只有标为远程投递失败且非认证/授权错误的故障可重试;本地数据库错误不会被伪装成网络重试无限循环。
已核对的实现 · 本地代码快照
Worker 的依赖与运行状态 backend · c6dbe05c
worker/internal/core/execution/application/runtime.go · 第 48–100 行
符号:Runtime · 核对日期 2026-10-02
来源与提交版本一致
// Runtime 是 Worker 的长期运行应用:注册会话、恢复本地状态、消费命令并可靠上报事件。
type Runtime struct {
registration Registration
control ControlPlane
commands CommandStream
store LocalStore
capabilityExecutor CapabilityExecutor
ids executionports.IDGenerator
observer executionports.Observation
now func() time.Time
retryInitial time.Duration
outboxWake chan struct{} // 合并事件发送唤醒,可靠数据始终保存在本地数据库。
outboxDrain chan chan error // 请求发送器完成一轮发送并返回结果。
admissionMu sync.Mutex // 串行化本地容量检查与尝试接收。
activeMu sync.Mutex // 保护正在执行的协程及其取消函数。
active map[string]context.CancelFunc // 本进程持有的运行中尝试。
canceled map[string]bool // 记录当前会话已收到的取消。
draining bool // 当前进程的排空状态,持久标记用于重启恢复。
resourceInventory func() executiondomain.NodeInventory
resourceUsage func(context.Context) executiondomain.ResourceUsage
}
func New(registration Registration, control ControlPlane, commands CommandStream, store LocalStore, capabilityExecutor CapabilityExecutor, ids executionports.IDGenerator, observer executionports.Observation) (*Runtime, error) {
if registration.InstanceID == "" || registration.Name == "" || registration.NodeID == "" {
return nil, executiondomain.Wrap(
executiondomain.Invalid,
"Worker 注册信息缺少实例、名称或 Node",
errors.New("worker instance ID, name, and node ID are required"),
)
}
if registration.Inventory.CPU.CapacityMillis < 1 {
return nil, executiondomain.Wrap(
executiondomain.Invalid,
"Worker CPU 资源采集不可用",
errors.New("worker CPU inventory is required"),
)
}
if control == nil || commands == nil || store == nil || capabilityExecutor == nil || ids == nil || observer == nil {
return nil, errors.New("control plane, command stream, local store, capability, ID generator, and observer are required")
}
return &Runtime{
registration: registration,
control: control,
commands: commands,
store: store,
capabilityExecutor: capabilityExecutor,
ids: ids,
observer: observer,
now: time.Now,
retryInitial: time.Second,
outboxWake: make(chan struct{}, 1),
outboxDrain: make(chan chan error),
active: make(map[string]context.CancelFunc),片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
单发送器及网络重试 backend · c6dbe05c
worker/internal/core/execution/application/outbox.go · 第 34–87 行
符号:sendEvents · 核对日期 2026-10-02
来源与提交版本一致
// sendEvents 独占当前会话的事件发送;能力协程只负责持久化并唤醒发送器。
// 临时网络失败保留事件重试,不取消其他运行中的尝试。
func (runtime *Runtime) sendEvents(ctx context.Context, session Session) error {
delay := time.Second
timer := time.NewTimer(0)
defer timer.Stop()
for {
var drained chan error
select {
case <-ctx.Done():
return nil
case <-runtime.outboxWake:
case drained = <-runtime.outboxDrain:
case <-timer.C:
}
err := runtime.flushEvents(ctx, session)
if drained != nil {
drained <- err
}
if err == nil {
delay = time.Second
timer.Reset(time.Second)
continue
}
if ctx.Err() != nil {
return nil
}
if !retryableDelivery(err) {
return err
}
runtime.observer.Event(ctx, "event delivery deferred", err)
// 退避期间不响应生产者唤醒;新增事件仍保留在本地数据库中。
retry := time.NewTimer(delay)
select {
case <-ctx.Done():
retry.Stop()
return nil
case <-retry.C:
}
timer.Reset(0)
delay = min(delay*2, 30*time.Second)
}
}
// 仅标记远程投递失败;本地数据库错误不能被当作网络故障无限重试。
type deliveryError struct{ error }
func (e *deliveryError) Unwrap() error { return e.error }
func retryableDelivery(err error) bool {
var delivery *deliveryError
return errors.As(err, &delivery) && domain.Kind(err) != domain.Unauthenticated && domain.Kind(err) != domain.Forbidden
}
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
Capability 可替换执行契约 backend · c6dbe05c
worker/internal/core/execution/ports/outbound/capability_executor.go · 第 9–27 行
符号:CapabilityExecutor · 核对日期 2026-10-02
来源与提交版本一致
// CapabilityExecutor 将通用执行命令委托给具体 Capability。
type CapabilityExecutor interface {
// Execute 根据命令中的能力名称和版本执行,并通过 reporter 上报进度。
Execute(context.Context, domain.Command, ProgressReporter) (*domain.AttemptEvent, error)
}
// Progress 是能力执行过程中上报给控制面的结构化进度。
type Progress struct {
Stage string `json:"stage"` // 当前处理阶段名称。
Status string `json:"stage_status"` // 阶段状态。
Progress int `json:"progress"` // 整体 0 到 100 的进度。
Completed int `json:"completed"` // 当前阶段已完成条目数。
Total int `json:"total"` // 当前阶段总条目数。
Failed int `json:"failed"` // 当前阶段失败条目数。
}
// ProgressReporter 将能力进度转换为本地检查点和控制面尝试心跳。
type ProgressReporter func(context.Context, Progress) error
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
SSE 连接只负责命令传输,PostgreSQL worker_commands 与 Worker SQLite Inbox 提供事实。删除 SQLite 不等于清缓存,而是放弃身份和未同步恢复状态。
失败与边界:文件、进程和远端服务
ZIP 清理要拒绝路径穿越、限制展开体积并定义产物提交方式;命令执行要保存进程或外部任务身份,不能因为新进程查不到旧 goroutine 就推断没有任务在运行。
人工 drain 与自动恢复是不同状态。容量缩减或能力变化时,应防止新任务进入,同时处理已开始的 Attempt。清理临时文件也要保留尚有恢复价值的数据。
迁移练习与参考答案
练习:执行器调用云渲染 API,创建成功但还未保存 job ID 就崩溃。如何避免重复渲染?
参考答案:创建前持久化稳定请求 ID,远端按该 ID 幂等创建或提供查找接口;重启查询后保存 job ID 再继续。若远端不支持这些能力,结果处于未知状态,需要对账或人工处理,不能仅靠“先写本地 RUNNING”保证不重复。