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

同一容器多次调用ContainerTestUtils.waitForAssignment及消费者重复问题咨询

Kafka监听器集成测试问题解决方案(外部Kafka容器场景)

问题1:固定groupId时,后续测试waitForAssignment报错

问题现象

使用固定groupId(如test.some-event),多个集成测试类继承包含ContainerTestUtils.waitForAssignment调用的基类时,第一个测试正常执行,第二个测试抛出java.lang.IllegalStateException: Expected 1 but got 0 partitions,且未使用@DirtiesContext,Spring上下文处于缓存状态。

原因

Spring上下文缓存导致Kafka消费者容器在第一个测试后未重启,而固定groupId的消费者已在Kafka集群中完成分区注册,第二个测试时容器尝试重新分配分区,但Kafka认为该groupId的消费者仍处于活跃状态,不会重新分配,最终waitForAssignment检测到0个分区。

解决方案

在每个测试的初始化和清理阶段,手动重启消费者容器,强制Kafka重新分配分区:

@ActiveProfiles("test")
@SpringBootTest
abstract class BaseKafkaIntegrationSpec extends Specification {

    @Autowired(required = true)
    private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry

    void setup() {
        kafkaListenerEndpointRegistry.getAllListenerContainers()
            .stream()
            .forEach { container ->
                // 停止运行中的容器,确保重新启动
                if (container.isRunning()) {
                    container.stop()
                }
                container.start()
                ContainerTestUtils.waitForAssignment(container, 1)
            }
    }

    void cleanup() {
        kafkaListenerEndpointRegistry.getAllListenerContainers()
            .stream()
            .forEach { container ->
                container.stop()
            }
    }
}

问题2:随机groupId时,多测试类创建重复消费者实例

问题现象

配置随机groupId(test.some-event-${random.uuid})后,尽管Spring缓存上下文,但每个测试类仍会创建新的SampleKafkaConsumer实例,多个不同groupId的消费者同时运行并消费事件,期望复用现有Bean/容器而非新建。

原因

${random.uuid}在Spring上下文初始化时生成,不同测试类触发上下文刷新(因配置属性变化),会生成新的上下文及对应的消费者Bean和容器。

解决方案

复用现有消费者容器,在每个测试前动态修改groupId,避免创建新Bean实例:

@ActiveProfiles("test")
@SpringBootTest
abstract class BaseKafkaIntegrationSpec extends Specification {

    @Autowired(required = true)
    private KafkaListenerEndpointRegistry kafkaListenerEndpointRegistry

    void setup() {
        // 生成当前测试唯一的groupId
        def uniqueGroupId = "test.some-event-${UUID.randomUUID()}"
        
        kafkaListenerEndpointRegistry.getAllListenerContainers()
            .stream()
            .forEach { container ->
                if (container.isRunning()) {
                    container.stop()
                }
                // 动态修改容器的groupId
                container.getContainerProperties().setGroupId(uniqueGroupId)
                container.start()
                ContainerTestUtils.waitForAssignment(container, 1)
            }
    }

    void cleanup() {
        kafkaListenerEndpointRegistry.getAllListenerContainers()
            .stream()
            .forEach { container ->
                container.stop()
            }
    }
}

该方案无需修改原SampleKafkaConsumer代码,直接通过操作容器实例实现groupId动态变更,同时复用Spring上下文中的Bean,避免重复创建消费者。


内容的提问来源于stack exchange,提问作者BartekN

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 04:02:08