Spring Boot Kafka Producer禁用Topic自动创建问题咨询
问题:禁用Kafka Producer自动创建Topic并在Topic不存在时抛出异常
我使用Spring Boot 2.7.6(对应spring-kafka版本2.8.11),需要禁用Kafka Producer在Topic不存在时自动创建Topic的功能,期望发送消息时若Topic不可用则抛出异常(当前会自动创建Topic)。
我已找到Kafka Consumer对应的配置项ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG,且该配置生效。
现有配置代码
Consumer配置
@Bean fun kafkaConsumerFactory(): ConsumerFactory<String?, String?> { val props: MutableMap<String, Any> = HashMap() props[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = bootstrapAddress props[ConsumerConfig.GROUP_ID_CONFIG] = groupId props[ConsumerConfig.AUTO_OFFSET_RESET_CONFIG] = "latest" props[ConsumerConfig.ALLOW_AUTO_CREATE_TOPICS_CONFIG] = false props[ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG] = ErrorHandlingDeserializer::class.java props[ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG] = ErrorHandlingDeserializer::class.java props[ErrorHandlingDeserializer.KEY_DESERIALIZER_CLASS] = StringDeserializer::class.java props[ErrorHandlingDeserializer.VALUE_DESERIALIZER_CLASS] = StringDeserializer::class.java return DefaultKafkaConsumerFactory(props) }
Producer配置
@Bean fun kafkaProducerFactory(): ProducerFactory<String?, String?> { val configProps: MutableMap<String, Any> = HashMap() configProps[ProducerConfig.BOOTSTRAP_SERVERS_CONFIG] = bootstrapAddress configProps[ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG] = StringSerializer::class.java configProps[ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG] = StringSerializer::class.java return DefaultKafkaProducerFactory(configProps) }
解决方案
这个需求可以实现,你遗漏了关键的配置逻辑,具体说明如下:
核心原理
Kafka Producer本身没有类似Consumer的ALLOW_AUTO_CREATE_TOPICS_CONFIG配置项,自动创建Topic的行为本质上是由Kafka Broker端的auto.create.topics.enable参数控制的。默认情况下该参数为true,当Producer发送消息到不存在的Topic时,Broker会自动创建Topic。
实现方式
方式一:修改Broker配置(推荐)
将Broker端的auto.create.topics.enable设置为false。此时Producer发送消息到不存在的Topic时,会直接抛出UnknownTopicOrPartitionException异常,完全符合你的需求。
方式二:Producer端自定义检查(无需修改Broker)
如果无法修改Broker配置,可以在发送消息前主动检查Topic是否存在,不存在则抛出异常:
- 配置AdminClient用于Topic检查
@Bean fun kafkaAdmin(): KafkaAdmin { val configs = mapOf( AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG to bootstrapAddress ) return KafkaAdmin(configs) } @Bean fun adminClient(): AdminClient { return AdminClient.create(kafkaAdmin().configurationProperties) }
- 在消息发送服务中添加检查逻辑
@Service class KafkaMessageService( private val kafkaTemplate: KafkaTemplate<String?, String?>, private val adminClient: AdminClient ) { fun sendMessage(topic: String, message: String) { // 同步检查Topic是否存在 val existingTopics = adminClient.listTopics().names().get() if (!existingTopics.contains(topic)) { throw IllegalStateException("目标Topic $topic 不存在") } // 发送消息 kafkaTemplate.send(topic, message) .addCallback( { result -> /* 发送成功处理逻辑 */ }, { ex -> /* 发送失败处理逻辑,比如重新抛出异常 */ throw ex } ) } }
注意事项
- 如果你设置了Producer的
retries参数(默认值为2147483647),即使Broker禁用了自动创建Topic,Producer会自动重试发送,直到重试次数耗尽才会抛出异常。如果需要快速抛出异常,可以调整retries为0,但需根据业务权衡可靠性。
内容的提问来源于stack exchange,提问作者Gobanit
相关产品推荐
相关产品推荐

