You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 14:07:03