跳转到内容

事务、缓存与队列

一个创建异步任务的 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-once
HTTP COMMANDIdempotency Key同一业务意图使用同一键
ONE DB TRANSACTIONTask + Idempotency Record + Outbox业务状态与待发布事件一起提交
BROKERAt-least-once Delivery发布和消费都可能重试
WORKERInbox / Effect Ledger先去重,再执行可验证副作用
UNKNOWN OUTCOME查询幂等记录超时后先查状态,不用新键盲目重放。
PARTIAL EFFECT对账或补偿补偿是新的业务动作,不是数据库回滚的远程版本。
POISON MESSAGE有限重试后隔离保留输入、错误分类和重放条件,不无限 requeue。
CACHE只作派生视图命中与否不能改变权威写入规则。
数据库事务只能原子地保护它覆盖的资源。跨数据库、消息代理和外部 API 的流程,需要幂等记录、Outbox、去重、对账与补偿共同收敛。

系统包含:

Client
│ POST /tasks + Idempotency-Key
API Service ── DB(transaction) ── Outbox row
Outbox Publisher ── Broker
Worker + Inbox
Artifact Store
Cache:只保存任务读模型,不参与权威写入判定

每个边界的可靠语义不同:

边界能可靠断言什么不能断言什么
单个数据库事务事务覆盖的行要么一起提交,要么一起回滚队列或远程服务也同步提交
Broker Publisher ConfirmBroker 按其协议确认接收发布消费者已处理、业务副作用完成
Consumer Ack消费者告诉 Broker 此次投递可移除外部副作用天然只执行一次
Cache在命中时返回某个版本的派生值值最新、键永不淘汰、写入已持久化
HTTP 2xx服务按 API 协议接受或完成操作客户端一定收到;超时一定没执行

可靠系统不是把这些边界想象成一个大事务,而是设计每个边界失败后怎样重新观察和推进。

2. 事务:原子性只覆盖事务里的资源

Section titled “2. 事务:原子性只覆盖事务里的资源”

事务(Transaction)事务Transaction把一组数据库读写作为一个提交或回滚边界处理的机制。打开术语条目 → 让数据库内一组修改以一个提交点出现。本课创建任务时在同一事务写三类记录:

  1. task:业务权威状态;
  2. idempotency_record:客户端业务意图与响应身份;
  3. outbox_event:待发布的 TaskCreated 事件。

如果事务回滚,三者都不存在;如果提交,发布器迟早可以从 Outbox 恢复。不要在事务中先发布消息再写数据库,因为 Broker 不能跟随数据库回滚。

examples/reliability-schema.sql
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 中的唯一约束是最后防线。应用层先查询可以返回更清晰错误,但两个并发请求仍可能同时看到“不存在”;只有数据库唯一约束能在提交点裁决。

隔离级别(Isolation Level)隔离级别Isolation Level规定并发事务能够观察哪些中间结果,以及哪些冲突由数据库或应用处理。打开术语条目 → 决定并发事务能观察到什么,以及冲突如何暴露。即使使用常见的 Read Committed,也不能用“先读余额/配额,再无条件写回”维护跨请求不变量。

错误模式:

T1读取 remaining=1
T2读取 remaining=1
T1写 remaining=0
T2写 remaining=0
两个任务都被接受,但配额只有1

可选修复:

UPDATE quota
SET remaining = remaining - 1
WHERE 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;也不要漏掉会改变副作用的字段。

  1. API 验证 Idempotency Key 格式、作用域和请求 Schema。
  2. 计算请求指纹,尝试插入 in_progress 幂等记录。
  3. 若插入成功,在同一事务创建 Task 和 Outbox,并写入可重放响应。
  4. 若唯一约束冲突,读取已有记录。
  5. 指纹不同:返回稳定的 IDEMPOTENCY_KEY_REUSED_WITH_DIFFERENT_REQUEST
  6. 指纹相同且已成功:返回首次响应,不再创建 Task。
  7. 指纹相同但处理中:返回当前 Task 身份或 409/202 风格的处理中结果;不要并发执行第二次。
  8. 锁过期也不能立即假设首次执行没发生;恢复器要检查 Task、Outbox 与副作用账本。

两个客户端重试同时携带键 key-7 和相同指纹 fp-A

T1 INSERT idempotency(scope=public-demo,key=key-7,fp-A) → 成功
T2 INSERT 同一主键 → 唯一冲突,等待T1提交
T1 INSERT task-101 + outbox-101 → COMMIT
T2 SELECT existing → state=succeeded, taskId=task-101

结果只有一个 Task。若 T2 使用 fp-B,它不能拿到 task-101 作为自己的成功响应;服务返回键复用错误,让调用方选择新键。

幂等记录必须至少覆盖客户端可能重试的窗口、队列延迟与故障恢复期。过早删除后旧请求会被当成新请求。无限保留则增加存储与隐私风险。协议要公开说明键作用域和保留边界;清理前还要保证关联任务已经终态、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。

