上一篇讲事务、锁、死锁和并发扣库存。那篇文章最后留了一个问题:事务里不要发 MQ,也不要调第三方接口。
但真实业务偏偏经常这么写:
tx := db.Begin()orderID := createOrder(tx)deductStock(tx)tx.Commit()mq.Publish("order.created", orderID)这段代码看着很合理。数据库提交成功以后,再发消息。
问题也藏在这里:如果 Commit 成功了,进程在 Publish 之前崩了呢?
订单已经创建,库存已经扣了,但后续消息没发出去。物流不知道,积分不知道,通知不知道,搜索索引也不知道。
这个坑不是 MQ 能单独解决的,也不是 PostgreSQL 事务能单独解决的。它发生在两者之间。
Outbox 模式解决的就是这个缝隙。
完整代码在 postgres-outbox-lab[1]。这篇文章里的核心代码都来自这个仓库。
很多人第一次听 Outbox,会把它理解成:
业务表旁边加一张 outbox 表,把消息先存进去。
这句话没错,但太浅了。
Outbox 真正要解决的是这三个问题:
• 业务数据和待投递事件必须一起提交。 • 事件投递失败之后必须能重试。 • 事件可能重复投递,下游必须能幂等。
也就是说,Outbox 不承诺 exactly-once。
它更诚实:我保证事件不会因为进程崩溃而凭空丢掉,但我允许事件被投递多次。
这也是工程里更容易落地的可靠性边界。
先看表结构
实验里只有三张表。
一张订单表:
CREATE TABLE outbox_lab_orders ( id BIGSERIAL PRIMARY KEY, order_no TEXT NOT NULL UNIQUE, amount_cents BIGINT NOT NULL CHECK (amount_cents > 0), status TEXT NOT NULL DEFAULT 'created', created_at TIMESTAMPTZ NOT NULL DEFAULT now());一张 outbox 表:
CREATE TABLE outbox_lab_messages ( id BIGSERIAL PRIMARY KEY, aggregate_type TEXT NOT NULL, aggregate_id BIGINT NOT NULL, topic TEXT NOT NULL, payload JSONB NOT NULL, status TEXT NOT NULL CHECK (status IN ('pending', 'processing', 'sent', 'failed')), attempts INTEGER NOT NULL DEFAULT 0, next_run_at TIMESTAMPTZ NOT NULL DEFAULT now(), locked_by TEXT, locked_at TIMESTAMPTZ, last_error TEXT, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), sent_at TIMESTAMPTZ);再加一个索引:
CREATE INDEX outbox_lab_messages_ready_idxON outbox_lab_messages (status, next_run_at, id);实验里还放了一张 published_events 表:
CREATE TABLE outbox_lab_published_events ( id BIGSERIAL PRIMARY KEY, message_id BIGINT NOT NULL, topic TEXT NOT NULL, payload JSONB NOT NULL, publisher TEXT NOT NULL, created_at TIMESTAMPTZ NOT NULL DEFAULT now());这张表不是生产环境的标准配置。它只是为了让实验可观察:我们把“发布到 MQ”模拟成“写入 published_events 表”,这样可以清楚看到某条消息到底被发布了几次。
真正上线时,这里可能是 Redis Streams、asynq、Kafka、RabbitMQ,或者你自己的事件总线。
第一步:业务数据和 outbox 同事务写入
Outbox 的核心不是 relay,而是写入边界。
业务数据和 outbox 消息必须在同一个 PostgreSQL 事务里提交。
代码在 CreateOrder:
func (s *Store) CreateOrder(ctx context.Context, in CreateOrderInput) (CreateOrderResult, error) { tx, err := s.pool.BeginTx(ctx, pgx.TxOptions{IsoLevel: pgx.ReadCommitted}) if err != nil { return CreateOrderResult{}, err } deferfunc() { _ = tx.Rollback(ctx) }() var orderID int64 if err := tx.QueryRow(ctx, `INSERT INTO outbox_lab_orders (order_no, amount_cents)VALUES ($1, $2)RETURNING id`, in.OrderNo, in.AmountCents).Scan(&orderID); err != nil { return CreateOrderResult{}, fmt.Errorf("insert order: %w", err) } payload, err := json.Marshal(map[string]any{ "event_id": fmt.Sprintf("order-%d-created", orderID), "order_id": orderID, "order_no": in.OrderNo, "amount_cents": in.AmountCents, }) if err != nil { return CreateOrderResult{}, fmt.Errorf("marshal payload: %w", err) } var messageID int64 if err := tx.QueryRow(ctx, `INSERT INTO outbox_lab_messages ( aggregate_type, aggregate_id, topic, payload, status)VALUES ('order', $1, 'order.created', $2::jsonb, 'pending')RETURNING id`, orderID, string(payload)).Scan(&messageID); err != nil { return CreateOrderResult{}, fmt.Errorf("insert outbox message: %w", err) } if err := tx.Commit(ctx); err != nil { return CreateOrderResult{}, err } return CreateOrderResult{OrderID: orderID, MessageID: messageID}, nil}这里有个关键点:CreateOrder 没有发消息。
它只做两件事:
• 写订单。 • 写 outbox 消息。
这两件事在同一个事务里。要么一起成功,要么一起回滚。
如果进程在 Commit 之前崩了,订单没提交,outbox 消息也没提交。
如果 Commit 成功了,订单和 outbox 消息都在数据库里。即使进程下一秒崩了,relay 也能之后把消息捞出来。
运行:
go run ./cmd/create-order输出类似:
order_id=1 outbox_message_id=1第二步:relay 只从 outbox 表取消息
业务请求不发消息,谁发?
独立 relay 发。
relay 可以是同一个服务里的后台 goroutine,也可以是单独进程。它只关心 outbox 表:
func (s *Store) RelayOnce(ctx context.Context, publisher Publisher, limit int, simulateCrashAfterPublish bool) (RelayResult, error) { if limit <= 0 { limit = 10 } messages, err := s.claim(ctx, limit) if err != nil { return RelayResult{}, err } result := RelayResult{Claimed: len(messages)} for _, msg := range messages { if err := publisher.Publish(ctx, msg); err != nil { result.Failed++ if markErr := s.markFailed(ctx, msg.ID, err); markErr != nil { return result, markErr } continue } result.Published++ if simulateCrashAfterPublish { result.SimulatedCrash = true return result, nil } if err := s.markSent(ctx, msg.ID); err != nil { return result, err } result.MarkedSent++ } return result, nil}这段代码的流程是:
claim pending messages -> publish -> publish failed: mark failed, next_run_at 往后推 -> publish success: mark sent注意:这里不是把“取消息 + 发布 + 标记 sent”全部塞进一个数据库事务。
发布 MQ 是外部副作用。你不能指望 PostgreSQL 事务回滚掉 Kafka、Redis Streams 或 HTTP 调用。
所以 relay 的关键不是“一个大事务包住所有事”,而是承认外部投递可能失败、可能重复,然后用状态机和幂等把系统拉回来。
FOR UPDATE SKIP LOCKED 怎么用
多个 relay 同时跑,不能都拿到同一条消息。
这时用 FOR UPDATE SKIP LOCKED:
rows, err := tx.Query(ctx, `SELECT id, aggregate_id, topic, payload::text, attemptsFROM outbox_lab_messagesWHERE status IN ('pending', 'failed') AND next_run_at <= now()ORDER BY idLIMIT $1FOR UPDATE SKIP LOCKED`, limit)FOR UPDATE 会锁住当前事务选中的行。
SKIP LOCKED 的意思是:如果某些行已经被其他事务锁住,不要等,直接跳过。
所以多个 relay 可以并行抢任务:
relay-a 锁住 message 1,2,3relay-b 查询时跳过 1,2,3,拿到 4,5,6这正是 PostgreSQL 官方文档里提到的 queue-like table 场景。
但这句话也要说完整:SKIP LOCKED 会给你一个“不完整视图”,不适合普通业务查询。它适合任务队列、outbox relay 这类“谁抢到谁处理”的场景。
选中消息后,实验代码把它们标记成 processing:
for _, msg := range messages { if _, err := tx.Exec(ctx, `UPDATE outbox_lab_messagesSET status = 'processing', attempts = attempts + 1, locked_by = $2, locked_at = now()WHERE id = $1`, msg.ID, s.relayID); err != nil { return nil, fmt.Errorf("claim outbox message %d: %w", msg.ID, err) }}然后提交事务。
到这里,claim 结束。真正发布发生在事务外。
跑一次正常 relay
运行:
go run ./cmd/relay输出类似:
claimed=3 published=3 marked_sent=3 failed=0orders=3 outbox_messages=3 published_events=3这说明:
• 创建了 3 笔订单。 • 同时写了 3 条 outbox。 • relay 抢到 3 条。 • 3 条都发布成功。 • 3 条都标记成 sent。
这条路径是最舒服的路径。
但工程里真正要看的不是舒服路径。
最重要的故障窗口:发布成功,但没标记 sent
Outbox 最容易被误解的地方在这里。
很多人会以为:消息在数据库里,所以不会重复。
不对。
看这个窗口:
relay 拿到 outbox messagerelay 发布到 MQ 成功relay 准备 UPDATE outbox SET status = 'sent'进程崩了这时 MQ 已经收到消息,但 outbox 表里还是 processing,或者之后被恢复成 pending。
下一轮 relay 会再次投递。
配套代码里专门模拟了这个窗口:
first, err := store.RelayOnce(ctx, publisher, 1, true)if err != nil { log.Fatal(err)}afterCrash, err := publisher.CountByMessage(ctx, created.MessageID)if err != nil { log.Fatal(err)}recovered, err := store.RecoverProcessing(ctx)if err != nil { log.Fatal(err)}second, err := store.RelayOnce(ctx, publisher, 1, false)if err != nil { log.Fatal(err)}afterRetry, err := publisher.CountByMessage(ctx, created.MessageID)if err != nil { log.Fatal(err)}运行:
go run ./cmd/crash-window输出:
message_id=1first_relay claimed=1 published=1 marked_sent=0 simulated_crash=truepublished_after_crash=1 recovered_processing=1second_relay claimed=1 published=1 marked_sent=1published_after_retry=2这就是 Outbox 的真实语义:at-least-once。
至少投递一次,不保证只投递一次。
如果下游不能接受重复消息,那不是 Outbox 坏了,是下游缺幂等。
下游幂等怎么做
payload 里我放了一个 event_id:
payload, err := json.Marshal(map[string]any{ "event_id": fmt.Sprintf("order-%d-created", orderID), "order_id": orderID, "order_no": in.OrderNo, "amount_cents": in.AmountCents,})消费端应该用这个 event_id 做幂等。
最简单的做法是事件处理表:
CREATE TABLE processed_events ( event_id TEXT PRIMARY KEY, topic TEXT NOT NULL, processed_at TIMESTAMPTZ NOT NULL DEFAULT now());消费时:
INSERT INTO processed_events (event_id, topic)VALUES ($1, $2)ON CONFLICT (event_id) DO NOTHING;如果插入成功,说明第一次处理。
如果 RowsAffected() == 0,说明这个事件处理过,直接 ACK。
更贴近业务的幂等也可以放在业务表上:
• 订单状态机只允许 created -> paid执行一次。• 积分流水用 event_id做唯一键。• 优惠券发放表用 (user_id, coupon_id, event_id)做唯一约束。
关键不是“哪里做幂等”,而是你必须承认重复会发生。
processing 卡住怎么办
实验代码里有一个简化版恢复逻辑:
func (s *Store) RecoverProcessing(ctx context.Context) (int64, error) { tag, err := s.pool.Exec(ctx, `UPDATE outbox_lab_messagesSET status = 'pending', locked_by = NULL, locked_at = NULL, next_run_at = now()WHERE status = 'processing'`) if err != nil { return 0, fmt.Errorf("recover processing messages: %w", err) } return tag.RowsAffected(), nil}这是为了实验容易看懂。
生产环境不要把所有 processing 一把扫回 pending。你至少要加一个超时条件:
UPDATE outbox_messagesSET status = 'pending', locked_by = NULL, locked_at = NULL, next_run_at = now()WHERE status = 'processing' AND locked_at < now() - interval '5 minutes';否则一个 relay 正在正常处理,另一个恢复任务把它的消息改回 pending,就会人为制造重复。
重复可以接受,但没必要主动放大。
失败重试别立刻打爆自己
发布失败时,实验代码这样标记:
func (s *Store) markFailed(ctx context.Context, messageID int64, publishErr error) error { nextRunAt := time.Now().Add(2 * time.Second) _, err := s.pool.Exec(ctx, `UPDATE outbox_lab_messagesSET status = 'failed', next_run_at = $2, locked_by = NULL, locked_at = NULL, last_error = $3WHERE id = $1`, messageID, nextRunAt, publishErr.Error()) if err != nil { return fmt.Errorf("mark failed: %w", err) } return nil}实验里固定 2 秒,是为了容易观察。
生产环境更常见的是指数退避:
1s -> 5s -> 30s -> 2m -> 10m还要考虑最大重试次数。
有些错误不是重试能解决的,比如 payload 格式不合法、topic 配错、下游永远拒绝。一直重试只会把 outbox 表变成垃圾堆。
通常我会给 outbox 加这些字段:
• attempts• next_run_at• last_error• dead_at• dead_reason
超过最大次数就进 dead-letter 状态,再由人工或补偿任务处理。
Outbox 和 asynq / Redis Streams 怎么接
Outbox 不替代队列。
它解决的是“本地事务成功以后,队列投递不能丢”。
真正投递到哪里,取决于你的系统:
type Publisher interface { Publish(ctx context.Context, msg Message) error}实验里用的是 DBPublisher:
func (p *DBPublisher) Publish(ctx context.Context, msg Message) error { _, err := p.pool.Exec(ctx, `INSERT INTO outbox_lab_published_events (message_id, topic, payload, publisher)VALUES ($1, $2, $3::jsonb, $4)`, msg.ID, msg.Topic, msg.Payload, p.name) if err != nil { return fmt.Errorf("publish message %d: %w", msg.ID, err) } return nil}换成 Redis Streams,大概就是:
func (p *StreamsPublisher) Publish(ctx context.Context, msg Message) error { return p.rdb.XAdd(ctx, &redis.XAddArgs{ Stream: msg.Topic, Values: map[string]any{ "message_id": msg.ID, "payload": msg.Payload, }, }).Err()}换成 asynq,大概就是:
func (p *AsynqPublisher) Publish(ctx context.Context, msg Message) error { task := asynq.NewTask(msg.Topic, []byte(msg.Payload)) _, err := p.client.EnqueueContext(ctx, task, asynq.TaskID(fmt.Sprintf("outbox:%d", msg.ID))) return err}这里故意用了 TaskID。因为 outbox relay 可能重复投递,同一个 outbox message 重试时,最好能给队列层一个稳定 ID。
但不要把“队列层去重”当成唯一防线。真正的幂等还是要落到消费端业务状态上。
这套代码上线前还要补什么
实验项目为了讲清楚,保持得很小。
如果要放到生产,我至少会补这些东西:
分批与限速。 relay 每次 claim 多少条、每秒最多投递多少条,要能配置。下游抖动时不能把自己和下游一起打穿。
stale processing 恢复。 只恢复 locked_at 超时的消息,不要恢复正在处理的消息。
dead-letter。 超过最大重试次数后进入 dead 状态,保留错误原因和人工重放入口。
指标。 pending 数量、failed 数量、最老 pending 年龄、投递耗时、重试次数、dead-letter 数量,这些都要能看见。
状态更新条件。 实验代码为了简洁,markSent 只按 id 更新。生产环境最好带上 status = 'processing' 和 locked_by = 当前 relay,避免误把不属于自己的消息标记成 sent。
幂等约束。 消费端必须有事件处理表、业务唯一键或状态机保护。
payload 版本。 事件 payload 迟早会变。建议带 event_type、event_version、event_id。
归档清理。sent 消息不能无限留在主表里。可以按时间归档,或者定期清理。
这些不是“架构洁癖”。Outbox 一旦成为可靠事件投递的底座,它出问题就是整条业务链路出问题。
最后看边界
Outbox 能解决:
• 本地业务写入和“待投递事件”原子提交。 • 进程崩溃后事件不丢。 • relay 多实例并发投递。 • 投递失败后的重试。
Outbox 不能解决:
• 下游业务天然幂等。 • 外部 MQ exactly-once。 • payload 设计混乱。 • 消费端处理一半崩溃产生的业务副作用。
这也是为什么我说 Outbox 不是“加张表”。
那张表只是入口。真正重要的是这条链路:
业务事务写 outbox -> relay 抢消息 -> 外部投递 -> 标记 sent / failed -> 失败重试 -> 消费端幂等把这条链路想清楚,你才是在做可靠事件投递。
引用链接
[1] postgres-outbox-lab: https://github.com/arixbit/postgres-outbox-lab
夜雨聆风