SpringBoot集成测试:EmbeddedKafka事件断言的优化方案咨询
问题场景
我正在为SpringBoot应用编写集成测试,需要验证系统发送了两条消息。为避免使用Mock,我选择采用@EmbeddedKafka实现,并编写了仅用于测试包的消费者:
@Component public class StudentCreatedTestConsumer { private ObjectMapper mapper = new ObjectMapper(); @Setter private Runnable afterReceive = () -> {}; @Getter private List<Map<String, String>> events = new ArrayList<>(); public void reset() { afterReceive = () -> {}; events = new ArrayList<>(); } @KafkaListener(topics = "student.created") public void receive(ConsumerRecord<String, String> evt) { events.add(parseHashMap(evt.value())); afterReceive.run(); } @SneakyThrows private HashMap<String, String> parseHashMap(String evt) { TypeReference<HashMap<String, String>> typeRef = new TypeReference<>() {}; return mapper.readValue(evt, typeRef); } public List<Map<String, String>> getEventsSortedBy(String key) { return events.stream() .sorted(comparing(e -> e.get(key))) .toList(); } }
我在测试中这样使用该消费者:
@Test void test() throws Exception { //given CountDownLatch latch = new CountDownLatch(2); // will wait for 2 events testConsumer.setAfterReceive(latch::countDown); // when mockMvc.perform(post("/api/students") .content(...) .contentType(contentType)) .andExpect(status().isOk()); //then latch.await(10, TimeUnit.SECONDS); List<Map<String, String>> events = testConsumer.getEventsSortedBy("name"); assertThat(events).hasSize(2); assertThat(events.get(0)) .containsKey("studentId") .containsEntry("name", "anna") .containsEntry("age", "27"); //and assertThat(events.get(1)) .containsKey("studentId") .containsEntry("name", "bobby") .containsEntry("age", "25"); }
目前我使用CountDownLatch等待前2个事件,但如果应用实际产生3个事件,当前方案会存在问题。请问是否有更优的测试方式?
优化方案
1. 改用CompletableFuture结合条件判断
放弃CountDownLatch,改用CompletableFuture实现“等待到指定事件数量或超时”的逻辑,即使多产生事件,也能准确等到目标数量后再执行断言:
先给测试消费者新增等待方法:
public CompletableFuture<Void> waitForEventCount(int expectedCount, Duration timeout) { CompletableFuture<Void> future = new CompletableFuture<>(); // 先检查当前是否已经满足数量 if (events.size() >= expectedCount) { future.complete(null); return future; } // 自定义回调,每次收到事件后检查数量 this.afterReceive = () -> { if (events.size() >= expectedCount) { future.complete(null); } }; // 超时处理 CompletableFuture.delayedExecutor(timeout.toMillis(), TimeUnit.MILLISECONDS) .execute(() -> future.completeExceptionally(new TimeoutException("Timeout waiting for " + expectedCount + " events"))); return future; }
测试代码修改为:
@Test void test() throws Exception { //given testConsumer.reset(); // when mockMvc.perform(post("/api/students") .content(...) .contentType(contentType)) .andExpect(status().isOk()); //then testConsumer.waitForEventCount(2, Duration.ofSeconds(10)).get(); List<Map<String, String>> events = testConsumer.getEventsSortedBy("name"); assertThat(events).hasSize(2); // 后续断言逻辑不变... }
2. 使用Spring Kafka Test的KafkaTestUtils
Spring Kafka Test提供的KafkaTestUtils工具类可以直接从嵌入式Kafka拉取消息,无需维护自定义消费者,逻辑更简洁灵活:
@Autowired private EmbeddedKafkaBroker embeddedKafkaBroker; @Test void test() throws Exception { // when mockMvc.perform(post("/api/students") .content(...) .contentType(contentType)) .andExpect(status().isOk()); //then // 拉取指定数量的消息,超时时间10秒 List<ConsumerRecord<String, String>> records = KafkaTestUtils.getRecords( KafkaTestUtils.createConsumer(embeddedKafkaBroker.getBrokersAsString()), Duration.ofSeconds(10), 2 ); assertThat(records).hasSize(2); // 解析并断言消息内容 ObjectMapper mapper = new ObjectMapper(); List<Map<String, String>> events = records.stream() .map(record -> { try { return mapper.readValue(record.value(), new TypeReference<HashMap<String, String>>() {}); } catch (JsonProcessingException e) { throw new RuntimeException(e); } }) .sorted(comparing(e -> e.get("name"))) .toList(); assertThat(events.get(0)) .containsKey("studentId") .containsEntry("name", "anna") .containsEntry("age", "27"); assertThat(events.get(1)) .containsKey("studentId") .containsEntry("name", "bobby") .containsEntry("age", "25"); }
这种方式不需要自定义测试消费者,直接从Kafka拉取需要的2条消息,即使生产者多发送消息也不会影响测试——因为我们只处理指定数量的消息。
3. 增强自定义消费者的等待逻辑
如果坚持使用自定义测试消费者,可以修改逻辑,让等待信号只在达到目标数量时触发一次,避免后续多余消息干扰:
@Component public class StudentCreatedTestConsumer { // ... 原有代码不变 private CompletableFuture<Void> pendingFuture; private int expectedCount; public synchronized CompletableFuture<Void> waitFor(int count, Duration timeout) { this.expectedCount = count; pendingFuture = new CompletableFuture<>(); if (events.size() >= count) { pendingFuture.complete(null); return pendingFuture; } // 超时任务 CompletableFuture.delayedExecutor(timeout.toMillis(), TimeUnit.MILLISECONDS) .execute(() -> { synchronized (this) { if (!pendingFuture.isDone()) { pendingFuture.completeExceptionally(new TimeoutException("Expected " + count + " events, got " + events.size())); } } }); return pendingFuture; } @Override @KafkaListener(topics = "student.created") public void receive(ConsumerRecord<String, String> evt) { synchronized (this) { events.add(parseHashMap(evt.value())); if (pendingFuture != null && !pendingFuture.isDone() && events.size() >= expectedCount) { pendingFuture.complete(null); } } } }
测试时直接调用waitFor方法即可,后续的第3条消息不会重复触发完成信号,确保断言逻辑不受干扰。
内容的提问来源于stack exchange,提问作者Emanuel Trandafir
相关产品推荐
相关产品推荐

