springboot中怎么实现kafa指定offset消费-创新互联
小编给大家分享一下springboot中怎么实现kafa指定offset消费,相信大部分人都还不怎么了解,因此分享这篇文章给大家参考一下,希望大家阅读完这篇文章后大有收获,下面让我们一起去了解一下吧!
kafka消费过程难免会遇到需要重新消费的场景,例如我们消费到kafka数据之后需要进行存库操作,若某一时刻数据库down了,导致kafka消费的数据无法入库,为了弥补数据库down期间的数据损失,有一种做法我们可以指定kafka消费者的offset到之前某一时间的数值,然后重新进行消费。
首先创建kafka消费服务
@Service@Slf4j//实现CommandLineRunner接口,在springboot启动时自动运行其run方法。public class TspLogbookAnalysisService implements CommandLineRunner { @Override public void run(String... args) { //do something }}
kafka消费模型建立
kafka server中每个主题存在多个分区(partition),每个分区自己维护一个偏移量(offset),我们的目标是实现kafka consumer指定offset消费。
在这里使用consumer-->partition一对一的消费模型,每个consumer各自管理自己的partition。
@Service@Slf4jpublic class TspLogbookAnalysisService implements CommandLineRunner { //声明kafka分区数相等的消费线程数,一个分区对应一个消费线程 private static final int consumeThreadNum = 9; //特殊指定每个分区开始消费的offset private List
kafka consumer对offset的处理
声明kafka consumer的配置类
private Properties buildKafkaConfig() { Properties kafkaConfiguration = new Properties(); kafkaConfiguration.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, ""); kafkaConfiguration.put(ConsumerConfig.GROUP_ID_CONFIG, ""); kafkaConfiguration.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, ""); kafkaConfiguration.put(ConsumerConfig.AUTO_COMMIT_INTERVAL_MS_CONFIG, ""); kafkaConfiguration.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ""); kafkaConfiguration.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ""); kafkaConfiguration.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG,""); kafkaConfiguration.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, ""); ...更多配置项 return kafkaConfiguration;}
创建kafka consumer,处理offset,开始消费数据任务#
private void startConsume(int partitionIndex) { //创建kafka consumer KafkaConsumer
消费数据逻辑,offset操作
private void kafkaRecordConsume(KafkaConsumer
以上是“springboot中怎么实现kafa指定offset消费”这篇文章的所有内容,感谢各位的阅读!相信大家都有了一定的了解,希望分享的内容对大家有所帮助,如果还想学习更多知识,欢迎关注创新互联行业资讯频道!
分享文章:springboot中怎么实现kafa指定offset消费-创新互联
浏览地址:http://ybzwz.com/article/dpodpo.html