【发布时间】:2017-12-17 21:54:18
【问题描述】:
当消费者组 A 的 Kafka 消费者连接到 Kafka 代理时,我想寻找所有分区的末尾,即使在代理端存储了偏移量。如果更多额外的消费者正在为同一个消费者组连接,他们应该获取最新存储的偏移量。 我正在执行以下操作:
consumer.poll(timeout)
consumer.seekToEnd(emptyList())
while(true) {
val records = consumer.poll(timeout)
if(records.isNotEmpty()) {
//print records
consumer.commitSync()
}
}
问题是当我连接消费者组 A 的第一个消费者 c1 时,一切都按预期工作,如果我连接消费者组 A 的另一个消费者 c2,则该组正在重新平衡,c1 将消耗跳过的偏移量。
有什么想法吗?
【问题讨论】:
标签: java apache-kafka kotlin