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

Go NSQ 基础:Topic、Channel、消费者与存量系统使用边界

本文以 Go 1.26.4、NSQ v1.3.0 和 github.com/nsqio/go-nsq v1.1.0 为基准。NSQ 是 Go 编写的分布式实时消息平台,核心模型很小:Producer 向 Topic 发布,不同 Channel 各自得到一份消息,同一 Channel 内多个 Consumer 分摊消息。它没有 Kafka 式分区日志和长期 offset,也没有 RabbitMQ 式 Exchange 路由;理解这些缺失比记住 API 更重要。

NSQ 在存量 Go 系统中仍有价值,但 v1.3.0 发布于 2023 年,升级频率和治理生态都明显弱于活跃的新平台。维护者应补齐幂等、死信、权限隔离、恢复演练和容量文档;新系统则应把项目活跃度、回放、顺序和安全需求纳入选型。

1. nsqd、nsqlookupd、nsqadmin 的职责

nsqd 接收、排队并投递消息,是数据平面;每个节点独立拥有自己的 Topic/Channel 数据,不靠共识协议复制到其他 nsqdnsqlookupd 保存节点注册和 Topic 拓扑,供消费者发现生产节点,是发现平面而非消息存储。nsqadmin 聚合 HTTP API,展示深度、速率、客户端和错误,并提供受控管理操作。

producer --------> nsqd-A ---- Topic ---- Channel ---- consumer group
       \----------> nsqd-B         \----- Channel ---- another service
                         \注册
                          nsqlookupd <----- consumer 查询与轮询
                          nsqadmin   <----- 运维观察

生产者通常直接持有一个或多个 nsqd 地址;消费者从 lookupd 发现所有承载该 Topic 的节点,并分别建立 TCP 连接。lookupd 故障不会立刻中断既有连接,但会阻止拓扑更新。因为没有内建跨节点复制,发送到单个 nsqd 后磁盘永久损坏可能丢消息,不能把多节点部署误读为自动三副本。

2. Topic、Channel 与扇出语义

Topic 是发布入口。Channel 第一次创建后开始独立排队,同一消息会复制到当时存在的每个 Channel;某个 Channel 没有在线 Consumer 时仍保留自己的积压。Channel 内多个客户端竞争消息,实现负载均衡。

临时 Topic/Channel 名以 #ephemeral 结尾,最后一个客户端断开后会删除,因此只适合允许丢失的在线流。普通 Channel 应在发布前通过部署或管理 API 创建,否则早于 Channel 创建的消息不会神奇地补发给它。

curl -fsS -X POST 'http://127.0.0.1:4151/topic/create?topic=article_events'
curl -fsS -X POST 'http://127.0.0.1:4151/channel/create?topic=article_events&channel=search_index'
curl -fsS --data-binary '{"event_id":"evt-42"}' \
  'http://127.0.0.1:4151/pub?topic=article_events'

Topic/Channel 名应低基数且稳定,不要为每个租户或订单动态创建。多租户隔离、schema version 与事件类型可以放进受控命名或消息信封;无限创建 Channel 会增加文件、内存和运维成本。

3. TCP 握手和连接生命周期

go-nsq 与 nsqd 建立长 TCP 连接,发送协议魔数后执行 IDENTIFY,协商心跳、压缩、TLS、消息超时等能力。服务端周期发送 _heartbeat_,客户端以 NOP 响应。Consumer 还要发送 SUB topic channelRDY n 才开始接收。

config := nsq.NewConfig()
config.DialTimeout = 3 * time.Second
config.ReadTimeout = 60 * time.Second
config.WriteTimeout = 5 * time.Second
config.HeartbeatInterval = 30 * time.Second
config.MaxInFlight = 64

consumer, err := nsq.NewConsumer("article_events", "search_index", config)
if err != nil {
	return fmt.Errorf("create nsq consumer: %w", err)
}

读超时必须大于心跳间隔并留出网络抖动。连接中断后 Consumer 会退避重连并继续从 lookupd 发现节点;在途未 FIN 的消息超时后会重投。停止时调用 consumer.Stop() 并等待 <-consumer.StopChan,使 RDY 归零并完成在途 handler,而不是直接结束进程。

