Spring Boot集成Testcontainers Kafka:如何验证生产者消息发送成功?
验证Testcontainers Kafka消息发送的几种方案(无需业务消费者或Mock)
方案1:测试中临时创建轻量消费者
直接在测试代码里创建一次性Kafka消费者,专门用于验证消息发送结果,完全不改动业务代码。示例代码:
@SpringBootTest @Testcontainers public class KafkaProducerTest { @Container private static final KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.4.0")); @Autowired private KafkaTemplate<String, String> kafkaTemplate; @Test void testMessageSentSuccessfully() throws InterruptedException, ExecutionException { // 发送测试消息 String testTopic = "test-topic"; String testMessage = "hello-testcontainers"; kafkaTemplate.send(testTopic, testMessage).get(); // 配置临时消费者参数 Properties consumerProps = new Properties(); consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers()); consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, "test-group-" + UUID.randomUUID()); consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); // 创建并使用消费者验证消息 try (Consumer<String, String> consumer = new KafkaConsumer<>(consumerProps)) { consumer.subscribe(Collections.singleton(testTopic)); // 等待消息(设置合理超时时间) ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(5)); // 断言验证 assertThat(records).isNotEmpty(); assertThat(records.iterator().next().value()).isEqualTo(testMessage); } } }
核心是仅在测试阶段临时创建消费者,用完即销毁,完全不侵入业务逻辑,也无需Mock生产者接口。
方案2:用Spring Kafka TestUtils简化验证逻辑
如果项目使用Spring Kafka,可直接借助TestUtils类快速获取消息,代码更简洁:
@SpringBootTest @Testcontainers public class KafkaProducerTest { @Container private static final KafkaContainer kafkaContainer = new KafkaContainer(DockerImageName.parse("confluentinc/cp-kafka:7.4.0")); @Autowired private KafkaTemplate<String, String> kafkaTemplate; @Autowired private ConsumerFactory<String, String> consumerFactory; @Test void testMessageSentWithSpringTestUtils() throws ExecutionException, InterruptedException { String testTopic = "test-topic"; String testMessage = "hello-spring-test"; kafkaTemplate.send(testTopic, testMessage).get(); // 创建临时消费者并获取消息 Consumer<String, String> consumer = consumerFactory.createConsumer("test-group-" + UUID.randomUUID(), "test-instance"); consumer.subscribe(Collections.singleton(testTopic)); ConsumerRecord<String, String> record = TestUtils.getSingleRecord(consumer, testTopic, 5000); assertThat(record.value()).isEqualTo(testMessage); consumer.close(); } }
TestUtils.getSingleRecord会自动等待直到获取到目标消息或超时,省去手动poll的繁琐代码。
方案3:通过AdminClient检查主题偏移量(轻量验证)
如果只需确认消息已写入主题、无需校验内容,可借助Kafka AdminClient检查分区偏移量变化:
@Test void testMessageOffsetIncreased() throws ExecutionException, InterruptedException { String testTopic = "test-topic"; // 初始化AdminClient AdminClient adminClient = AdminClient.create(Map.of(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, kafkaContainer.getBootstrapServers())); // 获取发送前的主题偏移量总和 TopicDescription topicDescription = adminClient.describeTopics(Collections.singleton(testTopic)).values().get(testTopic).get(); long initialOffset = topicDescription.partitions().stream() .map(PartitionInfo::partition) .mapToLong(p -> adminClient.listOffsets(Map.of(new TopicPartition(testTopic, p), OffsetSpec.latest())).all().get().get(new TopicPartition(testTopic, p)).offset()) .sum(); // 发送消息 kafkaTemplate.send(testTopic, "test-message").get(); // 获取发送后的偏移量总和 long afterSendOffset = topicDescription.partitions().stream() .map(PartitionInfo::partition) .mapToLong(p -> adminClient.listOffsets(Map.of(new TopicPartition(testTopic, p), OffsetSpec.latest())).all().get().get(new TopicPartition(testTopic, p)).offset()) .sum(); // 验证偏移量增加(说明有消息写入) assertThat(afterSendOffset).isEqualTo(initialOffset + 1); adminClient.close(); }
该方案适合只需要确认消息已投递的场景,代码更轻量,但无法验证消息具体内容。
内容的提问来源于stack exchange,提问作者testtt
相关产品推荐
相关产品推荐

