Spring Kafka @RetryableTopic异常:禁用自动创建仍生成主题及自定义主题问题
问题1:设置autoCreateTopics=false后仍自动创建主题
原因分析
@RetryableTopic的autoCreateTopics属性仅控制重试/死信主题的自动创建,而Kafka消费者客户端默认开启auto.create.topics.enable=true,当消费者尝试订阅不存在的主题时,会触发客户端自动创建主题。此外,你代码中autoCreateTopics="false"使用字符串形式,虽Spring支持该写法,但需配合消费者端配置才能彻底禁用自动创建。
解决方案
修正
@RetryableTopic配置
保持注解中autoCreateTopics="false"(或改为false,注解属性为String类型支持SpEL),确保明确关闭重试主题的自动创建:@RetryableTopic( attempts = "4", backoff = @Backoff(delay = 1000), autoCreateTopics = "false", kafkaTemplate = "consumerKafkaTemplate", include = RuntimeException.class )禁用消费者端自动创建主题
在配置文件(如application.yml)中添加:spring: kafka: consumer: properties: auto.create.topics.enable: false或通过代码配置
ConsumerFactory时显式设置:@Bean public ConsumerFactory<Object, Object> consumerFactory() { Map<String, Object> props = new HashMap<>(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "your-kafka-servers"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "your-group-id"); props.put(ConsumerConfig.AUTO_CREATE_TOPICS_ENABLE_CONFIG, false); // 其他消费者配置 return new DefaultKafkaConsumerFactory<>(props); }
问题2:自定义RetryTopicNamesProviderFactory出现未知主题异常
原因分析
启用自定义重试主题名称后,由于autoCreateTopics=false,Spring Kafka不会自动生成这些主题,若Kafka集群中不存在对应主题,消费者订阅时就会抛出UNKOWN_TOPIC_OR_PARTITION异常。同时需确保自定义名称提供者正确关联到@RetryableTopic。
解决方案
实现自定义RetryTopicNamesProvider
创建自定义主题命名规则类:@Component public class CustomRetryTopicNamesProvider implements RetryTopicNamesProvider { @Override public String getTopicName(String originalTopic, int retryAttempt) { // 自定义重试主题命名,例:bookStore-retry-1、bookStore-retry-2 return originalTopic + "-retry-" + retryAttempt; } @Override public String getDltTopicName(String originalTopic) { // 自定义死信主题命名,例:bookStore-dlt return originalTopic + "-dlt"; } }配置RetryTopicNamesProviderFactory Bean
将自定义提供者包装为Factory:@Bean public RetryTopicNamesProviderFactory customRetryTopicNamesProviderFactory(CustomRetryTopicNamesProvider provider) { return RetryTopicNamesProviderFactory.create(provider); }关联自定义Factory到@RetryableTopic
在注解中指定自定义Factory,同时修正@KafkaListener的语法错误:@RetryableTopic( attempts = "4", backoff = @Backoff(delay = 1000), autoCreateTopics = "false", kafkaTemplate = "consumerKafkaTemplate", include = RuntimeException.class, retryTopicNamesProviderFactory = "customRetryTopicNamesProviderFactory" ) @KafkaListener(topicPartitions = {@TopicPartition(topic = "bookStore")}) public void listen(Message message, Acknowledgement ack) { log.info("Message Recieved from Topic ::" + message.getHeaders().get("topic")); }手动创建Kafka主题
在Kafka集群中手动创建所有需要的主题(主主题、重试主题、死信主题),确保名称与自定义规则一致:# 创建主主题 kafka-topics.sh --create --topic bookStore --bootstrap-server your-kafka-servers --partitions 3 --replication-factor 1 # 创建重试主题 kafka-topics.sh --create --topic bookStore-retry-1 --bootstrap-server your-kafka-servers --partitions 3 --replication-factor 1 kafka-topics.sh --create --topic bookStore-retry-2 --bootstrap-server your-kafka-servers --partitions 3 --replication-factor 1 kafka-topics.sh --create --topic bookStore-retry-3 --bootstrap-server your-kafka-servers --partitions 3 --replication-factor 1 # 创建死信主题 kafka-topics.sh --create --topic bookStore-dlt --bootstrap-server your-kafka-servers --partitions 3 --replication-factor 1
额外注意事项
- 确保
consumerKafkaTemplate的默认主题配置不会干扰重试主题的消息路由。 - 重试主题的分区数、副本数建议与主主题保持一致,避免消息堆积。
内容的提问来源于stack exchange,提问作者kumar

