KafkaConsumer调用seekToEnd()无法从最新偏移量消费问题求解
seekToEnd()未生效的根因
你当前调用seekToEnd(emptyList())没有生效的核心原因是调用时机不对:
Kafka消费者调用subscribe()仅注册了主题订阅关系,并不会立刻触发分区分配,要等到第一次poll()触发重平衡完成分区分配后,seekToEnd()才会对已分配的分区生效。你在run()方法启动时直接调用seekToEnd(emptyList()),此时消费者还没有拿到任何分配的分区,相当于空操作,不会产生任何效果。
另外Kafka消费者的偏移量读取优先级为:已提交的偏移量 > 手动seek设置的偏移量 > auto.offset.reset配置。你崩溃前已经提交过偏移量,所以即使你配置了auto.offset.reset=latest,重启后会优先读取之前提交的旧偏移量,不会触发latest逻辑。
seekToEnd()的正确使用方法
你可以通过注册消费者重平衡监听器,在分区分配完成后执行seekToEnd,修改你createConsumer方法的订阅逻辑即可:
import org.apache.kafka.clients.consumer.ConsumerRebalanceListener import org.apache.kafka.common.TopicPartition fun <T> createConsumer(props: Properties): KafkaConsumer<String, T> { val consumer = KafkaConsumer<String, T>(props) consumer.subscribe(listOf(config.kafka.inputTopic), object : ConsumerRebalanceListener { override fun onPartitionsAssigned(partitions: MutableCollection<TopicPartition>) { // 分区分配完成后对所有新分配的分区执行seekToEnd,跳过所有历史消息 consumer.seekToEnd(partitions) } override fun onPartitionsRevoked(partitions: MutableCollection<TopicPartition>) { // 原有偏移量提交逻辑不变 } }) return consumer }
修改后你不需要再在run()方法里调用seekToEnd,每次消费者重启完成分区分配后都会自动执行seekToEnd(),直接跳过崩溃到重启期间的所有历史消息,从最新偏移量开始消费。
内容的提问来源于stack exchange,提问作者hermanjakobsen
相关产品推荐
相关产品推荐

