Skip to content

Outbox 与双写问题 ​

问题:订单已经保存,通知却没有发出 ​

最直觉的实现是保存订单后发送消息。数据库提交与消息发送之间崩溃,订单存在但消息丢失。调换顺序也无济于事:消息先到达,数据库后来回滚,接收端看到不存在的订单。

跨两个系统的两次写入无法靠调整顺序获得单库事务原子性。Outbox 把“稍后要发送这件事实”先变成同一个数据库里的可靠记录。

原理:把发送意图纳入事务 ​

一次事务同时写订单与 Outbox 事件。独立发送器认领可发送记录,发布到消息系统,收到确认后标记成功。进程重启后扫描未完成记录,继续发送。

这提供恢复依据,但不消除重复:消息系统接收后,发送器可能在记录成功之前崩溃。重启不知道消息已经送达,只能重投同一事件。消费端必须处理重复。

最小例子:事务里只存事实 ​

text
事务 A:INSERT order + INSERT outbox(event_id, payload) → COMMIT
发送器:认领 event_id → 发布 → 接收确认
事务 B:按 event_id + claim_id 标记 PUBLISHED

event_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、幂等与乱序、可观测性。

理解原理 · 分析取舍 · 用真实代码检验