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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 13:15:53