Spring Kafka consumer - IllegalStateException: Correlation id for response does not match request

Viewed 699

Encountering this issue on kafka consumers - request you to let me know if this is a broker or consumer issue? If a consumer issue is there a workaround for this with the below given versions?

Trying to debug on my end as well - will update with progress soon.

Update 1 - To add to the complexity even a consumer restart doesn;t solve this and it makes me think the underlying problem is with the broker due to this.

spring-kafka consumer using @KafkaListener

spring-kafka version - 2.2.14.RELEASE kafka client version - kafka-client - 2.0.1

kafka cluster running on 1.1.1

Update 2: I also see the exception - Node 597397927 was unable to process the fetch request with (sessionId=1939114665, epoch=2): INVALID_FETCH_SESSION_EPOCH

Consumer exception - cause: {} - java.lang.IllegalStateException: Correlation id for response (400801) does not match request (400737), request header: RequestHeader(apiKey=OFFSET_COMMIT, apiVersion=3, clientId=consumer-7, correlationId=400737)
at org.apache.kafka.clients.NetworkClient.correlate(NetworkClient.java:853)
at org.apache.kafka.clients.NetworkClient.parseStructMaybeUpdateThrottleTimeMetrics(NetworkClient.java:638)
at org.apache.kafka.clients.NetworkClient.handleCompletedReceives(NetworkClient.java:757)
at org.apache.kafka.clients.NetworkClient.poll(NetworkClient.java:519)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:271)
at org.apache.kafka.clients.consumer.internals.ConsumerNetworkClient.poll(ConsumerNetworkClient.java:242)
at org.apache.kafka.clients.consumer.KafkaConsumer.pollForFetches(KafkaConsumer.java:1247)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1187)
at org.apache.kafka.clients.consumer.KafkaConsumer.poll(KafkaConsumer.java:1154)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.pollAndInvoke(KafkaMessageListenerContainer.java:742)
at org.springframework.kafka.listener.KafkaMessageListenerContainer$ListenerConsumer.run(KafkaMessageListenerContainer.java:699)
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)
0 Answers
Related