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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.24 15:24:55