未提交偏移量的Kafka消息无法自动重消费,仅重启应用生效
Hey there! Let's break down why your unacknowledged Kafka messages aren't reprocessing automatically and fix it step by step.
The Root Cause
When using manual acknowledgment mode, Spring Kafka's listener container won't automatically retry failed messages out of the box. Without additional configuration, once an exception is caught and you don't explicitly tell Kafka to retry, the container will just move on to the next message. The reason messages reappear after a restart is that the consumer reinitializes and starts from the last committed offset (which never updated for the failed message).
Fix 1: Configure a SeekToCurrentErrorHandler
This is the key component you're missing. The SeekToCurrentErrorHandler tells the consumer to "seek back" to the failed message's offset when an exception occurs, so it gets reconsumed immediately without needing a restart.
Add this to your listener container factory configuration:
@Bean public ConcurrentKafkaListenerContainerFactory<String, Object> tempListenerContainerFactory(ConsumerFactory<String, Object> consumerFactory) { ConcurrentKafkaListenerContainerFactory<String, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.setConsumerFactory(consumerFactory); // Enable manual immediate acknowledgment (more responsive than plain MANUAL) factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL_IMMEDIATE); // Add the error handler to trigger automatic retries factory.setErrorHandler(new SeekToCurrentErrorHandler()); // Optional: Limit retries + route to dead letter queue after failures // factory.setErrorHandler(new SeekToCurrentErrorHandler( // new DeadLetterPublishingRecoverer(kafkaTemplate), // new FixedBackOff(1000L, 3L) // Retry 3 times with 1s delay // )); return factory; }
Fix 2: Adjust Your Consumer Code
You have two options to trigger the retry behavior:
- Let the exception propagate: Remove the catch block (or rethrow the exception) so the
SeekToCurrentErrorHandlercan intercept it. - Explicitly call
nack(): If you want to handle the exception locally first, useacknowledgment.nack()to tell Kafka to retry the message after a specified delay.
Here's the adjusted code with option 2:
@KafkaListener(topics = "temp_topic", groupId = "temp-consumer", containerFactory = "tempListenerContainerFactory") public void receive(ConsumerRecord<?, ?> consumerRecord, Acknowledgment acknowledgment) { try { log.info("Consuming message : {} ", consumerRecord.value().toString()); String message = consumerRecord.value().toString(); objectMapper.configure(DeserializationFeature.FAIL_ON_UNKNOWN_PROPERTIES, false); Payload payload = objectMapper.readValue(message,Payload.class); tempService.processEvent(payload); log.info("Message: {} successfully consumed",message); acknowledgment.acknowledge(); } catch (Exception e) { log.error("Error while listening to kafka event: ", e); // Tell Kafka to retry this message after 1 second acknowledgment.nack(1000L); } }
Key Notes
MANUAL_IMMEDIATEvsMANUAL:MANUAL_IMMEDIATEcommits the offset immediately when you callacknowledge(), whileMANUALbatches offset commits. For manual retry scenarios,MANUAL_IMMEDIATEis more reliable.- Dead Letter Queue (DLQ): If you don't want infinite retries, configure a DLQ using
DeadLetterPublishingRecovereras shown in the commented code. This routes messages that fail after X retries to a separate topic for later analysis.
内容的提问来源于stack exchange,提问作者Napstablook

