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

基于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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 16:25:55