资讯动态

Apache Kafka Streams 应用测试指南:TopologyTestDriver 与 MockProcessorContext 实战解析

发布时间:2026/9/10 14:48:11 来源:尧图企业网站定制
Apache Kafka Streams 应用测试指南TopologyTestDriver 与 MockProcessorContext 实战解析【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka导读本文以 Kafka 官方文档 Testing a Streams Application 为核心系统讲解如何为 Kafka Streams 应用编写单元测试与集成测试。你将掌握三个层次的测试能力使用kafka-streams-test-utils测试工具包中的TopologyTestDriver驱动整条 Topology 处理数据并校验输出使用TestInputTopic/TestOutputTopic便捷地注入与读取记录以及使用MockProcessorContext对单个 Processor 进行隔离式单元测试。文章同时结合本仓库streams/test-utils模块的源码实现说明测试工具的内部运行机制帮助你写出可复现、可维护、结果可精确断言的高质量 Streams 测试。一、导入测试工具kafka-streams-test-utils测试 Kafka Streams 应用的第一步是将官方提供的测试工具test-utils作为普通依赖引入测试代码。该构件与主库同源发布与kafka-clients、kafka-streams属于同一套版本体系本仓库中对应的模块源码位于 streams/test-utils当前仓库构建版本为4.5.0-SNAPSHOT见 gradle.properties。Maven 工程在pom.xml中按如下方式声明scope 为test仅用于测试期dependency groupIdorg.apache.kafka/groupId artifactIdkafka-streams-test-utils/artifactId version4.5.0-SNAPSHOT/version scopetest/scope /dependency该构件核心提供了两大测试设施TopologyTestDriver进程内模拟 Kafka Streams 运行时驱动整条拓扑处理数据是端到端式拓扑测试的主力MockProcessorContext一个可捕获转发数据、提交状态、定时器punctuator的模拟上下文用于对单个 Processor 做纯单元测试。两者分别对应下文第二、第三节。引入后需要哪些 Serde则参考 Datatypes 指南 中内置 Serdes 的用法。二、使用 TopologyTestDriver 测试 Streams 拓扑2.1 测试驱动器的设计原理TopologyTestDriver的核心思想是在单个 JVM 进程内、不使用真实 Kafka 集群模拟 Kafka Streams 的运行时行为——持续从输入主题拉取记录、沿拓扑遍历处理、把结果写入输出主题并将所有状态存储内嵌在驱动器中。从源码看它的内部维护了一个MockTime实例mockWallClockTime来模拟墙钟时间见 TopologyTestDriver.java一个outputRecordsByTopic队列映射按主题暂存所有输出记录见 TopologyTestDriver.java。由于不依赖真实 Broker测试速度极快且天然具备确定性——这正是对拓扑计算结果做精确断言的前提。2.2 组装待测拓扑与创建驱动器待测拓扑既可以用 Processor API 手工组装也可以用 DSLStreamsBuilder构建// Processor API 方式 Topology topology new Topology(); topology.addSource(sourceProcessor, input-topic); topology.addProcessor(processor, ..., sourceProcessor); topology.addSink(sinkProcessor, output-topic, processor); // 或者使用 DSL 方式 StreamsBuilder builder new StreamsBuilder(); builder.stream(input-topic).filter(...).to(output-topic); Topology topology builder.build(); // 创建测试驱动器 TopologyTestDriver testDriver new TopologyTestDriverBuilder(topology).build();需要说明的是仓库当前版本的TopologyTestDriver构造函数如TopologyTestDriver(topology)、TopologyTestDriver(topology, config)已标记为Deprecated(since 4.4)见 TopologyTestDriver.java官方推荐的入口是TopologyTestDriverBuilder。它支持链式传入配置与初始墙钟时间Properties props new Properties(); props.setProperty(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.setProperty(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass().getName()); TopologyTestDriver testDriver new TopologyTestDriverBuilder(topology) .withConfig(props) .build();即使不显式配置驱动器也会自动补齐必要项从源码看构造时会缺省注入dummy-bootstrap-host:0作为 bootstrap.servers、随机生成dummy-topology-test-driver-app-id-*作为 application.id并关闭窗口聚合的发射节流见 TopologyTestDriver.java。因此测试中你只需关心与业务相关的配置。2.3 用 TestInputTopic 注入数据TopologyTestDriver.createInputTopic()接收主题名与键/值对应的Serializer返回TestInputTopic对象TestInputTopicString, Long inputTopic testDriver.createInputTopic( input-topic, stringSerde.serializer(), longSerde.serializer()); inputTopic.pipeInput(key, 42L);TestInputTopic提供多种注入方式方法定义见 TestInputTopic.java方法作用pipeInput(K key, V value)写入单条记录时间戳为当前模拟时间pipeInput(V value)写入仅含 value 的记录pipeInput(TestRecordK,V record)写入一条完整记录可精确控制时间戳、头信息pipeKeyValueList(ListKeyValueK,V list)批量写入多条 KeyValue 记录pipeRecordList(ListTestRecordK,V list)批量写入多条 TestRecordadvanceTime(Duration advance)推进注入记录的时间戳默认情况下输入记录使用驱动器的当前模拟时间作为时间戳且不自动递增如需模拟时间流逝产生的新记录可使用带startTimestamp与autoAdvance参数的重载版createInputTopic见 TopologyTestDriver.java。2.4 用 TestOutputTopic 校验输出TestOutputTopic在初始化时配置主题名与对应的Deserializer并提供一系列读取辅助方法见 TestOutputTopic.javaTestOutputTopicString, Long outputTopic testDriver.createOutputTopic( output-topic, stringSerde.deserializer(), longSerde.deserializer()); assertEquals(new KeyValue(key, 42L), outputTopic.readKeyValue());方法作用readKeyValue()读取下一条记录为KeyValue忽略时间戳便于标准断言readRecord()读取下一条记录为完整TestRecord含时间戳、头readKeyValuesToList()把队列中全部记录读为ListKeyValuereadValuesToList()只读 value 列表readKeyValuesToMap()转为MapKey, ValuegetQueueSize()返回输出队列中剩余记录数isEmpty()判断输出队列是否为空如果只关心键和值而不关心记录时间戳使用readKeyValue()配合标准 JUnit 断言即可需要验证时间戳、头信息等元数据时再改用readRecord()。2.5 时间控制事件时间与墙钟时间的 Punctuation 测试TopologyTestDriver完整支持 punctuator周期性回调的测试事件时间STREAM_TIMEpunctuation基于已处理记录的时间戳自动触发无需额外操作墙钟时间WALL_CLOCK_TIMEpunctuation由驱动器内部模拟的墙钟时间驱动通过advanceWallClockTime()显式推进来触发testDriver.advanceWallClockTime(Duration.ofSeconds(20));从源码实现看advanceWallClockTime会推进内部MockTime并立即执行系统时间类型的maybePunctuateSystemTime()随后提交并完成剩余可处理的工作见 TopologyTestDriver.java。这让你可以精确控制经过 60 秒后发生了什么而无需真实等待。2.6 通过驱动器访问状态存储TopologyTestDriver暴露多种状态存储访问入口便于测试前后检查或预置数据KeyValueStoreString, Long store testDriver.getKeyValueStore(store-name);主要方法见 TopologyTestDriver.java方法用途getKeyValueStore(name)获取 KeyValue 类型状态存储L1044getWindowStore(name)获取窗口化状态存储L1172getSessionStore(name)获取 Session 状态存储L1273getStateStore(name)按名获取任意 StateStoreL952getAllStateStores()获取全部状态存储L918测试前访问存储可预置初始值注意需要关闭状态存储的日志记录 changelog 才允许预置见下文示例测试后访问存储可验证状态更新是否符合预期。2.7 资源释放测试驱动器内部持有任务、状态目录等资源务必在测试结束时关闭testDriver.close();最佳实践是放在 JUnit 的After/tearDown()方法中保证每个测试用例独立、无资源泄漏。三、完整示例基于状态存储与 Punctuation 的最大值聚合以下示例演示测试驱动器的完整用法拓扑按 key 用 KeyValueStore 计算最大值处理期间不产生输出仅当事件时间或墙钟时间触发 punctuation 时才把当前存储内容转发到下游。private TopologyTestDriver testDriver; private TestInputTopicString, Long inputTopic; private TestOutputTopicString, Long outputTopic; private KeyValueStoreString, Long store; private SerdeString stringSerde new Serdes.StringSerde(); private SerdeLong longSerde new Serdes.LongSerde(); Before public void setup() { Topology topology new Topology(); topology.addSource(sourceProcessor, input-topic); topology.addProcessor(aggregator, new CustomMaxAggregatorSupplier(), sourceProcessor); topology.addStateStore( Stores.keyValueStoreBuilder( Stores.inMemoryKeyValueStore(aggStore), Serdes.String(), Serdes.Long()).withLoggingDisabled(), // 必须禁用 changelog 日志否则无法预置存储 aggregator); topology.addSink(sinkProcessor, result-topic, aggregator); // 配置测试驱动器 Properties props new Properties(); props.setProperty(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); props.setProperty(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass().getName()); testDriver new TopologyTestDriverBuilder(topology).withConfig(props).build(); // 配置测试主题 inputTopic testDriver.createInputTopic(input-topic, stringSerde.serializer(), longSerde.serializer()); outputTopic testDriver.createOutputTopic(result-topic, stringSerde.deserializer(), longSerde.deserializer()); // 预置状态存储 store testDriver.getKeyValueStore(aggStore); store.put(a, 21L); } After public void tearDown() { testDriver.close(); } Test public void shouldFlushStoreForFirstInput() { inputTopic.pipeInput(a, 1L); assertEquals(new KeyValue(a, 21L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); } Test public void shouldNotUpdateStoreForSmallerValue() { inputTopic.pipeInput(a, 1L); assertEquals(21L, store.get(a)); assertEquals(new KeyValue(a, 21L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); } Test public void shouldUpdateStoreForLargerValue() { inputTopic.pipeInput(a, 42L); assertEquals(42L, store.get(a)); assertEquals(new KeyValue(a, 42L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); } Test public void shouldUpdateStoreForNewKey() { inputTopic.pipeInput(b, 21L); assertEquals(21L, store.get(b)); assertEquals(new KeyValue(a, 21L), outputTopic.readKeyValue()); assertEquals(new KeyValue(b, 21L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); } Test public void shouldPunctuateIfEventTimeAdvances() { final Instant recordTime Instant.now(); inputTopic.pipeInput(a, 1L, recordTime); assertEquals(new KeyValue(a, 21L), outputTopic.readKeyValue()); inputTopic.pipeInput(a, 1L, recordTime); assertTrue(outputTopic.isEmpty()); // 事件时间前进 10 秒触发 STREAM_TIME 类型 punctuation inputTopic.pipeInput(a, 1L, recordTime.plusSeconds(10L)); assertEquals(new KeyValue(a, 21L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); } Test public void shouldPunctuateIfWallClockTimeAdvances() { // 墙钟时间前进 60 秒触发 WALL_CLOCK_TIME 类型 punctuation testDriver.advanceWallClockTime(Duration.ofSeconds(60)); assertEquals(new KeyValue(a, 21L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); }配套的处理器实现如下——它在init()中同时注册了两种类型的 punctuatorflushStore()将存储中的全部键值对转发到下游public class CustomMaxAggregatorSupplier implements ProcessorSupplierString, Long { Override public ProcessorString, Long get() { return new CustomMaxAggregator(); } } public class CustomMaxAggregator implements ProcessorString, Long { ProcessorContext context; private KeyValueStoreString, Long store; SuppressWarnings(unchecked) Override public void init(ProcessorContext context) { this.context context; context.schedule(Duration.ofSeconds(60), PunctuationType.WALL_CLOCK_TIME, time - flushStore()); context.schedule(Duration.ofSeconds(10), PunctuationType.STREAM_TIME, time - flushStore()); store (KeyValueStoreString, Long) context.getStateStore(aggStore); } Override public void process(String key, Long value) { Long oldValue store.get(key); if (oldValue null || value oldValue) { store.put(key, value); } } private void flushStore() { KeyValueIteratorString, Long it store.all(); while (it.hasNext()) { KeyValueString, Long next it.next(); context.forward(next.key, next.value); } } Override public void close() {} }这个示例清晰地展示了测试驱动器的几大关键能力预置存储withLoggingDisabled()后写入初始值、精确注入pipeInput、输出队列消费与判空readKeyValueisEmpty以及双时间轴 punctuation 的确定性触发。仓库中TopologyTestDriverTest、TopologyTestDriverAtLeastOnceTest、TopologyTestDriverEosTest见 streams/test-utils/src/test/java/org/apache/kafka/streams进一步覆盖了 at-least-once 与 EOS 处理语义下的驱动行为是深入研读的绝佳样例。四、使用 MockProcessorContext 单元测试 ProcessorTopologyTestDriver适合整条拓扑的测试当使用 Processor API 编写单个Processor时往往只需要对它做轻量级单元测试。问题在于Processor通过ProcessorContext把结果转发出去而非直接返回因此单元测试需要一个能捕获转发数据的模拟上下文——这就是MockProcessorContext实现见 MockProcessorContext.java的职责。4.1 构造与初始化实例化待测 Processor用 mock context 初始化final Processor processorUnderTest ...; final MockProcessorContextString, Long context new MockProcessorContext(); processorUnderTest.init(context);如果需要向 Processor 传递配置或设置默认 Serde可以带配置构造final Properties props new Properties(); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.Long().getClass()); props.put(some.other.config, some config value); final MockProcessorContextString, Long context new MockProcessorContext(props);MockProcessorContext还提供带TaskId与状态目录的构造重载见 MockProcessorContext.java需要验证 task 相关行为时可用。4.2 捕获数据的断言Processor 每次context.forward(...)都会被 mock 捕获可逐一断言processorUnderTest.process(key, value); final IteratorCapturedForward? extends String, ? extends Long forwarded context.forwarded().iterator(); assertEquals(forwarded.next().record(), new Record(..., ...)); assertFalse(forwarded.hasNext()); // 重置转发记录便于构造更长的交互场景 context.resetForwards(); assertEquals(context.forwarded().size(), 0);若 Processor 转发到特定子处理器可按键名查询捕获结果final ListCapturedForward? extends String, ? extends Long captures context.forwarded(childProcessorName);4.3 提交commit行为的验证mock 还会捕获 Processor 是否调用了context.commit()assertTrue(context.committed()); // 提交状态也可重置 context.resetCommit(); assertFalse(context.committed());4.4 设置记录元数据如果 Processor 逻辑依赖记录的 topic、partition、offset 等元数据可以显式设置context.setRecordMetadata(topicName, /*partition*/ 0, /*offset*/ 0L);设置之后 context 会持续返回相同值直到再次设置新值见 MockProcessorContext.java。4.5 注册状态存储若 Processor或其 punctuator是有状态的mock context 支持注册状态存储。推荐使用合适类型的内存存储KeyValue、Windowed 或 Session因为 mock context 不管理 changelog、状态目录等设施final KeyValueStoreString, Integer store Stores.keyValueStoreBuilder( Stores.inMemoryKeyValueStore(myStore), Serdes.String(), Serdes.Integer() ) .withLoggingDisabled() // MockProcessorContext 不支持 changelog .build(); store.init(context, store); context.register(store, /*已废弃参数*/ false, /*mock 中不使用的参数*/ null);4.6 验证 PunctuatorProcessor 可以调度 punctuator 执行周期任务mock context不会自动执行它们但会捕获以便逐一单测final MockProcessorContext.CapturedPunctuator capturedPunctuator context.scheduledPunctuators().get(0); final long interval capturedPunctuator.getIntervalMs(); final PunctuationType type capturedPunctuator.getType(); final boolean cancelled capturedPunctuator.cancelled(); final Punctuator punctuator capturedPunctuator.getPunctuator(); punctuator.punctuate(/*timestamp*/ 0L);CapturedPunctuator暴露了间隔毫秒数、punctuation 类型、是否已取消以及实际的Punctuator对象你可以手动触发它并断言副作用。如果需要测试 punctuator 的自动触发由记录时间或墙钟时间驱动官方建议将 Processor 放入一个简单拓扑改用 TopologyTestDriver 来完成——两种工具各司其职配合使用即可覆盖全部测试场景。五、测试策略小结场景推荐工具关键能力验证整条拓扑DSL 或 Processor API 组装的数据流结果TopologyTestDriverTestInputTopic/TestOutputTopic进程内模拟运行时、确定性的时间控制、内嵌状态存储验证单条记录的时间敏感行为、事件时间/墙钟时间 punctuationTopologyTestDriverpipeInput带时间戳、advanceWallClockTime双时间轴精确推进对单个Processor做纯单元测试MockProcessorContext捕获 forward/commit、记录元数据、注册内存存储、捕获 punctuator验证自动触发的 punctuator 全流程简单拓扑 TopologyTestDriver由驱动器按模拟时间自动调度无论使用哪种工具都请牢记三条实践准则测试后关闭驱动器释放资源After中testDriver.close()预置状态存储前确认已禁用 changelog 日志withLoggingDisabled()用TestRecord/readRecord()校验时间戳与头信息用KeyValue/readKeyValue()只校验业务键值。基于kafka-streams-test-utils你可以完全不启动真实 Kafka 集群就能把 Streams 应用的核心逻辑测试得又快又稳。相关文档Processor API 指南 · Streams 数据类型与序列化 · Streams 配置 · 开发指南首页【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价