Skip to main content

任务状态同步系统生产设计:全链路、状态机、实时通知与高可用

Rainy
雨落无声,代码成诗 —— 致力于技术与艺术的极致平衡
Rainy
66 MIN READ... VIEWS

真正可靠的任务看板,不是“WebSocket 一直在线”,而是:状态先成为可恢复的事实,通知只负责让事实更快到达,任何一层短暂故障后都能重新收敛到正确状态。

文件解析、AI 生成、视频转码、报表导出这类异步任务,通常会经过多个服务和步骤。用户真正关心的往往不是复杂的工作流编排,而是几个朴素问题:

  • 我的新任务是否已经创建?
  • 当前是等待、运行、成功,还是失败?
  • 当前由哪个服务执行什么功能,不同功能的数据该如何展示?
  • 哪个步骤正在执行,已经用了多久?
  • 步骤失败、自动重试或执行器失联时,用户如何收到明确告知?
  • 页面断线、刷新或切到后台后,会不会永远停在旧状态?
  • 服务多副本、节点重启、网络闪断时,状态通知是否仍能恢复?

本文设计一套专注于任务状态、步骤状态和耗时展示的生产链路。它不负责启动下一步、不替代工作流引擎,也不把实时消息系统当数据库。

最终方案可以概括为:

Worker 状态上报
→ Task Status Service 事务校验与落库
→ 同事务写 MySQL Outbox
→ Outbox Relay 领取并调用 Centrifugo 发布
→ 用户摘要频道 + 任务详情频道
→ 前端按 version 合并
→ 恢复失败时重新读取数据库快照

一、先划清系统边界

1.1 这套系统负责什么

能力具体职责
当前状态保存任务与步骤当前处于哪个状态
耗时保存起止时间和最终耗时,支持运行中本地计时
上报校验校验任务归属、执行轮次、事件幂等和状态转换
实时展示将已经提交的状态近实时通知给授权用户
断线恢复短断线尝试消息恢复,恢复失败读取业务快照
多租户权限只让用户订阅自己有权查看的任务
可观测性监控状态写入、Outbox 积压、发布延迟、连接与恢复结果

1.2 它明确不负责什么

  • 不决定下一步执行哪个服务;
  • 不创建和调度 Worker;
  • 不因为“目前已知步骤都成功”就猜测总任务成功;
  • 不把频道历史当永久事件审计库;
  • 不保证 Worker 已完成业务但尚未上报的状态自动出现;
  • 不为页面上每秒变化的耗时发送一次消息。

如果系统需要 DAG 编排、补偿、人工审批、长周期定时器和跨服务重试,应交给 Temporal、Camunda、Argo Workflows 等编排能力。状态系统只观察和展示,不悄悄演化成第二套工作流引擎。


二、总体架构:事实层、通知层、展示层

任务状态同步系统生产级整体架构:MySQL 事实源、Outbox Relay、Centrifugo、Redis 与多端客户端

整体架构图:MySQL 保存权威事实,Centrifugo 与 Redis 负责实时分发和短期恢复;虚线表示实时链路无法证明连续时的快照兜底。

品牌标识来源:MySQL 官方 Logo 页面Centrifugo 官方项目Redis 官方项目

三层的职责不能互相替代:

  1. MySQL 是事实源:回答“任务现在到底是什么状态”。
  2. Centrifugo 是通知与短期恢复层:回答“如何让在线客户端尽快知道变化”。
  3. 前端状态仓库是视图:按版本合并快照与推送,不能反向成为权威状态。

Centrifugo 官方也把频道历史定义为有界缓存,而不是永久消息队列;历史不足时,应用应回到自己的数据库恢复状态。参见 Stream history and recovery


三、完整链路:从创建任务到页面显示完成

一次状态变化如何可靠到达当前用户的多个任务

核心流程图:写入成功只代表事实已经持久化;通知可重复、可乱序,客户端以版本合并,断线后采用 History 与 MySQL 快照两级恢复。

3.1 创建任务

任务创建时就应确定任务类型、功能定义版本和每个步骤的 operationCode,不能等 Worker 上报时才临时猜测:

{
"taskType": "DOCUMENT_PROCESS",
"statusDefinitionVersion": 3,
"operationDefinitionVersion": 2,
"steps": [
{"stepId": "01_prepare", "operationCode": "FILE.PREPARE", "producerId": "file-worker"},
{"stepId": "02_parse", "operationCode": "DOCUMENT.PARSE", "producerId": "document-parser"},
{"stepId": "03_store", "operationCode": "RESULT.PERSIST", "producerId": "result-service"}
]
}

状态服务根据功能注册表校验上述组合并落库,再为每个步骤签发受限执行上下文。后续上报只能引用这份已登记定义。

任务、已知步骤和通知命令必须在同一个事务里提交。这样不会出现下面的缺口:

task 已经提交
→ 进程崩溃
→ 还没有调用实时发布 API
→ 页面长期不知道任务存在

创建接口响应和实时消息可能乱序:推送可能先到,也可能 HTTP 响应先到。前端不能用“收到一条就追加一行”,而应以 taskId 定位对象、以 version 判断新旧。

3.2 步骤开始与完成

Worker 不直接写任务表,也不直接向浏览器发布。它向状态服务上报业务事实:

POST /internal/tasks/8a4d/events
Content-Type: application/json
Idempotency-Key: 019a-event-001
{
"eventId": "019a-event-001",
"taskId": "task-8a4d",
"taskType": "DOCUMENT_PROCESS",
"taskAttempt": 1,
"stepId": "02_parse",
"operationCode": "DOCUMENT.PARSE",
"stepAttempt": 1,
"producerId": "document-parser",
"sourceSeq": 2,
"eventType": "STEP_COMPLETED",
"businessStatus": "PARSE_COMPLETED",
"occurredAt": "2026-09-21T07:10:18Z",
"durationMs": 18000,
"data": {
"schema": "document.parse.completed.v2",
"value": {"pageCount": 36, "language": "zh-CN"}
}
}

服务端处理顺序如下:

这些防线分别解决不同问题:

字段或机制解决的问题规则
executionToken服务越权替其他任务或功能上报绑定任务、步骤、功能、轮次、生产者和允许事件类型
operationCode无法判断上报的具体功能必须匹配任务创建时登记的步骤定义
data.schema不同功能数据结构混用必须属于该功能和事件类型的允许 Schema
eventId请求超时后的重复提交同一次重试必须使用相同 ID 和相同内容
attempt上一次执行晚到并覆盖新一轮旧轮次只能被忽略,不能改变当前轮次
sourceSeq同一轮次内部乱序只接受比当前序号更大的状态变化
versionHTTP、推送和多端接收乱序服务端每次有效变更递增,客户端只接收新版本

3.3 先统一八类标识:任务、功能、来源和事件各自独立

同一个任务可以由文档服务、AI 服务和结果服务依次或并行处理。它们不能只传一个模糊的 status,而要携带足够上下文,让状态服务知道“这是哪类任务、哪个业务功能、哪个逻辑步骤、哪一次执行、由谁产生、在该来源中的第几个事件”。

标识生成方生命周期作用
taskIdCoordinator / Business API整个任务不变聚合不同服务产生的步骤与事件
taskType业务 API / 任务模板创建后不变区分 DOCUMENT_PROCESSVIDEO_EXPORT 等整类业务任务
stepId任务模板或 Coordinator同一逻辑步骤不变标识这一次任务图中的具体步骤,如 02_parse
operationCode功能注册表 / 任务模板步骤定义内不变标识步骤执行的稳定业务功能,如 DOCUMENT.PARSE
attemptCoordinator每次重试递增隔离旧执行实例的晚到事件
producerId服务身份系统服务实例或逻辑生产者说明谁在上报,但不等同于业务功能
eventId事件生产服务每次业务事件唯一HTTP 超时重试时保持相同,用于幂等
sourceSeq事件生产服务taskId + stepId + attempt + producerId 内递增判断同一生产者事件的先后顺序

serviceName/producerIdoperationCode 必须分开:一个 AI 服务可能同时实现摘要、翻译和分类三个功能;同一个 DOCUMENT.PARSE 功能以后也可能从旧解析服务迁移到新解析服务。用服务名判断功能,会让扩容、重构和灰度发布直接改变业务语义。

协调方启动步骤时,应把不可缺失的执行上下文传给服务,而不是让服务自己猜:

{
"taskId": "task-8a4d",
"taskType": "DOCUMENT_PROCESS",
"taskAttempt": 1,
"stepId": "02_parse",
"operationCode": "DOCUMENT.PARSE",
"stepAttempt": 1,
"producerId": "document-parser",
"statusDefinitionVersion": 3,
"operationDefinitionVersion": 2,
"executionToken": "short-lived-signed-token"
}

executionToken 由状态服务或协调方签发并绑定以上字段以及允许的 eventType。服务 A 不能拿服务 B 的 Token 给另一个步骤上报,也不能把 DOCUMENT.PARSE 改成 VIDEO.TRANSCODE,更不能把 attempt=1 改成 attempt=2。状态服务以数据库中的步骤定义和 Token 为准;请求体中的功能字段只是显式自描述,必须完全匹配,不能由 Worker 自由声明。

sourceSeq 也不能只放在进程内存中:来源服务应把“业务结果、本地待上报事件、下一序号”放在自己的持久化事务或可靠 Outbox 中。服务重启后继续同一 attempt 时恢复序号;如果是一次真正的重新执行,则由 Coordinator 创建新的 stepAttempt

3.3.1 不同功能下发不同数据:稳定信封 + 带版本的数据体

所有功能共享同一层稳定信封,客户端先依据 kindtaskTypeoperationCode 和版本完成路由与去重;只有 data 是功能特有内容:

{
"kind": "task.detail",
"schemaVersion": 1,
"taskId": "task-8a4d",
"taskType": "DOCUMENT_PROCESS",
"stepId": "02_parse",
"operationCode": "DOCUMENT.PARSE",
"eventType": "STEP_COMPLETED",
"taskVersion": "12",
"stepVersion": "4",
"data": {
"schema": "document.parse.completed.v2",
"value": {
"pageCount": 36,
"language": "zh-CN",
"resultRef": "result://document/task-8a4d/parse-v2"
}
}
}

视频转码功能可以使用完全不同的数据体,但不能改变公共信封:

