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

如何为带手动确认的@KafkaListener消费方法编写单元测试?

Kafka监听器单元测试实现方案

先搞定依赖

确保测试依赖里引入Spring Kafka Test:

<dependency>
    <groupId>org.springframework.kafka</groupId>
    <artifactId>spring-kafka-test</artifactId>
    <scope>test</scope>
</dependency>

测试类完整实现

核心思路是用@EmbeddedKafka启动内置Kafka集群,通过KafkaTemplate发送测试消息,再验证监听器的处理结果。

测试代码示例

import org.apache.avro.generic.GenericRecord;
import org.apache.avro.generic.GenericRecordBuilder;
import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.kafka.core.KafkaTemplate;
import org.springframework.kafka.test.context.EmbeddedKafka;
import org.springframework.test.context.TestPropertySource;

import java.util.concurrent.CountDownLatch;
import java.util.concurrent.TimeUnit;

import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.Mockito.verify;

@SpringBootTest
// 启动嵌入式Kafka,指定分区数和端口
@EmbeddedKafka(partitions = 1, brokerProperties = {"listeners=PLAINTEXT://localhost:9092", "port=9092"})
// 覆盖Kafka配置,用测试用的groupId和偏移量策略
@TestPropertySource(properties = {
        "spring.kafka.consumer.group-id=test-group",
        "spring.kafka.consumer.auto-offset-reset=earliest",
        "spring.kafka.producer.bootstrap-servers=localhost:9092"
})
public class YourKafkaListenerTest {

    @Autowired
    private KafkaTemplate<String, Object> kafkaTemplate;

    // 用来同步测试线程和监听器线程的计数器
    private final CountDownLatch processLatch = new CountDownLatch(1);

    @Autowired
    private YourKafkaListener targetListener;

    @Test
    public void testListenMethod() throws InterruptedException {
        // 构造符合业务格式的GenericRecord(这里以Avro为例,替换成你的实际Schema)
        GenericRecord testData = new GenericRecordBuilder(
                new org.apache.avro.Schema.Parser().parse("{\"type\":\"record\",\"name\":\"DemoRecord\",\"fields\":[{\"name\":\"userId\",\"type\":\"string\"}]}"))
                .set("userId", "test_123")
                .build();

        // 发送测试消息到指定主题
        kafkaTemplate.send("<topic>", "test-key", testData);

        // 等待5秒,直到监听器处理完成
        boolean processed = processLatch.await(5, TimeUnit.SECONDS);

        // 验证消息是否被处理
        assertTrue(processed, "监听器未在超时时间内处理消息");

        // 如果你需要验证业务逻辑,比如监听器调用了某个服务,可以加断言
        // 比如:assertEquals("test_123", targetListener.getLastProcessedUserId());
    }

    // 可选:不想修改原监听器的话,用MockBean替代业务依赖验证调用
    // @MockBean
    // private YourBizService bizService;
    // 测试时:verify(bizService, timeout(5000)).handleData(any(GenericRecord.class));
}

关键细节提示

  • @EmbeddedKafka:自动启动轻量Kafka集群,无需额外搭建外部环境,完全适配单元测试场景。
  • CountDownLatch:解决多线程同步问题,确保监听器处理完消息后再执行断言逻辑。
  • GenericRecord构造:如果使用Avro,必须匹配实际业务的Schema;仅测试消费流程的话,也可以用Mock模拟GenericRecord对象。
  • 手动提交偏移量:你的监听器使用ack.acknowledge()手动提交偏移量,测试时无需额外处理,嵌入式Kafka会正常处理偏移量逻辑。
  • 无侵入式测试:不想修改原监听器代码的话,用@MockBean替换监听器的业务依赖,通过验证依赖方法的调用来确认消费成功。

内容的提问来源于stack exchange,提问作者Kstackr

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 03:33:27