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
相关产品推荐
相关产品推荐

