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

如何在同一连接中重新消费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会一直等待新消息,导致无限阻塞。

实现方法

要在同一连接下重复消费未提交的消息,需要手动重置消费位置到目标消息的偏移量,强制客户端重新拉取该位置的消息。具体步骤如下:

  1. 第一次消费后,记录每个分区需要重置的偏移量

    • 对于拉取到的每条消息,其Offset就是该消息在分区中的位置,我们需要将对应分区的消费位置重置为这个Offset(因为Kafka的消费位置代表下一条要拉取的消息偏移量)
  2. 使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 16:30:04