kafka提交偏移量失敗導(dǎo)致重復(fù)消費(fèi)的解決
問題詳情
org.springframework.kafka.KafkaException: Seek to current after exception; nested exception is org.apache.kafka.clients.consumer.CommitFailedException: Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.
at org.springframework.kafka.listener.SeekToCurrentBatchErrorHandler.handle(SeekToCurrentBatchErrorHandler.java:92)
at org.springframework.kafka.listener.RecoveringBatchErrorHandler.handle(RecoveringBatchErrorHandler.java:124)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.handleConsumerException(KafkaMessageListenerContainer.java:1365)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1063)
at java.util.concurrent.Executors$RunnableAdapter.call(Executors.java:511)
at java.util.concurrent.FutureTask.run(FutureTask.java:266)
at java.lang.Thread.run(Thread.java:748)
Caused by: org.apache.kafka.clients.consumer.CommitFailedException: Offset commit cannot be completed since the consumer is not part of an active group for auto partition assignment; it is likely that the consumer was kicked out of the group.
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.sendOffsetCommitRequest(ConsumerCoordinator.java:1116)
at org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.commitOffsetsSync(ConsumerCoordinator.java:983)
at org.apache.kafka.clients.consumer.KafkaConsumer.commitSync(KafkaConsumer.java:1510)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.doCommitSync(KafkaMessageListenerContainer.java:2311)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.commitSync(KafkaMessageListenerContainer.java:2306)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.commitIfNecessary(KafkaMessageListenerContainer.java:2292)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.processCommits(KafkaMessageListenerContainer.java:2106)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:1097)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:1031)
... 3 common frames omitted
解決思路
kafka的好多配置,在spring-kafka中沒有明確的配置對應(yīng),但是預(yù)留了一個(gè)properties屬性,可以設(shè)置所有的kafka配置
spring.kafka.properties.session.timeout.ms=10000 // 單位:毫秒 spring.kafka.properties.max.poll.interval.ms=300000 // 單位:毫秒
kafka會有一個(gè)心跳線程來同步服務(wù)端,告訴服務(wù)端自己是正??捎玫?,默認(rèn)是3秒發(fā)送一次心跳,超過session.timeout.ms(默認(rèn)10秒)服務(wù)端沒有收到心跳就會認(rèn)為當(dāng)前消費(fèi)者失效。max.poll.interval.ms決定了獲取消息后提交偏移量的最大時(shí)間,超過設(shè)定的時(shí)間(默認(rèn)5分鐘),服務(wù)端也會認(rèn)為該消費(fèi)者失效。
Kafka配置max.poll.interval.ms參數(shù)
max.poll.interval.ms默認(rèn)值是5分鐘,如果需要加大時(shí)長就需要給這個(gè)參數(shù)重新賦值
這里解釋下自己為什么要修改這個(gè)參數(shù):因?yàn)榈谝淮谓邮誯afka數(shù)據(jù),需要加載一堆基礎(chǔ)數(shù)據(jù),大概執(zhí)行時(shí)間要8分鐘,而5分鐘后,kafka認(rèn)為我沒消費(fèi),又重新發(fā)送,導(dǎo)致我這邊收到許多重復(fù)數(shù)據(jù),所以我需要調(diào)大這個(gè)值,避免接收重復(fù)數(shù)據(jù)
大部分文章都是如下配置:
public static KafkaConsumer<String, String> createConsumer() {
Properties properties = new Properties();
properties.put(CommonClientConfigs.BOOTSTRAP_SERVERS_CONFIG, BOOTSTRAP_SERVER);
properties.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
properties.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
properties.put(ConsumerConfig.GROUP_ID_CONFIG, "group1");
properties.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
properties.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, 10000);
return new KafkaConsumer<>(properties);
}
或是:
max.poll.interval.ms = 300000
如果需要在yml文件中配置,應(yīng)該怎么寫呢?
spring:
kafka:
consumer:
max-poll-records: 500
properties:
max.poll.interval.ms: 600000
總結(jié)
以上為個(gè)人經(jīng)驗(yàn),希望能給大家一個(gè)參考,也希望大家多多支持腳本之家。
相關(guān)文章
Java設(shè)計(jì)模式之建造者模式(Builder模式)介紹
這篇文章主要介紹了Java設(shè)計(jì)模式之建造者模式(Builder模式)介紹,本文講解了為何使用建造者模式、如何使用建造者模式、Builder模式的應(yīng)用等內(nèi)容,需要的朋友可以參考下2015-03-03
詳解用Spring Boot Admin來監(jiān)控我們的微服務(wù)
這篇文章主要介紹了用Spring Boot Admin來監(jiān)控我們的微服務(wù),文中通過示例代碼介紹的非常詳細(xì),對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)學(xué)習(xí)吧2020-08-08
Java模擬實(shí)現(xiàn)HashMap算法流程詳解
在java開發(fā)中,HashMap是最常用、最常見的集合容器類之一,文中通過示例代碼介紹HashMap為啥要二次Hash,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價(jià)值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧2023-02-02
微信小程序訂閱消息推送實(shí)戰(zhàn)圖文教程(Java?Spring?Boot?+?Redis)
訂閱消息是微信小程序提供的一種消息推送方式,用戶可以訂閱某個(gè)公眾號或小程序的消息,當(dāng)有新消息時(shí),系統(tǒng)會自動(dòng)推送通知給用戶,這篇文章主要介紹了微信小程序訂閱消息推送(Java Spring Boot+Redis)的相關(guān)資料,需要的朋友可以參考下2026-04-04
Spring Boot+Nginx+MySQL容器化實(shí)戰(zhàn)指南
這篇文章主要介紹了Spring Boot+Nginx+MySQL容器化實(shí)戰(zhàn)指南,本文通過思路代碼給大家介紹的非常詳細(xì),對大家的學(xué)習(xí)或工作具有一定的參考借鑒價(jià)值,需要的朋友參考下吧2026-03-03