Publisher Confirm 表示 Broker 按其协议接受了发布。它不是消费者完成证明。发布器要记录:

  • outboxId 作为稳定 Message ID;
  • 发布尝试次数;
  • Broker 返回的确认或拒绝;
  • 最后错误类别;
  • 下次可重试时间;
  • 超过上限后的人工/隔离状态。

网络超时时,结果可能未知。不要生成新 outboxId;使用同一 Message ID 重发,让消费者 Inbox 去重。

简化 SQL:

SELECT outbox_id, payload
FROM outbox_event
WHERE published_at IS NULL
ORDER BY created_at, outbox_id
FOR UPDATE SKIP LOCKED
LIMIT 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 处理”。

Worker 收到消息时,先在数据库记录 Inbox:

(messageId, consumerName, payloadHash)

处理协议:

  1. 同 Message ID 不存在:插入 processing
  2. 同 ID、同 Hash、状态 completed:返回已完成摘要并 Ack;
  3. 同 ID、同 Hash、状态 processing:检查租约;未过期则不并发处理;
  4. 同 ID、不同 Hash:协议冲突,进入隔离;
  5. 执行业务状态转换与本地记录;
  6. 外部副作用按 Ledger 执行和验证;
  7. Inbox 标记 completed 后 Ack。

Ack 太早会丢工作,Ack 太晚会增加重复。可靠性不来自“恰到好处的 Ack 时机”,而来自重复发生时处理仍安全。

7. 外部副作用:完成、未知结果与验证

Section titled “7. 外部副作用:完成、未知结果与验证”

假设 Worker 把 Artifact 写入外部存储。调用返回前连接断开:

PUT artifact(idempotencyToken=effect-42) → timeout

timeout 只表示调用方没有在期限内观察到结果。副作用可能:

  • 未到达服务;
  • 已到达但未执行;
  • 已执行并提交,但响应丢失;
  • 正在执行;
  • 部分写入。

Effect Ledger 将状态设为 unknown。下一步优先查询同一业务身份:

HEAD artifact/effect-42
GET operation-status/effect-42

若外部服务支持幂等令牌,使用同一令牌重试;若不支持查询也不支持幂等,自动重试可能重复副作用,应转人工或设计本地代理层。

存储 API 返回 2xx 后,Ledger 先是 completed,随后读取元数据或内容,核对:

artifact.taskId == taskId
artifact.attempt == currentAttempt
artifact.sha256 == expectedSha256
artifact.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-101

Schema 升级时切换命名空间,避免旧序列化被新代码误读。值中保留 stateVersion,客户端可以判断流式事件是否比缓存新。

缓存键可能因内存策略被淘汰,不能把“缓存一定存在”作为协议。缓存击穿、雪崩和热键需要请求合并、随机化 TTL、容量监控或分片,但是否采用取决于负载证据。课程不提供虚构吞吐数字。

  1. 客户端生成一个业务意图级 Idempotency Key,并发送规范请求。
  2. API 计算指纹,在数据库事务中写幂等记录、Task 和 Outbox。
  3. 提交后返回 Task ID;响应丢失时客户端用同一键查询或重试。
  4. Outbox Publisher 用稳定 Message ID 发布,等待 Confirm;未知结果使用同一 ID 重发。
  5. Worker 在 Inbox 中登记消息身份和 Payload Hash。
  6. Worker 以版本条件把 Task 从 queued 推进到 running
  7. Worker 在 Ledger 中计划副作用,使用稳定令牌执行,并验证 Artifact。
  8. 在数据库事务中写 completed、Artifact 身份和后续 Outbox 事件。
  9. Inbox 标记完成并 Ack;重复投递返回首次结果。
  10. 投影器更新或失效缓存;失败时从权威状态重建。
  11. 对账器扫描 unknown、超时租约、未发布 Outbox 和状态差异。
  12. 无法自动收敛的记录进入明确人工处理队列,而不是无限重试。

11. 一个确定性 TypeScript 任务处理器

Section titled “11. 一个确定性 TypeScript 任务处理器”

下面代码是可运行的内存教学实现。它模拟数据库原子区段、Broker 重投和 Artifact Store;所有 ID、时钟与故障都由输入注入,因此测试结果确定。真实生产实现应把存储接口替换为数据库事务,并按所用 Broker 官方语义处理 Confirm/Ack。

examples/reliable-task-processor.ts
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 仍只有一个业务身份、重复消息不会再次推进状态。

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;不会创建第二份副作用。

事实:数据库事务可能已经提交。客户端必须用同一 Idempotency Key 查询或重试。使用新键会表达新意图,服务无法自动判断重复。

事实:Broker 可能已接收。发布器用同一 outboxId 重发;消费者按 Message ID 去重。不能因为超时创建新事件身份。

