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 必须同一事务
订单从 pending 到 paid 的状态转换和 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 完整技术体系学习路线:从语法、并发到框架、中间件与 AI
- 上一篇:Go MQTT 与 Paho:QoS、会话、保留消息和设备连接
- 下一篇:Go Asynq 异步任务:Redis 队列、重试、定时与唯一任务
- 延伸:Go RabbitMQ 实战:Exchange、Queue、确认、重试与死信
- 延伸:Go Kafka 完整指南:分区、消费者组、Offset 与 franz-go
- 延伸:Go 分布式事务实践:本地事务、Outbox、Saga、TCC 与 DTM
官方资料
本文依据 Go 官方规范、标准库文档和 Go 官方博客重新梳理;正文与示例由 WR BLOG 编写。

评论
0 条讨论