{
"taskType": "VIDEO_EXPORT",
"operationCode": "VIDEO.TRANSCODE",
"eventType": "STEP_PROGRESS",
"data": {
"schema": "video.transcode.progress.v1",
"value": {
"frame": 18420,
"totalFrames": 30000,
"fps": 48.2,
"preset": "1080p-h264"
}
}
}

下发数据还要按受众做投影,而不是把同一份大 JSON 广播给所有频道:

投影视图面向对象数据范围
USER_SUMMARY当前用户任务列表标题、功能摘要、总状态、进度、用户可读异常,不含步骤原始结果
TASK_DETAIL已授权且正在查看详情的客户端步骤、operationCode、功能特有 data、脱敏错误与操作建议
PROJECT_BROADCAST项目成员聚合进度和非敏感告警,不含个人输入、结果地址和诊断信息
SYSTEM_BROADCAST所有在线用户维护与服务状态,不携带任何业务任务数据

data.schema 必须登记在该 operationCode 允许的 Schema 列表中。兼容字段可以在同一版本追加;删除、改名或改变含义必须发布新版本。状态服务负责 JSON Schema 校验、字段白名单、大小限制和敏感字段裁剪,大结果只下发短期授权的 resultRef,不能直接塞进实时消息。

前端按功能注册渲染器,未知功能或未知 Schema 使用通用状态卡降级,不能因此丢弃任务的生命周期状态:

const renderer = operationRenderers[message.operationCode]?.[message.data?.schema];

return renderer
? renderer(message.data.value)
: renderGenericTaskState(message.lifecycleStatus, message.businessStatus);

3.4 一次多服务任务如何完整传递

多个服务在同一任务中上报状态、耗时并由状态服务分配版本

同一个 taskId 下,来源服务只维护自己的 sourceSeq;状态服务在 MySQL 事务里为所有已接受事件分配全局递增的 taskVersion。

假设任务有三个步骤:

01_prepare FILE.PREPARE 文件服务:准备文件
02_parse DOCUMENT.PARSE AI 服务:解析内容
03_store RESULT.PERSIST 结果服务:保存结果

完整过程如下:

  1. Coordinator 创建任务和三个步骤,状态服务写入 taskVersion=1,返回每个步骤的执行上下文;
  2. 文件服务上报 STEP_STARTED, sourceSeq=1,状态服务接受后分配 taskVersion=2
  3. 文件服务上报 STEP_COMPLETED, sourceSeq=2, durationMs=3100,分配 taskVersion=3
  4. AI 服务上报 STEP_STARTED, sourceSeq=1。它是另一个 stepId + producerId,序号可以重新从 1 开始,分配 taskVersion=4
  5. AI 服务可以上报节流后的 STEP_PROGRESS,例如从 35% 到 60%,每个被接受的变化获得新 taskVersion
  6. AI 服务完成后上报 STEP_COMPLETED, durationMs=18000,结果服务再以自己的 sourceSeq 上报;
  7. Coordinator 确认所有业务收尾完成,最后上报 TASK_COMPLETED,状态服务将生命周期置为 SUCCEEDED

不同服务不共享 sourceSeq 计数器,也不直接生成 taskVersion。它们只对自己的事件流负责,状态服务负责把跨服务并发事件排成一个任务级提交顺序。

这里传递的是状态事实,不是服务间业务结果。服务 A 的文件、模型输出或索引地址仍通过现有工作流、对象存储或业务数据库交给服务 B;状态事件只携带可展示的摘要和不敏感引用。服务之间不互相转发状态,也不经浏览器中转,所有服务都独立上报给 Task Status Service。

3.5 统一事件协议:状态、进度和耗时一起传

推荐使用一个事件入口,而不是每个服务自定义一套接口:

POST /internal/v1/task-events
Authorization: Bearer <executionToken>
Idempotency-Key: evt-01JQ...
Content-Type: application/json
{
"eventId": "evt-01JQ...",
"taskId": "task-8a4d",
"taskType": "DOCUMENT_PROCESS",
"taskAttempt": 1,
"stepId": "02_parse",
"operationCode": "DOCUMENT.PARSE",
"stepAttempt": 1,
"producerId": "document-parser",
"sourceSeq": 2,
"eventType": "STEP_COMPLETED",
"businessStatus": "PARSE_COMPLETED",
"occurredAt": "2026-09-21T07:10:18.120Z",
"durationMs": 18000,
"progress": {
"current": 36,
"total": 36,
"unit": "page"
},
"data": {
"schema": "document.parse.completed.v2",
"value": {
"pageCount": 36,
"language": "zh-CN",
"resultRef": "result://document/task-8a4d/parse-v2"
}
}
}

taskType 表达整类任务,operationCode 表达正在执行的功能,eventType 表达这次发生了什么,businessStatus 表达业务现在处于什么阶段,四者不要混用:

eventType是否改变状态常用字段说明
STEP_STARTEDoccurredAt第一次接受时确定 startedAt
STEP_PROGRESS可选progressmessage应节流、可覆盖,避免高频写入
STEP_COMPLETEDdurationMsdata映射到成功生命周期
STEP_FAILEDdurationMserror错误内容先脱敏再通知客户端
STEP_HEARTBEAToccurredAt只更新健康时间,不等同业务进度
STEP_HEALTH_CHANGEDhealthreason只允许健康检查器上报,不擅自改业务终态
TASK_COMPLETEDdurationMs只能由任务协调方上报
TASK_FAILEDerrorfailedStepId只能由任务协调方或注册的失败策略执行器上报

普通心跳只更新 lastHeartbeatAt,可以不增加 taskVersion、也不向客户端发布;只有健康状态从 HEALTHY 变成 STALE 或恢复时,Health Reconciler 才提交 STEP_HEALTH_CHANGED 并分配新版本。

状态服务的响应必须把最终接受结果和版本返回给调用方:

{
"accepted": true,
"duplicate": false,
"stale": false,
"taskId": "task-8a4d",
"stepId": "02_parse",
"operationCode": "DOCUMENT.PARSE",
"taskVersion": "12",
"stepVersion": "4",
"lifecycleStatus": "RUNNING",
"businessStatus": "PARSE_COMPLETED",
"committedAt": "2026-09-21T07:10:18.156Z"
}

请求超时后,生产服务使用相同 eventId 和相同内容重试。状态服务从幂等表返回第一次提交的同一个 taskVersion,不能再次递增。

3.6 七种序号和版本到底分别从哪里来

这里有七种容易混淆但用途完全不同的序号或版本:

字段产生或绑定方作用范围主要用途
schemaVersion实时协议发布方消息格式判断公共 JSON 信封如何解析
statusDefinitionVersion创建任务时由服务端绑定任务生命周期固定业务状态代码与生命周期的映射
operationDefinitionVersion创建任务时由服务端绑定任务生命周期固定功能、生产者、数据 Schema 与失败策略
data.schema来源服务从允许集合选择,服务端校验单类功能数据判断 data.value 的具体结构与兼容版本
sourceSeq来源服务生成单个步骤执行流拒绝同一来源的乱序旧事件
stepVersion状态服务生成单个 taskId + stepId合并某一步骤的快照与消息
taskVersion状态服务生成整个任务跨多个服务建立统一提交顺序

版本在 MySQL 事务中分配。伪代码如下:

START TRANSACTION;

SELECT version, attempt, lifecycle_status
FROM task
WHERE id = :task_id
FOR UPDATE;

SELECT version, attempt AS step_attempt, source_seq, lifecycle_status
FROM task_step
WHERE task_id = :task_id
AND step_id = :step_id
AND attempt = :step_attempt
FOR UPDATE;

-- 应用进程在持有行锁时完成:
-- 1. eventId 幂等检查
-- 2. executionToken / attempt / sourceSeq / 状态机校验
-- 3. new_task_version = current_task_version + 1
-- 4. new_step_version = current_step_version + 1

UPDATE task_step
SET version = :new_step_version,
source_seq = :source_seq,
lifecycle_status = :step_lifecycle,
business_status = :business_status,
updated_at = UTC_TIMESTAMP(3)
WHERE task_id = :task_id
AND step_id = :step_id
AND attempt = :step_attempt;

UPDATE task
SET version = :new_task_version,
updated_at = UTC_TIMESTAMP(3)
WHERE id = :task_id;

INSERT INTO task_event_log (..., task_version, step_version, ...)
VALUES (..., :new_task_version, :new_step_version, ...);

INSERT INTO task_outbox (..., aggregate_version, payload)
VALUES (..., :new_task_version, :complete_state_json);

COMMIT;

如果服务 A 和服务 B 同时上报,二者会竞争同一 task 行锁:先提交的事件得到 taskVersion=12,后提交的得到 taskVersion=13。这不是要求服务之间同步时钟,而是以数据库提交顺序为任务建立稳定总序。某个超热门任务可能因此成为单行热点,但普通异步任务的状态变化频率通常远低于业务数据吞吐;确有证据时再把步骤事件流与任务摘要版本拆分。

3.7 事件耗时如何得到,哪些时间可以相信

每个时间字段回答不同问题:

字段产生位置用途
occurredAt来源服务墙钟展示业务事件大致发生时间
receivedAt状态服务观察网络与排队延迟
committedAtMySQL 事务提交附近定义该版本何时成为事实
durationMs来源服务单调时钟精确表示该次执行耗时

步骤开始时,来源服务记录单调时钟起点,并上报 STEP_STARTED;结束时用同一进程的单调时钟计算 durationMs。不要用两台机器的 finishedAt - startedAt 计算精确执行耗时,因为 NTP 偏差、休眠和时钟回拨都会污染结果。

状态服务同时保存三类耗时:

queueDurationMs = startedAt - createdAt 排队等待
runDurationMs = Worker 上报的 durationMs 实际执行
totalDurationMs = task committed terminal - task createdAt

对于进程重启导致单调时钟起点丢失的情况,Worker 应从自己的持久化执行记录恢复;确实无法恢复时把 durationSource 标记为 WALL_CLOCK_ESTIMATE,而不是伪装成精确值。

3.8 客户端怎样获得版本和事件时间线

客户端有三种版本来源,取最大值合并:

创建任务响应:taskVersion=1
快照接口响应:taskVersion=11, 各 stepVersion
实时消息: taskVersion=12

