You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.01 09:53:19