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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 10:55:20