推荐快照接口:

GET /api/tasks/task-8a4d
GET /api/tasks/task-8a4d/events?afterTaskVersion=8&limit=100

快照返回当前完整状态,事件接口返回用于时间线展示的已提交事件。实时 task.detail 消息同时带 triggerEvent 和完整当前状态:事件区可以追加时间线,状态区只接受更高版本。

{
"kind": "task.detail",
"taskVersion": "12",
"triggerEvent": {
"eventId": "evt-01JQ...",
"eventType": "STEP_COMPLETED",
"stepId": "02_parse",
"operationCode": "DOCUMENT.PARSE",
"durationMs": 18000,
"data": {
"schema": "document.parse.completed.v2",
"value": {"pageCount": 36, "language": "zh-CN"}
},
"committedAt": "2026-09-21T07:10:18.156Z"
},
"state": {
"taskType": "DOCUMENT_PROCESS",
"lifecycleStatus": "RUNNING",
"businessStatus": "PARSE_COMPLETED",
"steps": []
}
}

如果实时消息从版本 9 直接跳到 12,不代表一定丢失了当前状态,因为消息携带完整快照;如果 UI 还要展示每一个中间事件,再调用 events?afterTaskVersion=9 补时间线。当前状态恢复和审计事件补齐是两个不同问题。

3.9 总任务完成

总任务终态应由现有业务协调方明确上报:

{
"eventId": "019a-task-final-001",
"taskId": "task-8a4d",
"taskAttempt": 1,
"producerId": "task-coordinator",
"sourceSeq": 7,
"eventType": "TASK_COMPLETED",
"businessStatus": "COMPLETED",
"durationMs": 42600
}

不要让状态服务根据“当前登记的步骤都成功”自动推导总任务成功。真实业务可能还在保存结果、生成索引、写审计记录,甚至下一步尚未登记。只有当工作流定义固定、完整且由状态服务权威管理时,才适合做状态聚合推导。


四、状态机:允许什么,拒绝什么

4.1 推荐状态集合

PENDING 已创建,尚未开始
RUNNING 正在执行
SUCCEEDED 成功终止
FAILED 失败终止
CANCELLED 取消终止

如果业务明确区分超时,可以增加 TIMED_OUT;如果只是展示“执行器可能失联”,更适合增加独立的健康字段,而不是把业务状态直接改成失败。

重试不是把同一个 attempt=1FAILED 改回 RUNNING,而是由协调方创建 attempt=2。这能保留“第一次失败、第二次成功”的事实,也让晚到的第一次完成消息无法污染第二轮。

4.2 终态不可被普通事件覆盖

同一轮次中,下面的变化都应拒绝:

SUCCEEDED → RUNNING
FAILED → SUCCEEDED
CANCELLED → RUNNING

如果业务确实需要人工纠错,不应伪装成普通 Worker 上报。应设计独立的管理员纠错命令,记录操作者、原因、旧值、新值和审计时间,并产生新的业务版本。

4.3 业务状态与健康状态分离

{
"lifecycleStatus": "RUNNING",
"health": "STALE",
"lastHeartbeatAt": "2026-09-21T07:11:00Z"
}
  • lifecycleStatus=RUNNING:最后一个已确认的业务事实是正在运行;
  • health=STALE:一段时间未观察到执行端心跳;
  • UI 显示“运行中,状态待确认”,而不是擅自改成“失败”。

连接断开、Worker 心跳超时和任务失败是三件不同的事。把它们混成一个状态,会制造大量假失败。

4.4 异常不是一个字符串:失败、失联和通知故障分别建模

异常至少分成四类,处理主体和用户文案都不同:

异常类型谁确认状态如何变化用户应该看到什么
业务执行失败当前步骤 Worker 上报 STEP_FAILED当前步骤 attempt 进入 FAILED“文档解析失败:文件格式不受支持”
可重试步骤失败Worker + Coordinator本 attempt 为 FAILED,任务可保持 RUNNING/RETRY_WAIT,随后创建新 attempt“解析失败,正在进行第 2 次重试”
Worker 心跳超时Health Reconcilerhealth=STALE,生命周期仍保留最后已确认值“仍在运行,但暂时无法确认执行器状态”
状态上报或实时通知故障来源 Outbox / Relay / 客户端恢复机制不制造业务 FAILED“实时更新暂不可用”,恢复后读取快照

步骤失败时建议使用结构化错误,而不是让每个服务随意拼接一段异常栈:

{
"eventType": "STEP_FAILED",
"operationCode": "DOCUMENT.PARSE",
"businessStatus": "PARSE_FAILED",
"durationMs": 18240,
"error": {
"category": "BUSINESS_INPUT",
"code": "UNSUPPORTED_FILE_FORMAT",
"retryable": false,
"userMessage": "暂不支持该文件格式,请转换为 PDF 后重试",
"diagnosticsRef": "trace://01JQ7M..."
}
}
  • code 是稳定的机器可读错误码,驱动前端操作按钮和统计;
  • userMessage 必须可安全展示且经过长度限制、脱敏和国际化映射;
  • diagnosticsRef 只给服务端排障使用,客户端不能凭它读取日志;
  • 原始堆栈、SQL、对象存储地址、密钥和第三方响应不能进入实时消息。

步骤失败不一定意味着总任务立即失败。每个 operationCode 在功能定义中声明 failurePolicy

FAIL_TASK 关键步骤失败,由 Coordinator 明确提交 TASK_FAILED
RETRY_THEN_FAIL_TASK 先创建新 attempt,耗尽重试后提交 TASK_FAILED
CONTINUE_WITH_WARNING 可选步骤失败,任务继续但 health=DEGRADED
WAIT_FOR_USER 任务保持 RUNNING,businessStatus=WAITING_USER_ACTION

状态服务可以校验策略,但总任务终态仍应由 Coordinator 或明确注册的策略执行器提交。它不能仅凭“看到一个步骤失败”猜测任务失败,因为该步骤可能可选、可补偿或正在重试。

4.5 允许业务自定义状态,但不要让它接管生命周期

不同产品需要不同的展示语义。例如文档任务可能经历“正在 OCR”“等待人工复核”,视频任务可能经历“正在渲染”“正在上传”。如果把这些值全部塞进通用状态机,每接入一个业务都要修改状态服务、客户端和告警规则。

更稳妥的模型是同时保存两套状态:

字段示例作用
lifecycleStatusRUNNING跨业务稳定,驱动终态判断、重试、耗时和告警
businessStatusWAITING_REVIEW业务自定义,驱动标题、说明、进度阶段和视觉呈现
statusDefinitionVersion3固定本轮任务使用的状态定义,防止运行中配置漂移
healthSTALE描述执行端健康,不改写已确认的业务事实

例如:

{
"taskId": "8a4d",
"lifecycleStatus": "RUNNING",
"businessStatus": "WAITING_REVIEW",
"statusDefinitionVersion": 3,
"health": "HEALTHY",
"version": "12"
}

状态上报方只能提交已经注册的 businessStatus 代码。服务端根据任务类型与定义版本解析出它对应的生命周期,并校验允许的来源状态:

DOCUMENT_PROCESS / v3

QUEUED → lifecycle=PENDING
PARSING → lifecycle=RUNNING
WAITING_REVIEW → lifecycle=RUNNING
PUBLISHING → lifecycle=RUNNING
COMPLETED → lifecycle=SUCCEEDED, terminal=true
REJECTED → lifecycle=FAILED, terminal=true

调用方不能自行声明 terminal=true,也不能决定颜色、严重级别或对应的生命周期,否则一个拼错或恶意字段就可能绕过终态保护。未知业务状态应拒绝写入;旧客户端遇到暂不认识的状态时,使用 lifecycleStatus + 原始 code 降级展示,不能崩溃或误判为成功。

4.6 自定义状态注册表

状态定义可以按租户、产品和任务类型配置,但定义变更必须版本化:

CREATE TABLE task_status_definition (
tenant_id CHAR(36) NOT NULL COMMENT '租户 ID',
product VARCHAR(64) NOT NULL COMMENT '产品或业务域标识',
task_type VARCHAR(128) NOT NULL COMMENT '任务类型',
definition_version INT UNSIGNED NOT NULL COMMENT '业务状态定义版本',
code VARCHAR(64) NOT NULL COMMENT '业务状态编码',
display_name VARCHAR(128) NOT NULL COMMENT '默认展示名称',
lifecycle_status VARCHAR(32) NOT NULL COMMENT '映射到的系统生命周期状态',
terminal BOOLEAN NOT NULL DEFAULT FALSE COMMENT '该业务状态是否表示终态',
severity VARCHAR(16) NOT NULL DEFAULT 'INFO' COMMENT '展示与告警严重级别',
sort_order INT NOT NULL DEFAULT 0 COMMENT '同类状态的展示顺序',
presentation JSON NULL COMMENT '颜色、图标和国际化文案键等展示元数据',
allowed_from JSON NOT NULL COMMENT '允许迁移到当前状态的前置业务状态集合',
enabled BOOLEAN NOT NULL DEFAULT TRUE COMMENT '状态定义是否允许用于新事件',
PRIMARY KEY (
tenant_id, product, task_type, definition_version, code
)
) ENGINE=InnoDB COMMENT='租户和任务类型维度的版本化业务状态定义';

presentation 只保存诸如 colorTokeniconKey 和本地化文案键等展示元数据,不保存任意 HTML。新任务固定使用创建时的 definition_version;发布新版定义不会悄悄改变正在运行或历史任务的含义。

4.7 功能定义注册表:限制谁能上报什么、下发什么

operationCode 不能只是一个没有约束的字符串。任务模板引用版本化的功能定义,状态服务据此校验生产者、事件类型、差异化数据 Schema 和失败策略:

CREATE TABLE task_operation_definition (
tenant_id CHAR(36) NOT NULL COMMENT '租户 ID',
product VARCHAR(64) NOT NULL COMMENT '产品或业务域标识',
task_type VARCHAR(128) NOT NULL COMMENT '任务类型',
definition_version INT UNSIGNED NOT NULL COMMENT '功能定义版本',
operation_code VARCHAR(128) NOT NULL COMMENT '稳定的业务功能编码',
display_name VARCHAR(128) NOT NULL COMMENT '功能默认展示名称',
producer_pattern VARCHAR(255) NOT NULL COMMENT '允许上报该功能的服务身份模式',
allowed_event_types JSON NOT NULL COMMENT '允许该功能提交的事件类型集合',
allowed_data_schemas JSON NOT NULL COMMENT '各事件允许使用的数据 Schema 集合',
failure_policy VARCHAR(32) NOT NULL COMMENT '失败后的任务级处理策略',
enabled BOOLEAN NOT NULL DEFAULT TRUE COMMENT '是否允许用于新任务',
PRIMARY KEY (
tenant_id, product, task_type, definition_version, operation_code
)
) ENGINE=InnoDB COMMENT='任务功能、上报权限、数据契约与失败策略定义';

例如 DOCUMENT.PARSE 可以允许 document.parse.progress.v1document.parse.completed.v2,但拒绝 video.transcode.progress.v1producerId=document-parser 也不能给 RESULT.PERSIST 上报。新任务在创建时固定 operation_definition_version,从而避免任务运行途中配置更新导致同一事件突然改变解释方式。


五、MySQL 数据模型:当前状态、幂等和 Outbox

第一版可以使用五张运行时核心表:任务、步骤、不可变事件、事件幂等和 Outbox,外加上一节的版本化状态定义表与功能定义表。下面以 MySQL 8.0 和 InnoDB 为基线,为了便于阅读使用 CHAR(36) 保存 UUID;高写入量场景可以改为 BINARY(16),但必须统一字节序和序列化规则。

CREATE TABLE task (
id CHAR(36) NOT NULL COMMENT '任务全局唯一 ID',
tenant_id CHAR(36) NOT NULL COMMENT '任务所属租户 ID',
owner_id CHAR(36) NOT NULL COMMENT '任务所属用户 ID',
product VARCHAR(64) NOT NULL COMMENT '产品或业务域标识',
task_type VARCHAR(128) NOT NULL COMMENT '任务类型',
title VARCHAR(255) NOT NULL COMMENT '面向用户展示的任务标题',
lifecycle_status VARCHAR(32) NOT NULL COMMENT '系统生命周期状态',
business_status VARCHAR(64) NOT NULL COMMENT '业务自定义状态编码',
status_definition_version INT UNSIGNED NOT NULL COMMENT '创建任务时绑定的业务状态定义版本',
operation_definition_version INT UNSIGNED NOT NULL COMMENT '创建任务时绑定的功能定义版本',
attempt INT UNSIGNED NOT NULL DEFAULT 1 COMMENT '任务级执行轮次',
coordinator_source_seq BIGINT UNSIGNED NOT NULL DEFAULT 0 COMMENT 'Coordinator 任务级事件序号',
version BIGINT UNSIGNED NOT NULL DEFAULT 1 COMMENT '状态服务分配的任务级单调版本',
health VARCHAR(16) NOT NULL DEFAULT 'HEALTHY' COMMENT 'HEALTHY、DEGRADED 或 STALE',
summary_schema VARCHAR(128) NULL COMMENT '任务列表差异化摘要数据的 Schema',
summary_data JSON NULL COMMENT '经过白名单裁剪的任务列表差异化摘要数据',
failure_code VARCHAR(128) NULL COMMENT '任务级稳定失败码',
failure_message VARCHAR(512) NULL COMMENT '面向用户且已脱敏的任务失败说明',
failed_step_id VARCHAR(128) NULL COMMENT '导致任务失败的步骤 ID',
created_at DATETIME(3) NOT NULL COMMENT '任务创建时间',
started_at DATETIME(3) NULL COMMENT '任务首次进入运行态的时间',
finished_at DATETIME(3) NULL COMMENT '任务进入权威终态的时间',
duration_ms BIGINT UNSIGNED NULL COMMENT '任务最终确认耗时,单位毫秒',
updated_at DATETIME(3) NOT NULL COMMENT '任务状态最后提交时间',
PRIMARY KEY (id),
KEY task_owner_updated_idx
(tenant_id, owner_id, updated_at DESC, id)
) ENGINE=InnoDB COMMENT='任务当前权威状态与任务级版本';

CREATE TABLE task_step (
task_id CHAR(36) NOT NULL COMMENT '所属任务 ID',
step_id VARCHAR(128) NOT NULL COMMENT '任务内稳定的步骤 ID',
operation_code VARCHAR(128) NOT NULL COMMENT '该步骤执行的稳定业务功能编码',
service_name VARCHAR(128) NOT NULL COMMENT '负责执行并上报状态的服务名',
attempt INT UNSIGNED NOT NULL COMMENT '步骤执行轮次',
source_seq BIGINT UNSIGNED NOT NULL DEFAULT 0 COMMENT '当前服务在本步骤轮次内的单调事件序号',
version BIGINT UNSIGNED NOT NULL DEFAULT 1 COMMENT '状态服务分配的步骤级单调版本',
lifecycle_status VARCHAR(32) NOT NULL COMMENT '步骤生命周期状态',
business_status VARCHAR(64) NOT NULL COMMENT '步骤业务自定义状态编码',
health VARCHAR(16) NOT NULL DEFAULT 'HEALTHY' COMMENT '步骤执行器的健康状态',
last_heartbeat_at DATETIME(3) NULL COMMENT '最近一次已接受心跳时间',
data_schema VARCHAR(128) NULL COMMENT '当前功能特有状态数据的 Schema',
state_data JSON NULL COMMENT '经过校验和裁剪的功能特有当前状态数据',
started_at DATETIME(3) NULL COMMENT '步骤开始时间',
finished_at DATETIME(3) NULL COMMENT '步骤进入终态的时间',
duration_ms BIGINT UNSIGNED NULL COMMENT '执行服务用单调时钟测得的耗时,单位毫秒',
error_category VARCHAR(64) NULL COMMENT '业务输入、依赖、系统等错误分类',
error_code VARCHAR(128) NULL COMMENT '稳定的机器可读错误码',
error_message TEXT NULL COMMENT '脱敏后的错误说明',
error_retryable BOOLEAN NULL COMMENT '该步骤失败是否允许重试',
diagnostics_ref VARCHAR(255) NULL COMMENT '仅供服务端排障关联的诊断引用',
updated_at DATETIME(3) NOT NULL COMMENT '步骤状态最后提交时间',
PRIMARY KEY (task_id, step_id, attempt),
CONSTRAINT task_step_task_fk
FOREIGN KEY (task_id) REFERENCES task(id)
) ENGINE=InnoDB COMMENT='每个任务步骤及执行轮次的当前权威状态';

CREATE TABLE task_event_log (
event_id CHAR(36) NOT NULL COMMENT '生产者生成的全局幂等事件 ID',
task_id CHAR(36) NOT NULL COMMENT '所属任务 ID',
task_version BIGINT UNSIGNED NOT NULL COMMENT '该事件提交后获得的任务级版本',
step_id VARCHAR(128) NULL COMMENT '步骤事件对应的步骤 ID,任务级事件为空',
operation_code VARCHAR(128) NULL COMMENT '步骤事件对应的业务功能编码',
step_version BIGINT UNSIGNED NULL COMMENT '该事件提交后获得的步骤级版本',
task_attempt INT UNSIGNED NOT NULL COMMENT '事件所属任务执行轮次',
step_attempt INT UNSIGNED NULL COMMENT '事件所属步骤执行轮次',
producer_id VARCHAR(128) NOT NULL COMMENT '事件生产服务或协调器标识',
source_seq BIGINT UNSIGNED NOT NULL COMMENT '生产者在对应执行流内的单调序号',
event_type VARCHAR(64) NOT NULL COMMENT '事件类型,例如 STEP_STARTED',
lifecycle_status VARCHAR(32) NOT NULL COMMENT '事件提交后的生命周期状态',
business_status VARCHAR(64) NOT NULL COMMENT '事件提交后的业务状态编码',
health VARCHAR(16) NOT NULL COMMENT '事件提交后的健康状态',
data_schema VARCHAR(128) NULL COMMENT '功能特有数据的 Schema',
error_code VARCHAR(128) NULL COMMENT '该事件携带的稳定错误码',
occurred_at DATETIME(3) NOT NULL COMMENT '事件在生产者侧发生的墙钟时间',
received_at DATETIME(3) NOT NULL COMMENT '状态服务接收到事件的时间',
committed_at DATETIME(3) NOT NULL COMMENT '事件在 MySQL 中提交的时间',
duration_ms BIGINT UNSIGNED NULL COMMENT '生产者测量的事件或步骤耗时,单位毫秒',
payload JSON NULL COMMENT '审计和回放所需的扩展事件数据',
PRIMARY KEY (event_id),
UNIQUE KEY task_event_version_uk (task_id, task_version),
KEY task_event_timeline_idx (task_id, committed_at, event_id),
UNIQUE KEY task_step_event_seq_uk
(task_id, step_id, step_attempt, source_seq)
) ENGINE=InnoDB COMMENT='已接受状态事件的不可变时间线与版本审计记录';

CREATE TABLE task_event_dedup (
event_id CHAR(36) NOT NULL COMMENT '生产者生成的全局幂等事件 ID',
task_id CHAR(36) NOT NULL COMMENT '事件所属任务 ID',
payload_hash CHAR(64) NOT NULL COMMENT '规范化事件内容的 SHA-256 摘要',
result JSON NOT NULL COMMENT '首次处理结果,包含已分配版本等响应字段',
created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) COMMENT '首次处理并记录幂等结果的时间',
PRIMARY KEY (event_id),
KEY task_event_task_idx (task_id, created_at)
) ENGINE=InnoDB COMMENT='事件幂等判定与首次响应结果缓存';

task.coordinator_source_seq 只排序 Coordinator 产生的任务级事件;每个服务的步骤事件分别使用对应 task_step.source_seq。不能拿一个任务级 source_seq 比较所有服务,否则文件服务的 seq=2 会错误地把 AI 服务合法的 seq=1 当成旧消息。

5.1 MySQL Outbox 表

Centrifugo 当前的内置异步消费者列表没有 MySQL。因此这里不让 Centrifugo 直接扫描数据库,而是由独立的 Outbox Relay 从 MySQL 领取记录,再调用 Centrifugo HTTP 或 gRPC Server API。官方支持的内置消费者及自定义 consumer 接入边界见 Built-in API command async consumers