4. 发布确认与生产者复用

Producer 应长期复用。Publish 把命令写到指定 nsqd,只有收到 OK 才返回 nil;若写入或响应阶段断线,结果可能未知,消息可能已经入队。重试必须保留相同 event ID,并依赖消费者幂等。

producer, err := nsq.NewProducer("127.0.0.1:4150", nsq.NewConfig())
if err != nil {
	return fmt.Errorf("create nsq producer: %w", err)
}
defer producer.Stop()

if err := producer.Ping(); err != nil {
	return fmt.Errorf("ping nsqd: %w", err)
}
if err := producer.Publish("article_events", payload); err != nil {
	return fmt.Errorf("publish article event: %w", err)
}

NSQ 不替生产者完成多节点复制。应用可按确定规则选择节点或通过外层代理分发,但向两个节点双发会产生两条独立消息。业务变更与发布之间仍有双写窗口,关键事件应使用数据库 Outbox;Relay 发布成功后标记已发送,崩溃造成的重复由 Inbox 消除。

MultiPublish 能批量减少往返,但批次越大,失败重试的重复范围越大。AsyncPublish 必须读取 response channel,限制并发 pending,并在进程退出前停止接收新任务、等待结果、再停止 Producer。

5. 消息帧、FIN、REQ 与 TOUCH

每条消息包含 16 字节 ID、时间戳、尝试次数和 Body。Consumer 收到后进入 in-flight。处理成功发送 FIN message_id;暂时失败发送 REQ message_id timeout,服务端延迟后重新排队;任务仍在执行时发送 TOUCH message_id 重置超时。

nsqd --MSG--> consumer
  |              |-- FIN ------> 完成并移出 in-flight
  |              |-- REQ ------> 延迟后重新入队,Attempts 增加
  |              `-- TOUCH ----> 延长当前处理期限
  `-- timeout -----------------> 自动重新入队

默认 Handler 返回 nil 时 go-nsq 自动 FIN,返回 error 时自动 REQ。错误发生在本地事务提交后,返回 error 会重复投递;错误发生在提交前却返回 nil 会丢业务效果。因此提交是 Ack 边界,handler 必须识别重复。TOUCH 只适合有明确上限的长处理,反复 TOUCH 会让卡死任务永远占据 in-flight。

6. 可运行的 Handler 与错误分类

普通 Handler 的返回值无法表达永久失败,因此复杂系统可关闭自动响应,显式选择 Finish 或 Requeue。下面保留自动模式,并把永久无效消息写入独立失败存储后返回 nil;若失败存储也不可用则返回错误,避免证据丢失。

consumer.AddHandler(nsq.HandlerFunc(func(message *nsq.Message) error {
	ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
	defer cancel()

	err := service.Apply(ctx, message.ID.String(), message.Body)
	if err == nil || errors.Is(err, ErrDuplicate) {
		return nil
	}
	if errors.Is(err, ErrInvalidEvent) {
		if archiveErr := failures.Save(ctx, message.ID.String(), message.Body, err); archiveErr != nil {
			return fmt.Errorf("archive invalid event: %w", archiveErr)
		}
		return nil
	}
	return fmt.Errorf("apply article event: %w", err)
}))

if err := consumer.ConnectToNSQLookupd("127.0.0.1:4161"); err != nil {
	return fmt.Errorf("connect to nsqlookupd: %w", err)
}

请求级 context 应由进程生命周期和处理期限派生。回调内部创建的 goroutine 若在返回后继续使用 Body,会破坏 Ack 和内存所有权;需要并发时使用 AddConcurrentHandlers 的固定并发数,或复制数据交给有界且可关闭的 worker。

7. RDY、MaxInFlight 与背压

NSQ 的关键流控是 RDY。每条连接告诉 nsqd 最多可再发送多少消息;收到消息会消耗一个 RDY 额度,FIN/REQ 后库再分配。MaxInFlight 是一个 Consumer 跨连接的总预算,go-nsq 根据连接数量重新分配 RDY。

