← blog 专题系列 · 2026-04-11

RL-Scheduler 中的 Master-Worker 机制可用性分析

Master,Worker维持任务的最终一致性分析

4 min read

Master-Worker 可用性评估(修正版)


一、3 个 Redis Key 的作用

任务分派时,Master 通过 tryDispatchQueuedTaskToWorker() 依次设置:

Key设置时机TTL作用
worker:{id}:task=taskItryPreemptWorker() Lua 脚本原子写入120sWorker 正在处理哪个任务
task:{taskId}:workerId=workerIdregisterTaskOwner()120s任务被哪个 Worker 认领
worker:{id}:hbWorker 心跳时续期30sWorker 是否存活

二、分派时序(正常流程)

  1. tryPreemptWorker() → 原子写入 worker:{id}:task = taskId (TTL 120s)
  2. registerTaskOwner() → 写入 task:{taskId}:workerId = workerId (TTL 120s)
  3. dispatchTask() → 通过 Netty 发给 Worker
  4. DB status → PENDING 更新为 RUNNING

三、Worker 宕机场景

graph TD
    A["Worker 崩溃"]
    B["不再发送心跳"]
    C["worker:{id}:hb 在 30s 后过期"]
    D["RunningTaskRecovery (每5s) 发现:<br/>taskOwnerKey 存在 ✓ / hb 不存在 → Worker 已死 / task 存在 → 任务孤立"]
    E["标记 PENDING + releaseTaskOwner()"]
    F["PendingTaskReconciler (每2s) 发现 idle worker → 重新分发"]
    A --> B --> C --> D --> E --> F
    note["恢复延迟:30s (心跳TTL) + 5s (扫描间隔) ≈ 35s"]

四、Master 宕机场景(核心)

场景 A:Master 在步骤 1-2 之间宕机(未设置 taskOwnerKey)

  宕机前只设置了 worker:{id}:task
  DB 状态 = RUNNING(已在 dispatchTask 前更新)
  重启后 reconstructWorkerTasksFromRedis():
    - 发现 worker:{id}:task = taskId
    - 发现 task:{taskId}:workerId 不存在或不匹配
    → 标记 PENDING + 入队

场景 B:Master 在步骤 2-3 之间宕机(已设置 taskOwnerKey 但未收到 Worker 响应)

  宕机前设置了 task:{taskId}:workerId + worker:{id}:task
  但 DB status = RUNNING 未变(dispatchTask 发送失败或未处理)
  重启后:
    - worker:{id}:task 存在
    - task:{taskId}:workerId 存在且匹配
    - RunningTaskRecovery: hb 存在则跳过(Worker 还在跑或刚完成)
    - PendingTaskReconciler: worker:{id}:task 存在则不视为 idle
  最终由 Worker 心跳或 checkAndFixStaleRunningTasks() 清理

场景 C:Master 在设置 taskOwnerKey 之前宕机(只设置了 Lua 脚本的 worker:{id}:task)

  宕机前只设置了 worker:{id}:task
  task:{taskId}:workerId 不存在
  DB status = RUNNING(dispatchTask 前已更新)
  重启后 reconstructWorkerTasksFromRedis():
    - worker:{id}:task 存在
    - task:{taskId}:workerId 不存在或不匹配
    → 标记 PENDING + 入队

五、关键机制总结

组件职责扫描频率
RunningTaskRecovery检测 Worker 宕机导致的任务孤立每 5s
reconstructWorkerTasksFromRedisMaster 重启时恢复被中断的分派@PostConstruct 一次
PendingTaskReconciler将 PENDING 任务重新入队给 idle worker每 2s
checkAndFixStaleRunningTasks心跳时检测 通知丢失导致的任务卡住每次心跳
心跳续期renewTaskOwnerTtl() 保持 taskOwnerKey 不过期每次心跳

六、3 Key 的最终一致性保证

无论 Worker 宕机还是 Master 宕机:

  Worker 宕机:
    worker:{id}:hb 过期 → RunningTaskRecovery 检测 → PENDING + 重新入队

  Master 宕机(未收到完成通知):
    1. 重启时 reconstructWorkerTasksFromRedis() 修正不一致状态
    2. 或者 Worker 心跳时 checkAndFixStaleRunningTasks() 发现 taskKey 没了但 ownerKey 还在 → 标记
  COMPLETED

七、当前设计的限制

  1. 故障检测延迟:Worker 宕机后约 35s 才能发现(30s TTL + 5s 扫描)
  2. Master 重启恢复:依赖 @PostConstruct 一次性重建,不保证 Master 运行时实时修复所有不一致
  3. channelInactive() 有 TODO 未完成:MasterHandler.java:59-63 未触发恢复逻辑,但依赖心跳 TTL 兜底
  4. 无真正零宕机切换:Master 宕机期间任务状态更新请求会丢失,依赖 Worker 侧重试或最终一致性

Sources

No external sources for this entry.

Related