CREATE TABLE task_outbox (
id BIGINT UNSIGNED NOT NULL AUTO_INCREMENT COMMENT 'Outbox 自增主键',
aggregate_id CHAR(36) NOT NULL COMMENT '聚合根 ID,此处为任务 ID',
aggregate_version BIGINT UNSIGNED NOT NULL COMMENT '消息对应的任务级版本',
channel VARCHAR(255) NOT NULL COMMENT '服务端计算出的 Centrifugo 目标频道',
projection VARCHAR(32) NOT NULL COMMENT 'USER_SUMMARY、TASK_DETAIL 或广播投影类型',
payload JSON NOT NULL COMMENT '包含完整任务状态的发布载荷',
status VARCHAR(16) NOT NULL DEFAULT 'PENDING' COMMENT 'PENDING、PROCESSING、PUBLISHED 或 DEAD',
available_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) COMMENT '允许 Relay 下一次领取的时间',
locked_by VARCHAR(128) NULL COMMENT '当前领取该记录的 Relay 实例 ID',
locked_at DATETIME(3) NULL COMMENT 'Relay 获得领取租约的时间',
attempts INT UNSIGNED NOT NULL DEFAULT 0 COMMENT '累计发布尝试次数',
last_error VARCHAR(1024) NULL COMMENT '最近一次脱敏后的发布错误',
created_at DATETIME(3) NOT NULL DEFAULT CURRENT_TIMESTAMP(3) COMMENT 'Outbox 记录创建时间',
published_at DATETIME(3) NULL COMMENT '最近一次成功发布时间',
PRIMARY KEY (id),
UNIQUE KEY outbox_aggregate_version_channel_uk
(aggregate_id, aggregate_version, channel, projection),
KEY outbox_claim_idx (status, available_at, id),
KEY outbox_lock_timeout_idx (status, locked_at)
) ENGINE=InnoDB COMMENT='状态变化的事务发布命令与 Relay 租约状态';

状态更新和通知命令在同一个 InnoDB 事务中提交:

START TRANSACTION;

SELECT version
FROM task
WHERE id = :task_id
FOR UPDATE;

SELECT version, source_seq, operation_code
FROM task_step
WHERE task_id = :task_id
AND step_id = :step_id
AND attempt = :step_attempt
FOR UPDATE;

UPDATE task_step
SET lifecycle_status = :step_lifecycle_status,
business_status = :business_status,
health = :step_health,
data_schema = :data_schema,
state_data = CAST(:state_data_json AS JSON),
error_category = :error_category,
error_code = :error_code,
error_message = :error_message,
error_retryable = :error_retryable,
diagnostics_ref = :diagnostics_ref,
source_seq = :source_seq,
version = :new_step_version,
updated_at = UTC_TIMESTAMP(3)
WHERE task_id = :task_id
AND step_id = :step_id
AND attempt = :step_attempt
AND source_seq < :source_seq;

UPDATE task
SET lifecycle_status = :lifecycle_status,
business_status = :business_status,
health = :task_health,
summary_schema = :summary_schema,
summary_data = CAST(:summary_data_json AS JSON),
failure_code = :failure_code,
failure_message = :failure_message,
failed_step_id = :failed_step_id,
version = :new_version,
updated_at = UTC_TIMESTAMP(3)
WHERE id = :task_id
AND attempt = :task_attempt;

INSERT INTO task_outbox (
aggregate_id,
aggregate_version,
channel,
projection,
payload
) VALUES (
:task_id,
:new_version,
:channel,
:projection,
CAST(:complete_state_json AS JSON)
);

INSERT INTO task_event_log (
event_id,
task_id,
task_version,
step_id,
operation_code,
step_version,
task_attempt,
step_attempt,
producer_id,
source_seq,
event_type,
lifecycle_status,
business_status,
health,
data_schema,
error_code,
occurred_at,
received_at,
committed_at,
duration_ms,
payload
) VALUES (
:event_id,
:task_id,
:new_version,
:step_id,
:operation_code,
:new_step_version,
:task_attempt,
:step_attempt,
:producer_id,
:source_seq,
:event_type,
:lifecycle_status,
:business_status,
:health,
:data_schema,
:error_code,
:occurred_at,
:received_at,
UTC_TIMESTAMP(3),
:duration_ms,
CAST(:event_payload_json AS JSON)
);

INSERT INTO task_event_dedup (
event_id,
task_id,
payload_hash,
result
) VALUES (
:event_id,
:task_id,
:payload_hash,
CAST(:result_json AS JSON)
);

COMMIT;

实际实现应先 SELECT ... FOR UPDATE 锁定任务并读取当前版本,再完成状态机校验和更新;不要依赖上面简化 SQL 猜测 :new_version

同一已接受事件可以在事务内生成多条不同频道的 Outbox 记录:用户频道写 USER_SUMMARY 投影,任务频道写 TASK_DETAIL 投影。投影器从已校验的当前状态构建载荷,不能把 Worker 原始 JSON 原样转发;否则来源服务可能绕过字段裁剪,把只允许出现在详情中的数据泄漏到摘要或广播频道。

5.2 多副本 Relay 如何安全领取

多个 Relay 可以使用 MySQL 8.0 的 FOR UPDATE SKIP LOCKED 并发领取不同记录。MySQL 官方明确指出,SKIP LOCKED 会跳过已经被其他事务锁定的行,适合多个会话访问同一队列表;它不适合需要一致性视图的普通业务查询。参见 InnoDB Locking Reads

START TRANSACTION;

SELECT id
FROM task_outbox
WHERE status = 'PENDING'
AND available_at <= UTC_TIMESTAMP(3)
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED;

UPDATE task_outbox
SET status = 'PROCESSING',
locked_by = :relay_id,
locked_at = UTC_TIMESTAMP(3),
attempts = attempts + 1
WHERE id IN (:claimed_ids);

COMMIT;

Relay 在事务外调用 Centrifugo Server API,成功后将记录更新为 PUBLISHED。失败时按指数退避计算新的 available_at,回到 PENDING。如果 Relay 在“发布成功、标记完成之前”崩溃,记录会在锁租约过期后重试,因此客户端必须允许重复消息。

每条消息必须携带完整状态和 aggregate_version。这样即使多个 Relay 导致同一任务的通知晚到或乱序,客户端仍会忽略旧版本,正确性不依赖严格的全局发布顺序。

Outbox 只覆盖“已经到达状态服务”的事件

如果 Worker 已完成业务,却在状态上报前崩溃,状态服务的 Outbox 无法凭空知道这次完成。来源服务仍需要持久化重试、本地 Outbox,或由对账任务查询实际结果并修复状态。


六、通知设计:可选传输、有限订阅与受控广播

客户端可选 WebSocket 或 SSE,一条连接承载用户摘要、有限任务详情与受控广播

订阅模型图:一条物理连接可以复用多个逻辑频道;用户摘要常驻,任务详情按需且有限,项目、租户与系统广播只能由服务端授权和发布。

6.1 用户摘要频道

登录期间始终订阅:

task_user:{tenantId}/{userId}

它发送用户有权看到的任务摘要:

{
"kind": "task.summary",
"schemaVersion": 1,
"task": {
"id": "8a4d",
"product": "DOCUMENT",
"taskType": "DOCUMENT_PROCESS",
"title": "合同解析",
"lifecycleStatus": "RUNNING",
"businessStatus": "WAITING_REVIEW",
"health": "HEALTHY",
"version": "12",
"startedAt": "2026-09-21T07:10:00Z",
"finishedAt": null,
"elapsedMs": 18000,
"asOf": "2026-09-21T07:10:18Z",
"summaryData": {
"schema": "document.process.summary.v1",
"value": {"fileName": "contract.pdf", "pageCount": 36}
}
}
}

新任务 B、C、D 都归入同一用户频道,所以不需要重建连接,也不需要提前知道未来的任务 ID。

这也是“当前用户如何同时收到不同任务数据”的关键:不是给每个任务各建一条 WebSocket,而是所有任务摘要都发布到当前用户唯一的摘要频道。消息带上 taskIdkindversion,客户端再路由到不同本地对象:

function onPublication(message: TaskSummaryMessage) {
const current = taskStore.get(message.task.id);
if (!current || BigInt(message.task.version) > BigInt(current.version)) {
taskStore.set(message.task.id, message.task);
}
}

无论任务 A、B、C 谁先变化,都进入同一个连接;任务 D 刚刚创建,也会自动出现在频道中。Centrifugo 的服务器侧订阅正适合这种用户个人频道,且官方建议用个人频道承载用户更新,避免为大量对象逐个创建订阅。参见 Server-side subscriptionsFAQ: scalability considerations

6.2 任务详情频道

用户打开详情时按需订阅:

task_detail:{tenantId}/{taskId}

关闭详情时只取消这个订阅,应用级连接仍保留。详情消息携带任务本身和步骤的完整当前状态:

{
"kind": "task.detail",
"schemaVersion": 1,
"taskId": "8a4d",
"version": "12",
"taskType": "DOCUMENT_PROCESS",
"lifecycleStatus": "RUNNING",
"businessStatus": "PARSING",
"steps": [
{
"id": "01_prepare",
"operationCode": "FILE.PREPARE",
"lifecycleStatus": "SUCCEEDED",
"businessStatus": "PREPARED",
"durationMs": 3000
},
{
"id": "02_parse",
"operationCode": "DOCUMENT.PARSE",
"lifecycleStatus": "RUNNING",
"businessStatus": "PARSING",
"startedAt": "2026-09-21T07:10:03Z",
"data": {
"schema": "document.parse.progress.v1",
"value": {"currentPage": 18, "totalPages": 36}
}
},
{
"id": "03_store",
"operationCode": "RESULT.PERSIST",
"lifecycleStatus": "PENDING",
"businessStatus": "QUEUED"
}
]
}

6.3 为什么分两级

只用用户频道只用任务频道两级订阅
容易发现新任务,但所有步骤详情会广播过多详情精确,但新任务需要不断新增订阅列表自动发现,详情按需加载

用户有一万个历史任务时,摘要频道仍应只推发生变化的当前摘要;列表快照使用分页。详情频道只服务正在查看该任务的客户端。

6.4 客户端可选 WebSocket 或 SSE

传输方式不改变频道模型、鉴权或版本规则,只改变连接如何承载消息:

