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

使用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_CONFIG
    • ConsumerConfig.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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.02 18:57:03