必须用本地事务表写入消息,否则发出去的消息和业务状态一定不一致;因为先提交数据库事务再发消息时,若网络抖动、mq不可用或进程崩溃,消息丢失而数据库已落库,导致状态撕裂。

必须用本地事务表写入消息,否则发出去的消息和业务状态一定不一致。
为什么不能在事务提交后再发消息
常见错误是先 tx.Commit(),再调用 kafka.Producer.WriteMessages() 或 amqp.Publish()。一旦网络抖动、MQ 拒绝连接或进程 panic,消息就彻底丢失,而数据库已落库——下游永远收不到事件。
- 用
go func() { }()异步发也不行:goroutine 可能没跑完就被 OOM kill 或容器重启干掉 -
time.AfterFunc更危险:它依赖进程存活,重启后所有待触发逻辑全丢 - 哪怕用了 RabbitMQ 的
channel.Tx(),也只是保证“发到 Broker”,不防消费者宕机、不防网络分区、不提供补偿路径
outbox_events 表怎么建才安全
核心不是“有没有这张表”,而是字段设计和写入时机是否绑定业务事务。
Go语言(Golang)1.26.0版本提供 Go 官方 Windows amd64 MSI 安装包下载入口,版本号 1.26.0,可用于旧项目维护、兼容性测试和指定版本开发环境配置。
- 必须包含:
id(主键)、topic(如"order.created")、payload(JSON 字符串)、processed(BOOLEAN DEFAULT false)、created_at - 禁止用两个事务分别插入订单和消息:主从延迟下轮询可能读到未同步的
processed = false记录,导致重复投递 - 禁止在
defer tx.Commit()后写消息表——panic 时事务回滚,但消息已落库,状态撕裂 - 正确姿势:在同一个
*sql.Tx中,先INSERT INTO orders,再INSERT INTO outbox_events,最后tx.Commit()
轮询投递器要注意什么
轮询不是简单查表更新,要防并发冲突、DB 压力和静默失败。
- PostgreSQL 用
SELECT ... FOR UPDATE SKIP LOCKED;MySQL 建议加乐观锁字段(如version)并用UPDATE ... WHERE id = ? AND version = ? - 别用固定 100ms 轮询:指数退避更稳,比如首次
time.Second * 1,失败后*2、*4,上限设为 30 秒 - UPDATE 必须带
WHERE processed = false条件,否则多实例同时运行会重复处理同一条 - 投递失败的消息不能跳过,要进死信队列(DLQ),推荐用 Kafka compact topic 或 NATS JetStream 存储待重试事件
消费者端不幂等,前面全白做
Kafka 和 RabbitMQ 都不保证“恰好一次”。网络分区、手动 nack、broker 重启都会触发重投——你得靠自己拦住重复。
- 别只靠
UPDATE SET stock = stock - 1 WHERE order_id = ? AND status = 'pending':status 字段可能被并发更新覆盖 - 优先用数据库唯一约束:比如建
UNIQUE (order_id, event_type),插入失败就直接 return - 或者查状态再处理:
SELECT COUNT(*) FROM inventory_adjustments WHERE order_id = ? AND type = 'deduct' AND status = 'done',有结果就跳过 - 避免逻辑错位:恢复库存时写
WHERE status = 'frozen',而不是WHERE status != 'restored',后者可能误放行已修复记录
最容易被忽略的是对账——再严谨的流程也会漏:主从延迟导致轮询读旧数据、MQ 分区丢某批次、补偿任务被误删……这些不会报错,只会静默不一致。对账不是备选方案,是强制兜底动作。










