乐于分享
好东西不私藏

AI 时代重学 PostgreSQL 实战:Outbox 不是加张表,用 PostgreSQL + Go 实现可靠事件投递

AI 时代重学 PostgreSQL 实战:Outbox 不是加张表,用 PostgreSQL + Go 实现可靠事件投递

上一篇讲事务、锁、死锁和并发扣库存。那篇文章最后留了一个问题:事务里不要发 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 (&#x27;pending&#x27;, &#x27;failed&#x27;)  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 = &#x27;processing&#x27;,    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 = &#x27;sent&#x27;进程崩了

这时 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 = &#x27;pending&#x27;,    locked_by = NULL,    locked_at = NULL,    next_run_at = now()WHERE status = &#x27;processing&#x27;`)    if err != nil {        return 0, fmt.Errorf("recover processing messages: %w", err)    }    return tag.RowsAffected(), nil}

这是为了实验容易看懂。

生产环境不要把所有 processing 一把扫回 pending。你至少要加一个超时条件:

UPDATE outbox_messagesSET status = &#x27;pending&#x27;,    locked_by = NULL,    locked_at = NULL,    next_run_at = now()WHERE status = &#x27;processing&#x27;  AND locked_at < now() - interval &#x27;5 minutes&#x27;;

否则一个 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 = &#x27;failed&#x27;,    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_typeevent_versionevent_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