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

如何测试Spring Cloud Stream Kafka Binder?求简单单元测试示例

测试Spring Cloud Stream Kafka Streams Function的两种方案

你可以直接测试这个Function,也可以用Kafka Streams提供的TopologyTestDriver做更底层的拓扑测试,两种方式都不需要真实Kafka集群,只需要引入对应的测试依赖即可。下面是两种方案的具体示例:

一、直接测试Function(轻量快速)

这种方式直接调用你的Function Bean,借助Spring Cloud Stream的测试绑定器模拟消息输入输出,适合快速验证业务逻辑。

所需Maven依赖

<dependency>
    <groupId>org.springframework.cloud</groupId>
    <artifactId>spring-cloud-stream-test-support</artifactId>
    <scope>test</scope>
</dependency>
<dependency>
    <groupId>org.apache.kafka</groupId>
    <artifactId>kafka-streams-test-utils</artifactId>
    <scope>test</scope>
</dependency>

测试代码

import org.junit.jupiter.api.Test;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.boot.test.context.SpringBootTest;
import org.springframework.cloud.stream.binder.test.TestChannelBinderConfiguration;
import org.springframework.cloud.stream.binder.test.TestInputTopic;
import org.springframework.cloud.stream.binder.test.TestOutputTopic;
import org.springframework.context.annotation.Import;
import org.apache.kafka.streams.KeyValue;
import java.util.List;

@SpringBootTest
@Import(TestChannelBinderConfiguration.class) // 启用测试绑定器,无需真实消息中间件
public class WordCountFunctionTest {

    @Autowired
    private TestInputTopic inputTopic;

    @Autowired
    private TestOutputTopic outputTopic;

    @Autowired
    private Function<KStream<Object, String>, KStream<?, WordCount>> process;

    @Test
    void testWordCountProcessing() {
        // 发送测试文本
        inputTopic.send("Hello Hello World");

        // 接收输出结果
        List<KeyValue<?, WordCount>> results = outputTopic.receive(2, KeyValue.class);

        // 验证计数逻辑
        for (KeyValue<?, WordCount> result : results) {
            WordCount wc = result.value;
            if ("hello".equals(wc.getWord())) {
                assert wc.getCount() == 2;
            } else if ("world".equals(wc.getWord())) {
                assert wc.getCount() == 1;
            }
        }
    }
}

二、用TopologyTestDriver测试(贴近生产逻辑)

如果需要验证Kafka Streams拓扑的细节(比如窗口、状态存储),可以用TopologyTestDriver,之前失败大概率是拓扑构建或Serde配置问题,以下是正确示例:

测试代码

import org.junit.jupiter.api.AfterEach;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
import org.apache.kafka.common.serialization.Serdes;
import org.apache.kafka.streams.StreamsConfig;
import org.apache.kafka.streams.Topology;
import org.apache.kafka.streams.TopologyTestDriver;
import org.apache.kafka.streams.kstream.Produced;
import org.apache.kafka.streams.test.ConsumerRecordFactory;
import org.apache.kafka.streams.test.OutputVerifier;
import org.springframework.kafka.support.serializer.JsonSerde;
import java.util.Properties;
import java.util.Date;

public class WordCountTopologyTest {

    private TopologyTestDriver testDriver;
    private ConsumerRecordFactory<Object, String> recordFactory;

    @BeforeEach
    void setUp() {
        // 实例化你的Function Bean
        YourFunctionBean functionBean = new YourFunctionBean();
        Function<KStream<Object, String>, KStream<?, WordCount>> wordCountFunc = functionBean.process();

        // 构建Kafka Streams拓扑
        Topology topology = new Topology();
        topology.addSource("input-source", Serdes.Object().deserializer(), Serdes.String().deserializer(), "input-topic");
        KStream<Object, String> inputStream = topology.stream("input-source");
        KStream<?, WordCount> outputStream = wordCountFunc.apply(inputStream);
        outputStream.to("output-topic", Produced.with(Serdes.String(), new JsonSerde<>(WordCount.class)));

        // 配置测试驱动参数
        Properties props = new Properties();
        props.put(StreamsConfig.APPLICATION_ID_CONFIG, "test-word-count-app");
        props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "dummy:9092"); // 虚拟地址,无需真实集群
        props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.Object().getClass());
        props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass());

        // 初始化测试驱动
        testDriver = new TopologyTestDriver(topology, props);
        recordFactory = new ConsumerRecordFactory<>("input-topic", Serdes.Object().serializer(), Serdes.String().serializer());
    }

    @AfterEach
    void tearDown() {
        testDriver.close();
    }

    @Test
    void testWindowedWordCount() {
        // 发送测试消息
        testDriver.pipeInput(recordFactory.create("Hello Hello World"));

        // 验证hello的计数结果
        OutputVerifier.compareKeyValue(
                testDriver.readOutput("output-topic", Serdes.String().deserializer(), new JsonSerde<>(WordCount.class).deserializer()),
                "hello",
                new WordCount("hello", 2, new Date(0), new Date(5000))
        );
        // 验证world的计数结果
        OutputVerifier.compareKeyValue(
                testDriver.readOutput("output-topic", Serdes.String().deserializer(), new JsonSerde<>(WordCount.class).deserializer()),
                "world",
                new WordCount("world", 1, new Date(0), new Date(5000))
        );
    }
}

关键注意事项

  • WordCount类需要支持序列化,要么实现Serializable接口,要么使用JsonSerde(需引入Spring Kafka的Jackson依赖)。
  • 若要测试窗口关闭逻辑,需手动推进测试时间:testDriver.advanceWallClockTime(Duration.ofSeconds(5)),触发窗口计数输出。

内容的提问来源于stack exchange,提问作者Papi Abi

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 13:54:52