nsq需严格按序启动nsqlookupd、nsqd、nsqadmin并正确配置地址;producer必须全局复用、监听err()通道并手动重连;publishasync需配错误channel,connecttonsqlookupd才支持服务发现;消息不保证exactly-once。

NSQ 不是开箱即用的“自动重连+服务发现”消息队列,Go 客户端一连不上就 panic,发不出消息也不报错,连对了但收不到消息也常见——问题几乎全出在启动顺序、连接方式和错误处理这三处。
nsqlookupd、nsqd、nsqadmin 必须按序启动且地址配对
本地调试时 nsq.NewProducer("127.0.0.1:4150", cfg) 直接 panic 或报 connection refused,99% 是 nsqd 根本没起来,或它没连上 nsqlookupd。
- 必须严格按顺序执行:先
nsqlookupd(监听4160TCP /4161HTTP),再nsqd --lookupd-tcp-address=127.0.0.1:4160 --broadcast-address=127.0.0.1,最后nsqadmin --lookupd-http-address=127.0.0.1:4161 -
--broadcast-address=127.0.0.1不可省略:不设或设成内网 IP(如192.168.x.x),消费者通过 lookupd 拿到的就是错地址,本地必连不上 - 验证是否就绪:
curl -s http://localhost:4161/topics应返回空 JSON;http://localhost:4171能打开管理页且 Topics 列表无报错
producer.Publish 和 PublishAsync 的错误处理逻辑截然不同
很多人把 PublishAsync 当作“更快的 Publish”,结果消息静默丢失,日志里连 error 都没有。
Go 配置库,使用 spf13/viper — 分层优先级(flag > env >file > KV > default),提供 BindPFlag/BindPFlags、SetEnvPrefix + SetEnvKeyReplace 等功能。
-
producer.Publish("topic", []byte("msg"))是同步阻塞调用,返回error可直接if err != nil判断;适合低频、关键路径(如审计日志) -
producer.PublishAsync("topic", []byte("msg"), nil)若第三个参数传nil,失败会直接吞掉,消息永久丢失 - 高频生产务必用带缓冲的错误 channel:
errChan := make(chan *nsq.Error, 100),再起 goroutine 消费:go func() { for range errChan { /* 记日志 or 告警 */ } }() - 别在
PublishAsync回调里做 DB 写入、HTTP 请求或加锁——会阻塞 producer 内部 goroutine,导致后续消息堆积超时
ConnectToNSQD 和 ConnectToNSQLookupd 容易用反
这两个方法名极具迷惑性:名字带 NSQD 的是直连单节点,名字带 NSQLookupd 的才是走服务发现——生产环境必须用后者。
- 本地单机调试可用
consumer.ConnectToNSQD("127.0.0.1:4150"),简单省事 - 一旦上生产、多 nsqd 实例、要扩缩容,必须用
consumer.ConnectToNSQLookupd("127.0.0.1:4161")(注意是 lookupd 的 TCP 端口4160,不是 HTTP 端口4161) - 调用
ConnectToNSQLookupd后不会立刻收消息:注册是异步的,需等 lookupd 同步完成(通常毫秒级),首次消费前建议轮询/nodes接口确认 - 漏掉
consumer.AddHandler(&MyHandler{}),或HandleMessage返回非nil,消息都会被重发
Producer 必须全局复用并监听 Err() 通道
NSQ Producer 默认不重连,网络抖动或 nsqd 重启后,Publish() 可能一直返回 nil error 却发不出消息——因为底层 TCP 已断,但缓冲区还没刷完。
- 禁止在 HTTP handler 里反复
nsq.NewProducer():每个实例占一个 TCP 连接 + goroutine,短生命周期创建会快速耗尽本地端口 -
nsq.NewProducer()后必须显式调p.Connect(),否则Publish()可能阻塞或 panic - 必须立即启动 goroutine 监听
p.Err(),捕获"io: read/write timeout"、"use of closed network connection"等底层错误 - 重连不是调
p.Connect()再试一次,而是要p.Stop()→ 新建nsq.Producer实例 →Connect()→ 重新监听Err() - 高频场景务必设
config.MaxInFlight = 100(默认是 1),否则并发 publish 会被限流卡死
最常被忽略的是:NSQ 本身不提供 exactly-once 或强顺序保证,它只承诺最终一致性。线上静默失效往往不是代码写错了,而是连接断了没重连、错误回调为空、或者消费者连到了 lookupd 却没等同步完成就开始消费——这些地方一漏,消息就进黑洞,连日志都难查。
golang免费学习笔记(深入):立即使用
在学习笔记中,你将探索golang的核心概念和高级技巧!










