必须走消息队列而非伪异步,因rabbitmq默认publish为fire-and-forget,无确认机制易丢消息;需显式开启confirm、设persistent、序号管理、失败落库或转dlx,nats需幂等ack与过期键,kafka-go须设超时与分区键,asynq需taskid防重复。

直接结论:别用 goroutine 直接调 HTTP/gRPC 做“伪异步”,必须走消息队列(RabbitMQ / NATS JetStream / Kafka),否则根本不算解耦,只是把阻塞藏起来而已。
为什么 RabbitMQ 的 Publish 会丢消息?
默认 channel.Publish() 是 fire-and-forget,网络抖动、Broker 拒绝、磁盘写失败时,Go 程序完全感知不到——消息静默丢失。
- 必须显式开启发布确认:
channel.Confirm()+channel.NotifyPublish()监听成功/失败 - 每条消息要设
amqp.Publishing{DeliveryMode: amqp.Persistent},否则重启后消息消失 - 不要在 for 循环里连续
Publish(),得等上一条确认回调返回后再发下一条,或自己维护序号映射批量确认 - 失败时不能只打日志,得写入本地
outbox表,由后台任务重试;或转发到死信交换器(DLX)人工介入
NATS JetStream 消费者怎么避免重复处理?
JetStream 默认是 At-Least-Once,msg.Ack() 必须在业务逻辑执行完且落库成功后才调,否则超时自动重投——但 ACK 太晚又可能被重复消费。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
- 订阅时用
js.PullSubscribe(),传nats.MaxDeliver(3)和nats.BackOff(...)控制重试节奏 - Handler 开头就用
redis.SetNX(ctx, "evt:"+msg.Header.Get("Nats-Stream"), "1", 24*time.Hour)做幂等标记(带过期) - 更稳妥的是用业务主键 upsert,比如
INSERT INTO event_log (id, status) VALUES (?, 'processed') ON CONFLICT (id) DO NOTHING - 千万别在
Ack()前调远程服务(如发邮件),先落库再触发,否则 ACK 了但邮件没发出去,就不可逆丢失
Kafka-go 生产者卡死的常见原因
不是 Kafka 本身卡,而是 kafka-go 客户端默认超时无限等待,一次网络抖动就能让整个 goroutine 挂住。
- 必须显式设置
WriteTimeout和ReadTimeout,例如WriteTimeout: 5 * time.Second - PartitionKey 不设会导致同一
user_id散列到不同分区,下游无法保证顺序——用kafka.StringPartitioner或自定义 key - 主题名必须带版本,比如
events.user.updated.v1,消费者升级 v2 时能过滤旧格式 - 序列化只用
json.Marshal(),别碰gob,跨语言消费者(Python/Node.js)一解析就 panic
asynq 里 TaskID 为什么不能省?
asynq.NewTask() 不设 asynq.TaskID(...),相同 payload 多次进队列就会生成多个任务——支付回调、库存扣减这类场景,重复执行就是资损。
- ID 要基于业务唯一标识构造,比如
fmt.Sprintf("pay_%s", getBizID(body)),而不是用随机 UUID - payload 里只放 ID,敏感字段(密钥、token)绝对不进 Redis,worker 启动后查 DB 获取详情
- 用
asynq.Timeout(30*time.Second)防止 handler 卡死拖垮整个队列 - 别把
asynq当通用消息队列用——它本质是任务调度器,没有消息广播、分区、顺序语义,只适合单点可靠执行
真正难的不是发消息,而是当 Broker 不可用、消费者宕机、网络分区时,你的事件还能不能最终一致。所有中间件都只是工具,关键在幂等设计、状态可追溯、失败有兜底——这些没法靠配置解决,得写进 handler 里。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










