非@KafkaListener方式下Spring Kafka重试机制实现求助
手动轮询Kafka消费场景下的重试机制实现方案
问题背景
已通过@KafkaListener结合Spring Kafka的DefaultErrorHandler实现了针对指定异常的消费重试逻辑,但在自行编写的轮询消费代码中,无法为dataProcessor.readRecord(r, consumer)方法实现相同的重试机制。
现有@KafkaListener重试实现代码
@Component public class Consumer implements AcknowledgingMessageListener<String, String> { @Override @KafkaListener(topics = { "ff808081672c17c8016730733d020001.CBS_PROFILE"}) public void onMessage(ConsumerRecord<String, String> consumerRecord, Acknowledgment acknowledgment) { log.info("Computed "); throw new RuntimeException(); } } @Configuration @Slf4j public class KafkaConsumerConfig { public DefaultErrorHandler errorHandler() { System.out.println("Inside error handler"); List<Class<? extends Exception>> exceptionsToIgnore = List.of(JsonSchemaValidationException.class); List<Class<? extends Exception>> exceptionsToRetry = List.of(DataEngineNotAvailableException.class, ListenerExecutionFailedException.class); var fixedBackOff = new FixedBackOff(500L, 1); var errorHandler = new DefaultErrorHandler(fixedBackOff); errorHandler.setRetryListeners((consumerRecord, ex, deliveryAttempt) -> { log.info("Failed record {} in retry listeners, Exception : {}, delivery attempt : {}", consumerRecord, ex.getMessage(), deliveryAttempt); }); exceptionsToIgnore.forEach(errorHandler::addNotRetryableExceptions); exceptionsToRetry.forEach(errorHandler::addRetryableExceptions); return errorHandler; } @Bean public ConcurrentKafkaListenerContainerFactory<?, ?> kafkaListenerContainerFactory( ConcurrentKafkaListenerContainerFactoryConfigurer configurer, ConsumerFactory<Object, Object> kafkaConsumerFactory) { ConcurrentKafkaListenerContainerFactory<Object, Object> factory = new ConcurrentKafkaListenerContainerFactory<>(); configurer.configure(factory, kafkaConsumerFactory); factory.getContainerProperties().setAckMode(AckMode.MANUAL); factory.setConcurrency(3); factory.setCommonErrorHandler(errorHandler()); return factory; } }
手动轮询消费代码(待添加重试)
@PostConstruct public void onLoad() { configService.setConfigurations(); consumerRunnerV2.loadConsumers(); } public void loadConsumers() { consumerServiceV2.fillConsumerMap(); Map<String, LRConsumer> consumersMap = consumerServiceV2.getConsumersMap(); for (Map.Entry<String, LRConsumer> consumer : consumersMap.entrySet()) { consumerTaskExecutor.execute(() -> runConsumer(consumer.getValue())); } } private void runConsumer(LRConsumer consumer) { while (true) { consumer = consumerServiceV2.getConsumerById(consumer.getConsumerId()); if (consumer.getState() != ConsumerState.PAUSED) { log.info("Waiting to receive messages."); ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> r : records) { // 需要为该方法添加重试机制 dataProcessor.readRecord(r, consumer); } } else { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } log.info("Consumer is in paused state."); } } }
解决方案
方案1:复用Spring Kafka的DefaultErrorHandler
DefaultErrorHandler并非只能配合@KafkaListener使用,可直接在手动消费逻辑中调用它处理异常和重试,同时需手动管理offset提交。
- 保留原有
KafkaConsumerConfig中的errorHandler()方法(可改为@Bean或直接实例化) - 修改
runConsumer方法的消息处理逻辑:
// 注入或实例化DefaultErrorHandler private final DefaultErrorHandler errorHandler; // 构造方法注入 public ConsumerRunnerV2(DefaultErrorHandler errorHandler, DataProcessor dataProcessor) { this.errorHandler = errorHandler; this.dataProcessor = dataProcessor; } private void runConsumer(LRConsumer consumer) { while (true) { consumer = consumerServiceV2.getConsumerById(consumer.getConsumerId()); if (consumer.getState() != ConsumerState.PAUSED) { log.info("Waiting to receive messages."); ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> r : records) { try { dataProcessor.readRecord(r, consumer); // 处理成功后手动提交offset consumer.commitSync(Map.of(r.topicPartition(), new OffsetAndMetadata(r.offset() + 1))); } catch (Exception ex) { // 交给DefaultErrorHandler处理重试 boolean shouldRetry = errorHandler.handleOne(ex, r, consumer, null); if (!shouldRetry) { // 重试失败后,提交offset或发送到死信队列 consumer.commitSync(Map.of(r.topicPartition(), new OffsetAndMetadata(r.offset() + 1))); sendToDeadLetterQueue(r, ex); } } } } else { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } log.info("Consumer is in paused state."); } } }
方案2:使用Spring Retry框架
通过Spring Retry的RetryTemplate或@Retryable注解实现重试,灵活控制重试规则。
- 引入Spring Retry依赖(Maven示例):
<dependency> <groupId>org.springframework.retry</groupId> <artifactId>spring-retry</artifactId> </dependency> <dependency> <groupId>org.springframework</groupId> <artifactId>spring-aspects</artifactId> </dependency>
- 配置RetryTemplate:
@Configuration @EnableRetry public class RetryConfig { @Bean public RetryTemplate retryTemplate() { RetryTemplate retryTemplate = new RetryTemplate(); // 重试策略:针对指定异常重试,最多重试2次 SimpleRetryPolicy retryPolicy = new SimpleRetryPolicy(); Map<Class<? extends Throwable>, Boolean> retryableExceptions = new HashMap<>(); retryableExceptions.put(DataEngineNotAvailableException.class, true); retryableExceptions.put(ListenerExecutionFailedException.class, true); retryableExceptions.put(JsonSchemaValidationException.class, false); // 不重试该异常 retryPolicy.setRetryableExceptions(retryableExceptions); retryPolicy.setMaxAttempts(3); // 包括首次调用,共3次 // 退避策略:固定间隔1秒 FixedBackOffPolicy backOffPolicy = new FixedBackOffPolicy(); backOffPolicy.setBackOffPeriod(1000L); retryTemplate.setRetryPolicy(retryPolicy); retryTemplate.setBackOffPolicy(backOffPolicy); // 重试监听 retryTemplate.registerListener(new RetryListener() { @Override public <T, E extends Throwable> void onError(RetryContext context, RetryCallback<T, E> callback, Throwable throwable) { ConsumerRecord<String, String> record = (ConsumerRecord<String, String>) context.getAttribute("record"); log.info("Failed record {} in retry, Exception : {}, delivery attempt : {}", record, throwable.getMessage(), context.getRetryCount() + 1); } }); return retryTemplate; } }
- 在手动消费中使用RetryTemplate:
private final RetryTemplate retryTemplate; // 构造方法注入 public ConsumerRunnerV2(RetryTemplate retryTemplate, DataProcessor dataProcessor) { this.retryTemplate = retryTemplate; this.dataProcessor = dataProcessor; } private void runConsumer(LRConsumer consumer) { while (true) { consumer = consumerServiceV2.getConsumerById(consumer.getConsumerId()); if (consumer.getState() != ConsumerState.PAUSED) { log.info("Waiting to receive messages."); ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> r : records) { try { retryTemplate.execute(context -> { context.setAttribute("record", r); dataProcessor.readRecord(r, consumer); return null; }); consumer.commitSync(Map.of(r.topicPartition(), new OffsetAndMetadata(r.offset() + 1))); } catch (Exception ex) { // 重试失败后的处理 consumer.commitSync(Map.of(r.topicPartition(), new OffsetAndMetadata(r.offset() + 1))); sendToDeadLetterQueue(r, ex); } } } else { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } log.info("Consumer is in paused state."); } } }
方案3:手动实现重试逻辑
若不想引入额外依赖,可自行编写重试逻辑,控制重试次数、异常类型和退避时间。
private void runConsumer(LRConsumer consumer) { while (true) { consumer = consumerServiceV2.getConsumerById(consumer.getConsumerId()); if (consumer.getState() != ConsumerState.PAUSED) { log.info("Waiting to receive messages."); ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> r : records) { int retryCount = 0; boolean success = false; final int MAX_RETRIES = 2; final long BACKOFF_TIME = 1000L; while (retryCount <= MAX_RETRIES) { try { dataProcessor.readRecord(r, consumer); success = true; break; } catch (JsonSchemaValidationException e) { // 忽略该异常,直接终止重试 log.info("Ignoring exception for record {}, message: {}", r, e.getMessage()); break; } catch (DataEngineNotAvailableException | ListenerExecutionFailedException e) { retryCount++; log.info("Retry attempt {} for record {}, exception: {}", retryCount, r, e.getMessage()); if (retryCount > MAX_RETRIES) { log.error("Max retries reached for record {}", r); break; } // 退避等待 try { Thread.sleep(BACKOFF_TIME); } catch (InterruptedException ie) { Thread.currentThread().interrupt(); break; } } catch (Exception e) { // 其他异常不重试 log.error("Unexpected exception for record {}", r, e); break; } } // 无论成功还是重试失败,都提交offset(根据业务需求调整) if (success) { consumer.commitSync(Map.of(r.topicPartition(), new OffsetAndMetadata(r.offset() + 1))); } else { // 可选:发送到死信队列 sendToDeadLetterQueue(r, new RuntimeException("Retry failed")); consumer.commitSync(Map.of(r.topicPartition(), new OffsetAndMetadata(r.offset() + 1))); } } } else { try { Thread.sleep(1000); } catch (InterruptedException e) { Thread.currentThread().interrupt(); break; } log.info("Consumer is in paused state."); } } }
注意事项
- Offset管理:手动消费时必须明确控制offset提交时机,避免重复消费或丢失消息。
- 死信队列:重试失败的消息建议发送到死信队列,便于后续排查和处理。
- 线程安全:确保consumer实例在多线程环境下的安全访问,避免并发问题。
内容的提问来源于stack exchange,提问作者Muddassir Rahman
相关产品推荐
相关产品推荐

