EmbeddedKafka多主题测试分区数不匹配报错的解决办法
问题原因
你在@EmbeddedKafka中配置的partitions=6是每个主题的分区数量,当前有3个主题,所以Kafka会创建6×3=18个分区。而你的setUp方法里,调用ContainerTestUtils.waitForAssignment时传入的是embeddedKafkaBroker.getPartitionsPerTopic()(返回6),但你的监听器容器实际订阅了全部3个主题,会被分配18个分区,预期值和实际值不匹配,因此抛出错误。
解决方案
根据你的监听器订阅情况,选择以下一种方式修复:
方式一:计算总分区数(适用于单个容器订阅多个主题)
如果你的监听器容器同时订阅了topicA、topicB、topicC这三个主题,直接计算所有主题的总分区数传入方法:
@BeforeEach void setUp() { // 计算所有主题的总分区数:主题数量 × 每个主题的分区数 int totalPartitions = embeddedKafkaBroker.getTopics().size() * embeddedKafkaBroker.getPartitionsPerTopic(); for (MessageListenerContainer messageListenerContainer : endpointRegistry.getAllListenerContainers()) { ContainerTestUtils.waitForAssignment(messageListenerContainer, totalPartitions); } }
方式二:不指定预期分区数(通用方案)
ContainerTestUtils.waitForAssignment有一个重载方法,不需要传入预期分区数,它会自动等待容器完成分区分配,无需手动计算:
@BeforeEach void setUp() { for (MessageListenerContainer messageListenerContainer : endpointRegistry.getAllListenerContainers()) { ContainerTestUtils.waitForAssignment(messageListenerContainer); } }
方式三:针对单个主题的容器单独设置(适用于多个容器各订阅一个主题)
如果你的项目中有多个监听器容器,每个容器只订阅一个主题,那么可以为每个容器单独传入对应主题的分区数(6)。这种情况下需要确保每个容器对应一个主题,示例代码如下:
@BeforeEach void setUp() { int partitionsPerTopic = embeddedKafkaBroker.getPartitionsPerTopic(); for (MessageListenerContainer messageListenerContainer : endpointRegistry.getAllListenerContainers()) { ContainerTestUtils.waitForAssignment(messageListenerContainer, partitionsPerTopic); } }
(注:这种情况仅当每个容器只订阅单个主题时有效,否则仍会出现分区数不匹配的错误)
内容的提问来源于stack exchange,提问作者Anil
相关产品推荐
相关产品推荐

