使用@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
相关产品推荐
相关产品推荐