模式能力适用场景限制
WebSocket双向、动态订阅/取消、自动重连与 History 恢复Web、桌面和原生 App 的默认选择某些企业代理可能限制升级连接
SSE + HTTP emulation下行使用 SSE,上行命令使用 HTTP;JS SDK 仍可动态订阅WebSocket 不可用时的浏览器回退请求数和代理配置比 WebSocket 更复杂
原生单向 SSE / EventSource简单的服务器到客户端事件流固定用户摘要、只读大屏、轻量广播连接建立后不适合频繁动态增减详情订阅

在浏览器中,可以让 SDK 按顺序尝试 WebSocket 与双向 SSE emulation:

import { Centrifuge } from 'centrifuge';

const realtime = new Centrifuge(
[
{
transport: 'websocket',
endpoint: 'wss://realtime.example.com/connection/websocket',
},
{
transport: 'sse',
endpoint: 'https://realtime.example.com/connection/sse',
},
],
{
token: connectionToken,
getToken: refreshConnectionToken,
},
);

双向 SSE emulation 不是“只有下行的 EventSource”:SDK 通过 SSE 接收,通过 HTTP emulation 发送订阅等命令,因此仍能保留动态频道能力。原生单向 SSE 则适合连接时已经确定好的服务器侧订阅;如果页面需要频繁切换任务详情,优先使用 WebSocket 或双向 SSE emulation。参见 Transport overviewSSE bidirectional emulationUnidirectional SSE

6.5 任务详情订阅必须有限

“按需订阅”不等于“可以无限订阅”。推荐把限制分成三层:

  1. UI 只订阅当前打开、可见或固定关注的任务,页面离开后立即取消;
  2. 业务后端签发订阅 JWT 前检查任务 ACL、会话和应用级配额;
  3. 高风险场景启用订阅代理或由后端调用 Server API 动态订阅,把最终授权留在服务端。

例如产品可以配置 maxDetailSubscriptionsPerConnection = 20,这是应用策略,不是协议常量。超过上限时,客户端应先按 LRU 取消最早不可见的详情频道,再申请新频道;后端返回明确的 DETAIL_SUBSCRIPTION_LIMIT,不能静默放宽权限。

常驻:task_user:t1/u1 1 个
按需:task_detail:t1/task-a ... task-t 最多 20 个
总览:其余任务仍通过用户摘要频道持续收到变化

仅在客户端计数属于体验优化,不是安全边界。强制限制需要服务端维护连接或会话的订阅租约,并在签发 token、订阅代理或动态订阅时原子检查;短期订阅 JWT 还应绑定 sub + channel + exp。Centrifugo 的频道权限模式可参考 Channel permissions

6.6 广播支持:广播范围越大,载荷越小

系统可以支持广播,但必须先定义范围:

范围频道示例可以发送不应该发送
当前用户task_user:t1/u1该用户任务摘要其他用户数据
项目成员task_project:t1/p9项目级任务概览无权限成员的结果、密钥、完整错误栈
租户task_tenant:t1维护提示、租户级批处理进度单个用户的敏感任务详情
系统公共task_system:public版本通知、公开维护状态任何业务任务数据

同一个事件需要到达多个已授权频道时,Outbox Relay 可以调用 Centrifugo broadcast Server API;第一版也可以展开成多条 publish Outbox 记录,换取更简单的逐频道重试与审计。大规模广播应分批异步执行,避免一次命令过大阻塞普通任务通知。官方 Server API 支持对一组频道广播相同数据,见 Server API broadcast

{
"method": "broadcast",
"channels": [
"task_user:t1/u1",
"task_project:t1/p9"
],
"data": {
"kind": "task.summary",
"taskId": "8a4d",
"version": "12",
"lifecycleStatus": "RUNNING",
"businessStatus": "WAITING_REVIEW"
}
}

浏览器不能直接调用 Server API,也不能自己决定广播目标。业务服务先计算接收范围,事务内写入可审计的 Outbox 命令,Relay 才能执行发布。

6.7 一条连接为何能承载多个任务

Centrifugo 在一条传输连接上复用多个逻辑频道。连接身份保持不变,任务集合通过两条路径变化:

新任务出现
→ 用户摘要频道自动收到 task.summary
→ 不需要新增连接或订阅

用户打开某个任务
→ 申请该 task_detail 的短期订阅授权
→ 在现有连接上增加逻辑订阅
→ 关闭页面时取消,不影响其他任务

因此连接数大致随在线用户数增长,而不是随“在线用户数 × 历史任务数”增长;活跃详情订阅数则由页面行为和应用配额独立控制。

6.8 异常时如何让用户真正收到任务状态

异常状态不能只写进步骤表,更不能只发到“用户可能没有订阅”的详情频道。每次已确认异常至少生成两个投影:

  1. task_user:{tenantId}/{userId} 收到带异常摘要的完整 task.summary,保证用户只订阅一个常驻频道也能发现失败;
  2. task_detail:{tenantId}/{taskId} 收到步骤、功能数据和脱敏错误详情,供已经打开详情的页面展示。

用户摘要消息仍然是完整任务当前值,只额外包含一个可选 notice。这样客户端收到通知时可以同时更新列表、显示 Toast,并按 taskId + version 去重:

{
"kind": "task.summary",
"schemaVersion": 1,
"task": {
"id": "task-8a4d",
"taskType": "DOCUMENT_PROCESS",
"lifecycleStatus": "FAILED",
"businessStatus": "PARSE_FAILED",
"health": "HEALTHY",
"version": "13",
"failedStepId": "02_parse"
},
"notice": {
"id": "task-8a4d:13",
"type": "TASK_FAILED",
"level": "ERROR",
"title": "合同解析失败",
"message": "暂不支持该文件格式,请转换为 PDF 后重试",
"operationCode": "DOCUMENT.PARSE",
"errorCode": "UNSUPPORTED_FILE_FORMAT",
"retryable": false,
"requiresAction": true,
"actions": ["REUPLOAD_FILE", "VIEW_DETAILS"]
}
}

通知策略要由最终状态和异常类型驱动:

已确认事实用户摘要频道任务详情频道可选离线通知
STEP_FAILED,准备自动重试WARNING:“步骤失败,正在重试”完整失败 attempt 和下一轮信息通常不发
可选步骤失败、任务继续WARNING,任务标记 DEGRADED错误详情与受影响功能仅业务要求时发送
TASK_FAILEDERROR,必须发布且不能被进度合并吞掉最终错误、失败步骤和可执行操作可按用户偏好发送站内信、Push 或邮件
health=STALEWARNING:“状态暂时无法确认”最近心跳和检测时间通常不发,持续超阈值再升级
实时连接中断本地显示 RECOVERING/STALE重连后 History 或快照恢复不能把连接故障通知成任务失败

中间 STEP_PROGRESS 可以按短窗口合并,但 STEP_FAILEDTASK_FAILEDCANCELLEDSUCCEEDEDrequiresAction=true 的消息必须绕过普通进度合并。若 Relay 多次重试后进入 DEAD,这是运维告警,不是业务失败;客户端最终仍通过 MySQL 快照获得真实终态。

WebSocket/SSE 只覆盖在线连接。产品如果要求“用户离线也必须知道任务失败”,应把终态或需人工处理的事件同时写入独立、持久化的通知命令,由 Notification Service 生成站内未读记录并按用户偏好发送 Push/邮件。其幂等键可使用 userId + taskId + taskVersion + notificationType,避免 Relay 重试导致重复提醒。


七、认证与授权:频道名不是权限

7.1 连接 JWT

前端先用已有登录会话请求实时连接信息:

GET /api/realtime/session

后端从已验证身份中读取租户和用户,签发短期连接 JWT:

{
"sub": "t1/u1",
"channels": ["task_user:t1/u1"],
"aud": "task-realtime",
"iss": "task-status-service",
"exp": 1789981200
}

channels 会在连接建立时创建服务器侧订阅。官方特别说明,这个 claim 表示“自动订阅”,不是一张任意频道权限清单。参见 Client JWT authentication

7.2 订阅 JWT

打开任务详情时,前端向业务后端请求:

POST /api/realtime/subscription-token
Content-Type: application/json

{"taskId":"8a4d"}

后端必须重新执行任务权限检查:

验证登录身份
→ 查询任务 tenantId / ownerId / project ACL
→ 构造服务端认可的准确频道
→ 签发 channel + sub + exp
{
"sub": "t1/u1",
"channel": "task_detail:t1/8a4d",
"exp": 1789981200
}

订阅 JWT 的 sub 必须与连接用户一致,channel 必须是后端完成授权后构造的值。参见 Channel JWT authorization

下面这些做法都不安全:

  • 相信浏览器提交的 tenantIduserId 或完整频道名;
  • 把其他用户数据推到前端,再依赖 UI 隐藏;
  • 把 JWT 签名密钥或 Centrifugo Server API Key 发给浏览器;
  • 为解决 403,直接允许所有登录用户订阅任意任务频道。

八、客户端合并:快照和推送一定会乱序

8.1 每个对象只接受更高版本

type Versioned = {
id: string;
version: bigint;
};

function shouldApply(local: Versioned | undefined, incoming: Versioned) {
return !local || incoming.version > local.version;
}

版本建议以十进制字符串通过 JSON 传输,再在支持的客户端转换成 64 位整数或 bigint,避免 JavaScript number 的安全整数上限问题。

8.2 摘要与详情分别维护版本门

summaryVersion[taskId] = 12
detailVersion[taskId] = 10

收到 task.summary version=12,不能认为 task.detail version=12 已经处理,因为摘要并不包含步骤集合。摘要和详情要分别缓存、分别合并。

8.3 初始化采用“先订阅、后快照”

如果先查快照、再建立订阅,状态可能恰好在两者之间改变并被漏掉。订阅后缓冲、再加载快照、最后按版本重放,才能关闭这个竞态窗口。

8.4 重连采用两级恢复

Centrifugo 的双向 SDK 会保存频道 epoch/offset,在重订阅时尝试补齐短暂断线期间的消息;只有当历史连续可用时才返回恢复成功,否则应用应加载业务快照。官方恢复流程见 Stream history and recovery

对任务状态这种“当前值”系统,发布完整对象而不是差量补丁更容易恢复。即使漏掉:

PENDING → RUNNING → SUCCEEDED

