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

Spring-Kafka:如何为ConcurrentMessageListenerContainer配置RetryTopicConfiguration?

手动配置Kafka消费者关联重试主题与死信队列

步骤1:创建RetryTopicConfiguration Bean

定义重试规则的配置Bean,可适配泛型消息类型,按需调整重试策略:

import org.springframework.kafka.config.RetryTopicConfiguration;
import org.springframework.kafka.config.RetryTopicConfigurationBuilder;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.context.annotation.Bean;

@Bean
public <V> RetryTopicConfiguration genericRetryTopicConfig(KafkaProperties props, KafkaTemplate<String, V> kafkaTemplate) {
    return RetryTopicConfigurationBuilder.newInstance()
            .fixedBackOff(2000) // 固定退避间隔2秒
            .maxAttempts(3) // 最大重试次数(包含首次消费)
            .useSingleTopicForFixedDelays() // 固定延迟场景复用单个重试主题
            .deadLetterTopicSuffix("-dlt") // 死信队列后缀
            .retryTopicSuffix("-retry") // 重试主题后缀
            .create(kafkaTemplate);
}

步骤2:修改MyCustomListener,注入并应用重试配置

将重试配置注入Listener,替换原容器创建逻辑,让重试规则生效:

import org.springframework.kafka.config.RetryTopicConfiguration;
import javax.validation.constraints.NotNull;
import java.util.HashMap;
import org.apache.kafka.common.serialization.StringDeserializer;
import org.springframework.kafka.listener.ConcurrentMessageListenerContainer;
import org.springframework.kafka.listener.ContainerProperties;
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;

public class MyCustomListener<V> implements MessageListener<String, V> {
    private final String topicName;
    private final Class<V> recordType;
    private final String recordTypeName;
    private final KafkaProperties props;
    private final Logger logger;
    private final ConcurrentMessageListenerContainer<String, V> listenerContainer;
    private final RetryTopicConfiguration retryTopicConfiguration;

    // 构造方法新增重试配置注入
    public MyCustomListener(
            @NotNull String topicName,
            @NotNull Class<V> recordType,
            @NotNull KafkaProperties props,
            Logger logger,
            RetryTopicConfiguration retryTopicConfiguration
    ) {
        this.topicName = topicName;
        this.recordType = recordType;
        this.recordTypeName = recordType.getSimpleName();
        this.props = props;
        this.logger = logger;
        this.retryTopicConfiguration = retryTopicConfiguration;
        this.listenerContainer = setup();
    }

    protected ConcurrentMessageListenerContainer<String, V> setup() {
        logger.info("Setting up Kafka listener container for topic {}/{}", topicName, recordTypeName);

        var consumerFactory = new DefaultKafkaConsumerFactory<>(
                new HashMap<>() {{
                    put(BOOTSTRAP_SERVERS_CONFIG, props.getBootstrapServers());
                    put(GROUP_ID_CONFIG, props.getConsumer().getGroupId());
                }},
                new StringDeserializer(),
                new MyDataDeserializer<>()
        );

        var containerProps = new ContainerProperties(topicName);
        containerProps.setMessageListener(this);

        // 用重试配置初始化并配置容器,替代直接实例化
        return retryTopicConfiguration.configureContainer(
                new ConcurrentMessageListenerContainer<>(consumerFactory, containerProps),
                consumerFactory
        );
    }

    @Override
    public final void onMessage(ConsumerRecord<String, V> record) {
        logger.info(
                "(#{}) Received {} message in Kafka: ({})",
                Thread.currentThread().getId(),
                recordTypeName,
                record.key()
        );

        try {
            // 业务处理逻辑抽离,便于异常捕获触发重试
            processRecord(record);
        } catch (Exception e) {
            logger.error("Processing failed for record key: {}", record.key(), e);
            // 抛出异常触发重试机制,达到重试次数后消息会进入死信队列
            throw new RuntimeException("Record processing failed", e);
        }
    }

    // 业务处理方法,子类可重写
    protected void processRecord(ConsumerRecord<String, V> record) {
        // 你的业务逻辑实现
    }
}

关键注意事项

  • 重试触发条件:必须在业务处理中抛出异常,重试机制才会将消息转发到重试主题;达到最大重试次数后,消息自动进入死信队列。
  • 主题命名:默认重试主题为原主题名-retry,死信队列为原主题名-dlt,可通过RetryTopicConfigurationBuilder的方法自定义后缀。
  • 生命周期管理:建议将MyCustomListener注册为Spring Bean,交由Spring管理容器的启动/停止,避免手动调用container.start()引发的问题。
  • 多规则适配:若不同类型消息需不同重试策略,可创建多个RetryTopicConfiguration Bean,注入时通过@Qualifier指定对应Bean。

内容的提问来源于stack exchange,提问作者Higher-Kinded Type

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.23 08:45:28