【问题标题】:Kafka Streams Processor API clear state storeKafka Streams Processor API 清除状态存储
【发布时间】: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


    【解决方案1】:

    您只需要添加另一个输入主题 (add:Source) 并添加处理器来转换来自该主题的消息并基于它们从状态存储中删除人员。请注意,这些主题应该使用相同的键(因为分区)。

    【讨论】:

    • 谢谢,我们遇到的问题之一是第二个主题不能有相同的键。这些将有一些父级键,基于它们我可以迭代状态存储条目并清理其中一些。有没有一种方法可以让我们以某种方式从处理器的所有实例中获取来自第二个主题的所有消息,类似于全局 K-Table。
    • @HarshKumar,也许您应该基于 second topic 创建一些中间主题,并使用与第一个相同的键进行更改
    • 感谢 Bartosz,我在源应用程序中添加了发送相同密钥的逻辑。我必须在持久性存储中添加从父键到子键的查找,以找到与第一个主题中的键相同的子键。
    猜你喜欢
    • 1970-01-01
    • 2019-07-03
    • 1970-01-01
    • 2018-10-24
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-07-14
    • 2023-03-23
    相关资源
    最近更新 更多