Kafka关闭自动提交后如何手动提交消息偏移量?
我已经配置了关闭自动提交的Kafka消费者:
// Create Kafka consumer with auto-commit disabled kafkaConfig := &kafka.ConfigMap{} err := kafkaConfig.SetKey("enable.auto.commit", "false")
消费者将消息读取后发送到Worker处理:
func (c *consumer) start() { c.running = true for c.running { // read the message from Kafka (does NOT commit the offset) msg, err := c.k.ReadMessage(-1) if err == nil { // send the message to the channel for processing by workers c.msgChan <- msg
我希望Worker处理成功后提交消息偏移量,于是写了这段代码:
func worker(id int, c *consumer) { for c.running { select { case msg, ok := <-c.msgChan: if !ok { // Channel closed return } timestamp := strconv.FormatInt(time.Now().Unix()*1000, 10) m, err := newMessage(msg.Value, timestamp) if err != nil { log.Printf("Worker %d: Error creating new message: %v\n", id, err) break } payload, err := m.stringifyMessage() if err != nil { log.Printf("Worker %d: Error decoding message to JSON string: %v\n", id, err) break } key := uuid.New() status := c.r.Set(context.Background(), key.String(), string(payload), time.Hour*6) if status.Err() != nil { log.Printf("Worker %d: Failed to persist message %s. %v\n", id, key.String(), status.Err()) break } connID := m.getConnectionID() log.Printf("Worker %d:\t%s\tconnection:\t%s", id, key.String(), connID) // Commit the offset after successful processing _, commitErr := c.k.CommitMessage(msg) if commitErr != nil { log.Printf("Worker %d: Failed to commit message: %v (%v)\n", id, commitErr, msg) } case _ = <-c.sigChan: c.running = false return // Exit the worker gracefully } } }
但构建镜像时出现错误:
#0 8.386 ./worker.go:44:24: c.k.CommitMessage undefined (type Kafka has no field or method CommitMessage)
改成_, commitErr := c.CommitMessage(msg)后仍报错:
#0 9.260 ./worker.go:44:22: c.CommitMessage undefined (type *consumer has no field or method CommitMessage)
使用的依赖为github.com/confluentinc/confluent-kafka-go/kafka。
问题核心是你调用的提交方法不符合confluent-kafka-go库的API规范,该库的kafka.Consumer并没有CommitMessage这个方法,正确的偏移量提交方式如下:
方式一:直接调用CommitOffsets提交指定偏移量
从消息中提取Topic、Partition和Offset信息,构造TopicPartition对象后提交,注意要把偏移量设为当前消息的偏移量+1(告诉Kafka下一次从这条消息的下一条开始消费):
// 替换原来的CommitMessage调用代码 tp := kafka.TopicPartition{ Topic: msg.TopicPartition.Topic, Partition: msg.TopicPartition.Partition, Offset: msg.TopicPartition.Offset + 1, } commitErr := c.k.CommitOffsets([]kafka.TopicPartition{tp}) if commitErr != nil { log.Printf("Worker %d: Failed to commit message: %v (%v)\n", id, commitErr, msg) }
方式二:给自定义consumer结构体添加封装方法
如果想保留c.CommitMessage(msg)的调用形式,可以在你的consumer结构体中添加一个封装方法:
// 假设你的consumer结构体定义如下 type consumer struct { k *kafka.Consumer msgChan chan kafka.Message running bool sigChan chan os.Signal r redis.Client // 根据你的代码推测的Redis客户端字段 } // 添加CommitMessage封装方法 func (c *consumer) CommitMessage(msg kafka.Message) error { tp := kafka.TopicPartition{ Topic: msg.TopicPartition.Topic, Partition: msg.TopicPartition.Partition, Offset: msg.TopicPartition.Offset + 1, } return c.k.CommitOffsets([]kafka.TopicPartition{tp}) }
之后在Worker中就可以使用commitErr := c.CommitMessage(msg)来提交偏移量了。
内容的提问来源于stack exchange,提问作者Tjorriemorrie

