Go 基础体系 · 第 74/113 篇。示例统一基于 Go 1.26.4;核心片段可能省略 package 与 import,完整程序可直接按文中结构运行。

Go 可靠消息统一设计:Outbox、幂等、重试、顺序与死信

本文以 Go 1.26.4、PostgreSQL 17.6 和 github.com/jackc/pgx/v5 v5.7.6 为基准,Broker 接口接 NATS、Kafka、RabbitMQ。版本使事务与客户端行为可复现。

可靠消息是端到端属性。订单事务、Relay、Broker、消费事务和 Ack 都可能在动作完成、确认未到达时崩溃。可验证的目标通常是至少一次加业务幂等;脱离事务边界谈 Exactly Once,只是把重复藏到别处。

1. 先定义业务承诺和状态终点

设计前要明确延迟、重复、丢失、局部顺序、Broker 故障策略、人工处理和审计期限。支付、通知、搜索索引的答案不同。

HTTP 命令 -> 数据库业务事务 + Outbox
                         |
                         v
                    Relay 发布 -> Broker 持久化 -> Consumer
                                                      |
                                                      v
                                             Inbox + 业务事务 -> Ack

可靠链路的终点不是 Publish 返回 nil,而是业务状态可查询、失败可解释、重复可吸收、积压可恢复。关键事件通常选择至少一次,并为承诺配置指标和恢复手册。

2. 双写漏洞为何无法靠调整调用顺序消除

先提交数据库再发布,崩溃会让事件未发;先发布再提交,消费者可能看到回滚订单。把发布放进数据库事务也不能让普通 Broker 原子提交,反而延长锁并留下未知结果。

// 错误示意:任意顺序都存在数据库与 Broker 之间的崩溃窗口。
if err := orders.MarkPaid(ctx, orderID); err != nil {
	return err
}
if err := broker.Publish(ctx, event); err != nil {
	return err
}

XA 只有各组件都支持时才可能缩小窗口,同时引入协调者和 in-doubt 事务。多数服务选择 Outbox:业务与待发布意图共享本地事务,跨系统步骤靠重试和幂等收敛。

3. 消息信封是长期契约

Payload 不应只是一段没有身份的 JSON。事件信封至少包含全局唯一 event_id、事件类型、schema version、聚合 ID、聚合版本、发生时间、生产服务和 trace context。事件表示已经发生的事实,命令表示希望执行的动作,两者重试、过期和授权语义不同。

type Envelope struct {
	EventID         string          `json:"event_id"`
	EventType       string          `json:"event_type"`
	SchemaVersion   int             `json:"schema_version"`
	AggregateID     string          `json:"aggregate_id"`
	AggregateVersion int64          `json:"aggregate_version"`
	OccurredAt      time.Time       `json:"occurred_at"`
	TraceParent     string          `json:"trace_parent,omitempty"`
	Data            json.RawMessage `json:"data"`
}

event ID 首次请求时生成,Relay 重试保持不变。排序依赖业务版本而非设备时钟。破坏性 schema 变更发布新版本,并灰度兼容新旧格式;信封不携带无关敏感数据。

4. Outbox 表与约束设计

Outbox 与业务表位于同一数据库。行保存路由、序列化后的不可变 payload、状态、尝试次数和下次执行时间。event_id 唯一约束阻止同一业务请求生成两条身份不同的事件;按 (status, available_at) 建 Relay 扫描索引。

CREATE TABLE message_outbox (
    event_id uuid PRIMARY KEY,
    aggregate_id text NOT NULL,
    aggregate_version bigint NOT NULL,
    topic text NOT NULL,
    payload jsonb NOT NULL,
    status text NOT NULL CHECK (status IN ('pending', 'publishing', 'published', 'failed')),
    attempts integer NOT NULL DEFAULT 0,
    available_at timestamptz NOT NULL DEFAULT now(),
    lease_until timestamptz,
    published_at timestamptz,
    last_error text,
    created_at timestamptz NOT NULL DEFAULT now(),
    UNIQUE (aggregate_id, aggregate_version)
);
CREATE INDEX message_outbox_ready_idx
ON message_outbox (available_at, created_at)
WHERE status IN ('pending', 'publishing');

Payload 在事务内序列化,Relay 不再查询已变化的业务表。Outbox 持续增长,已发布行按审计要求归档;不能确认后立即删除,否则未知发布与对账缺少证据。

5. 业务变更与 Outbox 必须同一事务

