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

Spring Boot集成测试:如何验证Kafka Testcontainer中发送的消息?

验证Testcontainer Kafka消息发送的几种实用方法

1. 临时测试监听容器(Spring Kafka原生支持)

在集成测试类中定义一个临时的Kafka监听器,专门捕获目标Topic的消息并存入线程安全集合,直接断言内容即可。

@SpringBootTest
@Testcontainers
public class KafkaMessageIntegrationTest {
    @Container
    static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.4.0"));

    // 存储收到的消息,用线程安全集合避免并发问题
    private final ConcurrentLinkedQueue<YourMessageDto> receivedMessages = new ConcurrentLinkedQueue<>();

    // 测试专用监听器,autoStartup设为true自动启动
    @KafkaListener(topics = "your-target-topic", groupId = "test-unique-group-id", autoStartup = "true")
    public void listenTestTopic(YourMessageDto message) {
        receivedMessages.add(message);
    }

    @Test
    void testMessageSend() throws InterruptedException {
        // 调用业务方法触发消息发送
        yourService.generateAndSendKafkaMessage();

        // 等待消息被捕获(用CountDownLatch替代Thread.sleep会更严谨)
        Thread.sleep(1000);
        // 断言消息数量与内容
        assertThat(receivedMessages).hasSize(1);
        assertThat(receivedMessages.peek().getTargetField()).isEqualTo("expected-content");
    }
}

注意:groupId必须设置唯一值,避免和其他测试或服务的消费者冲突;如果需要精准等待,可配合CountDownLatch在监听器中触发计数。

2. 原生Kafka消费者拉取验证

直接使用Kafka原生Consumer API手动拉取消息,适合不想依赖Spring Kafka监听组件的场景:

@SpringBootTest
@Testcontainers
public class KafkaMessageIntegrationTest {
    @Container
    static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.4.0"));

    @Test
    void testMessageSend() {
        // 配置消费者属性
        Properties props = new Properties();
        props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers());
        props.put(ConsumerConfig.GROUP_ID_CONFIG, "test-consumer-group");
        props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest");
        props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
        props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, JsonDeserializer.class.getName());
        // 指定反序列化的目标消息类
        props.put(JsonDeserializer.VALUE_DEFAULT_TYPE, YourMessageDto.class.getName());

        // 创建消费者并订阅目标Topic
        try (KafkaConsumer<String, YourMessageDto> consumer = new KafkaConsumer<>(props)) {
            consumer.subscribe(Collections.singletonList("your-target-topic"));
            // 拉取消息,设置合理超时时间
            ConsumerRecords<String, YourMessageDto> records = consumer.poll(Duration.ofSeconds(2));
            
            // 断言结果
            assertThat(records.count()).isEqualTo(1);
            YourMessageDto message = records.iterator().next().value();
            assertThat(message.getTargetField()).isEqualTo("expected-content");
        }
    }
}

3. ProducerInterceptor拦截记录消息

自定义ProducerInterceptor,在消息发送时拦截并记录,既可以用于验证,也能把消息持久化供接收端测试使用:

第一步:实现拦截器

public class TestMessageInterceptor implements ProducerInterceptor<String, Object> {
    // 静态集合存储发送的消息,测试时可直接访问
    public static final ConcurrentLinkedQueue<ProducerRecord<String, Object>> sentMessages = new ConcurrentLinkedQueue<>();

    @Override
    public ProducerRecord<String, Object> onSend(ProducerRecord<String, Object> record) {
        sentMessages.add(record);
        return record;
    }

    // 以下方法默认实现即可
    @Override
    public void onAcknowledgement(RecordMetadata metadata, Exception exception) {}
    @Override
    public void close() {}
    @Override
    public void configure(Map<String, ?> configs) {}
}

第二步:在测试配置中启用拦截器

在application-test.yml中添加配置:

spring:
  kafka:
    producer:
      properties:
        interceptor.classes: com.yourpackage.TestMessageInterceptor

第三步:测试中验证

@SpringBootTest
@Testcontainers
public class KafkaMessageIntegrationTest {
    @Container
    static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.4.0"));

    @Test
    void testMessageSend() {
        // 清空之前的测试记录
        TestMessageInterceptor.sentMessages.clear();
        
        // 调用业务方法发送消息
        yourService.generateAndSendKafkaMessage();

        // 断言发送的消息
        assertThat(TestMessageInterceptor.sentMessages).hasSize(1);
        ProducerRecord<String, Object> record = TestMessageInterceptor.sentMessages.peek();
        YourMessageDto message = (YourMessageDto) record.value();
        assertThat(message.getTargetField()).isEqualTo("expected-content");
        assertThat(record.topic()).isEqualTo("your-target-topic");
    }
}

这种方式还可以把sentMessages中的内容序列化到文件,直接供接收端服务的测试用例读取使用。

4. @SpyBean增强KafkaTemplate

用@SpyBean包装真实的KafkaTemplate,既保留发送到Testcontainer Kafka的真实逻辑,又能验证调用细节:

@SpringBootTest
@Testcontainers
public class KafkaMessageIntegrationTest {
    @Container
    static KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.4.0"));

    @SpyBean
    private KafkaTemplate<String, YourMessageDto> kafkaTemplate;

    @Autowired
    private YourService yourService;

    @Test
    void testMessageSend() throws ExecutionException, InterruptedException {
        // 调用业务方法触发消息发送
        yourService.generateAndSendKafkaMessage();

        // 验证发送次数
        verify(kafkaTemplate, times(1)).send(anyString(), any(YourMessageDto.class));

        // 捕获发送的消息内容
        ArgumentCaptor<ProducerRecord<String, YourMessageDto>> recordCaptor = ArgumentCaptor.forClass(ProducerRecord.class);
        verify(kafkaTemplate).send(recordCaptor.capture());
        
        ProducerRecord<String, YourMessageDto> sentRecord = recordCaptor.getValue();
        assertThat(sentRecord.topic()).isEqualTo("your-target-topic");
        assertThat(sentRecord.value().getTargetField()).isEqualTo("expected-content");
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.22 05:53:19