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

如何使用Sarama在Apache Kafka中实现多主题订阅

如何用Sarama实现Apache Kafka多主题订阅

嘿,我来帮你搞定这个问题!你当前用的ConsumePartition方法是用来订阅单个主题的单个分区的,所以没法直接实现多主题订阅。下面给你两种常用的解决方案,按需选择:

方案1:使用普通消费者订阅多主题(所有分区)

如果不需要指定特定分区,只想订阅三个主题的所有分区,可以用Sarama的Consume方法,它支持传入多个主题:

import (
    "fmt"
    "os"
    "os/signal"
    "syscall"

    "github.com/Shopify/sarama"
)

func main() {
    // 配置Sarama客户端
    config := sarama.NewConfig()
    config.Consumer.Return.Errors = true
    // 根据你的Kafka版本设置对应的Version,比如这里用2.0.0
    config.Version = sarama.V2_0_0_0

    // 初始化消费者
    master, err := sarama.NewConsumer([]string{"localhost:9092"}, config)
    if err != nil {
        panic(err)
    }
    defer func() {
        if err := master.Close(); err != nil {
            panic(err)
        }
    }()

    // 定义要订阅的主题列表
    topics := []string{"Payments", "System", "Orders"}

    // 订阅所有主题的所有分区
    consumer, err := master.Consume(topics)
    if err != nil {
        panic(err)
    }

    // 监听中断信号,优雅退出
    sigchan := make(chan os.Signal, 1)
    signal.Notify(sigchan, syscall.SIGINT, syscall.SIGTERM)

    // 处理消息和错误
    doneCh := make(chan struct{})
    go func() {
        for {
            select {
            case msg := <-consumer.Messages():
                fmt.Printf("主题: %s, 分区: %d, 偏移量: %d, 消息: %s\n",
                    msg.Topic, msg.Partition, msg.Offset, string(msg.Value))
            case err := <-consumer.Errors():
                fmt.Printf("消费错误: %v\n", err)
            case <-sigchan:
                fmt.Println("收到中断信号,开始退出...")
                doneCh <- struct{}{}
                return
            }
        }
    }()

    <-doneCh
    fmt.Println("消费者已退出")
}

方案2:使用消费者组(推荐生产环境使用)

如果你的场景需要负载均衡、容错处理,或者要和其他消费者共享消费组,那更推荐用Sarama的消费者组功能。这种方式是Kafka官方推荐的消费模式:

import (
    "context"
    "fmt"
    "os"
    "os/signal"
    "syscall"

    "github.com/Shopify/sarama"
)

// 自定义消费者组处理器
type ConsumerGroupHandler struct{}

// Setup 在消费者组启动前执行(比如初始化资源)
func (h *ConsumerGroupHandler) Setup(_ sarama.ConsumerGroupSession) error {
    return nil
}

// Cleanup 在消费者组关闭后执行(比如清理资源)
func (h *ConsumerGroupHandler) Cleanup(_ sarama.ConsumerGroupSession) error {
    return nil
}

// ConsumeClaim 处理每个分区的消息
func (h *ConsumerGroupHandler) ConsumeClaim(session sarama.ConsumerGroupSession, claim sarama.ConsumerGroupClaim) error {
    for msg := range claim.Messages() {
        fmt.Printf("消费组: %s, 主题: %s, 分区: %d, 偏移量: %d, 消息: %s\n",
            session.MemberID(), msg.Topic, msg.Partition, msg.Offset, string(msg.Value))
        // 标记消息已处理(提交偏移量)
        session.MarkMessage(msg, "")
    }
    return nil
}

func main() {
    // 配置消费者组
    config := sarama.NewConfig()
    config.Consumer.Return.Errors = true
    config.Version = sarama.V2_0_0_0
    // 设置消费起始偏移量,这里用最旧的未消费消息
    config.Consumer.Offsets.Initial = sarama.OffsetOldest

    // 消费者组名称、Kafka地址、要订阅的主题
    groupID := "my-consumer-group"
    brokers := []string{"localhost:9092"}
    topics := []string{"Payments", "System", "Orders"}

    // 创建消费者组
    consumerGroup, err := sarama.NewConsumerGroup(brokers, groupID, config)
    if err != nil {
        panic(err)
    }
    defer func() {
        if err := consumerGroup.Close(); err != nil {
            panic(err)
        }
    }()

    // 监听中断信号
    sigchan := make(chan os.Signal, 1)
    signal.Notify(sigchan, syscall.SIGINT, syscall.SIGTERM)

    // 启动消费循环
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()

    doneCh := make(chan struct{})
    go func() {
        for {
            // 订阅主题并开始消费
            if err := consumerGroup.Consume(ctx, topics, &ConsumerGroupHandler{}); err != nil {
                fmt.Printf("消费者组错误: %v\n", err)
            }
            // 如果上下文被取消,退出循环
            if ctx.Err() != nil {
                doneCh <- struct{}{}
                return
            }
        }
    }()

    <-sigchan
    fmt.Println("收到中断信号,开始退出...")
    cancel()
    <-doneCh
    fmt.Println("消费者组已退出")
}

注意事项

  • 确保你的Sarama版本和Kafka集群版本兼容,config.Version要设置正确的Kafka版本号。
  • 如果使用普通消费者,记得在退出时关闭消费者实例,避免资源泄漏。
  • 消费者组模式下,同一个消费组的多个消费者会自动分摊主题的分区负载,适合多实例部署的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 07:40:41