切换主题
Outbox 与双写问题
问题:订单已经保存,通知却没有发出
最直觉的实现是保存订单后发送消息。数据库提交与消息发送之间崩溃,订单存在但消息丢失。调换顺序也无济于事:消息先到达,数据库后来回滚,接收端看到不存在的订单。
跨两个系统的两次写入无法靠调整顺序获得单库事务原子性。Outbox 把“稍后要发送这件事实”先变成同一个数据库里的可靠记录。
原理:把发送意图纳入事务
一次事务同时写订单与 Outbox 事件。独立发送器认领可发送记录,发布到消息系统,收到确认后标记成功。进程重启后扫描未完成记录,继续发送。
这提供恢复依据,但不消除重复:消息系统接收后,发送器可能在记录成功之前崩溃。重启不知道消息已经送达,只能重投同一事件。消费端必须处理重复。
最小例子:事务里只存事实
text
事务 A:INSERT order + INSERT outbox(event_id, payload) → COMMIT
发送器:认领 event_id → 发布 → 接收确认
事务 B:按 event_id + claim_id 标记 PUBLISHEDevent_id 在最初创建时确定,重投不换 ID。claim_id 标识当前发送所有者,租约过期后新的发送器可以拿走,旧发送器不能无条件覆盖新认领。
交互 02
在两次写入之间崩溃
订单未创建
待发记录无
投递次数0
业务效果0
教学简化:认领租约自动视为已到期,消费端的事件 ID 与业务效果同事务提交。真实租约、失败退避与外部副作用见正文;模型不承诺跨系统恰好一次。
请分别尝试“提交后崩溃”和“消息送达后崩溃”。观察订单、待发记录、投递次数与业务效果,解释为什么 Outbox 保留恢复依据,而去重解决另一类问题。
方案比较:轮询、唤醒与 CDC
| 方案 | 优势 | 成本与边界 |
|---|---|---|
| 直接双写 | 最少组件 | 需接受丢失或建设额外对账 |
| 表轮询 Outbox | 明确状态、容易恢复与诊断 | 扫描、索引和数据库负载 |
| 通知唤醒+表扫描 | 降低空轮询与响应延迟 | 通知不能代替可靠表记录 |
| CDC 读取提交日志 | 可减少主动扫描 | 运维、Schema 演进与消费恢复复杂 |
Outbox 不是事件溯源:它可以只是可靠传输记录,业务查询仍来自状态表。也不等于无限重试:永久无效事件必须进入可观测的失败终态。
真实案例:认领、退避与死信
TDP 的查询使用 FOR UPDATE SKIP LOCKED 选择待发或过期记录,并在同条语句中设置 PUBLISHING、claim_id、过期时间和递增 attempts。发布成功后的更新要求匹配 claim_id。
当前 Dispatcher 租约为一分钟,达到 20 次尝试的发布失败标记 DEAD;退避按 2^min(attempts, 10) 秒计算。第一次认领会先递增 attempts,所以不能把初次失败等待时间误写成固定一秒。这些数字是案例选择,通用系统应根据 SLO 和依赖能力确定。
已核对的实现 · 本地代码快照
Outbox 认领与所有权确认 backend · c6dbe05c
server/db/queries/runtime.sql · 第 4–50 行
符号:ClaimOutboxEvents / MarkOutboxEventPublished · 核对日期 2026-10-02
来源与提交版本一致
-- name: ClaimOutboxEvents :many
WITH candidates AS (
SELECT id
FROM outbox_events
WHERE (outbox_events.status IN ('PENDING', 'FAILED') AND outbox_events.available_at <= sqlc.arg(now))
OR (outbox_events.status = 'PUBLISHING' AND outbox_events.claim_expires_at <= sqlc.arg(now))
ORDER BY outbox_events.available_at, outbox_events.created_at
FOR UPDATE SKIP LOCKED
LIMIT sqlc.arg(batch_limit)
)
UPDATE outbox_events AS event
SET status = 'PUBLISHING',
attempts = event.attempts + 1,
claim_id = sqlc.arg(claim_id),
claim_expires_at = sqlc.arg(claim_expires_at)
FROM candidates
WHERE event.id = candidates.id
RETURNING event.id, event.aggregate_id, event.event_type, event.sequence, event.payload,
event.created_at, event.attempts, event.claim_id;
-- name: GetNextOutboxAttemptAt :one
SELECT MIN(candidate.due_at)::timestamptz AS due_at
FROM (
(SELECT available_at AS due_at
FROM outbox_events
WHERE status IN ('PENDING', 'FAILED')
ORDER BY available_at
LIMIT 1)
UNION ALL
(SELECT claim_expires_at AS due_at
FROM outbox_events
WHERE status = 'PUBLISHING' AND claim_expires_at IS NOT NULL
ORDER BY claim_expires_at
LIMIT 1)
) AS candidate;
-- name: MarkOutboxEventPublished :execrows
UPDATE outbox_events
SET status = 'PUBLISHED', published_at = sqlc.arg(published_at), last_error = NULL,
claim_id = NULL, claim_expires_at = NULL
WHERE id = sqlc.arg(id) AND status = 'PUBLISHING' AND claim_id = sqlc.arg(claim_id);
-- name: MarkOutboxEventFailed :execrows
UPDATE outbox_events
SET status = sqlc.arg(status), available_at = sqlc.arg(available_at),
last_error = sqlc.arg(last_error), claim_id = NULL, claim_expires_at = NULL
WHERE id = sqlc.arg(id) AND status = 'PUBLISHING' AND claim_id = sqlc.arg(claim_id);片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
发布、认领与退避 backend · c6dbe05c
server/internal/core/execution/application/outbox_dispatcher.go · 第 12–81 行
符号:DispatchOutboxBatch · 核对日期 2026-10-02
来源与提交版本一致
const (
outboxClaimTTL = time.Minute
outboxMaxAttempts = 20
)
// OutboxDispatcher 将领域事务中的 Outbox 事件可靠发布到消息系统。
type OutboxDispatcher struct {
store executionports.OutboxStore
publisher executionports.EventPublisher
clock executionports.Clock
}
func NewOutboxDispatcher(store executionports.OutboxStore, publisher executionports.EventPublisher, clock executionports.Clock) (*OutboxDispatcher, error) {
if store == nil || publisher == nil || clock == nil {
return nil, errors.New("outbox store, event publisher, and clock are required")
}
return &OutboxDispatcher{
store: store,
publisher: publisher,
clock: clock,
}, nil
}
func (dispatcher *OutboxDispatcher) DispatchOutboxBatch(ctx context.Context, limit int) (int, []error) {
if limit < 1 {
return 0, []error{errors.New("outbox batch limit must be positive")}
}
now := dispatcher.clock.Now().UTC()
events, err := dispatcher.store.ClaimOutboxEvents(ctx, now, outboxClaimTTL, limit)
if err != nil {
return 0, []error{fmt.Errorf("claim outbox events: %w", err)}
}
failures := make([]error, 0)
for _, event := range events {
err := dispatcher.publisher.PublishExecutionEvent(ctx, executionports.PublishedEvent{
ID: event.ID,
Type: event.EventType,
CorrelationID: event.AggregateID,
Sequence: event.Sequence,
OccurredAt: event.CreatedAt,
Payload: event.Payload,
})
if err == nil {
err = dispatcher.store.MarkOutboxEventPublished(ctx, event.ID, event.ClaimID, dispatcher.clock.Now().UTC())
if err != nil {
failures = append(failures, fmt.Errorf("mark outbox event %s published: %w", event.ID, err))
}
continue
}
status := "FAILED"
if event.Attempts >= outboxMaxAttempts {
status = "DEAD"
}
availableAt := dispatcher.clock.Now().UTC().Add(time.Second * time.Duration(1<<min(event.Attempts, 10)))
if markErr := dispatcher.store.MarkOutboxEventFailed(ctx, event.ID, event.ClaimID, status, availableAt, err.Error()); markErr != nil {
err = errors.Join(err, markErr)
}
failures = append(failures, fmt.Errorf("publish outbox event %s: %w", event.ID, err))
}
return len(events), failures
}
func (dispatcher *OutboxDispatcher) NextOutboxAttemptAt(ctx context.Context) (time.Time, bool, error) {
next, ok, err := dispatcher.store.NextOutboxAttemptAt(ctx)
if err != nil {
return time.Time{}, false, fmt.Errorf("read next outbox attempt: %w", err)
}
return next, ok, nil
}
片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
发布成功但标记失败只记录失败,未完成记录后续可再次认领,因此需要稳定事件 ID 和消费幂等。SKIP LOCKED 不保证业务事件全局有序;顺序要求需要单独设计。
失败与边界:可恢复不代表一定完成
死信、持续基础设施不可用、保留期误删都可能阻止最终投递。监控待发数量、最老记录年龄、失败次数及 DEAD 数量,提供按事件查看和受控补发能力。
接收端去重和业务修改必须原子提交。仅在缓存存“收到过”无法保护数据库修改;对第三方发邮件、支付等效果,还需下游幂等或对账,不能笼统承诺恰好一次。
迁移练习与参考答案
练习:支付成功后需通知商家,允许重试但不能重复增加商家收入。怎么设计?
参考答案:支付状态和通知事件同事务,事件有稳定 ID。发送器认领并重投;商家同事务插入 Inbox 与入账流水,入账业务键有唯一约束。超时可能已处理,发送方不换 ID。第三方账务无幂等契约时,需要查询支付引用对账,不能用本地 Inbox 单独保证远端效果。
继续阅读:Inbox、幂等与乱序、可观测性。