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

如何使用Micronaut监听器实现Kafka消费失败消息的重试

Micronaut Kafka 消费重试配置方案

Micronaut Kafka默认的消费异常处理逻辑是记录错误后直接提交offset,你可以通过以下两种方式实现自定义次数的重试:

方案1:全局/消费者组维度配置固定重试次数

直接在application.yml中添加重试配置即可对指定消费者生效:

kafka:
  consumers:
    default: # 填写你要配置的消费者group-id,填default表示全局所有消费者生效
      retry:
        attempts: 3 # 重试总次数,包含首次消费
        delay: 1000 # 每次重试的间隔时间,单位为毫秒
      enable-auto-commit: false # 建议关闭自动提交,由框架控制offset提交时机
      ack-mode: RECORD # 单条消息消费成功后再提交offset

如果需要仅针对特定异常触发重试,可以自定义全局错误处理器:

import io.micronaut.configuration.kafka.exceptions.KafkaListenerException;
import io.micronaut.configuration.kafka.exceptions.KafkaListenerExceptionHandler;
import io.micronaut.context.annotation.Replaces;
import jakarta.inject.Singleton;
import org.apache.kafka.common.errors.RetriableException;

@Singleton
@Replaces(KafkaListenerExceptionHandler.class)
public class CustomKafkaErrorHandler implements KafkaListenerExceptionHandler {

    @Override
    public void handle(KafkaListenerException exception) {
        // 非重试异常直接跳过,提交offset
        if (!(exception.getCause() instanceof RetriableException)) {
            exception.getConsumer().seek(
                exception.getConsumerRecord().topic(),
                exception.getConsumerRecord().partition(),
                exception.getConsumerRecord().offset() + 1
            );
            return;
        }
        // 重试异常抛出即可触发框架内置重试逻辑
        throw exception;
    }
}

方案2:为单个监听方法配置独立重试规则

如果不同的topic/消费逻辑需要不同的重试次数,可以直接在监听方法上加@Retryable注解,优先级高于全局配置:

import io.micronaut.retry.annotation.Retryable;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import io.micronaut.configuration.kafka.annotation.KafkaListener;
import io.micronaut.configuration.kafka.annotation.Topic;

@KafkaListener(groupId = "your-group-id")
public class CustomKafkaListener {

    @Topic("your-business-topic")
    // 配置重试5次,间隔2秒,仅针对RetriableException子类触发重试
    @Retryable(attempts = "5", delay = "2s", includes = {RetriableException.class})
    public void consume(ConsumerRecord<String, String> record) {
        // 你的业务消费逻辑
        System.out.printf("收到消息:topic=%s, value=%s%n", record.topic(), record.value());
    }
}

注意事项

  • 重试期间不会提交offset,直到重试次数耗尽且消费成功后才会提交
  • 如果重试次数耗尽后仍然抛出异常,默认仍会提交offset,如需转存死信队列可添加如下配置:
kafka:
  consumers:
    default:
      dlq:
        enabled: true
        topic: your-dead-letter-topic # 死信队列topic名

内容的提问来源于stack exchange,提问作者Solci

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 15:36:03