如何在同一连接中重新消费Kafka分区的未提交消息?
同一franz-go连接中重复消费未提交偏移量的消息问题
使用franz-go作为Kafka客户端库,场景如下:向topic生产消息后,第一次消费该消息但不提交偏移量,希望无需重新连接就能再次消费该消息(重新连接时可正常获取)。但当前代码中第二次调用PollFetches时陷入无限等待,请问如何实现同一连接下的重复消费?
原代码示例
package main import ( "context" "github.com/twmb/franz-go/pkg/kgo" ) func main() { cl, err := kgo.NewClient( kgo.SeedBrokers("localhost:9092"), kgo.ConsumerGroup("some_consumer_group"), kgo.ConsumeTopics("topic"), kgo.DefaultProduceTopic("topic"), kgo.DisableAutoCommit(), ) if err != nil { panic(err) } record := &kgo.Record{ Key: []byte("some_key"), Value: []byte("some value"), } ctx := context.Background() produceResult := cl.ProduceSync(ctx, record) if produceResult.FirstErr() != nil { panic(produceResult.FirstErr()) } fetches := cl.PollFetches(ctx) if fetches.Err() != nil { panic(fetches.Err()) } //未提交偏移量 doSomething(fetches) //希望再次获取同一消息并处理,但此处陷入无限等待 fetches = cl.PollFetches(ctx) if fetches.Err() != nil { panic(fetches.Err()) } doSomething(fetches) err = cl.CommitRecords(ctx, fetches.Records()...) if err != nil { panic(err) } } func doSomething(result kgo.Fetches) { //处理Kafka消息时出现异常 }
问题原因
franz-go的消费者在成功拉取消息后,会更新本地的消费位置游标(即使未提交偏移量到broker)。默认情况下,客户端会基于这个本地游标向broker请求下一批消息,不会主动重复拉取已经fetch过的消息,因此第二次PollFetches会一直等待新消息,导致无限阻塞。
实现方法
要在同一连接下重复消费未提交的消息,需要手动重置消费位置到目标消息的偏移量,强制客户端重新拉取该位置的消息。具体步骤如下:
第一次消费后,记录每个分区需要重置的偏移量
- 对于拉取到的每条消息,其
Offset就是该消息在分区中的位置,我们需要将对应分区的消费位置重置为这个Offset(因为Kafka的消费位置代表下一条要拉取的消息偏移量)
- 对于拉取到的每条消息,其
使用
SeekPartitions方法重置消费位置- 构造
kgo.PartitionSeek数组,指定要重置的topic、分区和目标偏移量 - 调用客户端的
SeekPartitions方法生效
- 构造
修改后的代码示例
package main import ( "context" "github.com/twmb/franz-go/pkg/kgo" ) func main() { cl, err := kgo.NewClient( kgo.SeedBrokers("localhost:9092"), kgo.ConsumerGroup("some_consumer_group"), kgo.ConsumeTopics("topic"), kgo.DefaultProduceTopic("topic"), kgo.DisableAutoCommit(), ) if err != nil { panic(err) } defer cl.Close() record := &kgo.Record{ Key: []byte("some_key"), Value: []byte("some value"), } ctx := context.Background() produceResult := cl.ProduceSync(ctx, record) if produceResult.FirstErr() != nil { panic(produceResult.FirstErr()) } // 第一次消费消息 fetches := cl.PollFetches(ctx) if fetches.Err() != nil { panic(fetches.Err()) } doSomething(fetches) // 记录需要重置的分区偏移量 var seeks []kgo.PartitionSeek fetches.EachPartition(func(p kgo.FetchTopicPartition) { p.EachRecord(func(r *kgo.Record) { // 将消费位置重置为当前消息的Offset,触发重新拉取 seeks = append(seeks, kgo.NewOffsetSeek(r.Topic, r.Partition, r.Offset)) }) }) // 重置消费位置 cl.SeekPartitions(ctx, seeks...) // 第二次消费,此时会拉取到同一条消息 fetches = cl.PollFetches(ctx) if fetches.Err() != nil { panic(fetches.Err()) } doSomething(fetches) // 提交最终偏移量 err = cl.CommitRecords(ctx, fetches.Records()...) if err != nil { panic(err) } } func doSomething(result kgo.Fetches) { //处理Kafka消息时出现异常 result.EachRecord(func(r *kgo.Record) { println("处理消息:", string(r.Value)) }) }
补充说明
- 如果只是处理逻辑异常需要重试,更高效的方式是直接保留第一次拉取的消息,重复调用处理函数,无需重新拉取消息,减少Kafka交互开销。
SeekPartitions方法会覆盖客户端的本地消费位置,后续的拉取请求会从指定的偏移量开始。
内容的提问来源于stack exchange,提问作者Kukwuka
相关产品推荐
相关产品推荐

