Spring Cloud Stream Kafka Streams 3.0 Binder API能否使用@KafkaStreamsStateStore注解?
首先直接给结论:在Spring Cloud Stream Kafka Streams 3.0的函数式API(比如你定义的BiFunction<KStream<String, Long>, KStream<String, Long>, KStream<String, Long>>这种Bean)上,@KafkaStreamsStateStore注解是不生效的。
为什么会报错?
你遇到的TopologyException: Invalid topology: StateStore my-store is not added yet异常,核心原因是@KafkaStreamsStateStore注解是为旧版的@StreamListener风格API设计的——框架会自动识别这个注解,并把声明的状态存储添加到@StreamListener方法对应的拓扑中。但函数式API是直接基于Kafka Streams原生的StreamsBuilder来构建拓扑,Spring Cloud Stream不会自动处理这个注解,所以你的状态存储根本没被添加到拓扑里,处理器自然找不到它。
正确的状态存储声明方式(函数式API)
针对函数式API,有几种可靠的方式来声明和使用状态存储:
1. 在DSL中直接创建并添加状态存储
这种方式最直接,适合只在当前流中使用的存储:
@Bean public BiFunction<KStream<String, Long>, KStream<String, Long>, KStream<String, Long>> myStream() { return (input1, input2) -> { // 定义窗口状态存储的Supplier WindowStoreSupplier<String, Long> storeSupplier = Stores.windowStoreBuilder( Stores.persistentWindowStore( "my-store", Duration.ofMillis(300000), // 窗口长度 Duration.ofMillis(300000), // 滑动间隔 false // 是否保留过期数据 ), Serdes.String(), Serdes.Long() ); // 将存储添加到当前拓扑的StreamsBuilder中 input1.processors().forEach(processorNode -> processorNode.context().getStreamsBuilder().addStateStore(storeSupplier) ); // 使用状态存储 return input1.transformValues(() -> new MyTransformer<>(), "my-store"); }; }
2. 通过StreamsBuilderFactoryBeanCustomizer添加全局状态存储
如果你的状态存储需要被多个函数式流共享,可以用这种方式全局注册:
@Bean public StreamsBuilderFactoryBeanCustomizer stateStoreCustomizer() { return factoryBean -> { WindowStoreSupplier<String, Long> storeSupplier = Stores.windowStoreBuilder( Stores.persistentWindowStore("my-store", Duration.ofMillis(300000), Duration.ofMillis(300000), false), Serdes.String(), Serdes.Long() ); // 把存储添加到全局的StreamsBuilder factoryBean.getStreamsBuilder().addStateStore(storeSupplier); }; } // 然后在函数式Bean中直接使用 @Bean public BiFunction<KStream<String, Long>, KStream<String, Long>, KStream<String, Long>> myStream() { return (input1, input2) -> input1.transformValues(() -> new MyTransformer<>(), "my-store"); }
3. 声明StateStore Bean并绑定到拓扑
你也可以把状态存储作为独立的Bean定义,然后在函数式方法中引用并添加到拓扑:
@Bean public WindowStore<String, Long> myStore() { return Stores.windowStoreBuilder( Stores.persistentWindowStore("my-store", Duration.ofMillis(300000), Duration.ofMillis(300000), false), Serdes.String(), Serdes.Long() ).build(); } @Bean public BiFunction<KStream<String, Long>, KStream<String, Long>, KStream<String, Long>> myStream(WindowStore<String, Long> myStore) { return (input1, input2) -> { // 将Bean形式的存储添加到拓扑 input1.processors().forEach(processorNode -> processorNode.context().getStreamsBuilder().addStateStore(myStore) ); return input1.transformValues(() -> new MyTransformer<>(), "my-store"); }; }
总结
函数式API更贴近Kafka Streams原生的拓扑构建逻辑,所以需要你手动管理状态存储的注册,而依赖@KafkaStreamsStateStore这种框架封装的注解是行不通的。上面的几种方式都能解决你遇到的拓扑异常问题,根据你的存储使用场景选择即可。
内容的提问来源于stack exchange,提问作者Sergey Shcherbakov

