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

Spring Boot中自定义Container Factory的Kafka Listener测试问题

问题原因

你自定义的containerFactory硬编码了ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG为127.0.0.1:9092,但Testcontainers启动的Kafka容器使用随机映射端口,通过DynamicPropertySource注入的spring.kafka.bootstrap-servers并没有被这个自定义工厂读取,导致监听器连接固定的127.0.0.1:9092,而非容器实际提供的地址。

解决方案

有两种可行的解决方式:

方式一:让自定义工厂读取配置属性

修改TestConfig中的containerFactory方法,不再硬编码bootstrap地址,而是通过配置属性注入获取动态地址:

@TestConfiguration
class TestConfig {
    @Autowired
    private lateinit var kafkaProperties: KafkaProperties

    @Bean
    fun containerFactory(): ConcurrentKafkaListenerContainerFactory<String, String> {
        val factory = ConcurrentKafkaListenerContainerFactory<String, String>()
        // 从KafkaProperties中读取完整的消费者配置,包括动态注入的bootstrap地址
        val consumerProps = kafkaProperties.buildConsumerProperties()
        val defaultConsumer = DefaultKafkaConsumerFactory(consumerProps, StringDeserializer(), StringDeserializer())
        factory.consumerFactory = defaultConsumer

        return factory
    }
}

如果不需要完整的KafkaProperties,也可以单独注入bootstrap地址:

@TestConfiguration
class TestConfig {
    @Value("\${spring.kafka.bootstrap-servers}")
    private lateinit var bootstrapServers: String

    @Bean
    fun containerFactory(): ConcurrentKafkaListenerContainerFactory<String, String> {
        val factory = ConcurrentKafkaListenerContainerFactory<String, String>()
        val props = HashMap<String, Any?>()
        props[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = bootstrapServers
        props[ConsumerConfig.AUTO_OFFSET_RESET_CONFIG] = "earliest"
        val defaultConsumer = DefaultKafkaConsumerFactory(props, StringDeserializer(), StringDeserializer())
        factory.consumerFactory = defaultConsumer

        return factory
    }
}

方式二:在测试中动态替换工厂配置

如果不想修改工厂原有定义,可在测试类中注入containerFactory,动态更新其bootstrap地址:

@Testcontainers
@SpringBootTest(properties = ["spring.kafka.consumer.auto-offset-reset=earliest"])
@EnableAutoConfiguration(exclude = [MongoAutoConfiguration::class])
@Import(KafkaListener::class, TestConfig::class)
class KafkaShippingGroupCreationListenerAdapterTest {

    companion object {
        @Container
        private val kafkaContainer = KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.3.3"))

        @JvmStatic
        @DynamicPropertySource
        fun overrideProperties(registry: DynamicPropertyRegistry) {
            registry.add("spring.kafka.bootstrap-servers", kafkaContainer::getBootstrapServers);
        }
    }

    @MockBean
    private lateinit var useCase: UseCase

    @Autowired
    private lateinit var kafkaTemplate: KafkaTemplate<String, Any>

    @Autowired
    private lateinit var containerFactory: ConcurrentKafkaListenerContainerFactory<String, String>

    @BeforeEach
    fun setUp() {
        val consumerFactory = containerFactory.consumerFactory as DefaultKafkaConsumerFactory<String, String>
        val updatedProps = consumerFactory.configuration.toMutableMap()
        updatedProps[ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG] = kafkaContainer.bootstrapServers
        consumerFactory.updateConfigs(updatedProps)
    }

    @Test
    fun `should invoke the use case`() {
        kafkaTemplate.send("topic", "message")

        verify(useCase, timeout(5000).times(1)).foo()
    }
}
说明

第一种方式更推荐,它让自定义工厂遵循Spring Boot配置规范,在测试和生产环境都能正确读取配置的Kafka地址;第二种方式适合临时调整测试配置,无需修改工厂原有代码。

内容的提问来源于stack exchange,提问作者Mauricio Avendaño

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.09 21:05:54