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

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是否存在,不存在则抛出异常:

  1. 配置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)
}
  1. 在消息发送服务中添加检查逻辑
@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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 09:05:23