如何用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) } } } }
关键改动说明
替换流保留策略
LimitsPolicy:消息会被保留直到达到你设置的存储限制(如最大消息数、最大存储大小、最长保留时长),适合需要长期存储消息的场景。- 如果不需要长期存储,可改用
InterestPolicy:当所有消费组都消费完某条消息后,该消息会被自动删除。
消费组的实现
- 每个消费组使用唯一的
Durable名称:JetStream会为每个Durable维护独立的消费游标,因此不同消费组可以各自读取全量消息,互不干扰。 - 同一
Durable的多个消费者实例:启动多个使用相同durableName的消费者时,JetStream会自动在实例间做负载均衡,消费者上下线时会触发自动重平衡,满足你第一个需求中的负载分配和重平衡要求。
- 每个消费组使用唯一的
精确一次投递保障
- 通过
AckExplicitPolicy配置显式确认机制:只有当你调用msg.Ack()确认消息处理完成后,JetStream才会标记该消息为已消费,避免重复投递。
- 通过
内容的提问来源于stack exchange,提问作者Yura
相关产品推荐
相关产品推荐