订单从 pendingpaid 的状态转换和 Outbox 插入由一个 PostgreSQL 事务提交。并发更新使用状态条件和版本约束,零行更新映射为明确业务冲突。任何一步失败都回滚。

func (s *Store) MarkPaid(ctx context.Context, orderID string, event Envelope) error {
	tx, err := s.pool.Begin(ctx)
	if err != nil {
		return fmt.Errorf("begin order transaction: %w", err)
	}
	defer tx.Rollback(ctx)

	result, err := tx.Exec(ctx, `
		UPDATE orders SET state = 'paid', version = version + 1
		WHERE order_id = $1 AND state = 'pending'`, orderID)
	if err != nil {
		return fmt.Errorf("mark order paid: %w", err)
	}
	if result.RowsAffected() != 1 {
		return ErrOrderStateConflict
	}
	payload, err := json.Marshal(event)
	if err != nil {
		return fmt.Errorf("marshal order event: %w", err)
	}
	if _, err := tx.Exec(ctx, `
		INSERT INTO message_outbox(event_id, aggregate_id, aggregate_version, topic, payload, status)
		VALUES ($1, $2, $3, $4, $5, 'pending')`,
		event.EventID, event.AggregateID, event.AggregateVersion, event.EventType, payload); err != nil {
		return fmt.Errorf("insert outbox event: %w", err)
	}
	if err := tx.Commit(ctx); err != nil {
		return fmt.Errorf("commit order transaction: %w", err)
	}
	return nil
}

Commit 返回网络错误时结果可能未知。API 使用请求幂等键或 order ID 查询最终状态,不能直接重建新 event ID 再执行。defer Rollback 的错误在已成功 Commit 后可忽略;业务函数不同时记录并返回同一错误。

6. Relay 抢占、租约和并发安全

多个 Relay 实例需要安全分工。PostgreSQL 可用 FOR UPDATE SKIP LOCKED 在短事务中领取批次,设置 publishing 与租约后立即提交,再在事务外访问 Broker。不能持有数据库锁等待网络。

WITH picked AS (
    SELECT event_id
    FROM message_outbox
    WHERE (status = 'pending' OR (status = 'publishing' AND lease_until < now()))
      AND available_at <= now()
    ORDER BY created_at
    FOR UPDATE SKIP LOCKED
    LIMIT 100
)
UPDATE message_outbox AS o
SET status = 'publishing',
    lease_until = now() + interval '30 seconds',
    attempts = attempts + 1
FROM picked
WHERE o.event_id = picked.event_id
RETURNING o.event_id, o.topic, o.payload, o.attempts;

租约长于发布 P99,过期后可由其他实例领取。批量、发布并发和数据库连接都有界,并控制同一 aggregate 的并发顺序。

7. 发布确认与不可消除的重复窗口

Relay 发布时把 event ID 映射到 Broker 的幂等键或消息头,要求持久发布确认。收到确认后再把 Outbox 标记 published。如果 Broker 已持久化而 Relay 在更新数据库前崩溃,租约过期后必然重复发布;这是 Outbox 算法允许重复以避免丢失的关键窗口。

type Publisher interface {
	Publish(ctx context.Context, message Message) (Confirmation, error)
}

func relayOne(ctx context.Context, publisher Publisher, row OutboxRow) error {
	confirmation, err := publisher.Publish(ctx, Message{
		ID:      row.EventID,
		Topic:   row.Topic,
		Payload: row.Payload,
	})
	if err != nil {
		return fmt.Errorf("publish outbox event %q: %w", row.EventID, err)
	}
	if err := markPublished(ctx, row.EventID, confirmation); err != nil {
		return fmt.Errorf("mark event %q published: %w", row.EventID, err)
	}
	return nil
}

确认信息保存 Stream sequence、partition/offset 或 Broker message ID。发布超时视为未知;重试保持 ID。Broker 短期去重不能替代 Inbox。

8. Inbox 将去重和业务副作用放进同一事务

消费者收到消息后,先尝试插入 (consumer_name, event_id) 唯一键,并在同一事务执行实际业务。若插入冲突说明此前事务已提交,可安全 Ack。Inbox 与业务更新分成两个事务会留下“已登记但未处理”的丢失窗口。

