Spring Boot集成Kafka Streams启动报错求助
Spring Boot集成Kafka Streams启动报错:Topology无流线程/全局线程
问题描述
我在开发Spring Boot应用,首次尝试集成Kafka Streams,配置类如下:
@Autowired NumberDetectionService numberDetectionService; private static final String INPUT_TOPIC = "nn_input"; private static final String OUTPUT_TOPIC = "nn_output"; private static final String STORE_NAME = "ImageDtoStore"; @Bean public Topology getTopology() { final StreamsBuilder builder = new StreamsBuilder(); StoreBuilder<KeyValueStore<String, ImageDto>> storeBuilder = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore(STORE_NAME), Serdes.String(), ImageDtoSerde.serde() ); builder.addStateStore(storeBuilder); builder.stream(INPUT_TOPIC, Consumed.with(Serdes.String(), ImageDtoSerde.serde())).transform(() -> new ImageDtoTransformer(numberDetectionService), STORE_NAME) .to(OUTPUT_TOPIC, Produced.with(Serdes.String(), ImageDtoSerde.serde())); System.out.println(builder.build().describe()); return builder.build(); } @Bean public KafkaStreams kafkaStreams() { //final StreamsBuilder builder = streamsBuilder(); KafkaStreams streams = new KafkaStreams(getTopology() , streamsConfiguration()); streams.cleanUp(); streams.start(); // attach shutdown handler to catch control-c Runtime.getRuntime().addShutdownHook(new Thread(streams::close)); return streams; } private Properties streamsConfiguration() { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "neural-streams"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 2); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); return props; } @Bean public NewTopic nnTopicIn() { return TopicBuilder.name(INPUT_TOPIC) .partitions(1).replicas(1).config(TopicConfig.RETENTION_MS_CONFIG,"3600000").build(); } @Bean public NewTopic nnTopicOut() { return TopicBuilder.name(OUTPUT_TOPIC).partitions(1).replicas(1).config(TopicConfig.RETENTION_MS_CONFIG,"3600000").build(); }
启动应用时抛出以下错误:
org.springframework.context.ApplicationContextException: Failed to start bean 'defaultKafkaStreamsBuilder'; nested exception is org.springframework.kafka.KafkaException: Could not start stream: ; nested exception is org.apache.kafka.streams.errors.TopologyException: Invalid topology: Topology has no stream threads and no global threads, must subscribe to at least one source topic or global table. at org.springframework.context.support.DefaultLifecycleProcessor.doStart(DefaultLifecycleProcessor.java:181) ~[spring-context-5.3.14.jar:5.3.14] at org.springframework.context.support.DefaultLifecycleProcessor.access$200(DefaultLifecycleProcessor.java:54) ~[spring-context-5.3.14.jar:5.3.14] at org.springframework.context.support.DefaultLifecycleProcessor$LifecycleGroup.start(DefaultLifecycleProcessor.java:356) ~[spring-context-5.3.14.jar:5.3.14] at java.base/java.lang.Iterable.forEach(Iterable.java:75) ~[na:na] at org.springframework.context.support.DefaultLifecycleProcessor.startBeans(DefaultLifecycleProcessor.java:155) ~[spring-context-5.3.14.jar:5.3.14] at org.springframework.context.support.DefaultLifecycleProcessor.onRefresh(DefaultLifecycleProcessor.java:123) ~[spring-context-5.3.14.jar:5.3.14] at org.springframework.context.support.AbstractApplicationContext.finishRefresh(AbstractApplicationContext.java:935) ~[spring-context-5.3.14.jar:5.3.14] at org.springframework.context.support.AbstractApplicationContext.refresh(AbstractApplicationContext.java:586) ~[spring-context-5.3.14.jar:5.3.14] at org.springframework.boot.web.servlet.context.ServletWebServerApplicationContext.refresh(ServletWebServerApplicationContext.java:145) ~[spring-boot-2.5.8.jar:2.5.8] at org.springframework.boot.SpringApplication.refresh(SpringApplication.java:765) ~[spring-boot-2.5.8.jar:2.5.8] at org.springframework.boot.SpringApplication.refreshContext(SpringApplication.java:445) ~[spring-boot-2.5.8.jar:2.5.8] at org.springframework.boot.SpringApplication.run(SpringApplication.java:338) ~[spring-boot-2.5.8.jar:2.5.8] at org.springframework.boot.SpringApplication.run(SpringApplication.java:1354) ~[spring-boot-2.5.8.jar:2.5.8] at org.springframework.boot.SpringApplication.run(SpringApplication.java:1343) ~[spring-boot-2.5.8.jar:2.5.8] at MultilayerPerceptron.MultiLayerPerceptronApplication.main(MultiLayerPerceptronApplication.java:20) ~[classes/:na] at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke0(Native Method) ~[na:na] at java.base/jdk.internal.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62) ~[na:na] at java.base/jdk.internal.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43) ~[na:na] at java.base/java.lang.reflect.Method.invoke(Method.java:566) ~[na:na] at org.springframework.boot.devtools.restart.RestartLauncher.run(RestartLauncher.java:49) ~[spring-boot-devtools-2.5.8.jar:2.5.8] Caused by: org.springframework.kafka.KafkaException: Could not start stream: ; nested exception is org.apache.kafka.streams.errors.TopologyException: Invalid topology: Topology has no stream threads and no global threads, must subscribe to at least one source topic or global table. at org.springframework.kafka.config.StreamsBuilderFactoryBean.start(StreamsBuilderFactoryBean.java:359) ~[spring-kafka-2.8.0.jar:2.8.0] at org.springframework.context.support.DefaultLifecycleProcessor.doStart(DefaultLifecycleProcessor.java:178) ~[spring-context-5.3.14.jar:5.3.14] ... 19 common frames omitted Caused by: org.apache.kafka.streams.errors.TopologyException: Invalid topology: Topology has no stream threads and no global threads, must subscribe to at least one source topic or global table. at org.apache.kafka.streams.KafkaStreams.getNumStreamThreads(KafkaStreams.java:948) ~[kafka-streams-2.8.0.jar:na] at org.apache.kafka.streams.KafkaStreams.<init>(KafkaStreams.java:854) ~[kafka-streams-2.8.0.jar:na] at org.apache.kafka.streams.KafkaStreams.<init>(KafkaStreams.java:711) ~[kafka-streams-2.8.0.jar:na] at org.springframework.kafka.config.StreamsBuilderFactoryBean.start(StreamsBuilderFactoryBean.java:337) ~[spring-kafka-2.8.0.jar:2.8.0] ... 20 common frames omitted
请问是缺少额外配置项吗?
解决方案
错误原因
Spring Kafka会自动初始化defaultKafkaStreamsBuilder这个Bean(对应StreamsBuilderFactoryBean),但你手动创建了独立的KafkaStreams实例并启动,同时Spring容器仍在尝试启动默认的StreamsBuilderFactoryBean——而这个默认Builder没有定义任何拓扑逻辑,因此触发"无流线程/全局线程"的报错。
修改方案
移除手动创建的KafkaStreams Bean,改用Spring Kafka的自动生命周期管理,让Spring负责KafkaStreams的启动、关闭:
方式一:定义Topology Bean
@Autowired NumberDetectionService numberDetectionService; private static final String INPUT_TOPIC = "nn_input"; private static final String OUTPUT_TOPIC = "nn_output"; private static final String STORE_NAME = "ImageDtoStore"; @Bean public Topology topology() { final StreamsBuilder builder = new StreamsBuilder(); StoreBuilder<KeyValueStore<String, ImageDto>> storeBuilder = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore(STORE_NAME), Serdes.String(), ImageDtoSerde.serde() ); builder.addStateStore(storeBuilder); builder.stream(INPUT_TOPIC, Consumed.with(Serdes.String(), ImageDtoSerde.serde())) .transform(() -> new ImageDtoTransformer(numberDetectionService), STORE_NAME) .to(OUTPUT_TOPIC, Produced.with(Serdes.String(), ImageDtoSerde.serde())); return builder.build(); } @Bean public StreamsConfig streamsConfig() { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "neural-streams"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 2); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); return new StreamsConfig(props); } @Bean public NewTopic nnTopicIn() { return TopicBuilder.name(INPUT_TOPIC) .partitions(1).replicas(1).config(TopicConfig.RETENTION_MS_CONFIG,"3600000").build(); } @Bean public NewTopic nnTopicOut() { return TopicBuilder.name(OUTPUT_TOPIC).partitions(1).replicas(1).config(TopicConfig.RETENTION_MS_CONFIG,"3600000").build(); }
方式二:直接注入StreamsBuilder构建拓扑
@Autowired NumberDetectionService numberDetectionService; private static final String INPUT_TOPIC = "nn_input"; private static final String OUTPUT_TOPIC = "nn_output"; private static final String STORE_NAME = "ImageDtoStore"; @Bean public void buildTopology(StreamsBuilder builder) { StoreBuilder<KeyValueStore<String, ImageDto>> storeBuilder = Stores.keyValueStoreBuilder( Stores.persistentKeyValueStore(STORE_NAME), Serdes.String(), ImageDtoSerde.serde() ); builder.addStateStore(storeBuilder); builder.stream(INPUT_TOPIC, Consumed.with(Serdes.String(), ImageDtoSerde.serde())) .transform(() -> new ImageDtoTransformer(numberDetectionService), STORE_NAME) .to(OUTPUT_TOPIC, Produced.with(Serdes.String(), ImageDtoSerde.serde())); } @Bean public StreamsConfig streamsConfig() { Properties props = new Properties(); props.put(StreamsConfig.APPLICATION_ID_CONFIG, "neural-streams"); props.put(StreamsConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(StreamsConfig.NUM_STREAM_THREADS_CONFIG, 2); props.put(StreamsConfig.DEFAULT_KEY_SERDE_CLASS_CONFIG, Serdes.String().getClass()); props.put(StreamsConfig.DEFAULT_VALUE_SERDE_CLASS_CONFIG, Serdes.String().getClass()); return new StreamsConfig(props); } @Bean public NewTopic nnTopicIn() { return TopicBuilder.name(INPUT_TOPIC) .partitions(1).replicas(1).config(TopicConfig.RETENTION_MS_CONFIG,"3600000").build(); } @Bean public NewTopic nnTopicOut() { return TopicBuilder.name(OUTPUT_TOPIC).partitions(1).replicas(1).config(TopicConfig.RETENTION_MS_CONFIG,"3600000").build(); }
额外检查
- 确保
ImageDtoSerde正确实现了Serde接口,序列化/反序列化逻辑无异常,这会影响拓扑的有效性。 - 无需手动添加关闭钩子,Spring会自动处理KafkaStreams的生命周期,包括优雅关闭。
内容的提问来源于stack exchange,提问作者Peter_G
相关产品推荐
相关产品推荐