只要最终收到更高版本的完整 SUCCEEDED 状态,客户端就能收敛;如果业务需要每个中间事件都可审计,应额外建设不可变事件表,而不是依赖实时频道历史。

8.5 状态合并和用户提醒要使用不同的去重门

任务状态按 taskId + version 合并;Toast、声音或系统通知则按 notice.id 去重。二者不能共用一个布尔值:页面可能已经通过快照得到 version=13,随后又收到同版本实时消息,此时不应回退状态,也不应重复弹出已读提醒。

function handleSummary(message: TaskSummaryMessage) {
mergeTaskWhenVersionIsNewer(message.task);

if (message.notice && !noticeStore.hasSeen(message.notice.id)) {
showNotice(message.notice);
noticeStore.markSeen(message.notice.id);
}
}

若“未读”需要跨设备同步,就把已读状态保存到持久化通知中心;浏览器本地集合只能防止单页面重复弹窗。重新加载历史快照时默认更新任务卡片,不重新弹出旧终态;只有持久化通知中心仍标记为未读时才展示未读入口。


九、耗时:后端确认,前端平滑显示

9.1 运行中不需要每秒推送

后端在步骤开始时发送:

{
"lifecycleStatus": "RUNNING",
"businessStatus": "PARSING",
"elapsedMs": 37000,
"asOf": "2026-09-21T07:10:37Z"
}

前端记录收到消息时的单调时钟:

const receivedAt = performance.now();
const baseElapsedMs = task.elapsedMs;

function displayedElapsedMs() {
return baseElapsedMs + performance.now() - receivedAt;
}

页面使用一个共享定时器刷新可见行即可。这不会产生网络请求、推送或数据库写入。

9.2 终态使用后端最终值

{
"lifecycleStatus": "SUCCEEDED",
"businessStatus": "COMPLETED",
"durationMs": 42000
}

收到终态后停止本地累加。浏览器从休眠或后台恢复、页面重新获得可见性、系统网络切换后,应重新获取当前状态校准。

9.3 明确耗时口径

指标推荐定义
排队耗时startedAt - createdAt
当前尝试耗时当前 attempt 的开始到结束
任务总耗时任务正式开始到权威终态的墙钟时间
Worker 执行耗时Worker 使用单调时钟测量并在终态上报
并行步骤耗时分别展示,不能直接相加为用户等待时间

跨机器时间戳适合表达时间线,但不适合在时钟未可靠同步时精确计算 Worker 内部耗时。精确执行时长应由同一进程的单调时钟测量。


十、通知可靠性:不要说“消息绝不丢”

这条链路包含不同层次的保证:

链路目标语义处理方式
Worker → 状态服务幂等重试eventId + payloadHash
状态服务 → MySQL原子提交当前状态、不可变事件、版本、去重结果、Outbox 同一个 InnoDB 事务
MySQL Outbox → Relay多副本竞争领取SKIP LOCKED、锁租约、失败退避
Relay → Centrifugo至少一次发布失败重试,崩溃窗口允许重复
Centrifugo → 在线客户端实时分发 + 有界恢复history、positioning、SDK recovery
客户端 → 最终视图最终收敛完整对象、version 去重、快照兜底

Relay 必须自己区分可重试错误和永久错误:网络超时、4295xx 进入带抖动的指数退避;无效频道或无法序列化的消息进入 DEAD 并告警,不能无限热循环。不能因为使用了 Outbox 就宣称整条链路“绝不丢”;官方也建议应用处理异步发布中的重复与晚到消息,详见 async consumers

真正可对外承诺的语义应该是:

状态服务成功响应,表示状态已持久化;实时通知可能重复、乱序或短暂中断,但客户端会通过版本和权威快照重新收敛。


十一、高可用:逐层消除单点

11.1 生产拓扑

11.2 Task Status Service

状态服务保持无状态并运行至少两个副本:

  • 登录身份和服务身份从可验证 Token 或 mTLS 获取;
  • 不把幂等结果、任务锁、订阅权限只保存在单机内存;
  • 每个请求有 deadline,数据库连接池有上限;
  • 就绪检查必须验证进程是否能够服务,但不要把瞬时下游抖动放大成全部 Pod 同时摘除;
  • 发布版本时先停止接收新流量,再等待在途事务完成。

并发更新同一个任务时,用任务行锁、乐观版本或按 taskId 分区的有序处理保证单任务串行;不要使用覆盖整个平台的全局锁。

11.3 MySQL

MySQL 保存唯一权威状态,因此它的 RPO/RTO 比实时网关更重要:

  • 使用 InnoDB、托管 Multi-AZ,或经过演练的 Group Replication / InnoDB Cluster;
  • 自动备份与时间点恢复不能省略;
  • 应用通过稳定 writer endpoint 连接,而不是写死节点 IP;
  • 故障切换期间,写入可以短暂失败,调用方按同一个 eventId 重试;
  • 定期做恢复演练,证明备份真的能还原,而不是只证明“备份任务显示成功”。

Outbox 与业务表在同一个数据库事务中,数据库切换后它们一起保留或一起回滚,不会出现跨两个存储提交一半的状态。

如果自建高可用,MySQL 官方说明 Group Replication 可以运行在 single-primary 模式并自动选举新主,但客户端仍需要 MySQL Router、代理或连接器切换到新节点;它本身不替应用完成连接故障转移。参见 Group Replication。时间点恢复则依赖完整备份后的 binary log,参见 Point-in-Time Recovery

11.4 Outbox Relay

Relay 至少运行两个副本,并把正确性放在 MySQL 中,而不是本机内存:

  • 使用 FOR UPDATE SKIP LOCKED 领取互不重叠的批次;
  • locked_at 是租约而不是永久所有权,超时记录可被其他副本接管;
  • 发布 API 使用短 deadline,失败后释放记录并延迟重试;
  • 单条永久错误进入死信状态,不阻塞后续记录;
  • 滚动发布时停止领取新批次,处理完已领取记录再退出;
  • 定期清理已发布记录,但保留足够窗口用于排障和审计。

Relay 不需要强制全局有序,因为消息包含完整状态与单调版本。若某个下游真的要求严格的单任务事件顺序,则需要再按 aggregate_id 做稳定分区并为每个分区维护单一租约。

11.5 Centrifugo

单节点验证可以使用 memory engine,但生产多副本需要共享 Engine。Centrifugo 官方说明,Redis engine 会在节点之间传播发布,并保存频道 history/presence;客户端连接任意节点都能接收目标频道消息。参见 Engines and scalability

示意配置如下:

{
"engine": {
"type": "redis",
"redis": {
"address": "redis+sentinel://redis-sentinel-a:26379?sentinel_master_name=realtime&addr=redis-sentinel-b:26379"
}
}
}

Relay 通过任意健康 Centrifugo 节点的 Server API 发布。共享 Redis Engine 会把消息转发到真正持有目标订阅的节点,因此 Relay 不需要知道某个浏览器连接在哪个 Pod。

11.6 Redis

两个 Centrifugo Pod 配一个单机 Redis,仍然只有一个新的单点。可选方式包括:

  • 托管高可用 Redis;
  • Redis Sentinel 主从切换;
  • Redis Cluster;
  • 规模更大时按频道进行应用级分片。

Sentinel 切换存在检测与选主窗口,因此高可用不等于零中断。连接侧必须允许重连,业务侧必须允许回快照。

11.7 负载均衡与连接排空

负载均衡器需要支持 WebSocket 升级和足够长的空闲超时。滚动发布时:

Pod 进入 terminating
→ readiness 变为 false
→ 不再接收新连接
→ 给存量连接一个排空窗口
→ 客户端重连到其他健康节点
→ SDK recovery 或业务快照补齐

不要求客户端永远粘在同一个节点;共享 Redis Engine 让新节点仍能访问同一频道分发与短期历史。粘性会话可以减少迁移,但不能成为正确性的前提。


十二、故障矩阵:每一层坏掉时会发生什么

故障用户看到什么系统如何恢复不应发生什么
一个 Status API Pod 崩溃某次请求失败或重试LB 切到其他 Pod,同 eventId 重试重试生成两次状态变化
MySQL 主库短暂切换上报延迟、查询短暂失败writer endpoint 恢复后重试,事务原子性保持返回成功但状态未落库
Outbox Relay 暂停数据库状态正确,页面更新延迟Relay 恢复后领取并清理积压为追求实时绕过事实库
一个 Centrifugo Pod 崩溃连接断开并重连连接其他 Pod,先 recovery,失败则快照把任务标成失败
Redis 主节点切换部分发布/恢复短暂失败Centrifugo 重连 Redis,客户端快照兜底依赖 history 作为唯一事实
浏览器休眠耗时暂时不刷新恢复可见时重新校准用错误本地时钟写回后端
消息重复无可见变化version 忽略同版或旧版重复新增任务行
消息乱序无可见回退仅应用更高版本终态退回 RUNNING
步骤业务失败列表收到警告或失败摘要,详情显示具体功能与错误Coordinator 按 failurePolicy 重试、继续或提交 TASK_FAILED只在详情频道发送,导致列表用户完全不知情
Worker 心跳超时显示“状态待确认”Health Reconciler 提交 health=STALE,恢复心跳后再提交 HEALTHY未确认事实时直接把任务写成 FAILED
Worker 完成后上报前崩溃长时间 RUNNING / STALE来源 Outbox、持久重试或对账修复认为中心 Outbox 能覆盖此缺口

高可用的目标不是让所有故障不可见,而是让故障有界、可观察、可恢复,并且不会破坏已确认事实


十三、降级策略:实时层不可用时仍然可用

13.1 客户端连接状态

UI 应区分:

LIVE 实时连接正常
RECOVERING 正在重连和恢复
STALE 暂时无法确认最新状态

不要在连接断开时清空任务列表。保留最后一次已确认状态,显示“正在恢复实时更新”,并在恢复失败时触发批量快照。

13.2 低频批量校验

如果部署阶段尚未启用可靠 history,或希望防止某些静默丢失,可对当前可见且未终止的任务做低频批量校验,例如页面活跃时每 30~60 秒一次:

POST /api/tasks/status:batchGet

{"taskIds":["8a4d","9b5e"],"knownVersions":{"8a4d":"12","9b5e":"3"}}

