Spring-Kafka:自定义重试与死信主题名称的RetryTopicConfiguration配置问题
问题a:让自定义命名提供者被Spring识别
你需要在构建RetryTopicConfiguration时,显式关联自定义的RetryTopicNamesProviderFactory——Spring不会自动绑定你定义的这个Bean,必须通过构建器方法手动注入。
修改你的kafkaRetryTopicConfig方法,注入自定义命名工厂Bean并添加到构建链中:
@Bean public RetryTopicConfiguration kafkaRetryTopicConfig(RetryTopicNamesProviderFactory retryTopicNamingProviderFactory, ...) { return RetryTopicConfigurationBuilder .newInstance() .fixedBackOff(...) .maxAttempts(..) .useSingleTopicForFixedDelays() .doNotRetryOnDltFailure() .listenerFactory(factory) .retryTopicNamesProviderFactory(retryTopicNamingProviderFactory) // 关键:关联自定义命名工厂 .create(template); }
完成后Spring就会使用你的自定义命名规则,而非默认的.retry/.deadLetter后缀。
问题b:区分重试主题与死信主题的命名
可以通过DestinationTopic.Properties对象的getDestinationType()方法判断目标主题类型,该方法返回DestinationTopic.Type枚举,包含RETRY、DLT、SCHEDULED_RETRY等类型标识。
修改你的命名提供者代码,根据类型返回对应配置的主题名称:
// 先注入配置文件中的主题名称 @Value("${kafka.retry.topic}") private String retryTopicName; @Value("${kafka.dlt.topic}") private String dltTopicName; @Bean public RetryTopicNamesProviderFactory retryTopicNamingProviderFactory() { return new RetryTopicNamesProviderFactory() { @Override public RetryTopicNamesProvider createRetryTopicNamesProvider(DestinationTopic.Properties properties) { DestinationTopic.Type topicType = properties.getDestinationType(); return new SuffixingRetryTopicNamesProvider(properties) { @Override public String getTopicName(String originalTopic) { if (DestinationTopic.Type.DLT.equals(topicType)) { return dltTopicName; } else if (DestinationTopic.Type.RETRY.equals(topicType) || DestinationTopic.Type.SCHEDULED_RETRY.equals(topicType)) { return retryTopicName; } return originalTopic; } }; } }; }
内容的提问来源于stack exchange,提问作者Higher-Kinded Type
相关产品推荐
相关产品推荐

