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

Kafka Go消费者未按预期触发重平衡问题求助

Kafka消费者重平衡异常问题排查

问题描述

测试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     

原因分析

  1. 无缓冲通道引发goroutine阻塞
    chano是无缓冲通道,在HandleMessage中发送消息时,由于同一goroutine内没有即时接收者,会导致该goroutine阻塞,无法执行后续的consumer.Close()逻辑。消费者无法正常退出消费组,Kafka集群会认为该消费者仍活跃,不会触发重平衡,原本分配给它的分区无人消费。

  2. 未处理重平衡核心事件
    当前代码仅处理了消息和错误事件,忽略了kafka.AssignedPartitions和kafka.RevokedPartitions事件。即使消费者正常退出,存活的消费者无法接收新的分区分配指令,不会开始消费新分配的分区。

  3. 自动提交偏移量时序问题
    默认自动提交是在每次poll时提交上一批消息的偏移量,处理"20"消息后立即退出,该消息的偏移量可能未提交,导致后续重平衡后分区可能重复消费,但这不是消费停止的直接原因。

解决方案

  1. 修改为缓冲通道
    将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) // 改为缓冲通道
        // 后续逻辑不变
    }
    
  2. 添加重平衡事件处理
    在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
    }
    
  3. 可选:手动控制偏移量提交
    如果需要精确控制偏移量,可改为手动提交,在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 01:37:30