轮询是兜底,不是每秒驱动 UI 的主链路。间隔应通过业务容忍度和容量测试确定,并加入随机抖动,避免所有客户端整点同时请求。

13.3 过载时优先保终态

中间 RUNNING 状态可以短窗口合并,但以下消息不能被普通合并规则吞掉:

  • SUCCEEDEDFAILEDCANCELLED
  • 新任务创建;
  • 新 attempt 开始;
  • 权限撤销或任务删除;
  • 需要用户操作的错误。

慢客户端积压时,可以丢弃可覆盖的旧摘要,只保留同一任务最高版本的完整状态;超过上限后主动断开,让客户端走恢复流程,不能无限增加内存队列。


十四、可观测性:监控正确性,而不只监控在线数

14.1 服务端指标

必须监控的指标
状态 API上报 QPS、P50/P95/P99、幂等命中、冲突、旧 attempt、功能身份不匹配、Schema 拒绝、非法转换
MySQL事务延迟、InnoDB 行锁等待、连接池等待、复制延迟、切换事件
Outbox Relay最老未发布年龄、PENDING/PROCESSING/DEAD 数量、领取与发布速率、租约超时、失败分类
Centrifugo在线连接、订阅数、发布速率、节点间 Broker 错误、断开原因
RedisCPU、内存、连接、Pub/Sub、复制延迟、故障切换
客户端连接成功率、重连次数、recovered=false 比例、快照耗时、未知功能/Schema、异常通知展示延迟

比“当前在线 10 万连接”更有价值的两个端到端指标是:

state_commit_to_publish_seconds
状态提交到实时发布的延迟

state_commit_to_visible_seconds
状态提交到客户端可见的延迟

14.2 日志关联

一次状态变化至少应能用下面的字段串起来:

traceId
eventId
taskId
stepId
taskType
operationCode
dataSchema
attempt
sourceSeq
taskVersion
errorCode
outboxId
channel

日志中不要记录完整 JWT、签名密钥、Server API Key 或敏感任务结果。

14.3 告警优先级

建议先对“用户会长期看到错误状态”的问题告警:

  1. Outbox 最老记录年龄持续增长;
  2. 状态提交成功但发布延迟超出 SLO;
  3. recovered=false 比例突增;
  4. 未终止任务超过最大合理运行时长;
  5. 同一任务行锁等待和版本冲突异常升高;
  6. TASK_FAILED 已提交但用户摘要投影未发布;
  7. 某个 operationCode 的 Schema 拒绝或错误码突然激增;
  8. Redis 或 MySQL 故障切换失败。

连接数下降可能只是业务低谷,Outbox 延迟上升却很可能意味着用户正在看到旧状态。


十五、容量模型:连接、变化和写入分开计算

假设:

在线连接 100,000
活跃任务 20,000
每个任务平均 30 秒一次状态变化
每次变化平均投递给 2 个客户端
平均消息体 800 B

粗略估算:

状态写入:20,000 / 30 ≈ 667 次/秒
客户端投递:667 × 2 ≈ 1,334 次/秒
纯消息体流量:1,334 × 800 B ≈ 1.07 MB/秒

这还没有计入 TLS、协议头、心跳、重连和快照流量。连接数、状态变化量、数据库写入量和热门频道扇出是四个不同维度,不能用“支持百万 WebSocket”替代系统容量测试。

至少要压测以下场景:

  • 大量空闲连接;
  • 稳态状态上报;
  • 单个热门项目的大扇出;
  • Outbox 积压后追赶;
  • 一个 Centrifugo 节点滚动下线;
  • Redis 主从切换;
  • MySQL 主备切换;
  • 10%~30% 客户端同时重连;
  • history 不足后集中读取快照。

十六、分阶段落地,不要第一天就堆满组件

阶段 1:单机验证

现有业务后端
+ MySQL 8.0
+ 单实例 Outbox Relay
+ 单实例 Centrifugo memory engine

先验证:

  • 状态机与幂等;
  • taskType + operationCode + eventType 的功能路由与权限校验;
  • 各功能 data.schema 校验、差异化投影和未知 Schema 降级;
  • 多服务各自维护 sourceSeq,状态服务统一分配 taskVersion
  • 并发事件的任务版本唯一、递增,重复事件返回第一次分配的版本;
  • Outbox 原子性;
  • 用户摘要和任务详情频道;
  • 失败、自动重试、STALE 与终态异常通知;
  • 订阅权限;
  • 快照与实时消息合并;
  • 前端本地计时。

阶段 2:生产高可用

Status API 多副本
+ MySQL Multi-AZ
+ Outbox Relay 多副本
+ Centrifugo 多副本
+ Redis HA Engine
+ LB / Ingress

加入滚动发布、连接排空、短期 history、恢复失败回快照、端到端告警和故障演练。

阶段 3:规模化

当证据表明 MySQL Outbox 扫描、Relay 发布或 Redis 成为瓶颈,再考虑:

  • 增加 Outbox 分区;
  • 将消费链路升级到 Kafka、NATS JetStream、Redis Streams 或云消息服务;
  • 按租户或 taskId 分片;
  • 管理看板只订阅聚合指标,不接收所有任务明细;
  • 将审计历史与当前状态查询分离。

升级依据应该是积压、CPU、锁等待、恢复风暴和成本数据,而不是“系统看起来应该用 Kafka”。


十七、生产验收清单

正确性

  • 相同 eventId 重试不会产生第二次状态变化;
  • 相同 eventId 携带不同内容会被拒绝;
  • Token、步骤定义与请求中的 taskType/operationCode/producerId 不一致时会被拒绝;
  • 功能数据只接受注册过的 data.schema,未知 Schema 不会污染当前状态;
  • 同一功能的摘要、详情和广播投影遵守各自字段白名单;
  • 旧 attempt、旧 sourceSeq、旧 version 都不能覆盖新状态;
  • 同一任务中,服务 A 的 sourceSeq=2 不会让服务 B 的 sourceSeq=1 被误判为旧事件;
  • 两个服务并发上报时得到不同且连续递增的 taskVersion
  • 重复 eventId 返回第一次提交的 taskVersion,不会再次增加版本;
  • 终态不会退回运行态;
  • 未注册业务状态会被拒绝,业务状态只能映射到服务端定义的生命周期;
  • 状态定义升级不会改变运行中和历史任务的语义;
  • 当前状态、不可变事件、版本、幂等结果与 Outbox 在同一个事务里提交;
  • 事件时间线可以从 afterTaskVersion 断点续取,并能识别已被清理的历史窗口;
  • Worker 上报的单调时钟耗时不会被跨机器墙钟差值覆盖;
  • 摘要和详情分别维护版本门;
  • 初始化期间的新消息不会被旧快照覆盖。
  • 客户端不认识某个 operationCodedata.schema 时仍能展示通用生命周期状态。
  • 用户未订阅任务详情时,仍能从用户摘要频道收到任务失败或需操作提醒。
  • STALE 只改变健康状态,不会被误报成业务 FAILED
  • STEP_FAILED、终态和需操作通知不会被进度消息合并策略吞掉。

安全

  • 连接 JWT 来自已验证登录身份;
  • 任务详情 Token 在服务端重新检查权限;
  • 浏览器无法自行指定其他租户频道;
  • JWT 密钥、Server API Key 不出现在前端;
  • 内部上报有服务身份认证、限流和审计。
  • 任务详情订阅数在服务端强制限额,Token 短期有效;
  • 项目、租户和系统广播不包含越权任务详情。

高可用

  • Status API 任一副本退出后请求可重试;
  • Centrifugo 任一节点退出后客户端能重连其他节点;
  • Redis 与 MySQL 切换经过真实演练;
  • Outbox 暂停后可以追赶且不会回退客户端状态;
  • history 不足时客户端会读取快照;
  • 重连风暴有退避、抖动和容量余量;
  • 备份做过恢复验证。

可观测性

  • 能查出某个 eventId 对应的任务版本和 Outbox 记录;
  • 能看到 Outbox 最老积压年龄;
  • 能区分发布失败、权限失败、恢复失败和快照失败;
  • 有状态提交到客户端可见的端到端延迟指标;
  • 对长期 RUNNING / STALE 任务有对账和修复流程。

十八、最终结论

任务状态同步系统最重要的设计原则不是选 SSE 还是 WebSocket,而是建立清晰的可靠性边界:

状态服务成功响应
= 状态已经持久化

Outbox 发布
= 尽力把已提交状态快速送到实时层

实时推送
= 更快的用户体验,不是唯一事实

version + snapshot
= 重复、乱序、断线和节点切换后的最终正确性

在这个边界之上,Centrifugo 解决连接、频道、重连与多节点分发,MySQL 负责权威状态和事务,Outbox Relay 负责可靠发布,Redis 负责实时集群协调,前端负责按版本合并和本地计时。任何一层短暂不可用,都不会要求重新定义业务状态。

最终可以把方案收敛为十三句话:

  1. 状态先落库,再通知;
  2. 同一事务写当前状态、不可变事件、版本、幂等结果和 Outbox;
  3. 每个服务只维护自己的 sourceSeq,状态服务统一分配任务级 taskVersion
  4. taskType 区分整类任务,operationCode 区分具体功能,eventType 区分发生的动作;
  5. 公共信封保持稳定,功能差异通过版本化 data.schema + data.value 表达;
  6. 稳定生命周期保证正确性,版本化业务状态负责产品表达;
  7. 失败、失联和通知故障分别建模,不能统一伪装成 FAILED
  8. 用户摘要发现多个新任务和异常,任务详情按需且有限订阅;
  9. 一条连接复用多个频道,客户端按环境选择 WebSocket 或 SSE;
  10. 连接 JWT 识别用户,订阅 JWT 授权具体任务;
  11. 广播范围由服务端计算,范围越大,载荷越小;
  12. 消息可以重复和乱序,客户端只接受更高版本;
  13. 短断线用 history 恢复,恢复失败读取权威快照;高可用的目标是断线后仍能安全收敛。

做到这些,系统才真正同时具备轻量、近实时、可扩展和可恢复,而不是只在演示环境里“看起来很实时”。

Logo
RainLib

Exploring the frontiers of technology, design, and distributed systems. Building tools for the future developers.

Suggestions & Feedback

© 2026 RainLib. Built for the Future.
All rights reserved.
System Normal