【问题标题】:how to get last committed offset from read_committed Kafka Consumer如何从 read_committed Kafka Consumer 获取最后提交的偏移量
【发布时间】:2018-07-01 22:21:11
【问题描述】:

我正在使用事务性 KafkaProducer 向主题发送消息。这工作正常。我使用具有 read_committed 隔离级别的 KafkaConsumer,但我对 seek 和 seekToEnd 方法有疑问。根据文档, seek 和 seekToEnd 方法给了我 LSO(最后稳定偏移)。但这有点令人困惑。因为它总是给我相同的价值,即主题的结束。无论最后一个条目是提交(由生产者)还是中止事务的一部分。 例如,在我中止最后 5 次尝试插入 20_000 条消息后,消费者不应读取最后 100_000 条记录。但是在 seekToEnd 期间,它会移动到主题的末尾(包括 100_000 条消息)。但是 poll() 不会返回它们。

我正在寻找一种方法来检索上次提交的偏移量(即生产者最后一次成功提交的消息)。似乎没有合适的 API 方法。那我需要自己动手吗?

选项是返回并轮询直到没有更多记录被检索到,这将导致最后提交的消息。但我会假设 Kafka 提供了这种方法。

我们使用 Kafka 1.0.0。

【问题讨论】:

  • 你能提供你的完整配置吗?另外,你可以试试seek-3 吗? -3 是代表last stable offset 的标记值。

标签: apache-kafka kafka-consumer-api kafka-producer-api


【解决方案1】:

要获取主题分区的最后提交偏移量,您可以使用KafkaConsumer.committed(TopicPartition partition) 函数。

TopicPartition topicPartition = new TopicPartition(record.topic(), record.partition());
Long committedOffset = consumer.committed(topicPartition).offset();
System.out.println("last committed offset: " + committedOffset);  

【讨论】:

    【解决方案2】:

    KafkaConsumer 类有一些不错的方法,例如:partitionForbegginingOffsetsendOffsets 以及 commitedposition

    检查哪一个适合您的需求。尤其要仔细考虑所有 4 种与偏移相关的方法。 partitionFor 方法返回完整的元数据对象和其他信息,但对于丰富日志记录很有用。

    【讨论】:

    • 我提出了一个 JIRA @ Kafka。 Seek 和 SeekToEnd、endOffsets 等移动到主题上的最后一条消息,无论该消息是否属于已提交事务的一部分......
    • 其实我们需要类似 seekToLastCommittedMessage
    • 您能否具体说明您是如何获得 Producer 上一条成功提交的消息的?据我了解, seekToEnd() 和 endOffsets() 将带您到主题的结尾。就我而言,我在一个主题中有 100 条消息。 Offset#96 是生产者最后提交的消息。 Offset#96-100 是未提交的消息。如何检索最后提交的消息?
    • @CoenDamen 你能链接kafka JIRA吗?谢谢!
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2020-12-28
    • 1970-01-01
    • 1970-01-01
    • 2021-01-17
    • 2020-09-06
    • 1970-01-01
    • 2021-11-15
    相关资源
    最近更新 更多