基于Go的Watermill Kafka库多订阅者实现问题咨询
问题:Watermill Kafka订阅者重复Subscribe后仅一个goroutine接收消息
使用Watermill库实现Kafka多订阅者读取指定Topic时,发现同一个消费者组下两次调用Subscribe,仅第二个goroutine能处理消息,第一个goroutine完全收不到数据。
问题代码
package main import ( "context" "fmt" "log" "time" "github.com/IBM/sarama" "github.com/ThreeDotsLabs/watermill" "github.com/ThreeDotsLabs/watermill-kafka/v2/pkg/kafka" "github.com/ThreeDotsLabs/watermill/message" ) func main() { saramaSubConfig := kafka.DefaultSaramaSubscriberConfig() saramaSubConfig.Consumer.Offsets.Initial = sarama.OffsetOldest brokers := []string{"127.0.0.1:9092", "127.0.0.1:9093"} subscriber, err := kafka.NewSubscriber(kafka.SubscriberConfig{ Brokers: brokers, Unmarshaler: kafka.DefaultMarshaler{}, OverwriteSaramaConfig: saramaSubConfig, ConsumerGroup: "TestGroup", }, watermill.NewStdLogger(false, false), ) if err != nil { panic(err) } messages, err := subscriber.Subscribe(context.Background(), "test") if err != nil { panic(err) } go process(messages, "one") messages2, err := subscriber.Subscribe(context.Background(), "test") if err != nil { panic(err) } go process(messages2, "two") pubisher, err := kafka.NewPublisher(kafka.PublisherConfig{ Brokers: brokers, Marshaler: kafka.DefaultMarshaler{}, }, watermill.NewStdLogger(false, false), ) if err != nil { panic(err) } publishMessage(pubisher) } func publishMessage(publisher message.Publisher) { counter := 0 for { msg := message.NewMessage(watermill.NewUUID(), []byte(fmt.Sprintf("Hello world %d", counter))) if err := publisher.Publish("test", msg); err != nil { panic(err) } counter++ time.Sleep(1 * time.Millisecond) } } func process(messages <-chan *message.Message, goroutineId string) { for msg := range messages { log.Printf("goroutine %s recived message: %s, payload: %s", goroutineId, msg.UUID, string(msg.Payload)) msg.Ack() } }
问题原因
同一个Watermill Kafka Subscriber实例针对相同消费者组+Topic重复调用Subscribe时,内部不会创建新的消费者连接,而是复用已有的消费流。第二次调用会覆盖第一个返回的消息通道的订阅关系,导致第一个goroutine的通道不再接收任何消息。
解决方案
方案1:单消息通道多goroutine并行消费(同消费者组分摊消息)
如果仅需要在同一个消费者组内用多个goroutine并行处理消息,无需重复调用Subscribe,直接将同一个消息通道交给多个goroutine即可:
func main() { // ... 原有初始化代码保持不变 ... messages, err := subscriber.Subscribe(context.Background(), "test") if err != nil { panic(err) } // 启动多个goroutine共享同一个消息通道 go process(messages, "one") go process(messages, "two") // ... 发布者代码保持不变 ... }
基于Go通道的特性,多个接收者会轮询获取消息,Kafka的Topic分区消息会被均匀分发到不同goroutine处理,实现并行消费。
方案2:创建独立Subscriber实例(不同消费者组全量消费)
如果需要多个独立订阅者各自接收Topic的全量消息,需为每个订阅者创建独立的Subscriber实例,并指定不同的消费者组:
func main() { saramaSubConfig := kafka.DefaultSaramaSubscriberConfig() saramaSubConfig.Consumer.Offsets.Initial = sarama.OffsetOldest brokers := []string{"127.0.0.1:9092", "127.0.0.1:9093"} // 第一个订阅者:消费者组TestGroup1 subscriber1, err := kafka.NewSubscriber(kafka.SubscriberConfig{ Brokers: brokers, Unmarshaler: kafka.DefaultMarshaler{}, OverwriteSaramaConfig: saramaSubConfig, ConsumerGroup: "TestGroup1", }, watermill.NewStdLogger(false, false)) if err != nil { panic(err) } messages1, err := subscriber1.Subscribe(context.Background(), "test") if err != nil { panic(err) } go process(messages1, "one") // 第二个订阅者:消费者组TestGroup2 subscriber2, err := kafka.NewSubscriber(kafka.SubscriberConfig{ Brokers: brokers, Unmarshaler: kafka.DefaultMarshaler{}, OverwriteSaramaConfig: saramaSubConfig, ConsumerGroup: "TestGroup2", }, watermill.NewStdLogger(false, false)) if err != nil { panic(err) } messages2, err := subscriber2.Subscribe(context.Background(), "test") if err != nil { panic(err) } go process(messages2, "two") // ... 发布者代码保持不变 ... }
这种方式下,两个订阅者属于不同消费者组,都会收到Topic的全量消息,适合独立消费链路的场景。
关键原理
- Kafka消费者组机制:同一消费者组内的消费者分摊Topic分区消息(每个分区仅被组内一个消费者消费);不同消费者组独立消费全量消息。
- Watermill Subscriber设计:同一个Subscriber实例针对相同消费者组+Topic仅维护一个消费流,重复
Subscribe会复用现有资源,导致之前的消息通道失效。
内容的提问来源于stack exchange,提问作者Tom
相关产品推荐
相关产品推荐

