Apache Kafka Streams 测试指南:TopologyTestDriver 与 MockProcessorContext 从入门到实战
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
本文基于 Apache Kafka 仓库中 docs/streams/developer-guide/testing.md 编写,系统讲解 Kafka Streams 应用的两级测试体系:面向完整Topology的TopologyTestDriver端到端驱动测试,以及面向单个Processor的MockProcessorContext单元测试。读完本文,你将掌握如何引入kafka-streams-test-utils测试依赖、通过输入/输出测试主题管道化数据、控制事件时间与墙钟时间、查询与预填充状态存储、验证 punctuator 调度行为,从而在不启动真实 Kafka 集群的情况下快速、确定性地验证流处理拓扑的正确性。
为什么需要专门的测试工具
Kafka Streams 应用本质上是一个由 Processor API 或 DSL 组装而成的有向无环拓扑(Topology):它持续从输入主题拉取记录、沿拓扑逐级处理、最终写入输出主题,并依赖底层的状态存储与 punctuator 定时器。直接对着真实集群写测试会遇到三座大山:环境搭建成本高、测试缓慢且不稳定、时间与分区等执行细节难以精确控制。
为此,Kafka 在streams/test-utils模块中提供了kafka-streams-test-utils测试构件,内置两类核心工具:
TopologyTestDriver:在本地 JVM 内模拟 Kafka Streams 运行时,把记录手动管道进拓扑并捕获输出,同时内嵌可控制的时钟;MockProcessorContext:为 Processor API 编写单元测试而设计的ProcessorContext替身,捕获 forward、commit、schedule 等一切行为。
两者的源码都位于仓库的 streams/test-utils/src/main/java/org/apache/kafka/streams 目录下,公开 API 均标注为@InterfaceAudience.Public,是官方支持的测试基础设施。
引入测试依赖
测试工具作为一个独立构件发布,与kafka-streams本身版本号保持一致,只需作为test 作用域依赖加入构建即可。使用 Maven 时,在pom.xml中加入:
<dependency> <groupId>org.apache.kafka</groupId> <artifactId>kafka-streams-test-utils</artifactId> <version>4.5.0-SNAPSHOT</version> <scope>test</scope> </dependency>版本号务必与你的
kafka-streams依赖版本一致。文档示例中写的是4.3.0;当前仓库gradle.properties中声明的版本为4.5.0-SNAPSHOT,你在使用时请替换为实际依赖的 Kafka 版本。
使用 Gradle 时等价写法为:
dependencies { testImplementation 'org.apache.kafka:kafka-streams-test-utils:4.5.0-SNAPSHOT' }使用 TopologyTestDriver 测试完整拓扑
TopologyTestDriver是整个测试工具包的核心。它的原理是模拟 Kafka Streams 的库运行时:测试驱动内部持续从输入主题"拉取"记录,沿拓扑遍历处理,并把结果记录捕获到输出主题。你不需要 broker、不需要网络,所有逻辑在单线程、单 JVM 内确定性地执行。
构造测试驱动
TopologyTestDriver接受一个Topology,这个拓扑既可以是用 Processor API 手工拼装的,也可以是用 DSL 经StreamsBuilder构建的:
// Processor API Topology topology = new Topology(); topology.addSource("sourceProcessor", "input-topic"); topology.addProcessor("processor", ..., "sourceProcessor"); topology.addSink("sinkProcessor", "output-topic", "processor"); // or // using DSL StreamsBuilder builder = new StreamsBuilder(); builder.stream("input-topic").filter(...).to("output-topic"); Topology topology = builder.build(); // create test driver TopologyTestDriver testDriver = new TopologyTestDriverBuilder(topology).build();注意当前仓库已引入流式构造器TopologyTestDriverBuilder.java(TopologyTestDriver的旧有构造函数仍可用但已标记为 deprecated),它提供三个可链式调用的方法:
| 方法 | 作用 | 默认值 |
|---|---|---|
withConfig(Properties config) | 传入驱动配置(如默认 Serde、application.id等) | 空Properties |
withInitialWallClockTime(Instant) | 设定驱动内部模拟墙钟时间的初始值 | 当前系统时间 |
build() | 构造并返回就绪的TopologyTestDriver | — |
例如传入配置并固定初始时钟:
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) .withInitialWallClockTime(Instant.ofEpochMilli(0)) .build();通过 TestInputTopic 管道输入
有了测试驱动后,用createInputTopic为每个输入主题创建TestInputTopic,创建时必须给出主题名以及对应的 key/value 序列化器:
TestInputTopic<String, Long> inputTopic = testDriver.createInputTopic("input-topic", stringSerde.serializer(), longSerde.serializer()); inputTopic.pipeInput("key", 42L);从 TestInputTopic.java 的源码可以看到,pipeInput提供了一组丰富的重载,覆盖"仅 value""key + value""value + 时间戳""key + value + 时间戳(long毫秒或Instant)"以及整条TestRecord等各种形态;此外还支持批量管道化:
| 方法 | 说明 |
|---|---|
pipeInput(V value) | 只发送 value,key 为 null |
pipeInput(K key, V value) | 发送 key-value 对 |
pipeInput(K key, V value, long timestampMs) | 携带毫秒时间戳发送 |
pipeInput(K key, V value, Instant timestamp) | 携带Instant时间戳发送 |
pipeInput(TestRecord<K, V> record) | 发送完整TestRecord(含 headers) |
pipeKeyValueList(List<KeyValue<K, V>>) | 批量发送 KeyValue 列表 |
pipeValueList(List<V> values) | 批量发送 value 列表 |
pipeRecordList(List<TestRecord<K, V>>) | 批量发送 TestRecord 列表 |
advanceTime(Duration advance) | 推进该输入主题内部跟踪的事件时间 |
批量方法的时间戳会根据构造时的起始时间自动递增;未显式指定时间戳的记录则使用输入主题当前跟踪的事件时间(源码中getTimestampAndAdvance()负责"取当前时间并自动推进")。
通过 TestOutputTopic 验证输出
与之对称,TestOutputTopic在初始化时配置主题与反序列化器,然后按需读取结果。如果你只关心 key 和 value 而不关心时间戳,可以像下面这样直接对读出的KeyValue做标准断言:
TestOutputTopic<String, Long> outputTopic = testDriver.createOutputTopic("output-topic", stringSerde.deserializer(), longSerde.deserializer()); assertEquals(new KeyValue<>("key", 42L), outputTopic.readKeyValue());TestOutputTopic.java 提供了从"最轻量"到"最完整"的多档读取粒度:
| 方法 | 返回 | 适用场景 |
|---|---|---|
readValue() | V | 只关心 value |
readKeyValue() | KeyValue<K, V> | 关心 key 与 value,忽略时间戳与 headers |
readRecord() | TestRecord<K, V> | 需要 key、value、时间戳、headers 全部字段 |
readRecordsToList() | List<TestRecord<K, V>> | 把结果视为流,读取全部记录 |
readKeyValuesToList() | List<KeyValue<K, V>> | 同上但丢弃时间戳 |
readKeyValuesToMap() | Map<K, V> | 把结果视为表,只关心每个 key 的最后一次更新(若 key 为 null 会抛出IllegalStateException) |
isEmpty() | boolean | 校验输出是否已读尽 |
控制时间:事件时间与墙钟时间
流处理中的时间语义(事件时间、墙钟时间)很难在真实环境中精确复现,TopologyTestDriver对此提供了决定性控制:
- 事件时间(event-time)punctuation:基于已处理记录的时间戳自动触发。只要你在
pipeInput时传入递增的时间戳,事件时间 punctuator 就会被自动调度执行; - 墙钟时间(wall-clock-time)punctuation:驱动在内部模拟了墙钟时间,你可以手动拨快它来触发对应类型的 punctuator:
testDriver.advanceWallClockTime(Duration.ofSeconds(20));仓库测试 TopologyTestDriverTest.java 中有大量此类用法,例如连续多次advanceWallClockTime(Duration.ofMillis(...))来逐段验证墙钟时间 punctuator 的触发边界。advanceWallClockTime与驱动构造时通过withInitialWallClockTime设定的初始值配合,可以构造任意精度的"时间剧本"。
访问与预填充状态存储
TopologyTestDriver允许在测试前、后访问嵌入的状态存储。测试前访问存储,可以预置初始值(例如绕过 changelog 恢复,直接给定 store 的种子数据);数据处理后访问存储,可以校验预期的更新结果:
KeyValueStore store = testDriver.getKeyValueStore("store-name");需要说明的是,getKeyValueStore只能获取 key-value 类存储;窗口(windowed)、会话(session)等类型的存储需要调用对应的getWindowStore/getSessionStore(测试驱动在类型不匹配时会抛出明确的类型错误提示,见测试文件中的断言信息)。
正确关闭测试驱动
TopologyTestDriver内部持有定时任务、状态存储与各类资源,测试结束时务必调用close()以确保资源被正确释放,避免跨测试用例的状态泄漏:
testDriver.close();实践中通常放在@After/tearDown方法里统一执行。
完整示例:基于状态存储的每键最大值聚合
下面这个示例完整演示了TopologyTestDriver及其辅助类的典型用法(内容取自 testing.md 原文档)。它构建的拓扑用 key-value store 计算每个 key 的最大值:处理过程中不产生任何输出,只更新状态存储;输出只由事件时间与墙钟时间两种 punctuator 触发,把 store 中全部数据 flush 到下游。
private TopologyTestDriver testDriver; private TestInputTopic<String, Long> inputTopic; private TestOutputTopic<String, Long> outputTopic; private KeyValueStore<String, Long> store; private Serde<String> stringSerde = new Serdes.StringSerde(); private Serde<Long> 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(), // need to disable logging to allow store pre-populating "aggregator"); topology.addSink("sinkProcessor", "result-topic", "aggregator"); // setup test driver 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(); // setup test topics inputTopic = testDriver.createInputTopic("input-topic", stringSerde.serializer(), longSerde.serializer()); outputTopic = testDriver.createOutputTopic("result-topic", stringSerde.deserializer(), longSerde.deserializer()); // pre-populate store 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 shouldNotUpdateStoreForLargerValue() { 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 shouldPunctuateIfEvenTimeAdvances() { 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()); inputTopic.pipeInput("a", 1L, recordTime.plusSeconds(10L)); assertEquals(new KeyValue<>("a", 21L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); } @Test public void shouldPunctuateIfWallClockTimeAdvances() { testDriver.advanceWallClockTime(Duration.ofSeconds(60)); assertEquals(new KeyValue<>("a", 21L), outputTopic.readKeyValue()); assertTrue(outputTopic.isEmpty()); } public class CustomMaxAggregatorSupplier implements ProcessorSupplier<String, Long> { @Override public Processor<String, Long> get() { return new CustomMaxAggregator(); } } public class CustomMaxAggregator implements Processor<String, Long> { ProcessorContext context; private KeyValueStore<String, 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 = (KeyValueStore<String, 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() { KeyValueIterator<String, Long> it = store.all(); while (it.hasNext()) { KeyValue<String, Long> next = it.next(); context.forward(next.key, next.value); } } @Override public void close() {} }这个例子蕴含了几个值得注意的工程细节:
withLoggingDisabled()的意义:示例注释明确说明,需要禁用 store 的 changelog 日志记录,才能在测试前置阶段直接store.put(...)预填充数据;- 两种 punctuator 的独立验证:事件时间推进通过给
pipeInput传入递增的Instant实现;墙钟时间推进则通过advanceWallClockTime(Duration.ofSeconds(60))实现——它们分别对应PunctuationType.STREAM_TIME与PunctuationType.WALL_CLOCK_TIME的调度; - 断言的确定性:由于驱动是单线程、确定性的,
outputTopic.readKeyValue()的返回顺序完全可预测,配合isEmpty()可以精确断言"该有的都有、不该有的没有"。
单元测试 Processor:MockProcessorContext
当你编写了自定义Processor(参见 processor-api 开发者指南)时,往往希望以最小的粒度对它做单元测试。问题在于:Processor并不返回结果,而是把结果转发给ProcessorContext。因此单测需要一个能捕获转发数据的 mock 上下文——这正是test-utils提供的MockProcessorContext的职责(源码见 processor/api/MockProcessorContext.java)。
从源码结构看,MockProcessorContext内部用capturedForwards列表记录所有forward(...)调用、用punctuators列表捕获schedule(...)注册的定时器、用stateStores映射管理注册的状态存储,并维护一个committed布尔标志——它只"记录目睹的一切",不采取任何自动动作(例如不会自动触发已调度的 punctuator)。
构造与初始化
实例化被测 processor,并用 mock 上下文初始化它:
final Processor processorUnderTest = ...; final MockProcessorContext<String, 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 MockProcessorContext<String, Long> context = new MockProcessorContext<>(props);源码中还有带TaskId与File stateDir的完整构造函数,可供需要访问taskId()/stateDir()的场景使用;无参构造会自动补齐application.id与bootstrap.servers两个占位配置。
断言捕获的转发数据
mock 会捕获 processor 转发的所有数据,你可以对其做断言:
processorUnderTest.process("key", "value"); final Iterator<CapturedForward<? extends String, ? extends Long>> forwarded = context.forwarded().iterator(); assertEquals(forwarded.next().record(), new Record<>(..., ...)); assertFalse(forwarded.hasNext()); // you can reset forwards to clear the captured data. This may be helpful in constructing longer scenarios. context.resetForwards(); assertEquals(context.forwarded().size(), 0);若 processor 转发给特定的子节点,可以按子节点名查询捕获数据:
final List<CapturedForward<? extends String, ? extends Long>> captures = context.forwarded("childProcessorName");mock 还会记录 processor 是否调用了commit():
assertTrue(context.committed()); // commit captures can also be reset. context.resetCommit(); assertFalse(context.committed());设置记录元数据
当 processor 的逻辑依赖记录元数据(主题、分区、偏移量)时,可以手工在上下文上设置:
context.setRecordMetadata("topicName", /*partition*/ 0, /*offset*/ 0L);设置一次之后,上下文会持续返回相同值,直到你再次设置新值。源码还提供了setRecordTimestamp、setCurrentSystemTimeMs、setCurrentStreamTimeMs等配套方法,用于精确控制处理器可见的时间视角。
注册状态存储
如果 processor 或它的 punctuator 有状态,mock 上下文允许注册状态存储。官方建议使用对应类型(KeyValue、Windowed 或 Session)的简单内存存储即可,因为 mock 上下文不会管理 changelog、状态目录等运行时设施:
final KeyValueStore<String, Integer> store = Stores.keyValueStoreBuilder( Stores.inMemoryKeyValueStore("myStore"), Serdes.String(), Serdes.Integer() ) .withLoggingDisabled() // Changelog is not supported by MockProcessorContext. .build(); store.init(context, store); context.register(store, /*deprecated parameter*/ false, /*parameter unused in mock*/ null);withLoggingDisabled()在这里不是可选项而是必须项:注释明确指出MockProcessorContext不支持 changelog。注册后,processor 内通过context.getStateStore("myStore")即可拿到同一实例;若需要把某个 store 加入 mock 而不走register,还可以使用addStateStore(...)方法。
验证 Punctuator
Processor 可以调度 punctuator 执行周期任务。mock 上下文不会自动执行punctuator,但会捕获它们,从而让你可以对定时器本身做单元测试:
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完整地保留了调度信息(起始时间getStartTime、间隔getInterval、类型getType、回调getPunctuator),并支持cancel()与cancelled()来验证取消语义。
如果测试需要自动触发已调度的 punctuator,而不是手动调用,官方建议:把你的 processor 放进一个最小的 source-processor-sink 拓扑,改用上文介绍的TopologyTestDriver来驱动——这正是两种工具的分工边界:MockProcessorContext做"显微镜级"的行为捕获,TopologyTestDriver做"运行时级"的自动调度。
两种工具的选择策略
| 维度 | TopologyTestDriver | MockProcessorContext |
|---|---|---|
| 测试对象 | 完整Topology(DSL 或 Processor API 组装) | 单个Processor |
| 运行方式 | 模拟库运行时,自动拉取/遍历/调度 | 纯手工驱动,仅捕获行为 |
| 时间控制 | 事件时间随记录推进、墙钟时间可手动拨快 | 需通过setCurrentSystemTimeMs等手工设定 |
| Punctuator | 自动触发 | 捕获后手动执行 |
| 状态存储 | 完整支持(含类型检查与预填充) | 仅支持内存 store 注册 |
| 适用场景 | 集成级行为验证、跨节点数据流校验 | 快速单测、边界分支覆盖 |
选择的原则很简单:能覆盖的行为越靠上(完整拓扑)越真实,越靠下(单处理器)越轻快。业务逻辑的边界条件建议先用MockProcessorContext快速覆盖,再把关键链路用TopologyTestDriver做整体验证。
深入阅读
- 本文主体对应官方指南 docs/streams/developer-guide/testing.md,配套的 Processor API 开发指南 解释了自定义处理器的编写规范;
- 测试工具全部源码位于 streams/test-utils/src/main/java/org/apache/kafka/streams,其中 TopologyTestDriverBuilder.java、TestInputTopic.java、TestOutputTopic.java 与 MockProcessorContext.java 是最值得精读的四个类;
- 想知道"官方自己怎么写测试",可直接阅读仓库测试 TopologyTestDriverTest.java,其中包含了输入/输出主题、时间推进、各类状态存储访问的完整断言示例;
- 若需了解 Kafka Streams 的整体架构与概念,可继续阅读 docs/streams/_index.md 与 开发者指南目录。
【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考