Kafka Go消费者未按预期触发重平衡问题求助
问题描述
测试Kafka消费者重平衡功能时,启动同属一个消费组的两个消费者,预期当某消费者处理到值为"20"的消息时,该消费者停止并关闭,触发重平衡后另一个消费者接管所有分区继续消费。但实际处理完"20"消息后,所有消费行为都停止了。
消费者代码如下:
package main import ( "errors" "os" "github.com/confluentinc/confluent-kafka-go/kafka" "github.com/sirupsen/logrus" ) type Callback func(*kafka.Message) error type CallbackMap map[string]Callback type Consumer struct { consumer *kafka.Consumer name string callbackMap CallbackMap } func NewKafkaConsumer(topics []string, name string) (*Consumer, error) { bootstrap := os.Getenv("KAFKA_HOSTNAME") consumer, err := kafka.NewConsumer( &kafka.ConfigMap{ "bootstrap.servers": bootstrap, "group.id": "testing", "security.protocol": "plaintext", "go.events.channel.enable": true, }, ) if err != nil { logrus.WithError(err).Fatal("failed to create consumer") return nil, err } if err := consumer.SubscribeTopics(topics, nil); err != nil { logrus.WithError(err).Fatal("failed to subscribe to topics") return nil, err } return &Consumer{ consumer: consumer, name: name, callbackMap: make(CallbackMap), }, nil } func (c *Consumer) Subscribe(key string, callback Callback) { c.callbackMap[key] = callback } func (c *Consumer) UnSubscribe(key string) { delete(c.callbackMap, key) } func (c *Consumer) HandleMessage(message *kafka.Message, chano chan error) { if message.TopicPartition.Error != nil { logrus.WithError(message.TopicPartition.Error).Warning() return } logrus.Errorf("CONSUMER NAME IS =========== %s", c.name) logrus.Errorf("MESSAGE is ============: %s", message) logrus.Errorf("VALUE IS =============: %s", string(message.Value)) logrus.Errorf("PARTITION IS =========== %d", message.TopicPartition.Partition) if string(message.Value) == "20" { logrus.Errorf("Exited consumer: %s", c.name) chano <- errors.New("err") return } } func (c *Consumer) Consume() { logrus.Errorf("entering consumer: %s", c.name) chano := make(chan error) for event := range c.consumer.Events() { select { case _ = <-chano: err := c.consumer.Close() if err != nil { logrus.Errorf("Failed to close consumer") return } return default: switch e := event.(type) { case *kafka.Message: c.HandleMessage(e, chano) case kafka.Error: logrus.WithError(e).Warning("an error occurred while reading from topic") default: // Ignore other event types } } } } func main() { kafkaConsumer1, err := NewKafkaConsumer([]string{"test_topic"}, "one") kafkaConsumer2, err := NewKafkaConsumer([]string{"test_topic"}, "two") logrus.Errorf("started") if err != nil { logrus.Errorf("err %s", err) return } go kafkaConsumer1.Consume() go kafkaConsumer2.Consume() select {} }
处理"20"消息的日志如下:
ERRO[0078] CONSUMER NAME IS =========== one ERRO[0078] MESSAGE is ============: test_topic[3]@947 ERRO[0078] VALUE IS =============: 20 ERRO[0078] PARTITION IS =========== 3 ERRO[0078] Exited consumer: one
原因分析
无缓冲通道引发goroutine阻塞
chano是无缓冲通道,在HandleMessage中发送消息时,由于同一goroutine内没有即时接收者,会导致该goroutine阻塞,无法执行后续的consumer.Close()逻辑。消费者无法正常退出消费组,Kafka集群会认为该消费者仍活跃,不会触发重平衡,原本分配给它的分区无人消费。未处理重平衡核心事件
当前代码仅处理了消息和错误事件,忽略了kafka.AssignedPartitions和kafka.RevokedPartitions事件。即使消费者正常退出,存活的消费者无法接收新的分区分配指令,不会开始消费新分配的分区。自动提交偏移量时序问题
默认自动提交是在每次poll时提交上一批消息的偏移量,处理"20"消息后立即退出,该消息的偏移量可能未提交,导致后续重平衡后分区可能重复消费,但这不是消费停止的直接原因。
解决方案
修改为缓冲通道
将chano := make(chan error)改为chano := make(chan error, 1),避免发送操作阻塞goroutine,确保消费者能正常执行关闭逻辑:func (c *Consumer) Consume() { logrus.Errorf("entering consumer: %s", c.name) chano := make(chan error, 1) // 改为缓冲通道 // 后续逻辑不变 }添加重平衡事件处理
在Consume函数的事件分支中处理重平衡相关事件,确保存活消费者能接收并处理新的分区分配:switch e := event.(type) { case *kafka.Message: c.HandleMessage(e, chano) case kafka.Error: logrus.WithError(e).Warning("an error occurred while reading from topic") case kafka.AssignedPartitions: if err := c.consumer.Assign(e.Partitions); err != nil { logrus.WithError(err).Error("failed to assign partitions") } case kafka.RevokedPartitions: if err := c.consumer.Unassign(); err != nil { logrus.WithError(err).Error("failed to unassign partitions") } default: // Ignore other event types }可选:手动控制偏移量提交
如果需要精确控制偏移量,可改为手动提交,在HandleMessage处理完消息后提交偏移量:func (c *Consumer) HandleMessage(message *kafka.Message, chano chan error) { // ... 原有逻辑 if string(message.Value) == "20" { // 先提交偏移量再退出 if _, err := c.consumer.CommitMessage(message); err != nil { logrus.WithError(err).Error("failed to commit message offset") } logrus.Errorf("Exited consumer: %s", c.name) chano <- errors.New("err") return } // 正常消息处理后也可提交 if _, err := c.consumer.CommitMessage(message); err != nil { logrus.WithError(err).Error("failed to commit message offset") } }同时在消费者配置中关闭自动提交:
&kafka.ConfigMap{ // ... 原有配置 "enable.auto.commit": false, }
内容的提问来源于stack exchange,提问作者Omri. B

