【发布时间】:2021-01-04 09:57:21
【问题描述】:
我正在使用 kafka 处理器 API 进行一些自定义计算。由于一些复杂的处理,DSL 并不是最合适的。流代码如下所示。
KeyValueBytesStoreSupplier storeSupplier = Stores.persistentKeyValueStore("storeName");
StoreBuilder<KeyValueStore<String, StoreObject>> storeBuilder = Stores.keyValueStoreBuilder(storeSupplier,
Serdes.String(), storeObjectSerde);
topology.addSource("SourceReadername", stringDeserializer, sourceSerde.deserializer(), "sourceTopic")
.addProcessor("processor", () -> new CustomProcessor("store"), FillReadername)
.addStateStore(storeBuilder, "processor") // define store for processor
.addSink("sinkName", "outputTopic", stringSerializer, resultSerde.serializer(),
Fill_PROCESSOR);
我需要根据来自单独主题的事件从状态存储中清除一些项目。我无法找到正确的方法来使用处理器 API 或其他方式加入另一个流来监听另一个主题中的事件,以便能够触发 CustomProcessor 类中的清理代码。 有没有办法我们可以在处理器 API 的另一个主题中获取事件?或者可能将 DSL 与处理器 API 混合,以便能够将两者结合起来,并将任何主题中的事件发送到 Process 方法,以便在清理主题中收到事件时运行清理代码?
谢谢
【问题讨论】:
标签: apache-kafka apache-kafka-streams