【问题标题】:Reliably get the last (already produced) message from Kafka topic可靠地获取来自 Kafka 主题的最后一条(已经产生的)消息
【发布时间】:2021-09-01 01:02:10
【问题描述】:

我正在做类似下面的伪代码

var consumer = new KafkaConsumer();
consumer.assign(topicPartitions);
var beginOff = consumer.beginningOffsets(topicPartitions);
var endOff = consumer.endOffsets(topicPartitions);
var lastOffsets = Math.max(beginOff, endOff - 1));
lastOffsets.forEach(consumer::seek);
lastMessages = consumer.poll(1 sec);
// do something with the received messages
consumer.close();

在我所做的简单测试中,这是可行的,但我想知道是否存在偏移量不是单调增加 1 的情况,例如生产者崩溃等?在这种情况下,我是否必须及时返回 seek(),或者我可以从 Kafka 获取最后一条已生成消息的消息偏移量?

我没有使用事务,所以我们不需要担心已提交的和未提交的消息。

编辑: 一个偏移不连续的例子是在日志压缩之后。但是,日志压缩应始终保留最后一条消息,因为它 - 显然 - 比所有先前的消息(相同或不同的键)更新。但是理论上可以压缩最后一条消息之前的偏移量。

【问题讨论】:

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


【解决方案1】:

kafka.apache.org/10/javadoc/中明确提到,consumer.endOffsets

Get the last offset for the given partitions. The last offset of a partition is the offset of the upcoming message, i.e. the offset of the last available message + 1.

因此,当您获得 endOff - 1 时,它是您获取该主题分区时最后一个可用的 Kafka 记录。因此,生产者的担忧不会因此受到影响。

还有一件事,Offset 不是由制片人决定的。由该主题分区的分区领导者决定。所以,它总是单调加一。

【讨论】:

  • 自 Kafka 1.0 以来发生了很多事情,但 2.8 的文档状态类似:In the default read_uncommitted isolation level, the end offset is the high watermark (that is, the offset of the last successfully replicated message plus one)。这听起来很确定,我希望这是真的。但我记得必须处理非连续的偏移量,但现在我不是 100%,如果这可能不是由于 read_committed... 将尝试检查。
  • @EvgeniyBerezovsky 我可以根据您的评论更新答案吗?
  • 您的回答很有说服力,但似乎是错误的——或者某处存在错误。我在这里记录了这样一个案例:stackoverflow.com/q/69186310
猜你喜欢
  • 2018-04-20
  • 2021-09-14
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2021-10-26
  • 2021-01-18
相关资源
最近更新 更多