【发布时间】: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