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

使用@RetryableTopic时主重试主题同名报错,求非阻塞重试解决方案

解决方案:实现同一主题的非阻塞重试与DLT流转

问题根源在于Spring Kafka的DefaultDestinationTopicResolver默认以主题名作为内部Map的key,当主主题与重试主题同名时,会尝试将两个不同的DestinationTopicHolder实例放入同一个key下,触发Duplicate key冲突。要实现「同一主题循环重试+指定DLT」的需求,需要通过自定义组件调整重试主题的命名逻辑与内部解析规则。


1. 修正自定义RetryTopicNamesProviderFactory

确保主端点返回原主题,重试阶段复用同一主题,DLT正确替换后缀:

@Override
@NonNull
public RetryTopicNamesProviderFactory retryTopicNamesProviderFactory() {
    return properties -> new SuffixingRetryTopicNamesProviderFactory.SuffixingRetryTopicNamesProvider(properties) {
        private static final String DEFAULT_RETRY_SUFFIX = "-retry";
        private static final String DEFAULT_DLT_SUFFIX = "-dlt";

        @Override
        @NonNull
        public String getTopicName(@NonNull String topic) {
            // 主监听端点返回原重试主题
            if (properties.isMainEndpoint()) {
                return topic;
            }
            // 死信主题替换后缀为-dlt
            if (properties.isDltTopic()) {
                return topic.replace(DEFAULT_RETRY_SUFFIX, DEFAULT_DLT_SUFFIX);
            }
            // 重试阶段直接复用当前主题,不生成新主题
            return topic;
        }

        @NonNull
        @Override
        public String getGroupId(@NonNull KafkaListenerEndpoint endpoint) {
            return "my-test-group";
        }
    };
}

2. 自定义DestinationTopicResolver解决重复key冲突

重写解析逻辑,为重试端点生成唯一key,避免与主端点的key重复:

@Component
public class SameTopicRetryDestinationResolver extends DefaultDestinationTopicResolver {

    @Override
    protected Map<String, DestinationTopicHolder> resolveDestinationTopics(RetryTopicConfiguration configuration) {
        Map<String, DestinationTopicHolder> originalMap = super.resolveDestinationTopics(configuration);
        Map<String, DestinationTopicHolder> correctedMap = new HashMap<>();

        originalMap.forEach((key, holder) -> {
            // 主端点保留原key,重试端点添加标识生成唯一key
            String targetKey = holder.getEndpoint().isMainEndpoint() 
                ? key 
                : key + "-retry-" + holder.getTopicSuffix();
            correctedMap.put(targetKey, holder);
        });

        return correctedMap;
    }
}

3. 关联自定义组件到重试配置

在重试配置中指定自定义的解析器与命名工厂:

@Configuration
public class KafkaRetryConfig {

    @Bean
    public RetryTopicConfiguration retryTopicConfiguration(KafkaTemplate<String, Object> kafkaTemplate,
                                                           SameTopicRetryDestinationResolver resolver) {
        return RetryTopicConfigurationBuilder
                .newInstance()
                .destinationTopicResolver(resolver)
                .retryTopicNamesProviderFactory(retryTopicNamesProviderFactory())
                .maxAttempts(100)
                .fixedBackOff(2000)
                .singleTopicStrategy()
                .create(kafkaTemplate);
    }

    // 放入第一步的retryTopicNamesProviderFactory方法
}

4. 确保重试元数据正常传递

Spring Kafka通过RetryTopicHeaders.RETRY_ATTEMPTS等头信息维护重试计数,需保证消费者容器工厂配置了支持头信息传递的HeaderMapper(默认SimpleKafkaHeaderMapper已支持),避免重试次数丢失导致无法触发DLT流转。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.12 01:40:21