如何测试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
相关产品推荐
相关产品推荐