如果连接数大于 MaxInFlight,一部分连接会周期轮换 RDY,避免少数节点永久饥饿。提高 MaxInFlight 只增加在途并发,不会提高慢数据库的容量;过大会放大内存、锁竞争和超时重投。并发数应小于等于下游可承受的连接/请求预算,并通过处理 P99 与 msg_timeout 校准。

config.MaxInFlight = 32
config.MsgTimeout = 30 * time.Second
config.MaxBackoffDuration = 2 * time.Minute
config.RDYRedistributeInterval = 5 * time.Second
consumer.AddConcurrentHandlers(handler, 16)

当 handler 连续失败,go-nsq 的 BackoffStrategy 会降低 RDY、暂停消费后探测恢复。这能保护下游,但不能替代失败分类;一个永久毒消息会反复触发退避并拖慢整个 Channel,应在最大尝试次数后归档。

8. 重试、最大尝试与自建死信

NSQ 没有功能完整的内建 DLQ 路由。MaxAttempts 到达后,go-nsq 默认调用失败日志器并放弃继续交给普通 Handler;业务若需要可查询、可修复、可重放,必须实现失败处理器,把原始 Body、消息 ID、Topic、Channel、Attempts、错误分类和时间存入持久系统。

consumer.SetLoggerLevel(nsq.LogLevelWarning)
consumer.AddHandler(handler)
consumer.SetBehaviorDelegate(&delegate{failureStore: failures})

业务永久错误不应经历全部指数退避;瞬时网络错误可重试;限流错误按服务端建议延迟;未知提交结果先查询业务状态。重试需要上限和抖动。死信重放生成新的操作记录并保留原 event ID,先验证 schema 与目标状态,分批限速,不能把整批失败消息直接灌回线上形成第二次事故。

9. 幂等消费与事务边界

NSQ 提供至少一次语义,重复是正常路径。消费者使用稳定 event ID,不应只用 NSQ 消息 ID:生产者在未知结果后重新 Publish 会得到新的 NSQ ID。推荐消息信封包含业务 event_idaggregate_idversionoccurred_atpayload

BEGIN;
INSERT INTO message_inbox(consumer_name, event_id, received_at)
VALUES ('search_index', $1, now())
ON CONFLICT DO NOTHING;
-- 只有确实插入 inbox 时才更新业务状态。
UPDATE article_index_state
SET version = $2, title = $3
WHERE article_id = $4 AND version < $2;
COMMIT;

Inbox 插入与业务更新必须在同一事务;事务提交后 Handler 才返回 nil。若只在缓存记录 ID,缓存淘汰后重复仍会穿透;若先写 Inbox 再在另一事务更新业务,二者之间崩溃会把未完成事件误认为成功。

10. 顺序语义及其限制

NSQ 不保证严格顺序。多 nsqd 发布、多个 Consumer、RDY 重分配、REQ 和超时重投都会打乱完成顺序。即使单节点、单 Channel、单 handler,也可能因失败重入队让后续消息先完成。

需要实体顺序时让事件携带递增版本:当前版本为 7,只接受 8;小于等于 7 视为重复,大于 8 暂缓并触发缺口修复。也可按聚合 ID 哈希到固定 Topic/Channel 分片并单线程消费,但扩缩分片会改变映射,热点键仍会阻塞。真正依赖可回放分区顺序的系统通常更适合 Kafka、Pulsar 或明确支持该模型的平台。

11. 内存队列、磁盘队列与恢复

每个 Topic/Channel 有内存队列;超过 mem-queue-size 后消息落到磁盘队列。设置为 0 可让普通消息尽快走磁盘,但仍有协议与操作系统缓冲。磁盘队列是本节点持久化,不是复制日志;正常退出会同步状态,崩溃恢复读取队列元数据并继续投递。

nsqd \
  --data-path=/var/lib/nsq \
  --mem-queue-size=0 \
  --max-msg-size=1048576 \
  --msg-timeout=30s \
  --max-msg-timeout=15m \
  --lookupd-tcp-address=nsqlookupd-1:4160