CREATE TABLE message_inbox (
    consumer_name text NOT NULL,
    event_id uuid NOT NULL,
    aggregate_id text NOT NULL,
    processed_at timestamptz NOT NULL DEFAULT now(),
    PRIMARY KEY (consumer_name, event_id)
);
func (s *Store) Consume(ctx context.Context, consumer string, event Envelope) error {
	tx, err := s.pool.Begin(ctx)
	if err != nil {
		return fmt.Errorf("begin inbox transaction: %w", err)
	}
	defer tx.Rollback(ctx)

	result, err := tx.Exec(ctx, `
		INSERT INTO message_inbox(consumer_name, event_id, aggregate_id)
		VALUES ($1, $2, $3) ON CONFLICT DO NOTHING`,
		consumer, event.EventID, event.AggregateID)
	if err != nil {
		return fmt.Errorf("insert inbox event: %w", err)
	}
	if result.RowsAffected() == 0 {
		return tx.Commit(ctx)
	}
	if err := applyProjection(ctx, tx, event); err != nil {
		return fmt.Errorf("apply projection: %w", err)
	}
	if err := tx.Commit(ctx); err != nil {
		return fmt.Errorf("commit inbox transaction: %w", err)
	}
	return nil
}

不同消费者使用不同 consumer name。Inbox 保留时间至少覆盖 Broker 重放、备份恢复和人工重放周期,避免旧消息再次产生副作用。

9. Ack、offset 与本地提交的正确顺序

处理成功提交后才向 Broker Ack 或提交 offset。先 Ack 后提交,进程在中间退出会永久丢失;先提交后 Ack,Ack 丢失只会重复,而 Inbox 可以吸收。可靠设计通常主动选择后者。

receive -> validate -> BEGIN -> inbox + business update -> COMMIT -> ACK
                           |             |
                           `---失败------`--> rollback + retry
COMMIT 后 ACK 前崩溃 -------------------------------------> redelivery + inbox 去重

Kafka 不能让后完成消息越过前一 offset;RabbitMQ/NATS 也要控制在途。长处理需延长可见性且可取消;关闭先停止拉取,再等待事务与 Ack。

10. 错误分类、退避和重试预算

瞬时网络/锁错误可以重试;schema 或业务前置条件错误通常永久失败;Commit 或发布响应丢失则结果未知。三类不能统一重试。

func retryDelay(attempt int, random *rand.Rand) time.Duration {
	const maxDelay = 2 * time.Minute
	delay := 500 * time.Millisecond * time.Duration(1<<min(attempt, 8))
	if delay > maxDelay {
		delay = maxDelay
	}
	return time.Duration(random.Int64N(int64(delay) + 1))
}

重试受次数、总时间和过期时间约束。Broker、handler 和 HTTP 重试不能层层放大;通常由一层持久调度负责,限流退避仍受业务截止时间限制。

11. 死信不是垃圾桶而是恢复队列

超过预算或永久失败的消息进入失败表,保存信封、来源、消费者、错误代码、尝试次数和时间;Payload 仍受加密和保留策略约束。

CREATE TABLE message_failure (
    failure_id bigserial PRIMARY KEY,
    consumer_name text NOT NULL,
    event_id uuid NOT NULL,
    payload jsonb NOT NULL,
    error_code text NOT NULL,
    error_detail text NOT NULL,
    attempts integer NOT NULL,
    first_failed_at timestamptz NOT NULL,
    last_failed_at timestamptz NOT NULL,
    replayed_at timestamptz
);

运营入口支持查询、修复、小批重放和审计。重放前检查版本、目标状态与容量,保留 event ID;修改 payload 必须形成新的受审计事件。

12. 顺序、聚合版本与缺口恢复

多数 Broker 只保证 partition 或队列内局部顺序。以 aggregate ID 路由可让同一订单进入同一分区,但热点键会限制并行度。

消费者保存 aggregate_version:旧版本忽略,当前加一才应用,更大表示缺口并触发回补。不能丢弃未来版本;可交换增量则可设计成无序合并。

并发 Relay 也可能反转顺序,可按 aggregate 领取、加锁或用 Broker key 分区,再由版本校正。仅为真实业务不变量承担严格顺序成本。

13. 背压与积压容量

Relay 批次、Consumer prefetch、worker 和下游请求都必须有界。下游变慢时暂停拉取,让积压留在 Broker 或 Outbox,而不是进程堆。

恢复时间 = 当前积压 / (可持续消费速率 - 新增速率)
Little's Law: 在途数量约等于吞吐率 × 平均处理时间

