状态、事件与复杂度
一个异步任务表最初可能只有三列:id、status、result。创建时写 pending,Worker 开始时写 running,结束时写 completed。后来加入取消、超时、重试、人工审核和恢复,代码逐渐出现下面的写法:
await db.task.update({ id, status: 'running' });// 另一个地方await db.task.update({ id, status: 'cancelled' });// 失败处理器await db.task.update({ id, status: 'pending' });// 管理命令await db.task.update({ id, status: requestedStatus });每一行单独看都合理,组合起来却没有一个位置回答:已取消任务还能不能完成?两个 Worker 同时启动时谁获胜?completed 后收到迟到的超时事件怎么办?重试是新任务还是同一任务的新尝试?恢复进程如何区分重复事件与缺失事件?
这不是“状态字段太多”的问题,而是写入没有协议。状态机(State Machine)状态机State Machine用有限状态和明确转换规则描述系统如何变化。打开术语条目 → 把所有合法变化集中到转换函数:当前状态与输入事件决定下一状态,非法组合返回明确拒绝。不变量(Invariant)不变量Invariant在所有合法状态与状态转换中都必须成立、可以被程序检查的条件。打开术语条目 → 则定义任何路径都不能破坏的业务事实。再配合 事件标识(Event ID)事件标识Event ID标识事件在一个有序事件流中的位置,用于断线续传、去重和回放。打开术语条目 →、版本检查、持久化事件和 快照(Snapshot)快照Snapshot某个已知版本上的完整状态副本,用于缩短恢复路径,但不能替代增量事件和校验。打开术语条目 →,任务生命周期才能测试、回放和恢复。
本课构建一个可测试任务状态机。它不把 Event Sourcing 当成所有系统的默认架构,而是借助事件记录解释状态变化,并明确什么时候只保存当前状态更简单。
1. 具体工程问题:迟到的完成覆盖了取消
Section titled “1. 具体工程问题:迟到的完成覆盖了取消”任务 task-42 的正常路径:
queued → running → completed用户在运行中取消:
queued → running → cancelling → cancelled真实系统里,取消不是瞬时的。控制服务先记录取消意图,向 Worker 发送信号;Worker 可能在信号到达前已经完成外部工作。于是两个请求并发发生:
请求 A:CancelTask(expectedVersion=8)请求 B:CompleteTask(expectedVersion=8, artifactHash=...)如果两个处理器只是执行 UPDATE tasks SET status = ? WHERE id = ?,后提交者覆盖前提交者。最终状态取决于时序,而不是业务规则。更糟的是,completed 或 cancelled 都可能缺少与之匹配的副作用证据。
我们先写出必须保持的事实:
completed、failed、cancelled是终态,不能被普通命令重新打开;completed必须携带经过验证的 Artifact 身份;cancelled表示执行已停止或已通过对账确认不会再产生新副作用;只记录取消请求时应处于cancelling;- 每个已提交事件都有唯一 Event ID;同一事件重放不能重复改变状态;
- 命令必须声明它读取的
expectedVersion;过期命令不能静默覆盖新状态; - 状态版本每接受一个新事件恰好增加一;拒绝和重复事件不增加版本。
状态列表应从这些不变量推导,而不是先画一堆方框再找含义。
命令、事件、状态与不变量
转换函数是唯一写入口2. 命令、事件、状态和观察不是同一种数据
Section titled “2. 命令、事件、状态和观察不是同一种数据”2.1 Command 表达意图,可能被拒绝
Section titled “2.1 Command 表达意图,可能被拒绝”StartTask、RequestCancellation、CompleteTask 是命令。它们来自用户、调度器、Worker 或恢复器,表达“希望系统发生什么”。命令可能因为权限、当前状态、版本冲突、参数无效或预算耗尽而被拒绝。
命令命名用祈使语义:
const command = { type: 'StartTask', taskId: 'task-42', expectedVersion: 3 } as const;不要把 TaskStarted 作为输入命令名称;这会把“请求开始”和“已经开始”混在一起。
2.2 Event 表达已经提交的事实
Section titled “2.2 Event 表达已经提交的事实”事件在事务提交后才成立:
const event = { eventId: 'evt-0004', type: 'TaskStarted', taskId: 'task-42', sequence: 4, attempt: 1,} as const;事件使用过去式。它不保证外部世界所有事情都成功,只保证系统定义的事实已经被权威存储接受。例如 CancellationRequested 只证明取消意图已登记,不证明 Worker 已停止;真正停止后再产生 TaskCancelled。
2.3 State 是事件按顺序折叠后的当前结论
Section titled “2.3 State 是事件按顺序折叠后的当前结论”状态可以被定义为:
Sₙ = fold(S₀, [E₁, E₂, …, Eₙ])其中 S0 是初始状态,E_i 是按任务序号排序的已接受事件。转换函数:
δ(Sᵢ₋₁, Eᵢ) → Sᵢ 或 Reject(reason)这里的 Reject 不能被忽略。若事件存储里出现无法应用的事件,说明历史损坏、版本不兼容或读取顺序错误;恢复器应停止并报告,而不是跳过后继续生成一个“尽量合理”的状态。
2.4 Observation 是可能过期的读模型
Section titled “2.4 Observation 是可能过期的读模型”列表页显示的进度、队列长度、预计完成时间通常是派生观察。它们可以来自缓存或异步投影,允许短暂滞后。命令处理器不能拿一个可能过期的列表页状态直接决定写入;它要读取权威版本,或在数据库条件更新中验证版本。
| 数据 | 时间语义 | 能否被拒绝 | 是否可重建 | 典型持久化 |
|---|---|---|---|---|
| Command | 希望发生 | 可以 | 不一定 | 请求日志、Inbox |
| Event | 已经提交 | 不再撤回,只能追加后续事实 | 是历史本身 | 追加事件表 |
| State | 截至当前的结论 | 不适用 | 可由事件重建,或直接存储 | 状态行、Snapshot |
| Observation | 面向查询的视图 | 不适用 | 可由权威状态重建 | 缓存、投影表 |
3. 从不变量推导状态集合
Section titled “3. 从不变量推导状态集合”本课任务状态:
export type TaskStatus = | 'queued' | 'running' | 'cancelling' | 'completed' | 'failed' | 'cancelled';为什么需要 cancelling?因为“收到取消请求”和“执行已停止”之间存在时间窗口。直接从 running 改成 cancelled 会制造错误事实:外部工具可能仍在写文件或发请求。
为什么没有 retrying?重试在本设计中是同一任务的下一次 attempt,它先通过 TaskFailed(retryable=true) 回到 queued,再接受新的 TaskStarted。如果产品需要展示退避等待,可以引入 waiting_retry 并保存 notBefore;状态应服务于可执行规则,而不是为了显示更丰富。
为什么 failed 是终态?本课把“可重试失败”建模为 AttemptFailed 事件,并回到 queued;TaskFailed 表示任务预算耗尽或错误不可恢复。使用不同事件避免“failed 有时终态、有时不是”的歧义。
3.1 状态转移表
Section titled “3.1 状态转移表”| 当前状态 | 输入事件 | Guard | 下一状态 | 必须附带的证据 |
|---|---|---|---|---|
queued | TaskStarted | attempt 等于上次 + 1 | running | workerId、leaseUntil、attempt |
queued | CancellationRequested | 无执行中副作用 | cancelled | requester、reason |
running | CancellationRequested | 请求身份有效 | cancelling | requester、reason、signalId |
running | TaskCompleted | Artifact 已验证;attempt 匹配 | completed | artifactHash、bytes、attempt |
running | AttemptFailed | retryable 且剩余预算 > 0 | queued | errorCode、attempt、nextAttempt |
running | TaskFailed | 不可重试或预算耗尽 | failed | errorCode、attempt |
cancelling | TaskCancelled | Worker 已停止或完成对账 | cancelled | stoppedAtSequence、effectCheck |
cancelling | TaskCompleted | 由业务策略决定 | completed 或拒绝 | Artifact、取消竞争裁决 |
| 任一终态 | 任一普通事件 | 无 | 拒绝 | TERMINAL_STATE |
cancelling + TaskCompleted 没有通用答案。如果完成的副作用不可撤回且 Artifact 已经产生,系统可能接受完成并记录取消过晚;如果取消承诺比完成优先,系统则拒绝完成并进入对账。规则必须在状态机中明确,不能留给数据库提交顺序决定。本课选择:只要 TaskCompleted 的 attempt 匹配,并且 Artifact 已验证,接受 completed,同时在 Trace 中记录 completionWonCancellationRace=true。
4. 状态转换函数:保持纯函数边界
Section titled “4. 状态转换函数:保持纯函数边界”核心转换函数不访问数据库、不读取当前时间、不调用工具。输入是完整当前状态和一个候选事件,输出是新状态或结构化拒绝:
export type TransitionResult = | { accepted: true; next: TaskState; evidence: TransitionEvidence } | { accepted: false; reason: RejectReason; current: TaskState };纯函数有三个直接收益:
- 可以枚举所有状态 × 事件组合;
- 可以对同一历史稳定回放;
- 数据库并发、消息去重和业务规则能分开测试。
副作用流程是“先判断可接受,再原子持久化事件和新状态”,而不是在纯函数里直接写库。
5. 一个可手算的状态历史
Section titled “5. 一个可手算的状态历史”初始状态:
status=queued, version=0, attempt=0按序事件:
E1 TaskStarted(attempt=1)E2 CancellationRequested(signalId=sig-9)E3 TaskCompleted(attempt=1, artifactHash=h1)E4 TaskCancelled(stoppedAtSequence=3)逐步折叠:
| 步骤 | 输入状态 | 事件 | 结果 |
|---|---|---|---|
| 0 | queued/v0/a0 | — | 初始状态 |
| 1 | queued/v0/a0 | TaskStarted(a1) | running/v1/a1 |
| 2 | running/v1/a1 | CancellationRequested | cancelling/v2/a1 |
| 3 | cancelling/v2/a1 | TaskCompleted(a1,h1) | completed/v3/a1,完成赢得竞争 |
| 4 | completed/v3/a1 | TaskCancelled | 拒绝:TERMINAL_STATE,版本保持 3 |
注意 E4 即使有更大的消息到达时间,也不能覆盖 E3。若 E4 已经被错误写入事件存储,重放器应报告非法历史;若它只是队列中的候选消息,Inbox 处理器记录拒绝后确认消息即可。
再看重复事件:E3 因 Worker 没收到确认而重投,Event ID 仍是同一个。处理器先查 Inbox,发现已处理,返回第一次结果的摘要;不再次调用转换函数,不增加版本。若 Worker 使用新 Event ID 重发同一业务事实,系统还需要业务去重键,例如 (taskId, type, attempt) 唯一约束。Event ID 去重只防止同一消息重复,不能防止发送方制造语义重复。
6. 乐观并发:过期写入必须失败
Section titled “6. 乐观并发:过期写入必须失败”乐观并发控制(Optimistic Concurrency Control)乐观并发控制Optimistic Concurrency Control提交写入时检查读取版本是否仍然有效,冲突时拒绝并要求重新计算。打开术语条目 → 让命令声明它基于哪个版本做决定。数据库更新可以写成:
UPDATE task_stateSET status = $1, version = version + 1, state_json = $2WHERE task_id = $3 AND version = $4;受影响行数为 0 时,不代表数据库坏了;它表示状态已被其他事务推进。处理器应重新读取,然后重新评估命令。不能直接把 expectedVersion 改成最新值再强行写入,因为新状态下原命令可能已非法。
以取消与完成竞争为例:
A读取 running/v8,准备 CancellationRequestedB读取 running/v8,准备 TaskCompletedA先提交:cancelling/v9B条件更新 version=8,影响0行B重新读取 cancelling/v9B按明确策略重新计算,接受 TaskCompleted → completed/v10最终结果由 cancelling + TaskCompleted 的规则决定,而不是“最后提交者覆盖”。
6.1 为什么事务隔离级别仍然重要
Section titled “6.1 为什么事务隔离级别仍然重要”版本条件更新避免丢失更新,但命令可能还读取其他行,例如剩余配额、租约、Artifact 记录。若业务不变量跨多行,必须让检查与写入处于合适事务边界,并处理序列化冲突或唯一约束冲突。不能因为用了 version 就声称整个业务操作线性一致。
本课来源中的 PostgreSQL 隔离文档用于核对各隔离级别的数据库语义;实际实现要针对所用数据库验证,不应把示例 SQL 直接外推到不同存储。
7. 事件乱序与两种 Sequence
Section titled “7. 事件乱序与两种 Sequence”分布式系统常同时出现两种顺序:
- 任务内权威序号:事件被权威存储接受时分配的
sequence;用于重放; - 消息到达顺序:消费者看到消息的先后;可能重复、延迟或乱序。
如果消息 sequence=12 先于 sequence=11 到达,消费者有三种策略:
| 策略 | 适用条件 | 风险 |
|---|---|---|
| 拒绝并稍后重试 | 上游保证缺失消息最终会到 | 长时间阻塞、重复重试 |
| 暂存 Gap Buffer | 乱序窗口有明确上限 | Buffer 膨胀,需要超时和补拉 |
| 从权威存储按序补拉 | 有按任务查询历史的接口 | 增加读取,需防止权限与版本错误 |
本课处理器遇到 sequence > current.version + 1 时返回 SEQUENCE_GAP,由恢复器从事件表补拉。遇到 sequence <= current.version 时不能直接判定重复,要先查 Event ID:
- Event ID 已处理:幂等返回;
- Event ID 未处理但序号已占用:历史冲突,进入隔离;
- 同一 Event ID 对应不同 Payload Hash:协议破坏,不能选择其中一个继续。
8. 完整 TypeScript 状态机
Section titled “8. 完整 TypeScript 状态机”下面是可运行的核心实现,不含具体数据库适配器。存储接口、内存实现和断言测试都给出,输入输出明确,不执行任意代码。
import assert from 'node:assert/strict';import { createHash } from 'node:crypto';
export type TaskStatus = | 'queued' | 'running' | 'cancelling' | 'completed' | 'failed' | 'cancelled';
export interface TaskState { taskId: string; status: TaskStatus; version: number; attempt: number; workerId: string | null; leaseUntilEpochMs: number | null; cancellationSignalId: string | null; artifactSha256: string | null; terminalErrorCode: string | null;}
interface EventBase { eventId: string; taskId: string; sequence: number;}
export type TaskEvent = | (EventBase & { type: 'TaskStarted'; attempt: number; workerId: string; leaseUntilEpochMs: number; }) | (EventBase & { type: 'CancellationRequested'; signalId: string; requester: 'user' | 'system'; reason: string; }) | (EventBase & { type: 'TaskCompleted'; attempt: number; artifactSha256: string; artifactVerified: true; }) | (EventBase & { type: 'AttemptFailed'; attempt: number; errorCode: string; retryable: true; nextAttempt: number; }) | (EventBase & { type: 'TaskFailed'; attempt: number; errorCode: string; }) | (EventBase & { type: 'TaskCancelled'; stoppedAtSequence: number; effectsReconciled: true; });
export type RejectReason = | 'TASK_ID_MISMATCH' | 'DUPLICATE_OR_OLD_SEQUENCE' | 'SEQUENCE_GAP' | 'TERMINAL_STATE' | 'ILLEGAL_TRANSITION' | 'ATTEMPT_MISMATCH' | 'INVALID_NEXT_ATTEMPT' | 'ARTIFACT_NOT_VERIFIED' | 'EFFECTS_NOT_RECONCILED';
export type TransitionResult = | { accepted: true; next: TaskState; evidence: { previousStatus: TaskStatus; nextStatus: TaskStatus; eventId: string; completionWonCancellationRace: boolean; }; } | { accepted: false; reason: RejectReason; current: TaskState };
export function initialState(taskId: string): TaskState { return { taskId, status: 'queued', version: 0, attempt: 0, workerId: null, leaseUntilEpochMs: null, cancellationSignalId: null, artifactSha256: null, terminalErrorCode: null, };}
function reject(current: TaskState, reason: RejectReason): TransitionResult { return { accepted: false, reason, current };}
export function transition(current: TaskState, event: TaskEvent): TransitionResult { if (event.taskId !== current.taskId) return reject(current, 'TASK_ID_MISMATCH'); if (event.sequence <= current.version) return reject(current, 'DUPLICATE_OR_OLD_SEQUENCE'); if (event.sequence !== current.version + 1) return reject(current, 'SEQUENCE_GAP'); if (['completed', 'failed', 'cancelled'].includes(current.status)) { return reject(current, 'TERMINAL_STATE'); }
const base = { ...current, version: event.sequence }; let next: TaskState | null = null; let completionWonCancellationRace = false;
switch (event.type) { case 'TaskStarted': { if (current.status !== 'queued') return reject(current, 'ILLEGAL_TRANSITION'); if (event.attempt !== current.attempt + 1) return reject(current, 'ATTEMPT_MISMATCH'); next = { ...base, status: 'running', attempt: event.attempt, workerId: event.workerId, leaseUntilEpochMs: event.leaseUntilEpochMs, cancellationSignalId: null, }; break; } case 'CancellationRequested': { if (current.status === 'queued') { next = { ...base, status: 'cancelled', cancellationSignalId: event.signalId, workerId: null, leaseUntilEpochMs: null, }; } else if (current.status === 'running') { next = { ...base, status: 'cancelling', cancellationSignalId: event.signalId }; } else { return reject(current, 'ILLEGAL_TRANSITION'); } break; } case 'TaskCompleted': { if (current.status !== 'running' && current.status !== 'cancelling') { return reject(current, 'ILLEGAL_TRANSITION'); } if (event.attempt !== current.attempt) return reject(current, 'ATTEMPT_MISMATCH'); if (event.artifactVerified !== true) return reject(current, 'ARTIFACT_NOT_VERIFIED'); completionWonCancellationRace = current.status === 'cancelling'; next = { ...base, status: 'completed', artifactSha256: event.artifactSha256, workerId: null, leaseUntilEpochMs: null, }; break; } case 'AttemptFailed': { if (current.status !== 'running') return reject(current, 'ILLEGAL_TRANSITION'); if (event.attempt !== current.attempt) return reject(current, 'ATTEMPT_MISMATCH'); if (event.nextAttempt !== current.attempt + 1) return reject(current, 'INVALID_NEXT_ATTEMPT'); next = { ...base, status: 'queued', workerId: null, leaseUntilEpochMs: null, terminalErrorCode: null, }; break; } case 'TaskFailed': { if (current.status !== 'running' && current.status !== 'cancelling') { return reject(current, 'ILLEGAL_TRANSITION'); } if (event.attempt !== current.attempt) return reject(current, 'ATTEMPT_MISMATCH'); next = { ...base, status: 'failed', workerId: null, leaseUntilEpochMs: null, terminalErrorCode: event.errorCode, }; break; } case 'TaskCancelled': { if (current.status !== 'cancelling') return reject(current, 'ILLEGAL_TRANSITION'); if (event.effectsReconciled !== true) return reject(current, 'EFFECTS_NOT_RECONCILED'); next = { ...base, status: 'cancelled', workerId: null, leaseUntilEpochMs: null, }; break; } }
return { accepted: true, next, evidence: { previousStatus: current.status, nextStatus: next.status, eventId: event.eventId, completionWonCancellationRace, }, };}
function canonicalEventHash(event: TaskEvent): string { const keys = Object.keys(event).sort(); const normalized = Object.fromEntries(keys.map((key) => [key, event[key as keyof TaskEvent]])); return createHash('sha256').update(JSON.stringify(normalized)).digest('hex');}
interface StoredEvent { event: TaskEvent; payloadSha256: string;}
export interface TaskStore { load(taskId: string): Promise<{ state: TaskState; events: StoredEvent[] }>; findProcessedEvent(eventId: string): Promise<StoredEvent | null>; appendAtomically(input: { expectedVersion: number; event: TaskEvent; payloadSha256: string; nextState: TaskState; }): Promise<'committed' | 'version-conflict' | 'event-id-conflict'>;}
export type ApplyResult = | { status: 'committed'; state: TaskState } | { status: 'duplicate'; state: TaskState } | { status: 'rejected'; reason: RejectReason; state: TaskState } | { status: 'conflict'; code: 'VERSION_CONFLICT' | 'EVENT_ID_PAYLOAD_CONFLICT'; state: TaskState };
export async function applyEvent(store: TaskStore, event: TaskEvent): Promise<ApplyResult> { const payloadSha256 = canonicalEventHash(event); const processed = await store.findProcessedEvent(event.eventId); if (processed) { const current = await store.load(event.taskId); return processed.payloadSha256 === payloadSha256 ? { status: 'duplicate', state: current.state } : { status: 'conflict', code: 'EVENT_ID_PAYLOAD_CONFLICT', state: current.state }; }
const current = await store.load(event.taskId); const result = transition(current.state, event); if (!result.accepted) { return { status: 'rejected', reason: result.reason, state: current.state }; }
const committed = await store.appendAtomically({ expectedVersion: current.state.version, event, payloadSha256, nextState: result.next, }); if (committed === 'version-conflict') { return { status: 'conflict', code: 'VERSION_CONFLICT', state: (await store.load(event.taskId)).state }; } if (committed === 'event-id-conflict') { return { status: 'conflict', code: 'EVENT_ID_PAYLOAD_CONFLICT', state: (await store.load(event.taskId)).state }; } return { status: 'committed', state: result.next };}
export class InMemoryTaskStore implements TaskStore { private readonly states = new Map<string, TaskState>(); private readonly eventsByTask = new Map<string, StoredEvent[]>(); private readonly eventsById = new Map<string, StoredEvent>();
constructor(taskIds: string[]) { for (const taskId of taskIds) this.states.set(taskId, initialState(taskId)); }
async load(taskId: string): Promise<{ state: TaskState; events: StoredEvent[] }> { const state = this.states.get(taskId); if (!state) throw new Error(`unknown task: ${taskId}`); return { state: structuredClone(state), events: structuredClone(this.eventsByTask.get(taskId) ?? []), }; }
async findProcessedEvent(eventId: string): Promise<StoredEvent | null> { return structuredClone(this.eventsById.get(eventId) ?? null); }
async appendAtomically(input: { expectedVersion: number; event: TaskEvent; payloadSha256: string; nextState: TaskState; }): Promise<'committed' | 'version-conflict' | 'event-id-conflict'> { const current = this.states.get(input.event.taskId); if (!current || current.version !== input.expectedVersion) return 'version-conflict'; const existing = this.eventsById.get(input.event.eventId); if (existing) return 'event-id-conflict';
const stored = { event: structuredClone(input.event), payloadSha256: input.payloadSha256 }; this.eventsById.set(input.event.eventId, stored); this.eventsByTask.set(input.event.taskId, [...(this.eventsByTask.get(input.event.taskId) ?? []), stored]); this.states.set(input.event.taskId, structuredClone(input.nextState)); return 'committed'; }}
export function replay(taskId: string, events: TaskEvent[]): TaskState { let state = initialState(taskId); const seen = new Set<string>(); for (const event of [...events].sort((a, b) => a.sequence - b.sequence)) { if (seen.has(event.eventId)) continue; seen.add(event.eventId); const result = transition(state, event); if (!result.accepted) { throw new Error(`invalid history at ${event.eventId}: ${result.reason}`); } state = result.next; } return state;}
async function selfTest(): Promise<void> { const store = new InMemoryTaskStore(['task-42']); const started: TaskEvent = { eventId: 'evt-1', taskId: 'task-42', sequence: 1, type: 'TaskStarted', attempt: 1, workerId: 'worker-a', leaseUntilEpochMs: 1000, }; assert.equal((await applyEvent(store, started)).status, 'committed'); assert.equal((await applyEvent(store, started)).status, 'duplicate');
const cancel: TaskEvent = { eventId: 'evt-2', taskId: 'task-42', sequence: 2, type: 'CancellationRequested', signalId: 'sig-9', requester: 'user', reason: 'no longer needed', }; assert.equal((await applyEvent(store, cancel)).state.status, 'cancelling');
const completed: TaskEvent = { eventId: 'evt-3', taskId: 'task-42', sequence: 3, type: 'TaskCompleted', attempt: 1, artifactSha256: 'a'.repeat(64), artifactVerified: true, }; assert.equal((await applyEvent(store, completed)).state.status, 'completed');
const lateCancel: TaskEvent = { eventId: 'evt-4', taskId: 'task-42', sequence: 4, type: 'TaskCancelled', stoppedAtSequence: 3, effectsReconciled: true, }; const rejected = await applyEvent(store, lateCancel); assert.equal(rejected.status, 'rejected'); if (rejected.status === 'rejected') assert.equal(rejected.reason, 'TERMINAL_STATE');
const loaded = await store.load('task-42'); assert.deepEqual(replay('task-42', loaded.events.map((item) => item.event)), loaded.state);}
await selfTest();内存存储只用于确定性教学测试。真实数据库适配器必须用一个事务同时:验证版本、插入 Event ID 唯一记录、追加事件、更新当前状态。若这些步骤分开提交,就可能出现事件已存但状态未更新,或状态已更新但事件丢失。
9. 数据库表与原子提交
Section titled “9. 数据库表与原子提交”简化 PostgreSQL Schema:
CREATE TABLE task_state ( task_id text PRIMARY KEY, status text NOT NULL, version bigint NOT NULL CHECK (version >= 0), attempt integer NOT NULL CHECK (attempt >= 0), state_json jsonb NOT NULL);
CREATE TABLE task_event ( task_id text NOT NULL, sequence bigint NOT NULL, event_id text NOT NULL UNIQUE, event_type text NOT NULL, payload jsonb NOT NULL, payload_sha256 text NOT NULL, created_at timestamptz NOT NULL DEFAULT now(), PRIMARY KEY (task_id, sequence), FOREIGN KEY (task_id) REFERENCES task_state(task_id));提交伪代码:
BEGIN;
SELECT versionFROM task_stateWHERE task_id = $task_idFOR UPDATE;
-- 应用层再次检查version与转换结果。INSERT INTO task_event (...)VALUES (...);
UPDATE task_stateSET status = $next_status, version = $next_version, attempt = $next_attempt, state_json = $next_stateWHERE task_id = $task_id AND version = $expected_version;
-- 若UPDATE影响0行,ROLLBACK并返回VERSION_CONFLICT。COMMIT;使用 FOR UPDATE 和版本条件都可以控制竞争;具体选型取决于事务长度、冲突率与数据库。即使用行锁,也要保留唯一约束来防 Event ID 重复。应用层预检查提升错误质量,数据库约束负责最后防线。
10. Snapshot 与回放
Section titled “10. Snapshot 与回放”快照(Snapshot)快照Snapshot某个已知版本上的完整状态副本,用于缩短恢复路径,但不能替代增量事件和校验。打开术语条目 → 保存某个事件序号处的状态,减少每次恢复需要折叠的事件数:
interface TaskSnapshot { taskId: string; throughSequence: number; state: TaskState; stateSha256: string; reducerVersion: string;}恢复算法:
- 找到不超过目标序号的最新可信 Snapshot。
- 验证 Snapshot 的 Schema、Hash、taskId、
state.version === throughSequence。 - 从
throughSequence + 1开始按序读取事件。 - 对每个事件执行同一个纯转换函数;遇到 Gap 或非法转换立即停止。
- 将重建状态与当前状态行比较;不同则生成差异报告,不自动选择“看起来更新”的一方。
Snapshot 不是事实来源的替代品。若只保留 Snapshot 并删除事件,就失去完整回放能力;这可能仍是合理存储策略,但必须明确恢复与审计边界。事件保留也不是免费:Schema 演进、隐私删除、存储增长和重放时间都要管理。
10.1 Reducer 演进
Section titled “10.1 Reducer 演进”旧事件不能因为代码升级就改变语义。常见策略:
- 事件携带版本,Reducer 兼容旧版本;
- 在读取边界把旧事件上转换成内部新格式;
- 进行一次可验证迁移,保留迁移前后 Hash 和映射报告;
- Snapshot 记录
reducerVersion,不兼容时从更早历史重建。
不要在 Event Payload 上就地修改历史然后称为“同一事件”。如果因合规要求必须删除字段,应记录明确的擦除或重写流程,并说明之后哪些审计能力不再存在。
11. 复杂度不是状态数量本身
Section titled “11. 复杂度不是状态数量本身”“有六个状态,所以系统简单”是错误判断。真正影响实现和测试的至少有:
- 合法转换数量;
- 每条转换的 Guard 数量;
- 并发命令的竞争组合;
- 外部副作用与内部状态的一致性边界;
- 事件乱序窗口;
- 超时、重试和人工操作是否增加正交维度;
- 状态与事件 Schema 的版本数量。
可以用一个工程计数帮助发现组合爆炸,但不要把它当成普适复杂度定律。假设状态集合为 S,事件类型集合为 E,完整枚举上界是:
N_pairs = |S| × |E|本课 6 个状态、6 个事件类型,至少有 36 个状态—事件组合需要明确“接受或拒绝”。如果每条接受路径还要测试“正确版本、过期版本、重复 Event ID、序号 Gap”,测试维度进一步扩展。
更实用的指标:
transition coverage = 已测试的状态-事件分类 / 已声明的状态-事件分类rejection coverage = 已测试的拒绝原因 / 全部拒绝原因race coverage = 已测试的竞争对 / 风险清单中的竞争对这些指标只说明测试覆盖了模型声明,不证明模型声明完整。新故障仍可能揭示缺失状态。
11.1 何时拆成正交状态
Section titled “11.1 何时拆成正交状态”如果把任务状态写成:
running-with-cancel-requested-and-artifact-uploading很快会得到笛卡尔积。可以拆成:
type OrthogonalTaskState = { lifecycle: 'queued' | 'running' | 'terminal'; cancellation: 'none' | 'requested' | 'confirmed'; artifact: 'none' | 'writing' | 'verified' | 'failed';};但拆分后必须有跨区域不变量,例如 lifecycle=terminal.completed 要求 artifact=verified。正交状态减少名称爆炸,却增加组合约束。选择依据是各维度是否真的可以独立变化,以及是否能写出完整不变量。
12. 恢复策略:先观察,再推进
Section titled “12. 恢复策略:先观察,再推进”任务卡在 running 且租约过期时,恢复器不能直接写 queued。它至少要回答:
- Worker 是否还活着?
- 外部副作用是否已经发生?
- Artifact 是否存在且属于当前 attempt?
- 最后一个权威事件是什么?
- 是否有未发布或乱序消息?
恢复命令可以产生不同事件:
LeaseExpiredObservedArtifactRecoveredAttemptFailed(retryable=true)TaskFailed(errorCode=UNKNOWN_EFFECT)HumanReviewRequested本课核心状态机没有把所有恢复事件展开,是为了保持示例可读;真实产物应根据可观察事实增加事件,而不是让恢复器直接修改 status。
13. 故障诊断矩阵
Section titled “13. 故障诊断矩阵”| 故障 | 检测证据 | 状态机响应 | 不应做什么 | 恢复策略 |
|---|---|---|---|---|
| 同一 Event ID 重投 | Event ID 与 Payload Hash 均相同 | 返回 duplicate,版本不变 | 再次应用转换 | 返回首次处理摘要并确认消息 |
| 同一 Event ID、不同 Payload | Event ID 相同,Hash 不同 | EVENT_ID_PAYLOAD_CONFLICT | 任选一个 Payload | 隔离发送方与消息,人工调查协议破坏 |
| 事件序号跳跃 | sequence > version + 1 | SEQUENCE_GAP | 跳过缺失序号继续 | 从权威事件存储补拉;超时后告警 |
| 过期命令 | 条件更新影响 0 行 | VERSION_CONFLICT | 把 expectedVersion 改新后强写 | 重新读取并重新执行业务 Guard |
| 终态收到迟到事件 | 当前状态属于终态 | TERMINAL_STATE | 把状态重新打开 | 记录拒绝;需要新工作时创建新任务或显式恢复命令 |
| Worker 租约过期 | 当前时间超过租约,仅是观察 | 不直接转状态 | 立即重试外部副作用 | 查询 Worker、Artifact、Ledger 后生成恢复事件 |
| Snapshot 与事件重放不一致 | 同序号状态 Hash 不同 | 停止恢复 | 选择字段更多的一方 | 固定 Reducer 版本,重放并生成差异报告 |
| 事件存储成功、投影失败 | 权威历史已推进,读模型滞后 | 权威状态不回滚 | 重新提交业务事件 | 重放投影或从 Snapshot 重建 |
14. 至少三个故障注入实验
Section titled “14. 至少三个故障注入实验”实验一:重复 Event ID
Section titled “实验一:重复 Event ID”- 对初始任务提交
evt-1 TaskStarted; - 再提交完全相同事件;
- 断言第二次结果为
duplicate; - 断言状态仍为
running/version=1/attempt=1; - 把同一
eventId的workerId改掉后提交; - 断言返回
EVENT_ID_PAYLOAD_CONFLICT; - 断言事件表中仍只有一条
evt-1。
验收:字节语义相同的重投幂等返回,身份相同但内容不同的消息被隔离。
实验二:取消与完成并发
Section titled “实验二:取消与完成并发”- 把任务推进到
running/version=8; - 构造两个基于 version 8 的命令;
- 让取消事务先提交到
cancelling/version=9; - 完成事务条件更新失败;
- 重新读取 version 9,再按
cancelling + TaskCompleted规则计算; - 断言最终为
completed/version=10,且 Transition Evidence 标记竞争; - 再切换策略为“取消优先”,确认同一场景稳定拒绝完成,而不是由时序随机决定。
验收:两种策略都必须显式、确定、可测试。
实验三:乱序与 Gap
Section titled “实验三:乱序与 Gap”- 正常提交 sequence 1;
- 先送达 sequence 3;
- 断言返回
SEQUENCE_GAP,状态不变; - 补拉并提交 sequence 2;
- 再提交 sequence 3;
- 断言状态依次推进,没有丢弃任何事件;
- 构造 sequence 2 但新 Event ID,断言为历史冲突而不是普通重复。
验收:到达顺序不决定权威顺序,缺口不能被静默跳过。
实验四:Snapshot 损坏
Section titled “实验四:Snapshot 损坏”- 在 version 50 保存 Snapshot;
- 修改 Snapshot 中
attempt,但不更新 Hash; - 恢复器应在重放前拒绝 Snapshot;
- 再更新 Hash 模拟“内部一致但语义错误”的 Snapshot;
- 从事件 1 全量重放,与 Snapshot 对比;
- 差异应进入报告,不自动覆盖事件历史;
- 固定
reducerVersion后重建可信 Snapshot。
验收:Snapshot 是加速器,不是跳过历史验证的许可证。
15. 决策表:当前状态表、事件历史还是二者都要
Section titled “15. 决策表:当前状态表、事件历史还是二者都要”| 需求 | 只存当前状态 | 存事件 + 当前状态 | 主要代价 |
|---|---|---|---|
| 简单 CRUD,历史无业务价值 | 通常足够 | 可能过度设计 | 事件 Schema 与迁移成本 |
| 需要解释“为什么到这里” | 需要额外审计表 | 事件天然提供因果记录 | 日志保留和隐私治理 |
| 需要重建多个读模型 | 需要另写变更日志 | 可由历史投影 | 重放时间和版本兼容 |
| 需要撤销 | 保存前值或补偿信息 | 仍需显式反向事件 | 事件历史不会自动可逆 |
| 高频计数更新 | 条件更新更直接 | 事件量可能很大 | 存储与聚合压力 |
| 跨服务流程 | 每个服务维护本地状态 | 事件可连接流程,但不提供全局事务 | 顺序、幂等、对账复杂度 |
Event Sourcing 不是“有事件消息就算”。如果数据库只保存最终状态、消息只是通知,仍然是状态存储加事件发布;这完全可以是更合适的设计。关键是明确事实来源和恢复路径。
16. 产物验收条件
Section titled “16. 产物验收条件”可测试任务状态机应满足:
- □ 命令、事件、状态、观察有不同类型和命名语义;
- □ 所有状态 × 事件组合都有接受或拒绝分类;
- □ 终态不能被普通事件重新打开;
- □
completed强制要求已验证 Artifact; - □
cancelled与cancelling的事实边界明确; - □ 每个事件具有 Event ID、任务内 Sequence 与 Payload Hash;
- □ 相同 Event ID 的相同 Payload 幂等返回,不增加版本;
- □ 相同 Event ID 的不同 Payload 被分类为协议冲突;
- □ 版本条件写入与事件追加处于一个事务;
- □ Gap、过期命令、迟到事件和终态事件都有测试;
- □ 回放器遇到非法历史会停止,不跳过;
- □ Snapshot 记录序号、Hash 与 Reducer 版本;
- □ 至少完成三个故障注入实验并保存 Transition Evidence;
- □ Trace 中不包含密钥、真实账号数据或无必要的输入全文。
17. 学完后应该能回答什么
Section titled “17. 学完后应该能回答什么”- Command 与 Event 的时间语义为什么不同?
- 为什么
CancellationRequested不能直接等同于TaskCancelled? - 如何从业务不变量推导状态,而不是先列状态名称?
δ(state,event)返回拒绝为什么与返回下一状态同样重要?- Event ID 去重和业务语义去重分别防什么?
- 为什么过期写入不能简单更新 expectedVersion 后重试?
- 任务内 Sequence 与消息到达顺序有什么区别?
- 遇到 Sequence Gap 有哪些策略,各自边界是什么?
- 为什么 Snapshot 必须记录 Reducer 版本和 Hash?
- 什么时候只存当前状态更合适,什么时候需要完整事件历史?
- 状态数量为什么不能单独代表复杂度?
- 租约超时后,恢复器为什么要先观察副作用再推进状态?
18. 来源边界
Section titled “18. 来源边界”- Statecharts 论文用于理解层级、并发和事件驱动状态模型的表达能力;本课使用的是较小的有限状态机,不声称实现完整 Statecharts 语义。
- Event Sourcing 材料用于区分事件历史、当前状态和投影;是否采用完整事件溯源仍取决于审计、恢复、隐私和运维成本。
- PostgreSQL 事务隔离文档用于核对示例数据库的并发边界;不同数据库的锁、唯一约束、序列化失败和事务语义必须按其官方文档验证。
- 所有 taskId、workerId、Event ID、时间和 Hash 都是抽象教学数据,没有真实账号或私有系统事实。