【问题标题】:Kafka Streams Persistent Store cleanupKafka Streams 持久存储清理
【发布时间】:2021-07-30 18:35:55
【问题描述】:

是否需要进行一些明确的清理以防止每个持久性存储的大小增长过多?我目前正在使用它来计算 DSL API 中的聚合。

【问题讨论】:

  • 不确定你的意思。你能详细说明一下吗?
  • 基本上,如果不检查持久状态存储的大小会产生什么影响,因为数据只是添加到其中但从未删除。是否有推荐的方法来清理此类国有商店。我可以考虑设置一个调度程序来删除旧记录(比如超过 1 个月),但想知道是否有更好的方法。
  • 您也许可以使用窗口商店(取决于您的需要)。或者你可以schedule()Punctuation手动删除旧记录。

标签: apache-kafka apache-kafka-streams


【解决方案1】:

我们遇到了类似的问题,我们只是安排了一个清洁商店的工作 在我们的处理器/变压器中。 只需实现您的 isDataOld(nextValue) 就可以了。

@Override
public void init(ProcessorContext context) {
this.kvStore = (KeyValueStore<Key, Value>) this.context.getStateStore("KV_STORE_NAME");
this.context.schedule(60000, PunctuationType.STREAM_TIME, (timestamp) -> {
    KeyValueIterator<Key, Value> iterator = kvStore.all();
    while (iterator.hasNext()){
    KeyValue<Key,Value> nextValue = iterator.next();
    if isDataOld(nextValue)
       kvStore.delete(nextValue.key);
    }

});
}

【讨论】:

  • 是的,我正在考虑做类似的事情。我现在只使用 DSL,但似乎只有处理器 API 有调度。我想我可以创建一个处理器节点来清除我的状态存储。
  • 您可以将 DSL API 与处理器一结合起来,在 DSL api 中使用方法:transformprocess
  • 是的,拥有一个处理器节点不会对您的 DSL 代码造成任何损害 :)
  • 在遍历state store的时候调用delete会有问题吗?
  • 注意:我在 Slack 聊天中询问了 KStreams,它似乎在迭代时删除对于状态存储来说是可以的。
【解决方案2】:

怎么样:

final var store =  Stores.persistentSessionStore(MATERIAL_STORE_NAME, Duration.ofDays(10));
    return Stores.sessionStoreBuilder(store,
            Serdes.Long(),
            keySpecificAvroSerde
    );

【讨论】:

    猜你喜欢
    • 2023-03-23
    • 2021-01-04
    • 2020-06-02
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-07-03
    • 1970-01-01
    相关资源
    最近更新 更多