【问题标题】:How to use interactive query within kafka process topology in spring-cloud-stream?如何在 spring-cloud-stream 中的 kafka 进程拓扑中使用交互式查询?
【发布时间】:2019-09-19 17:21:57
【问题描述】:

是否可以在 Spring Cloud Stream 中使用带有 @EnableBinding 注释的类或在带有 @StreamListener 的方法中使用交互式查询(InteractiveQueryService)?我尝试在提供的KStreamMusicSampleApplication 类和处理方法中实例化 ReadOnlyKeyValueStore,但它始终为空。

我的@StreamListener 方法正在侦听一堆 KTables 和 KStreams,并且在流程拓扑(例如过滤)期间,我必须检查 KStream 中的密钥是否已存在于特定 KTable 中。

我试图弄清楚如何扫描传入的 KTable 以检查密钥是否已经存在但没有运气。然后我遇到了 InteractiveQueryService,它的 get() 方法可用于检查 KTable 中的 state store materializedAs 中是否存在密钥。问题是我无法从流程拓扑(@EnableBinding 或@StreamListener)访问它。它只能从这些注释之外访问,例如 RestController。

有没有办法扫描传入的 KTable 以检查键或值是否存在?如果没有,我们可以在流程拓扑中访问 InteractiveQueryService 吗?

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams spring-cloud-stream


    【解决方案1】:

    Spring Cloud Stream 中的InteractiveQueryService 不能在StreamListener 的实际拓扑中使用。正如您所提到的,它应该在您的主要拓扑之外使用。但是,对于您描述的用例,您仍然可以使用主流程中的状态存储。例如,如果您有一个传入的KStream 和一个具体化为状态存储的KTable,那么您可以在KStream 上调用process 并以这种方式访问​​状态存储。这是一个粗略的代码来实现这一点。您需要将其转换为适合您的特定用例,但这是一个想法。

    ReadOnlyKeyValueStore<Object, String> store;
    
     input.process(() -> new Processor<Object, Product>() {
    
                    @Override
                    public void init(ProcessorContext processorContext) {
                        store = (ReadOnlyKeyValueStore) processorContext.getStateStore("my-store");
    
    
                    }
    
                    @Override
                    public void process(Object key, Object value) {
                        //find the key
                        store.get(key);
                    }
    
                    @Override
                    public void close() {
                        if (state != null) {
                            state.close();
                        }
                    }
                }, "my-store");
    

    【讨论】:

    • 感谢您的解决方案。我的下一个问题是来自另一个主题的消息在从另一个主题实现的状态存储完全填充(完成)之前到达。我的 application.yml 包含 spring.cloud.stream.kafka.streams.bindings.input.consumer.materializedAs: my-store。在第一次使用 application.yml 文件中的 materializedAs 启动应用程序时,如何确保我的流处理仅在从 kafka 主题完全填充状态存储后开始。谢谢
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2021-09-27
    • 2019-09-19
    • 2020-08-31
    • 2021-09-28
    • 2019-06-05
    • 2020-07-11
    • 1970-01-01
    相关资源
    最近更新 更多