新增 800/s、消费 1000/s、积压 720 万条,理想恢复需 10 小时;扩 worker 仅在下游有余量时有效。过期命令进入明确终态,大 Payload 放对象存储并携带校验引用。

14. 补偿、Saga 与事务边界

Outbox/Inbox 保证消息最终传播和消费幂等,不会自动让多个服务的业务状态同时提交。跨服务流程用 Saga 把每一步建成可重试的本地事务,并为已完成步骤定义补偿。补偿是新业务动作,不是时间倒流:退款可能失败,库存释放时商品状态可能已变化,通知无法“撤回已阅读”。

订单确认 -> 预占库存 -> 扣款 -> 发货
              |           |
失败补偿 <----释放库存 <---退款请求

Saga 保存全局 ID、状态机、步骤幂等键、截止时间和审计。补偿前先查询参与者,避免响应丢失时反向操作;不可自动补偿则进入对账和人工处理,并阻止流程继续产生冲突动作。

15. 对账修复是可靠性的最后一层

定时对账从业务真相源查询“已支付但无发布确认”“消费者长期无结果”“Saga 卡住”等不变量,按业务 ID 和版本核对,而非只比较消息数。修复通过同一幂等命令执行,记录操作者、原因和前后状态;批量操作先 dry-run、限速和小样本观察。禁止直接改库绕过 Outbox。

16. 故障矩阵与恢复推理

至少推演:业务提交后 Relay 前崩溃,Outbox 稍后被领取;Broker 确认丢失或标记 published 前崩溃,Relay 同 ID 重发;消费事务提交后 Ack 前崩溃,Inbox 去重;DLQ 写入失败时原消息不能确认。

故障点                         可见状态                       恢复
DB commit 后、publish 前        pending outbox                 Relay 重试
Broker commit 后、响应前        outbox 未确认、Broker 可能已有   同 ID 重发
Consumer commit 后、Ack 前      inbox 已有、Broker 未确认        重投后去重
补偿提交后、响应前              Saga 状态未知                   查询参与者

网络超时表示“未知”。恢复依据持久状态查询,不根据错误字符串猜测未执行;服务重启、Broker leader 切换、数据库切换和网络分区都要演练。

17. 监控、追踪和审计

生产端监控 Outbox pending、最老年龄、确认延迟和失败率;Broker 监控积压、重投、磁盘与副本;消费者监控处理耗时、Inbox 冲突、DLQ 和版本缺口。端到端 SLO 以事件发生到业务提交的年龄衡量。日志带 event ID、aggregate version、Broker position、consumer 和错误代码;敏感 payload 不落日志。消费 trace 使用新 span/link,告警应直接对应恢复动作。

18. 测试、竞态与故障注入

单测覆盖信封校验、退避、版本状态机和 Inbox 重复;真实 PostgreSQL 17.6 集成测试验证唯一约束、SKIP LOCKED 多 Relay 和回滚;Broker 测试验证确认、重投和排空。

gofmt -w .
go test ./...
go test -race ./...
go test -run TestRelayLeaseRecovery -count=100 ./...
go test -bench=BenchmarkInbox -benchmem ./...

故障注入在 DB Commit、Publish Ack、Inbox Commit 与 Broker Ack 边界杀进程,并注入延迟、断连和磁盘写满。断言最终状态、重复副作用、恢复证据和 goroutine 退出,而不只检查返回值。

19. 性能、分区和存储治理

批量领取和发布降低往返,却扩大重试范围;批次来自 P99 和 Broker 限制。Outbox 按时间分区,ready 索引只覆盖未完成行;Inbox 分区不能破坏去重范围。基准包含真实 payload、数据库 fsync、Broker 副本和下游事务。不得靠跳过确认或降低耐久参数换取吞吐。

20. 安全、部署与选型检查

Relay 只获 Outbox、租约和指定 Topic 权限;消费者只订阅自己的 Topic、写必要业务表。数据库和 Broker 使用 TLS/mTLS 或短期凭据,schema 校验大小、类型和租户。先发布兼容新旧 schema 的消费者,再发布生产者;滚动退出先停止拉取并排空,数据库迁移采用 expand/migrate/contract。

选 Broker 时比较持久确认、回放、局部顺序、死信、集群故障语义和运维经验。Outbox 解决生产事务边界,Inbox 吸收重复,Saga/补偿推动跨服务收敛,对账覆盖机制外异常;四层共同落地,可靠消息才是可证明、可恢复的系统能力。


系列导航与关联阅读

官方资料

本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。