You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

非@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提交。

  1. 保留原有KafkaConsumerConfig中的errorHandler()方法(可改为@Bean或直接实例化)
  2. 修改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注解实现重试,灵活控制重试规则。

  1. 引入Spring Retry依赖(Maven示例):
<dependency>
    <groupId>org.springframework.retry</groupId>
    <artifactId>spring-retry</artifactId>
</dependency>
<dependency>
    <groupId>org.springframework</groupId>
    <artifactId>spring-aspects</artifactId>
</dependency>
  1. 配置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;
    }
}
  1. 在手动消费中使用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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.26 00:42:37