磁盘满、文件损坏或节点永久丢失会越过这一保证。数据盘要有容量告警和备份策略,但复制正在变化的队列文件不一定形成一致快照,应按官方工具和恢复演练验证。若业务要求单节点损坏不丢消息,应在生产端保留 Outbox 或选择原生复制 Broker。

12. 发现、拓扑变化与故障模式

Consumer 周期查询所有 lookupd,合并 nsqd 地址并连接新节点。多个 lookupd 彼此不复制,但 nsqd 可同时注册到多个;Consumer 也应配置多个。lookupd 短暂不可用时已有消息连接继续工作,新节点和新 Topic 发现延迟。

典型故障包括:发布响应丢失导致重复;handler 提交后未 FIN 导致重投;处理超过 timeout 导致同一消息并发执行;节点磁盘满停止入队;下游变慢造成 depth、in-flight 和 requeue 上升;错误配置让临时 Channel 在断线后删除。每种故障都要写清可观察信号、恢复动作和数据核对方式。

curl -fsS 'http://127.0.0.1:4151/stats?format=json'
curl -fsS 'http://127.0.0.1:4161/nodes'
curl -fsS 'http://127.0.0.1:4171/api/topics'

不要依赖手工 curl 作为长期监控,采集 depth、memory_depth、backend_depth、in_flight_count、deferred_count、requeue_count、timeout_count、message_count 和客户端状态到指标系统。

13. 诊断、测试与性能验证

积压升高时先分解发布速率、完成速率、失败重入队和 handler P99;只加消费者可能把数据库压垮。timeout 上升要检查 GC 暂停、网络、处理时长和 MsgTimeout。单节点 backend depth 异常常表示生产分配不均或该节点消费者连接缺失。

nsq_to_file --topic=article_events --channel=audit \
  --lookupd-http-address=127.0.0.1:4161 --output-dir=/var/spool/nsq-audit
go test ./...
go test -race ./...
go test -run TestDuplicateDelivery -count=100 ./...

集成测试使用真实 NSQ v1.3.0,至少覆盖 Consumer 在提交后 FIN 前退出、REQ 延迟、超过消息 timeout、lookupd 重启、nsqd 重启和磁盘逼近上限。压力测试使用真实 Body 大小与下游延迟,观察稳定状态而非数秒峰值;分别测 Publish、MultiPublish 和多 Channel 扇出,因为每加一个 Channel 都增加排队与磁盘成本。

14. 安全与生产部署

NSQ 原生安全治理较弱。支持 TLS 的链路应启用证书校验,可结合 --auth-http-address 做授权;HTTP 管理端口和广播地址只暴露在受控网络,通过防火墙、反向代理身份认证和审计限制访问。不要把公网设备直接连到 nsqd,也不要在 Body 中无保护地放密钥和敏感个人数据。

nsqd --tls-required=true \
  --tls-cert=/run/secrets/nsqd.crt \
  --tls-key=/run/secrets/nsqd.key \
  --tls-root-ca-file=/run/secrets/ca.crt \
  --auth-http-address=http://nsq-auth.internal:4180

每个 nsqd 使用独立持久盘,限制文件描述符和消息大小,滚动维护前先把生产流量移走并等待 Channel 清空。升级前验证 go-nsq、压缩和 TLS 协商兼容,保留回滚路径。

15. 使用边界与迁移判断

NSQ 适合模型简单、低运维复杂度、允许至少一次、无需长期回放且已有成熟运行经验的异步分发。它不适合依赖严格分区顺序、跨节点复制、按 offset 长期回放、细粒度 ACL、事务生产或庞大连接设备生态的系统。

若主要是事件日志和流处理,比较 Kafka/Pulsar;若需要复杂路由、死信与队列治理,比较 RabbitMQ;若追求轻量服务通信并需要可选持久化,比较 NATS JetStream;设备弱网连接则比较 MQTT Broker。迁移时同时运行旧 NSQ Relay 和新 Broker,用相同 event ID 做消费者幂等与结果核对,先迁移读取再停止写入。选型不是追逐新技术,而是确保平台的故障语义覆盖业务承诺,并且团队能持续维护。


系列导航与关联阅读

官方资料

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