Spring Boot中EmbeddedKafkaBroker消息生产与监听集成测试问题
解决方案
1. 修复消费者配置类核心问题
你的KafkaConsumerConfig存在两个关键问题:
- 继承
KafkaProperties.Consumer但未正确初始化父类属性,且混用@Value和@ConfigurationProperties导致配置绑定异常 - 自定义的
kafkaConsumerCustomizer为容器设置KafkaAdmin,可能干扰EmbeddedKafka的自动配置
修改后的KafkaConsumerConfig如下:
@Configuration @ConfigurationProperties(prefix = "demo.kafka.consumer") public class KafkaConsumerConfig { private String clusterId; private String bootstrapServers; private boolean observationEnabled; private String autoOffsetReset = "earliest"; // 生成getter/setter,让@ConfigurationProperties自动绑定属性 public String getClusterId() { return clusterId; } public void setClusterId(String clusterId) { this.clusterId = clusterId; } public String getBootstrapServers() { return bootstrapServers; } public void setBootstrapServers(String bootstrapServers) { this.bootstrapServers = bootstrapServers; } public boolean isObservationEnabled() { return observationEnabled; } public void setObservationEnabled(boolean observationEnabled) { this.observationEnabled = observationEnabled; } public String getAutoOffsetReset() { return autoOffsetReset; } public void setAutoOffsetReset(String autoOffsetReset) { this.autoOffsetReset = autoOffsetReset; } @Bean public ConsumerFactory<String, String> kafkaConsumerFactory() { Map<String, Object> configProps = new HashMap<>(); configProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrapServers); configProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, autoOffsetReset); configProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); configProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class); if (StringUtils.hasText(clusterId)) { configProps.put("clusterId", clusterId); } return new DefaultKafkaConsumerFactory<>(configProps); } @Bean public ConcurrentKafkaListenerContainerFactory<String, String> kafkaContainerFactory( ConsumerFactory<String, String> kafkaConsumerFactory) { ConcurrentKafkaListenerContainerFactory<String, String> factory = new ConcurrentKafkaListenerContainerFactory<>(); factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.MANUAL); factory.setConsumerFactory(kafkaConsumerFactory); factory.getContainerProperties().setObservationEnabled(observationEnabled); return factory; } }
2. 优化监听器类(添加测试验证机制)
给监听器添加CountDownLatch,方便测试时验证消息是否被接收,同时补全手动提交偏移量的逻辑:
@Component @Slf4j public class KafkaListener { // 修正类名首字母大写,符合Java规范 public final KafkaTemplate<String, String> kafkaTemplate; public final ObjectMapper objectMapper; public final CountDownLatch messageLatch = new CountDownLatch(1); public KafkaListener(KafkaTemplate<String, String> kafkaTemplate, ObjectMapper objectMapper) { this.kafkaTemplate = kafkaTemplate; this.objectMapper = objectMapper; } @KafkaListener(topics = "${demo.topic-name}", groupId = "${demo.consumer-group-name}", containerFactory = "kafkaContainerFactory") public void listenFromTopic(ConsumerRecord<String, String> consumerRecord, Acknowledgment acknowledgment) { log.debug("Received message from topic: {}, key: {}, partition: {}, offset: {}", consumerRecord.topic(), consumerRecord.key(), consumerRecord.partition(), consumerRecord.offset()); messageLatch.countDown(); acknowledgment.acknowledge(); // 手动提交偏移量,避免重复消费 } }
3. 编写正确的集成测试类
使用@EmbeddedKafka自动启动嵌入式Kafka,通过@TestPropertySource覆盖生产者和消费者的bootstrap地址为嵌入式Kafka的动态地址,同时开启监听器的DEBUG日志:
@SpringBootTest @EmbeddedKafka(partitions = 1, topics = "${demo.topic-name}", kraft = false) @TestPropertySource(properties = { "demo.kafka.consumer.bootstrap-servers=${spring.embedded.kafka.brokers}", "spring.kafka.producer.bootstrap-servers=${spring.embedded.kafka.brokers}", "logging.level.your.package.KafkaListener=DEBUG" // 替换为监听器所在的实际包路径 }) public class KafkaListenerIntegrationTest { @Autowired private KafkaTemplate<String, String> kafkaTemplate; @Autowired private KafkaListener kafkaListener; @Value("${demo.topic-name}") private String topic; @Test void testListenerReceivesMessage() throws InterruptedException, ExecutionException { // 发送测试消息 String testKey = "test-key"; String testValue = "test-value"; kafkaTemplate.send(topic, testKey, testValue).get(10, TimeUnit.SECONDS); // 等待监听器处理消息,超时10秒 boolean messageReceived = kafkaListener.messageLatch.await(10, TimeUnit.SECONDS); assert messageReceived : "监听器未收到消息"; } }
4. 关键配置说明
${spring.embedded.kafka.brokers}是EmbeddedKafka自动暴露的变量,指向其实际运行的地址(随机端口,无需硬编码9092)- 测试时覆盖生产者和消费者的bootstrap地址,确保两者都连接到同一个嵌入式Kafka实例
- 开启监听器的DEBUG日志,确保能看到接收消息的日志输出
5. 额外检查点
- 确保你的
application.properties或application.yml中已配置demo.topic-name和demo.consumer-group-name - 如果使用Lombok,确认
@Slf4j注解已正确生成日志实例 - 测试时如果仍无日志输出,可临时将监听器的日志级别改为
INFO,避免DEBUG级别被日志框架过滤
内容的提问来源于stack exchange,提问作者geeta
相关产品推荐
相关产品推荐

