NATS-JetStream双消费者仅其一接收消息问题排查求助
NATS-JetStream 同Stream多消费者订阅不同Subject收不到消息排查
我用NATS-JetStream创建了两个持久化消费者,它们订阅不同的Subject、属于同一条Stream,且都配置了唯一的持久化名称。发布者单独部署在Kubernetes Pod中,两个消费者分别部署在同一命名空间的两个Pod里。现在的问题是,当发布者向两个Subject发送事件时,只有第一个消费者能收到消息,另一个完全收不到。我怀疑是Pod订阅时存在竞态条件,但还没定位到具体问题,求帮忙排查。
相关代码如下:
func StartNatsJetstream() error { fmt.Println("StartNatsJetstream runs") StartSubject := common.SoConfig.MessageBus.WorkflowRequestQueue + ".*" startConsumer := "tart-consumer" //Consumer name can be anything fmt.Printf("StartNatsJetstream runs%s", tartSubject) // Connect to the NATS server nc, js, err := natsjetstream.SetupJetStream(startSubject) if err != nil { log.Errorf("Nats server is not running", err) return err } defer nc.Close() fmt.Println("Connection got") // Subscribe to the stream and process messages _, err = js.Subscribe(StartSubject, func(msg *nats.Msg) { fmt.Printf("Received message: %s\n", msg.Data) // Acknowledge the message to prevent it from being sent again err := msg.Ack() if err != nil { fmt.Printf("Error acknowledging message: %v\n", err) return } // Process the message receivedPayload := msg.Data var msgData Message err = json.Unmarshal([]byte(receivedPayload), &msgData) if err != nil { fmt.Println("Error unmarshaling JSON:", err) return } //Decoded byte message understand to the user in logs decodedPayload, err := base64.StdEncoding.DecodeString(msgData.Payload) if err != nil { fmt.Println("Error decoding payload:", err) return } var payloadData PayloadData err = json.Unmarshal(decodedPayload, &payloadData) if err != nil { fmt.Println("Error unmarshaling JSON from decoded payload:", err) return } msgData.Payload = string(decodedPayload) fmt.Printf("Received message: %+v\n", msgData) //Send byte msgData to process workflow req. go processMessageNatsJetstream(msgData) }, nats.Durable(StartConsumer)) if err != nil { log.Errorf("Error while subscribe", err) } fmt.Println("Function completed") // To keep the goroutine running select {} }
排查方向
1. 检查Stream的Subject覆盖范围
- 用
nats stream info <stream-name>查看Stream的Subjects配置,确认两个消费者订阅的Subject(包括通配符匹配的情况)都在Stream的捕获范围内。如果Stream只配置了其中一个Subject,另一个消费者自然收不到消息。 - 验证发布者发送的消息是否被Stream存储:执行
nats stream view <stream-name>,确认两个Subject的消息都已写入Stream。
2. 核对消费者的订阅Subject是否正确
- 检查两个消费者的
StartSubject实际值:当前代码中用的是WorkflowRequestQueue + ".*",另一个消费者是否配置了对应的目标Subject(比如WorkflowResponseQueue + ".*")? - 修复代码中的笔误:
fmt.Printf("StartNatsJetstream runs%s", tartSubject)中的tartSubject未定义,应为StartSubject,修复后查看日志确认实际订阅的Subject是否符合预期。
3. 确认持久化消费者名称无冲突
- 务必保证两个消费者的
nats.Durable()参数值完全唯一。如果名称重复,后启动的消费者会覆盖前一个,导致只有一个能正常工作。 - 用
nats consumer info <stream-name> <durable-name>分别查询两个持久化消费者的状态,确认它们都存在且处于活跃状态。
4. 排查Kubernetes网络与日志
- 在消费者Pod内执行
nats ping,验证到NATS服务器的连通性是否正常。 - 查看消费者Pod的完整日志,确认是否存在连接失败、订阅失败的错误信息(比如代码中
log.Errorf("Error while subscribe", err)是否有输出)。
5. 检查SetupJetStream方法的逻辑
- 确认
SetupJetStream(startSubject)内部是否会重新创建Stream?如果第二个消费者启动时该方法覆盖了原有Stream配置,可能导致消息分发规则被修改。 - 确保Stream的创建逻辑是幂等的,仅在Stream不存在时创建,避免重复修改配置。
6. 修复订阅失败的处理逻辑
- 当前代码中订阅失败仅打印日志,程序仍会进入
select{},导致消费者看似启动但实际未订阅成功。建议在订阅失败时直接返回错误,终止程序,避免无意义的运行。
内容的提问来源于stack exchange,提问作者Aniruddha Kulkarni
相关产品推荐
相关产品推荐