事实:Task 可能仍是 running,Artifact 已存在。重投后先查 Ledger/Artifact;验证通过后推进状态。不要无条件重做副作用。

事实:同一 Message ID 对应不同 Hash。系统不能把它当普通重复,应隔离并告警,因为至少一个发送方违反协议或记录被篡改。

同一输入持续触发确定性错误。无限 requeue 会形成热循环。经过有限次数、稳定错误分类后,消息进入 死信队列(Dead-letter Queue)死信队列Dead-letter Queue隔离超过有限重试或无法处理消息的队列,便于诊断和受控重放。打开术语条目 → 或隔离表,同时保留原消息身份、失败代码、处理器版本和可重放条件。DLQ 不是垃圾桶;没有检查和回放流程的 DLQ 只是延迟丢失。

读模型显示 running/v1,权威状态已是 completed/v2。客户端收到流事件或重新读数据库时,应按版本选择更高状态。缓存过期不应触发第二次任务创建。

故障点当前事实不确定项检测证据收敛动作
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 对账,不盲重放

实验一:确认丢失导致重复发布

Section titled “实验一:确认丢失导致重复发布”
  1. 创建一个 Task 和 Outbox;
  2. 第一次 publish 把消息写入 Broker,但注入 publishConfirmLost=true
  3. Outbox 仍显示未发布;
  4. 第二次用同一 Message ID 发布并得到 Confirm;
  5. Broker 队列出现两次投递;
  6. Worker 第一次完成,第二次通过 Inbox 返回 duplicate
  7. 断言 Task 版本没有再次增加,Artifact Store 只有一个业务 Token。

验收:至少一次发布不产生重复业务副作用。

实验二:副作用已提交但调用超时

Section titled “实验二:副作用已提交但调用超时”
  1. 注入 artifactWriteTimeoutAfterCommit=true
  2. Store 实际保存 Artifact,调用返回 unknown
  3. Worker 不能立刻创建新 Token 重写;
  4. 使用相同 Token 读取并核对 Hash;
  5. Ledger 从 unknown 进入 verified
  6. Task 完成;
  7. 删除查询步骤后,测试应暴露重复或错误失败。

验收:未知结果通过观察收敛,而不是把 Timeout 等同于未执行。

实验三:Artifact 被破坏并触发补偿

Section titled “实验三:Artifact 被破坏并触发补偿”
  1. Store 写入后追加字节,注入 Hash 不匹配;
  2. Ledger 标记 compensation_required
  3. 执行 hide(token)
  4. 补偿验证后 Ledger 变为 compensated
  5. Task 进入 failed 并保存稳定错误码;
  6. 消息可以 Ack,避免确定性毒消息无限重试;
  7. 重放同 Message ID 返回已失败摘要。

验收:错误 Artifact 不可见,失败和补偿都有证据。

  1. intent-0001 创建输入 A;
  2. 用同键重试输入 A,返回同 Task;
  3. 用同键提交输入 B;
  4. 断言返回 IDEMPOTENCY_KEY_REUSED_WITH_DIFFERENT_REQUEST
  5. 断言没有创建第二个 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通过观察修复声明与现实差异扫描成本、权限、人工边界
  • □ 创建 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。
  1. 为什么数据库事务不能原子覆盖消息代理和远程 Artifact Store?
  2. Idempotency Key 为什么必须绑定请求指纹?
  3. 两个并发同键请求如何由唯一约束裁决?
  4. Transactional Outbox 解决什么双写问题,为什么仍会重复发布?
  5. Publisher Confirm、Consumer Ack 与业务完成分别证明什么?
  6. Event ID 去重与 Effect Token 去重有什么区别?
  7. Timeout 后为什么先查询或对账,而不是直接重试?
  8. completedverified 的副作用状态为什么要分开?
  9. Compensation 为什么不是远程事务回滚?
  10. Cache-aside 中为什么写后失效通常比缓存先写更易恢复?
  11. Poison Message 为什么不能无限 requeue?
  12. 如何用一个故障矩阵说明系统在部分失败后仍能收敛?
  • PostgreSQL 文档用于核对事务隔离和并发语义;示例 SQL 不代表所有数据库具有相同锁、冲突和约束行为。
  • RabbitMQ 文档用于说明 Publisher Confirm、Consumer Ack 与 Dead Lettering 的协议边界;课程不声称配置后自动获得业务 Exactly-once。
  • AWS Builders’ Library 的 Idempotent API 材料用于理解调用方请求身份和重试安全;本课 Schema、保留期和错误码是教学协议,需要按具体系统重新设计。
  • Redis Eviction 文档用于说明缓存键可能按策略淘汰;课程没有给出未经测量的容量或命中率承诺。
  • 所有任务、消息、Artifact、Token 和时间均为抽象教学数据,不包含真实服务、账号或私有项目事实。