KafkaStreams单元测试抛出IllegalArgumentException:未知主题问题求助
Kafka Streams TopologyTestDriver 测试报错:Unknown topic: input-topic 问题解决
问题描述
我有一个基于Kafka Streams的应用,用KStream读取数据、通过header过滤后写入KTable,拓扑构建代码如下:
public Topology buildTopology() { KStream<String,String> inputStream = builder.stream("topicname"); KStream<String,String> filteredStream = inputStream.transformValues(KSExtension::new) .filter((key,value) -> value!=null); kTable = filteredStream.groupByKey() .reduce(((value1, value2) -> value2), Materialized.as("ktable")); KafkaStreams streams = new KafkaStreams(builder.build(), streamsConfiguration); streams.start(); Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); return builder.build(); }
尝试用TopologyTestDriver写单元测试时抛出java.lang.IllegalArgumentException: Unknown topic: input-topic异常,测试代码如下:
private TopologyTestDriver td; private TestInputTopic<String, String> inputTopic; private TestOutputTopic<String, String> outputTopic; private Topology topology; private Properties streamConfig; @BeforeEach void setUp() { streamConfig = new Properties(); streamConfig.setProperty(StreamsConfig.APPLICATION_ID_CONFIG, "AppId"); streamConfig.setProperty(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "foo:1234"); streamConfig.setProperty(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); streamConfig.setProperty(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); topology = new Topology(); td = new TopologyTestDriver(topology, streamConfig); inputTopic = td.createInputTopic("input-topic", Serdes.String().serializer(), Serdes.String().serializer()); outputTopic = td.createOutputTopic("output-topic", Serdes.String().deserializer(), Serdes.String().deserializer()); } @Test void buildTopology(){ inputTopic.pipeInput("key1", "value1"); topology = app.buildTopology(); }
异常栈信息:
DEBUG org.apache.kafka.streams.processor.internals.InternalTopologyBuilder - No source topics using pattern subscription found, initializing consumer's subscription collection. java.lang.IllegalArgumentException: Unknown topic: input-topic at org.apache.kafka.streams.TopologyTestDriver.pipeRecord(TopologyTestDriver.java:582) at org.apache.kafka.streams.TopologyTestDriver.pipeRecord(TopologyTestDriver.java:945) at org.apache.kafka.streams.TestInputTopic.pipeInput(TestInputTopic.java:115) at org.apache.kafka.streams.TestInputTopic.pipeInput(TestInputTopic.java:137) at testclassname.buildTopology()
错误原因
- Topic名称不匹配:测试代码里创建的输入topic是
input-topic,但实际拓扑中读取的是topicname,TestDriver找不到对应的订阅topic。 - 测试流程颠倒:先调用
pipeInput发送数据,再构建拓扑,此时TestDriver使用的是初始化时的空Topology实例,根本没有订阅任何topic。 - 生产代码耦合问题:
buildTopology方法里直接启动了KafkaStreams并添加shutdown hook,单元测试中执行该方法会启动实际的Streams线程,不仅干扰测试,还不符合单一职责原则。
解决方案
1. 重构生产代码,分离拓扑构建与启动逻辑
把拓扑构建和Streams实例启动分开,避免测试时启动实际服务:
// 只负责构建拓扑,不启动Streams public Topology buildTopology() { KStream<String,String> inputStream = builder.stream("topicname"); KStream<String,String> filteredStream = inputStream.transformValues(KSExtension::new) .filter((key,value) -> value!=null); kTable = filteredStream.groupByKey() .reduce(((value1, value2) -> value2), Materialized.as("ktable")); return builder.build(); } // 单独的启动方法,用于实际运行时调用 public void startStreams() { KafkaStreams streams = new KafkaStreams(buildTopology(), streamsConfiguration); streams.start(); Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); }
2. 修正测试代码,调整流程并匹配Topic名称
确保先构建正确的拓扑,再初始化TestDriver,同时保持topic名称一致:
private TopologyTestDriver td; private TestInputTopic<String, String> inputTopic; private YourAppClass app; // 替换为你的应用类名 private Properties streamConfig; @BeforeEach void setUp() { streamConfig = new Properties(); streamConfig.setProperty(StreamsConfig.APPLICATION_ID_CONFIG, "AppId"); streamConfig.setProperty(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "foo:1234"); streamConfig.setProperty(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); streamConfig.setProperty(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass().getName()); app = new YourAppClass(); // 初始化你的应用实例 Topology topology = app.buildTopology(); // 先构建正确的拓扑 td = new TopologyTestDriver(topology, streamConfig); // 输入topic名称要和拓扑里的一致:"topicname" inputTopic = td.createInputTopic("topicname", Serdes.String().serializer(), Serdes.String().serializer()); } @Test void testFilterAndReduce() { // 先发送测试数据 inputTopic.pipeInput("key1", "value1"); inputTopic.pipeInput("key1", "value2"); // 测试reduce逻辑,保留最新值 // 验证KTable的结果,通过状态存储查询 ReadOnlyKeyValueStore<String, String> kTableStore = td.getKeyValueStore("ktable"); assertEquals("value2", kTableStore.get("key1")); }
补充说明
- KTable是状态存储组件,默认不会输出到topic,所以测试时不需要创建
outputTopic,而是通过TestDriver获取对应的状态存储来验证数据。 - 如果你的业务逻辑确实需要将KTable输出到某个topic,需要在拓扑中添加
to()操作,比如kTable.toStream().to("output-topic"),此时才能通过TestOutputTopic接收数据。
内容的提问来源于stack exchange,提问作者perplexedDev
相关产品推荐
相关产品推荐

