事务、缓存与队列
一个创建异步任务的 API 做了三件事:写数据库、把消息发到队列、返回任务 ID。正常路径只有十几行代码:
const task = await db.tasks.insert(input);await broker.publish({ taskId: task.id });return task;但任意两个语句之间都可能中断:数据库提交成功后进程崩溃,任务永远没有消息;消息发布成功后响应断线,客户端用新请求重试,生成两个任务;Worker 完成远程写入后确认消息失败,同一副作用再次执行;缓存还显示 queued,管理操作基于旧状态重复取消;毒消息无限 requeue,把正常任务饿死。
这些故障不能靠“再套一层 try/catch”解决。事务(Transaction)事务Transaction把一组数据库读写作为一个提交或回滚边界处理的机制。打开术语条目 → 只能原子保护它覆盖的资源;消息代理、缓存和外部服务通常不在同一事务里。可靠设计要先明确每个边界提供什么,再用 幂等键(Idempotency Key)幂等键Idempotency Key由调用方为同一业务意图稳定复用的身份,用于查找首次执行结果并阻止重复副作用。打开术语条目 →、事务性发件箱(Transactional Outbox)事务性发件箱Transactional Outbox把业务写入与待发布消息写进同一数据库事务,再由独立发布器可靠转发。打开术语条目 →、Inbox 去重、副作用账本(Side-effect Ledger)副作用账本Side-effect Ledger记录外部写入的意图、幂等身份、执行结果、验证状态和补偿状态的结构化账本。打开术语条目 →、对账收敛(Reconciliation)对账收敛Reconciliation比较期望状态与实际状态,并通过幂等动作持续缩小差异的恢复循环。打开术语条目 → 和 补偿(Compensation)补偿Compensation为已提交且无法直接回滚的业务副作用执行一个新的、可审计的纠正动作。打开术语条目 → 把部分失败收敛到可判断状态。
本课实现一个“文档转换任务处理器”。输入是抽象 JSON 文档,Worker 生成规范化 Artifact 并写入本地模拟 Artifact Store。所有外部副作用都是确定性模拟,不调用真实网络或收费 API。
事务边界之外,使用可重试协议收敛
不声称 Exactly-once1. 先画清四个权威边界
Section titled “1. 先画清四个权威边界”系统包含:
Client │ POST /tasks + Idempotency-Key ▼API Service ── DB(transaction) ── Outbox row │ ▼ Outbox Publisher ── Broker │ ▼ Worker + Inbox │ ▼ Artifact Store
Cache:只保存任务读模型,不参与权威写入判定每个边界的可靠语义不同:
| 边界 | 能可靠断言什么 | 不能断言什么 |
|---|---|---|
| 单个数据库事务 | 事务覆盖的行要么一起提交,要么一起回滚 | 队列或远程服务也同步提交 |
| Broker Publisher Confirm | Broker 按其协议确认接收发布 | 消费者已处理、业务副作用完成 |
| Consumer Ack | 消费者告诉 Broker 此次投递可移除 | 外部副作用天然只执行一次 |
| Cache | 在命中时返回某个版本的派生值 | 值最新、键永不淘汰、写入已持久化 |
| HTTP 2xx | 服务按 API 协议接受或完成操作 | 客户端一定收到;超时一定没执行 |
可靠系统不是把这些边界想象成一个大事务,而是设计每个边界失败后怎样重新观察和推进。
2. 事务:原子性只覆盖事务里的资源
Section titled “2. 事务:原子性只覆盖事务里的资源”事务(Transaction)事务Transaction把一组数据库读写作为一个提交或回滚边界处理的机制。打开术语条目 → 让数据库内一组修改以一个提交点出现。本课创建任务时在同一事务写三类记录:
task:业务权威状态;idempotency_record:客户端业务意图与响应身份;outbox_event:待发布的TaskCreated事件。
如果事务回滚,三者都不存在;如果提交,发布器迟早可以从 Outbox 恢复。不要在事务中先发布消息再写数据库,因为 Broker 不能跟随数据库回滚。
2.1 表结构
Section titled “2.1 表结构”CREATE TABLE task ( task_id text PRIMARY KEY, request_fingerprint text NOT NULL, status text NOT NULL CHECK (status IN ('queued', 'running', 'completed', 'failed', 'compensating', 'compensated')), version bigint NOT NULL DEFAULT 0, input_sha256 text NOT NULL, artifact_sha256 text, terminal_error_code text, created_at timestamptz NOT NULL DEFAULT now(), updated_at timestamptz NOT NULL DEFAULT now());
CREATE TABLE idempotency_record ( scope text NOT NULL, idempotency_key text NOT NULL, request_fingerprint text NOT NULL, state text NOT NULL CHECK (state IN ('in_progress', 'succeeded', 'failed_retryable', 'failed_final')), task_id text, response_status integer, response_body jsonb, locked_until timestamptz, expires_at timestamptz NOT NULL, PRIMARY KEY (scope, idempotency_key));
CREATE TABLE outbox_event ( outbox_id text PRIMARY KEY, aggregate_id text NOT NULL, aggregate_version bigint NOT NULL, event_type text NOT NULL, payload jsonb NOT NULL, payload_sha256 text NOT NULL, created_at timestamptz NOT NULL DEFAULT now(), published_at timestamptz, publish_attempts integer NOT NULL DEFAULT 0, last_error_code text, UNIQUE (aggregate_id, aggregate_version, event_type));
CREATE TABLE inbox_message ( consumer_name text NOT NULL, message_id text NOT NULL, payload_sha256 text NOT NULL, state text NOT NULL CHECK (state IN ('processing', 'completed', 'failed')), result_json jsonb, first_seen_at timestamptz NOT NULL DEFAULT now(), completed_at timestamptz, PRIMARY KEY (consumer_name, message_id));
CREATE TABLE effect_ledger ( effect_id text PRIMARY KEY, task_id text NOT NULL REFERENCES task(task_id), attempt integer NOT NULL, effect_kind text NOT NULL, idempotency_token text NOT NULL, target_identity text NOT NULL, state text NOT NULL CHECK (state IN ('planned', 'started', 'unknown', 'completed', 'verified', 'compensation_required', 'compensated', 'failed')), request_sha256 text NOT NULL, observed_result_sha256 text, error_code text, UNIQUE (effect_kind, idempotency_token));Schema 中的唯一约束是最后防线。应用层先查询可以返回更清晰错误,但两个并发请求仍可能同时看到“不存在”;只有数据库唯一约束能在提交点裁决。
2.2 隔离级别与丢失更新
Section titled “2.2 隔离级别与丢失更新”隔离级别(Isolation Level)隔离级别Isolation Level规定并发事务能够观察哪些中间结果,以及哪些冲突由数据库或应用处理。打开术语条目 → 决定并发事务能观察到什么,以及冲突如何暴露。即使使用常见的 Read Committed,也不能用“先读余额/配额,再无条件写回”维护跨请求不变量。
错误模式:
T1读取 remaining=1T2读取 remaining=1T1写 remaining=0T2写 remaining=0两个任务都被接受,但配额只有1可选修复:
UPDATE quotaSET remaining = remaining - 1WHERE scope = $1 AND remaining >= 1;检查影响行数;或使用行锁/更强隔离,并对序列化失败进行有限重试。选择哪种方式取决于不变量是否能表达为单条条件更新。不能笼统声称某隔离级别“解决并发”;要写出具体读写集合和冲突测试。
3. Idempotency Key:同一意图不重复创建资源
Section titled “3. Idempotency Key:同一意图不重复创建资源”幂等(Idempotency)幂等Idempotency同一个业务意图重复提交时,不会额外产生新的业务副作用。打开术语条目 → 意味着同一逻辑操作重复提交,系统状态不会被重复推进。幂等键(Idempotency Key)幂等键Idempotency Key由调用方为同一业务意图稳定复用的身份,用于查找首次执行结果并阻止重复副作用。打开术语条目 → 是调用方为一次业务意图选择的稳定身份。
最重要的规则:同一个键必须绑定同一个请求指纹。否则调用方误用键时,服务不知道是重试还是新意图。
请求指纹可以对规范化业务字段计算:
fingerprint = sha256(canonicalJson({ operation: 'create-document-task', inputSha256, outputFormat, policyVersion,}));不要把瞬时字段放进指纹,例如请求时间、Trace ID;也不要漏掉会改变副作用的字段。
3.1 并发创建协议
Section titled “3.1 并发创建协议”- API 验证 Idempotency Key 格式、作用域和请求 Schema。
- 计算请求指纹,尝试插入
in_progress幂等记录。 - 若插入成功,在同一事务创建 Task 和 Outbox,并写入可重放响应。
- 若唯一约束冲突,读取已有记录。
- 指纹不同:返回稳定的
IDEMPOTENCY_KEY_REUSED_WITH_DIFFERENT_REQUEST。 - 指纹相同且已成功:返回首次响应,不再创建 Task。
- 指纹相同但处理中:返回当前 Task 身份或
409/202风格的处理中结果;不要并发执行第二次。 - 锁过期也不能立即假设首次执行没发生;恢复器要检查 Task、Outbox 与副作用账本。
3.2 可手算例子
Section titled “3.2 可手算例子”两个客户端重试同时携带键 key-7 和相同指纹 fp-A:
T1 INSERT idempotency(scope=public-demo,key=key-7,fp-A) → 成功T2 INSERT 同一主键 → 唯一冲突,等待T1提交T1 INSERT task-101 + outbox-101 → COMMITT2 SELECT existing → state=succeeded, taskId=task-101结果只有一个 Task。若 T2 使用 fp-B,它不能拿到 task-101 作为自己的成功响应;服务返回键复用错误,让调用方选择新键。
3.3 保留期限不是随便清理
Section titled “3.3 保留期限不是随便清理”幂等记录必须至少覆盖客户端可能重试的窗口、队列延迟与故障恢复期。过早删除后旧请求会被当成新请求。无限保留则增加存储与隐私风险。协议要公开说明键作用域和保留边界;清理前还要保证关联任务已经终态、Outbox 已发布、未知副作用已对账。
4. Transactional Outbox:让数据库状态与待发布事件一起提交
Section titled “4. Transactional Outbox:让数据库状态与待发布事件一起提交”事务性发件箱(Transactional Outbox)事务性发件箱Transactional Outbox把业务写入与待发布消息写进同一数据库事务,再由独立发布器可靠转发。打开术语条目 → 解决“数据库已提交但消息没发”的双写问题:业务事务不直接依赖 Broker 成功,而是写一条 Outbox。独立发布器持续扫描未发布记录。
发布器循环:
select unpublished rows → publish(messageId=outboxId) → wait confirm → mark published仍然存在一个不可消除的窗口:Broker 已确认,但发布器在写 published_at 前崩溃。恢复后它会再次发布。因此 Outbox 通常提供 至少一次发布,消费者必须去重。它不创造 Exactly-once。
4.1 Publisher Confirm 的边界
Section titled “4.1 Publisher Confirm 的边界”Publisher Confirm 表示 Broker 按其协议接受了发布。它不是消费者完成证明。发布器要记录:
outboxId作为稳定 Message ID;- 发布尝试次数;
- Broker 返回的确认或拒绝;
- 最后错误类别;
- 下次可重试时间;
- 超过上限后的人工/隔离状态。
网络超时时,结果可能未知。不要生成新 outboxId;使用同一 Message ID 重发,让消费者 Inbox 去重。
4.2 批量锁定与并发 Publisher
Section titled “4.2 批量锁定与并发 Publisher”简化 SQL:
SELECT outbox_id, payloadFROM outbox_eventWHERE published_at IS NULLORDER BY created_at, outbox_idFOR UPDATE SKIP LOCKEDLIMIT 100;SKIP LOCKED 可让多个 Publisher 分工,但仍要设置租约或让事务足够短。不要在持有数据库事务时等待长网络调用;一种做法是先领取短租约、提交,再发布,最后条件更新。租约到期会产生重复发布,因此消费去重仍然必要。
5. Delivery Semantics:描述协议,不写口号
Section titled “5. Delivery Semantics:描述协议,不写口号”交付语义(Delivery Semantics)交付语义Delivery Semantics描述消息在故障和重试下可能不送达、送达一次或重复送达的协议边界。打开术语条目 → 常见表述:
| 表述 | 在消息层通常表示 | 业务层仍需处理 |
|---|---|---|
| At-most-once | 不重投或失败即丢弃 | 消息丢失、未处理任务 |
| At-least-once | 未确认会重投 | 重复消息、重复副作用 |
| Exactly-once | 在特定系统、范围和条件内去重/原子处理 | 外部系统、跨边界副作用和超出范围的故障 |
本课系统明确采用:Outbox 至少一次发布,Broker 至少一次投递,Inbox 对消息身份去重,Effect Ledger 对业务副作用收敛。不要把组合结果简写为“Exactly-once 处理”。
6. Inbox:消息去重与处理状态
Section titled “6. Inbox:消息去重与处理状态”Worker 收到消息时,先在数据库记录 Inbox:
(messageId, consumerName, payloadHash)处理协议:
- 同 Message ID 不存在:插入
processing; - 同 ID、同 Hash、状态
completed:返回已完成摘要并 Ack; - 同 ID、同 Hash、状态
processing:检查租约;未过期则不并发处理; - 同 ID、不同 Hash:协议冲突,进入隔离;
- 执行业务状态转换与本地记录;
- 外部副作用按 Ledger 执行和验证;
- Inbox 标记
completed后 Ack。
Ack 太早会丢工作,Ack 太晚会增加重复。可靠性不来自“恰到好处的 Ack 时机”,而来自重复发生时处理仍安全。
7. 外部副作用:完成、未知结果与验证
Section titled “7. 外部副作用:完成、未知结果与验证”假设 Worker 把 Artifact 写入外部存储。调用返回前连接断开:
PUT artifact(idempotencyToken=effect-42) → timeouttimeout 只表示调用方没有在期限内观察到结果。副作用可能:
- 未到达服务;
- 已到达但未执行;
- 已执行并提交,但响应丢失;
- 正在执行;
- 部分写入。
Effect Ledger 将状态设为 unknown。下一步优先查询同一业务身份:
HEAD artifact/effect-42GET operation-status/effect-42若外部服务支持幂等令牌,使用同一令牌重试;若不支持查询也不支持幂等,自动重试可能重复副作用,应转人工或设计本地代理层。
7.1 完成不等于验证
Section titled “7.1 完成不等于验证”存储 API 返回 2xx 后,Ledger 先是 completed,随后读取元数据或内容,核对:
artifact.taskId == taskIdartifact.attempt == currentAttemptartifact.sha256 == expectedSha256artifact.readable == true全部成立才标记 verified。这与上一课的 Artifact 证据一致:远程成功响应不是最终业务断言。
8. Compensation:新的业务动作,不是远程回滚
Section titled “8. Compensation:新的业务动作,不是远程回滚”补偿(Compensation)补偿Compensation为已提交且无法直接回滚的业务副作用执行一个新的、可审计的纠正动作。打开术语条目 → 用另一个可审计动作抵消已完成副作用。例如已发布临时 Artifact,但后续安全检查失败:
PublishArtifact(effect-42) 已完成SafetyValidation 失败HideArtifact(compensation-for-effect-42) 执行并验证补偿本身也可能失败、超时或重复,必须有自己的 Idempotency Token 和 Ledger 项。它不能保证世界恢复到“从未发生”:日志、通知、缓存读取或外部观察可能已经出现。因此补偿规则应表达业务可接受的收敛状态,而不是声称事务回滚。
| 原副作用 | 可能补偿 | 无法完全撤销的部分 |
|---|---|---|
| 创建可见 Artifact | 标记隐藏或删除 | 已被读取或缓存 |
| 预留配额 | 释放配额 | 期间阻塞过其他请求 |
| 发送通知 | 发送更正通知 | 首次通知已被看到 |
| 写外部索引 | 删除/覆盖索引项 | 搜索缓存和抓取快照 |
9. Cache:派生视图,不是权威状态捷径
Section titled “9. Cache:派生视图,不是权威状态捷径”旁路缓存(Cache-Aside)旁路缓存Cache-Aside应用先查缓存,未命中时读取权威存储并回填,写入时显式更新或失效缓存。打开术语条目 → 的典型读流程:
GET cache(taskId) ├─ hit → 返回带version的视图 └─ miss → 读DB → 写cache → 返回写流程可以先提交 DB,再删除 Cache。删除失败会留下旧值;因此读模型必须允许过期,关键写命令不能依赖 Cache 中的状态做最终 Guard。
9.1 为什么“先更新缓存再更新数据库”危险
Section titled “9.1 为什么“先更新缓存再更新数据库”危险”如果缓存更新成功、数据库事务失败,读者看到从未提交的状态。反过来先提交数据库、再失效缓存,最坏是短暂读旧值,权威事实仍在数据库。对于可接受短暂陈旧的任务列表,后者更容易恢复。
9.2 Cache Key 要包含版本或命名空间
Section titled “9.2 Cache Key 要包含版本或命名空间”task-view:v3:task-101Schema 升级时切换命名空间,避免旧序列化被新代码误读。值中保留 stateVersion,客户端可以判断流式事件是否比缓存新。
9.3 Eviction 是正常路径
Section titled “9.3 Eviction 是正常路径”缓存键可能因内存策略被淘汰,不能把“缓存一定存在”作为协议。缓存击穿、雪崩和热键需要请求合并、随机化 TTL、容量监控或分片,但是否采用取决于负载证据。课程不提供虚构吞吐数字。
10. 完整处理流程
Section titled “10. 完整处理流程”- 客户端生成一个业务意图级 Idempotency Key,并发送规范请求。
- API 计算指纹,在数据库事务中写幂等记录、Task 和 Outbox。
- 提交后返回 Task ID;响应丢失时客户端用同一键查询或重试。
- Outbox Publisher 用稳定 Message ID 发布,等待 Confirm;未知结果使用同一 ID 重发。
- Worker 在 Inbox 中登记消息身份和 Payload Hash。
- Worker 以版本条件把 Task 从
queued推进到running。 - Worker 在 Ledger 中计划副作用,使用稳定令牌执行,并验证 Artifact。
- 在数据库事务中写
completed、Artifact 身份和后续 Outbox 事件。 - Inbox 标记完成并 Ack;重复投递返回首次结果。
- 投影器更新或失效缓存;失败时从权威状态重建。
- 对账器扫描
unknown、超时租约、未发布 Outbox 和状态差异。 - 无法自动收敛的记录进入明确人工处理队列,而不是无限重试。
11. 一个确定性 TypeScript 任务处理器
Section titled “11. 一个确定性 TypeScript 任务处理器”下面代码是可运行的内存教学实现。它模拟数据库原子区段、Broker 重投和 Artifact Store;所有 ID、时钟与故障都由输入注入,因此测试结果确定。真实生产实现应把存储接口替换为数据库事务,并按所用 Broker 官方语义处理 Confirm/Ack。
import assert from 'node:assert/strict';import { createHash } from 'node:crypto';
type TaskStatus = 'queued' | 'running' | 'completed' | 'failed' | 'compensating' | 'compensated';
type Task = { taskId: string; requestFingerprint: string; status: TaskStatus; version: number; input: Record<string, unknown>; inputSha256: string; artifactSha256: string | null; terminalErrorCode: string | null;};
type IdempotencyRecord = { scope: string; key: string; fingerprint: string; state: 'in_progress' | 'succeeded' | 'failed_final'; taskId: string; response: { taskId: string; status: 'accepted' };};
type Outbox = { outboxId: string; aggregateId: string; aggregateVersion: number; eventType: 'TaskCreated' | 'TaskCompleted'; payload: Record<string, unknown>; published: boolean; publishAttempts: number;};
type Inbox = { consumer: string; messageId: string; payloadSha256: string; state: 'processing' | 'completed' | 'failed';};
type Effect = { effectId: string; taskId: string; token: string; state: 'planned' | 'started' | 'unknown' | 'completed' | 'verified' | 'compensation_required' | 'compensated' | 'failed'; expectedSha256: string; observedSha256: string | null; errorCode: string | null;};
type Message = { messageId: string; type: 'TaskCreated'; taskId: string; taskVersion: number;};
type Faults = { publishConfirmLost?: boolean; artifactWriteTimeoutAfterCommit?: boolean; artifactCorruptedAfterWrite?: boolean; ackLost?: boolean;};
function sha256(text: string): string { return createHash('sha256').update(text).digest('hex');}
function canonical(value: unknown): string { const normalize = (item: unknown): unknown => { if (Array.isArray(item)) return item.map(normalize); if (item !== null && typeof item === 'object') { return Object.fromEntries( Object.entries(item as Record<string, unknown>) .sort(([a], [b]) => a < b ? -1 : a > b ? 1 : 0) .map(([key, child]) => [key, normalize(child)]), ); } return item; }; return JSON.stringify(normalize(value));}
class InMemoryDatabase { readonly tasks = new Map<string, Task>(); readonly idempotency = new Map<string, IdempotencyRecord>(); readonly outbox = new Map<string, Outbox>(); readonly inbox = new Map<string, Inbox>(); readonly effects = new Map<string, Effect>();
transaction<T>(work: () => T): T { // 教学实现单线程执行;真实适配器必须由数据库提供原子性和约束。 return work(); }}
class SimulatedBroker { readonly queue: Message[] = []; private readonly acceptedIds = new Set<string>();
publish(message: Message, faults: Faults): 'confirmed' | 'unknown' { if (!this.acceptedIds.has(message.messageId)) { this.acceptedIds.add(message.messageId); this.queue.push(structuredClone(message)); } else { // 模拟至少一次发布:Broker或网络层仍可能再次投递同一身份。 this.queue.push(structuredClone(message)); } return faults.publishConfirmLost ? 'unknown' : 'confirmed'; }
deliver(): Message | null { return this.queue.shift() ?? null; }
redeliver(message: Message): void { this.queue.push(structuredClone(message)); }}
class SimulatedArtifactStore { readonly objects = new Map<string, string>();
put(token: string, bytes: string, faults: Faults): 'completed' | 'unknown' { if (!this.objects.has(token)) this.objects.set(token, bytes); if (faults.artifactCorruptedAfterWrite) this.objects.set(token, `${bytes}corrupted`); return faults.artifactWriteTimeoutAfterCommit ? 'unknown' : 'completed'; }
read(token: string): string | null { return this.objects.get(token) ?? null; }
hide(token: string): void { this.objects.delete(token); }}
function createTask( db: InMemoryDatabase, input: Record<string, unknown>, idempotencyKey: string, ids: { taskId: string; outboxId: string },): { status: 'accepted'; taskId: string } { if (!/^[a-zA-Z0-9_-]{6,80}$/.test(idempotencyKey)) { throw new Error('INVALID_IDEMPOTENCY_KEY'); } const scope = 'document-transform'; const fingerprint = sha256(canonical({ operation: scope, input })); const recordKey = `${scope}:${idempotencyKey}`;
return db.transaction(() => { const existing = db.idempotency.get(recordKey); if (existing) { if (existing.fingerprint !== fingerprint) { throw new Error('IDEMPOTENCY_KEY_REUSED_WITH_DIFFERENT_REQUEST'); } return existing.response; }
if (db.tasks.has(ids.taskId) || db.outbox.has(ids.outboxId)) { throw new Error('IDENTITY_CONFLICT'); } const task: Task = { taskId: ids.taskId, requestFingerprint: fingerprint, status: 'queued', version: 0, input: structuredClone(input), inputSha256: sha256(canonical(input)), artifactSha256: null, terminalErrorCode: null, }; const response = { status: 'accepted' as const, taskId: task.taskId }; db.tasks.set(task.taskId, task); db.outbox.set(ids.outboxId, { outboxId: ids.outboxId, aggregateId: task.taskId, aggregateVersion: 0, eventType: 'TaskCreated', payload: { taskId: task.taskId, taskVersion: 0 }, published: false, publishAttempts: 0, }); db.idempotency.set(recordKey, { scope, key: idempotencyKey, fingerprint, state: 'succeeded', taskId: task.taskId, response, }); return response; });}
function publishOutbox( db: InMemoryDatabase, broker: SimulatedBroker, outboxId: string, faults: Faults,): 'confirmed' | 'unknown' | 'already-published' { const row = db.outbox.get(outboxId); if (!row) throw new Error('OUTBOX_NOT_FOUND'); if (row.published) return 'already-published'; row.publishAttempts += 1; const outcome = broker.publish({ messageId: row.outboxId, type: 'TaskCreated', taskId: row.aggregateId, taskVersion: row.aggregateVersion, }, faults); if (outcome === 'confirmed') row.published = true; return outcome;}
function transform(input: Record<string, unknown>, taskId: string): string { return `${canonical({ schemaVersion: 1, taskId, normalized: input })}\n`;}
function processMessage( db: InMemoryDatabase, store: SimulatedArtifactStore, message: Message, faults: Faults,): { status: 'completed' | 'duplicate' | 'failed'; ack: boolean; code?: string } { const consumer = 'document-worker-v1'; const payloadSha256 = sha256(canonical(message)); const inboxKey = `${consumer}:${message.messageId}`; const existingInbox = db.inbox.get(inboxKey); if (existingInbox) { if (existingInbox.payloadSha256 !== payloadSha256) { return { status: 'failed', ack: false, code: 'MESSAGE_ID_PAYLOAD_CONFLICT' }; } if (existingInbox.state === 'completed') { return { status: 'duplicate', ack: !faults.ackLost }; } } else { db.inbox.set(inboxKey, { consumer, messageId: message.messageId, payloadSha256, state: 'processing', }); }
const task = db.tasks.get(message.taskId); if (!task) return { status: 'failed', ack: false, code: 'TASK_NOT_FOUND' }; if (task.status === 'completed') { db.inbox.get(inboxKey)!.state = 'completed'; return { status: 'duplicate', ack: !faults.ackLost }; } if (task.status !== 'queued' || task.version !== message.taskVersion) { return { status: 'failed', ack: false, code: 'TASK_VERSION_CONFLICT' }; }
db.transaction(() => { task.status = 'running'; task.version += 1; });
const artifactBytes = transform(task.input, task.taskId); const expectedSha256 = sha256(artifactBytes); const effectId = `${task.taskId}:artifact:attempt-1`; const token = effectId; const effect: Effect = db.effects.get(effectId) ?? { effectId, taskId: task.taskId, token, state: 'planned', expectedSha256, observedSha256: null, errorCode: null, }; db.effects.set(effectId, effect);
if (effect.state !== 'verified') { effect.state = 'started'; const outcome = store.put(token, artifactBytes, faults); effect.state = outcome === 'unknown' ? 'unknown' : 'completed';
// 无论调用返回完成还是未知,都通过查询同一token对账。 const observed = store.read(token); if (observed === null) { effect.state = 'failed'; effect.errorCode = 'ARTIFACT_NOT_FOUND_AFTER_WRITE'; return { status: 'failed', ack: false, code: effect.errorCode }; } const observedSha256 = sha256(observed); effect.observedSha256 = observedSha256; if (observedSha256 !== expectedSha256) { effect.state = 'compensation_required'; effect.errorCode = 'ARTIFACT_HASH_MISMATCH'; store.hide(token); effect.state = 'compensated'; task.status = 'failed'; task.version += 1; task.terminalErrorCode = effect.errorCode; db.inbox.get(inboxKey)!.state = 'failed'; return { status: 'failed', ack: true, code: effect.errorCode }; } effect.state = 'verified'; }
db.transaction(() => { task.status = 'completed'; task.version += 1; task.artifactSha256 = effect.observedSha256; db.inbox.get(inboxKey)!.state = 'completed'; }); return { status: 'completed', ack: !faults.ackLost };}
function reconcile(db: InMemoryDatabase, store: SimulatedArtifactStore): string[] { const actions: string[] = []; for (const effect of db.effects.values()) { if (effect.state !== 'unknown' && effect.state !== 'started') continue; const observed = store.read(effect.token); if (observed && sha256(observed) === effect.expectedSha256) { effect.observedSha256 = effect.expectedSha256; effect.state = 'verified'; actions.push(`${effect.effectId}:verified-by-observation`); } else if (observed) { effect.state = 'compensation_required'; actions.push(`${effect.effectId}:compensation-required`); } else { actions.push(`${effect.effectId}:safe-to-retry-with-same-token`); } } return actions;}
function selfTest(): void { const db = new InMemoryDatabase(); const broker = new SimulatedBroker(); const store = new SimulatedArtifactStore(); const input = { id: 'doc-1', title: 'demo' };
const first = createTask(db, input, 'intent-0001', { taskId: 'task-1', outboxId: 'outbox-1' }); const duplicate = createTask(db, input, 'intent-0001', { taskId: 'unused', outboxId: 'unused' }); assert.deepEqual(duplicate, first); assert.equal(db.tasks.size, 1);
assert.equal(publishOutbox(db, broker, 'outbox-1', { publishConfirmLost: true }), 'unknown'); assert.equal(publishOutbox(db, broker, 'outbox-1', {}), 'confirmed'); assert.equal(broker.queue.length, 2);
const message1 = broker.deliver(); assert.ok(message1); const execution = processMessage(db, store, message1, { artifactWriteTimeoutAfterCommit: true, ackLost: true }); assert.equal(execution.status, 'completed'); assert.equal(execution.ack, false);
const message2 = broker.deliver(); assert.ok(message2); const redelivery = processMessage(db, store, message2, {}); assert.equal(redelivery.status, 'duplicate'); assert.equal(db.tasks.get('task-1')?.status, 'completed'); assert.equal(store.objects.size, 1); assert.deepEqual(reconcile(db, store), []);}
selfTest();实现里 transaction 是单线程占位,代码注释已经说明不能冒充真实数据库原子性。教学重点是:重复发布、写入超时后已提交、Ack 丢失都发生时,Task 仍只有一个、Artifact 仍只有一个业务身份、重复消息不会再次推进状态。
12. 正常路径逐步验证
Section titled “12. 正常路径逐步验证”以 task-1 为例:
1. createTask(intent-0001) task=queued/v0 idempotency=succeeded → task-1 outbox=unpublished
2. publishOutbox(outbox-1) broker confirms outbox=published
3. Worker registers inbox(outbox-1) task queued/v0 → running/v1
4. Ledger plans artifact token=task-1:artifact:attempt-1 store.put returns completed read-back Hash matches effect=verified
5. DB transaction task running/v1 → completed/v2 inbox=completed
6. Ack message cache projection eventually sees completed/v2每个箭头都有可检查记录。若第 4 步完成、第 5 步失败,重投消息读取同一 Effect Token,发现 Artifact 已验证,再推进 Task;不会创建第二份副作用。
13. 至少五类失败路径
Section titled “13. 至少五类失败路径”13.1 API 响应丢失
Section titled “13.1 API 响应丢失”事实:数据库事务可能已经提交。客户端必须用同一 Idempotency Key 查询或重试。使用新键会表达新意图,服务无法自动判断重复。
13.2 Outbox 确认未知
Section titled “13.2 Outbox 确认未知”事实:Broker 可能已接收。发布器用同一 outboxId 重发;消费者按 Message ID 去重。不能因为超时创建新事件身份。
13.3 Worker 副作用完成后崩溃
Section titled “13.3 Worker 副作用完成后崩溃”事实:Task 可能仍是 running,Artifact 已存在。重投后先查 Ledger/Artifact;验证通过后推进状态。不要无条件重做副作用。
13.4 Payload 身份冲突
Section titled “13.4 Payload 身份冲突”事实:同一 Message ID 对应不同 Hash。系统不能把它当普通重复,应隔离并告警,因为至少一个发送方违反协议或记录被篡改。
13.5 Poison Message
Section titled “13.5 Poison Message”同一输入持续触发确定性错误。无限 requeue 会形成热循环。经过有限次数、稳定错误分类后,消息进入 死信队列(Dead-letter Queue)死信队列Dead-letter Queue隔离超过有限重试或无法处理消息的队列,便于诊断和受控重放。打开术语条目 → 或隔离表,同时保留原消息身份、失败代码、处理器版本和可重放条件。DLQ 不是垃圾桶;没有检查和回放流程的 DLQ 只是延迟丢失。
13.6 Cache Stale
Section titled “13.6 Cache Stale”读模型显示 running/v1,权威状态已是 completed/v2。客户端收到流事件或重新读数据库时,应按版本选择更高状态。缓存过期不应触发第二次任务创建。
14. 故障矩阵
Section titled “14. 故障矩阵”| 故障点 | 当前事实 | 不确定项 | 检测证据 | 收敛动作 |
|---|---|---|---|---|
| DB 提交后响应断线 | 客户端无响应 | Task 是否创建 | Idempotency Record | 同键查询/重试,返回首次 Task |
| DB 成功、Publisher 崩溃 | Outbox 未标记发布 | Broker 是否收到 | Outbox 状态 + Confirm 记录 | 同 outboxId 重发 |
| Broker 重投 | 同 Message ID 再到达 | 首次处理是否完成 | Inbox + Payload Hash | 幂等返回或恢复处理中记录 |
| Artifact PUT 超时 | 调用方未见结果 | Artifact 是否存在 | Effect Token 查询 + Hash | 对账;同 Token 重试或补偿 |
| Artifact Hash 错 | 路径存在 | 写入者/转换何处出错 | 期望与观察 Hash、Trace | 隔离、隐藏产物、任务失败 |
| Task 完成但 Ack 丢失 | Task/Inbox 已完成 | Broker 是否会再投递 | Inbox completed | 重投时返回 duplicate 并 Ack |
| Cache 删除失败 | DB 已新版本 | Cache 仍旧 | 版本字段、缓存年龄 | 下次读回源;异步重试失效 |
| Poison Message | 多次同类失败 | 修复后是否可处理 | 次数、稳定错误码、版本 | 隔离/DLQ,修复后受控重放 |
| 补偿超时 | 原副作用已发生 | 补偿是否完成 | 补偿 Ledger + 查询 | 同补偿 Token 对账,不盲重放 |
15. 故障注入实验
Section titled “15. 故障注入实验”实验一:确认丢失导致重复发布
Section titled “实验一:确认丢失导致重复发布”- 创建一个 Task 和 Outbox;
- 第一次
publish把消息写入 Broker,但注入publishConfirmLost=true; - Outbox 仍显示未发布;
- 第二次用同一 Message ID 发布并得到 Confirm;
- Broker 队列出现两次投递;
- Worker 第一次完成,第二次通过 Inbox 返回
duplicate; - 断言 Task 版本没有再次增加,Artifact Store 只有一个业务 Token。
验收:至少一次发布不产生重复业务副作用。
实验二:副作用已提交但调用超时
Section titled “实验二:副作用已提交但调用超时”- 注入
artifactWriteTimeoutAfterCommit=true; - Store 实际保存 Artifact,调用返回
unknown; - Worker 不能立刻创建新 Token 重写;
- 使用相同 Token 读取并核对 Hash;
- Ledger 从
unknown进入verified; - Task 完成;
- 删除查询步骤后,测试应暴露重复或错误失败。
验收:未知结果通过观察收敛,而不是把 Timeout 等同于未执行。
实验三:Artifact 被破坏并触发补偿
Section titled “实验三:Artifact 被破坏并触发补偿”- Store 写入后追加字节,注入 Hash 不匹配;
- Ledger 标记
compensation_required; - 执行
hide(token); - 补偿验证后 Ledger 变为
compensated; - Task 进入
failed并保存稳定错误码; - 消息可以 Ack,避免确定性毒消息无限重试;
- 重放同 Message ID 返回已失败摘要。
验收:错误 Artifact 不可见,失败和补偿都有证据。
实验四:同键不同请求
Section titled “实验四:同键不同请求”- 用
intent-0001创建输入 A; - 用同键重试输入 A,返回同 Task;
- 用同键提交输入 B;
- 断言返回
IDEMPOTENCY_KEY_REUSED_WITH_DIFFERENT_REQUEST; - 断言没有创建第二个 Task,也没有把 B 绑定到 A 的响应。
验收:Idempotency Key 不会成为覆盖请求差异的缓存键。
16. 决策表:什么时候使用什么机制
Section titled “16. 决策表:什么时候使用什么机制”| 问题 | 首选机制 | 为什么 | 仍需注意 |
|---|---|---|---|
| 客户端因断线重试创建命令 | Idempotency Key + 请求指纹 | 绑定同一业务意图和首次响应 | 作用域、保留期、并发唯一约束 |
| DB 写与消息发布 | Transactional Outbox | 让业务状态与待发布事实一起提交 | 仍会重复发布,消费者必须去重 |
| Broker 重复投递 | Inbox + Message ID/Hash | 稳定识别同一投递 | 语义重复还需业务键 |
| 外部副作用超时 | Effect Ledger + 查询/同 Token 重试 | 区分未知结果与确定失败 | 外部服务必须提供可查询或幂等接口 |
| 已完成副作用后流程失败 | Compensation | 用新业务动作收敛 | 不保证世界恢复到从未发生 |
| 高频读、允许短暂旧值 | Cache-aside | 减少权威存储读取 | 失效失败、Eviction、版本比较 |
| 确定性坏消息 | 有限重试 + 隔离/DLQ | 防止热循环阻塞正常消息 | 必须有检查、修复和重放流程 |
| 跨系统一致性检查 | Reconciliation | 通过观察修复声明与现实差异 | 扫描成本、权限、人工边界 |
17. 验收条件
Section titled “17. 验收条件”- □ 创建 API 要求 Idempotency Key,并校验作用域、格式和请求指纹;
- □ 相同键同请求返回首次响应,相同键不同请求稳定失败;
- □ Task、幂等记录和 Outbox 在同一数据库事务提交;
- □ Outbox 使用稳定 Message ID,Publisher Confirm 未知时不更换身份;
- □ Consumer 在业务处理前登记 Inbox 与 Payload Hash;
- □ 重复投递不会再次推进 Task 或重复执行副作用;
- □ 外部副作用具有稳定 Idempotency Token 与 Ledger;
- □ Timeout 进入
unknown,通过查询、验证或人工处理收敛; - □ Artifact 必须回读并核对 Task、Attempt、大小和 Hash;
- □ 补偿是独立、幂等、可验证的业务动作;
- □ Cache 只作派生读模型,关键写 Guard 读取权威状态;
- □ Poison Message 有有限重试、隔离、检查和受控重放路径;
- □ 至少执行三个故障注入实验,并保存 Task、Inbox、Outbox 和 Ledger 的最终快照;
- □ 文档与实现不声称跨数据库、Broker、缓存和外部服务的 Exactly-once。
18. 学完后应该能回答什么
Section titled “18. 学完后应该能回答什么”- 为什么数据库事务不能原子覆盖消息代理和远程 Artifact Store?
- Idempotency Key 为什么必须绑定请求指纹?
- 两个并发同键请求如何由唯一约束裁决?
- Transactional Outbox 解决什么双写问题,为什么仍会重复发布?
- Publisher Confirm、Consumer Ack 与业务完成分别证明什么?
- Event ID 去重与 Effect Token 去重有什么区别?
- Timeout 后为什么先查询或对账,而不是直接重试?
completed与verified的副作用状态为什么要分开?- Compensation 为什么不是远程事务回滚?
- Cache-aside 中为什么写后失效通常比缓存先写更易恢复?
- Poison Message 为什么不能无限 requeue?
- 如何用一个故障矩阵说明系统在部分失败后仍能收敛?
19. 来源边界
Section titled “19. 来源边界”- PostgreSQL 文档用于核对事务隔离和并发语义;示例 SQL 不代表所有数据库具有相同锁、冲突和约束行为。
- RabbitMQ 文档用于说明 Publisher Confirm、Consumer Ack 与 Dead Lettering 的协议边界;课程不声称配置后自动获得业务 Exactly-once。
- AWS Builders’ Library 的 Idempotent API 材料用于理解调用方请求身份和重试安全;本课 Schema、保留期和错误码是教学协议,需要按具体系统重新设计。
- Redis Eviction 文档用于说明缓存键可能按策略淘汰;课程没有给出未经测量的容量或命中率承诺。
- 所有任务、消息、Artifact、Token 和时间均为抽象教学数据,不包含真实服务、账号或私有项目事实。