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
相关产品推荐
相关产品推荐

