【问题标题】:Kafka Streams Rocksdb retention didn't remove old data with windowed functionKafka Streams Rocksdb 保留没有使用窗口函数删除旧数据
【发布时间】:2019-09-19 22:37:36
【问题描述】:

我正在运行具有窗口功能的 Kafka 流应用程序。但运行 24 小时后,本地磁盘使用量从 5G 增加到 20G 并不断增加。根据我的谷歌搜索,一旦我介绍了windowedBy,它应该会自动删除旧数据。

我的拓扑如下所示:

stream.selectKey(selectKey A)
.groupByKey(..)
.windowedBy(TimeWindows.of(Duration.ofMinutes(60)).grace(Duration.ZERO))
.reduce((value1,value2) -> value2)
.suppress.toStreams()
.selectKey(selectKey B).mapValues().filter()
.groupByKey().reduce.toStream().to()

我无法理解的是,从这个拓扑结构中,它将创建两个内部重新分区主题,如 repartition-03repartition-14 用于两个 groupBy 操作。从磁盘来看,所有正在执行repartition-03 任务的机器都具有很高的磁盘使用率,并且似乎永远不会删除旧数据,而正在运行repartition-14 任务的机器总是处于低磁盘使用率状态。

当我登录机器时,我发现这两台机器的路径不同,如下所示:

/tmp/kafka-streams/test-group/2_40/rocksdb/KSTREAM-REDUCE-STATE-STORE-0000000014
/tmp/kafka-streams/test-group/1_4/KSTREAM-REDUCE-STATE-STORE-0000000003/KSTREAM-REDUCE-STATE-STORE-0000000003.1568808000000

为什么他们有不同的路径? 2_40 用于repartition-14 任务,它的路径中有rocksdb,而另一个不包含rocksdb。同时,taks 1_4 保留了几个文件夹,如 KSTREAM-REDUCE-STATE-STORE-0000000003.1568808000000,但后缀不同。

虽然一旦我引入了windowedBy函数,rocksdb会在window过期时删除旧数据?为什么上述两个内部重新分区主题具有不同的路径和保留行为?

非常感谢任何帮助!谢谢!

【问题讨论】:

    标签: apache-kafka-streams rocksdb


    【解决方案1】:

    默认保留期为 24 小时。您可以通过

    减少它
    .reduce(..., Materialized.with(...).withRetention(...));
    

    【讨论】:

    • 嗨,Matthias,Materialized.as(null) 是什么?你能解释一下吗?它不允许我在这里将null 作为参数。谢谢!
    • 对不起。我的错。我认为这会起作用 - 更新了我的答案。试试with(...) 而不是as(...)
    • Materialized.withRetention 的更新,似乎 KStreams 提供了配置 WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG 来覆盖默认值 24HR。不知道是不是和Materialized.withRetention一样@
    • withRetention 设置存储保留时间,该保留时间也用于底层变更日志主题。配置参数WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG 仅适用于更改日志主题,并添加到保留时间。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2022-10-06
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多