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

如何用NATS构建类Kafka队列?多消费组配置报错求解

可以实现你的需求,问题出在流的保留策略选择上

你遇到的报错是因为使用了WorkQueuePolicy这个流保留策略——它的设计目标是消息被任意消费者确认后立刻从流中删除,因此只允许存在一个无过滤条件的消费者,天然无法支持多个消费组读取同一批消息。

解决方案:修改流的保留策略并调整消费者配置

要同时满足两个需求,你需要将流的保留策略改为LimitsPolicy(或InterestPolicy,根据存储需求选择),再通过不同的Durable名称区分消费组,每个消费组维护独立的消费进度。

修改后的完整代码

import (
    "context"
    "log"
    "time"

    "github.com/nats-io/nats.go"
    "github.com/nats-io/nats.go/jetstream"
)

func consumeFromNats(ctx context.Context, durableName string) error {
    nc, err := nats.Connect("localhost:4222")
    if err != nil {
        return err
    }
    defer nc.Close()

    js, err := jetstream.New(nc)
    if err != nil {
        return err
    }

    streamConfig := jetstream.StreamConfig{
        Name:         "TEST.STREAM",
        Retention:    jetstream.LimitsPolicy, // 替换为LimitsPolicy
        Subjects:     []string{"events"},
        MaxMessages:  10000, // 可选:设置流的最大消息存储量
        MaxAge:       7 * 24 * time.Hour, // 可选:设置消息保留时长
    }

    stream, err := js.CreateOrUpdateStream(ctx, streamConfig)
    if err != nil {
        return err
    }

    consumer, err := stream.CreateOrUpdateConsumer(ctx, jetstream.ConsumerConfig{
        Durable:       durableName, // 每个消费组传入唯一的durable名称
        AckPolicy:     jetstream.AckExplicitPolicy, // 显式确认,保证精确一次投递
        AckWait:       30 * time.Second, // 确认超时时间,按需调整
        MaxAckPending: 100, // 控制未确认消息的最大数量,防止负载过高
    })
    if err != nil {
        return err
    }

    for {
        select {
        case <-ctx.Done():
            return ctx.Err()
        default:
            msg, err := consumer.Next()
            if err != nil {
                if err == jetstream.ErrTimeout {
                    continue
                }
                return err
            }

            // 执行你的业务逻辑
            // ...

            // 显式确认消息,确保消息不会被重复投递
            if err := msg.Ack(); err != nil {
                log.Printf("failed to ack message: %v", err)
            }
        }
    }
}

关键改动说明

  1. 替换流保留策略

    • LimitsPolicy:消息会被保留直到达到你设置的存储限制(如最大消息数、最大存储大小、最长保留时长),适合需要长期存储消息的场景。
    • 如果不需要长期存储,可改用InterestPolicy:当所有消费组都消费完某条消息后,该消息会被自动删除。
  2. 消费组的实现

    • 每个消费组使用唯一的Durable名称:JetStream会为每个Durable维护独立的消费游标,因此不同消费组可以各自读取全量消息,互不干扰。
    • 同一Durable的多个消费者实例:启动多个使用相同durableName的消费者时,JetStream会自动在实例间做负载均衡,消费者上下线时会触发自动重平衡,满足你第一个需求中的负载分配和重平衡要求。
  3. 精确一次投递保障

    • 通过AckExplicitPolicy配置显式确认机制:只有当你调用msg.Ack()确认消息处理完成后,JetStream才会标记该消息为已消费,避免重复投递。

内容的提问来源于stack exchange,提问作者Yura

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.30 16:00:30