使用EmbeddedKafka做集成测试时waitForAssignment报期望1分区实际0错误
错误原因
- 代码存在多处拼写错误、变量名不统一问题,导致消费者配置不生效、无法正确连接嵌入式Kafka、配置读取失败
- 未提前创建消费者监听的目标Topic,EmbeddedKafka默认不会自动生成未声明的Topic,消费者监听不存在的Topic时自然无法获取到分区
- 部分配置参数写错导致消费者启动异常,无法正常向Kafka申请分区分配
解决方案
1. 修复所有语法/拼写错误
你需要先修正代码里的所有笔误:
- 注解修正:将
@EmbeddedKafka (partitions 1, ports = 9092)改为@EmbeddedKafka(partitions = 1, ports = 9092),补上partitions后的等号;将@SpringBoot Test改为@SpringBootTest,去掉中间空格 - 配置拼写修正:将
spring.kafka.bootstrap-servers=${spring.embedded.kakfa.brokers}里的kakfa改为kafka - 括号/引号修正:将
@TestPropertySource({"classpath:application.yaml"}}的末尾多余大括号删掉,改为@TestPropertySource({"classpath:application.yaml"});将@Value("${spring.kafka.group_id})补上末尾双引号,改为@Value("${spring.kafka.group_id}") - 变量名统一:你声明的队列变量是
records,初始化时写的consumerRecords = new LinkedBlockingQueue<>()要改为records = new LinkedBlockingQueue<>() - 消费者配置拼写修正:
Map<string, Object>改为Map<String, Object>- 方法参数里的
embeddedkafkaNroker改为embeddedKafkaBroker ConsumerConfig.BOOSTRAP_SERVER_CONFIG改为ConsumerConfig.BOOTSTRAP_SERVERS_CONFIGConsumerConfig.VALUE_DESERIALIZER_CLASS_COFIG改为ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG
2. 提前创建监听的目标Topic
在setUp方法最开头,先调用嵌入式Kafka的接口创建你要监听的Topic,确保消费者启动时Topic已经存在:
// 导入org.apache.kafka.clients.admin.NewTopic embeddedKafkaBroker.addTopics(new NewTopic(topicName, embeddedKafkaBroker.getPartitionsPerTopic(), (short) 1));
3. 可选:调整分区等待超时时间
如果嵌入式Kafka启动较慢,可以给waitForAssignment增加超时时间,避免启动慢导致的误判:
ContainerTestUtils.waitForAssignment(container, embeddedKafkaBroker.getPartitionsPerTopic(), 60);
内容的提问来源于stack exchange,提问作者Pavithra Srinivasan
相关产品推荐
相关产品推荐

