【问题标题】:Kafka kstream-kstream joins with sliding window memory usage grows over the time till OOMKafka kstream-kstream 加入滑动窗口内存使用量随着时间的推移而增长,直到 OOM
【发布时间】:2019-03-12 10:38:08
【问题描述】:

我在使用 kstream 连接时遇到问题。我所做的是从一个主题中分离出 3 种不同类型的消息到新的流中。 然后对两个流进行一次内连接,创建另一个流,最后我对新流和最后一个剩余流进行最后一次左连接。

连接的窗口时间为 30 秒。

这样做是为了过滤掉一些被其他消息覆盖的消息。

我在 kubernetes 上运行此应用程序,并且 pod 的磁盘空间无限增长,直到 pod 崩溃。

我意识到这是因为连接将数据本地存储在 tmp/kafka-streams 目录中。

目录被称为: KSTREAM-JOINTHIS... KSTREAM-OUTEROTHER..

它存储来自 RocksDb 的 sst 文件,并且这些文件无限增长。

我的理解是,因为我使用 30 秒的窗口时间,所以这些应该在特定时间后被刷新,但不是。

我还将 WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG 更改为 10 分钟,看看是否会发生变化,但事实并非如此。

我需要了解如何更改。

【问题讨论】:

    标签: java apache-kafka apache-kafka-streams rocksdb rocksdb-java


    【解决方案1】:

    窗口大小不决定您的存储要求,而是连接的宽限期。为了处理乱序记录,数据的存储时间比窗口大小要长。在较新的版本中,需要始终通过JoinWindows. ofTimeDifferenceAndGrace(...) 指定宽限期。在旧版本中,您可以通过JoinWindows.of(...).grace(...) 设置宽限期——如果未设置,则默认为 24​​ 小时。

    配置WINDOW_STORE_CHANGE_LOG_ADDITIONAL_RETENTION_MS_CONFIG 配置数据在集群中存储多长时间。因此,您可能也想减少它,但它无助于减少客户端存储需求。

    【讨论】:

    • 谢谢你这工作好多了!我现在可以看到它是稳定的 :) 只是另一个快速的:对于窗口和保留时间应该多长,是否有任何最佳实践?现在我已经尝试过 JoinWindow.of(30sec).until(6min)
    • 这取决于应用程序,您需要使用自己的推理。它还取决于您的数据的最大“延迟”,即“乱序程度”。
    • 好的,谢谢。在进行性能测试时,大小仍然很高,但一段时间后会稳定下来。有没有办法减少商店的记忆?还是我必须扩展应用程序和分区?
    • 好吧,对于连接,所有记录都必须在保留期内进行缓冲——因此知道您的数据速率、保留时间和记录大小,您可以估计所需的内存。除了减少保留时间,我想不出别的。当然,水平缩放应该允许您减少每个实例的内存占用。
    • @MatthiasJ.Sax 这是哪里(“您可以通过将 Materialized.as(null).withRetention(...) 传入 join(...) 运算符来减少保留时间。” ) 在 API 中?我看不到。
    猜你喜欢
    • 2018-09-18
    • 2017-06-02
    • 2020-04-21
    • 2022-10-24
    • 2019-09-28
    • 2018-08-13
    • 2017-01-08
    • 2018-08-30
    • 2020-05-02
    相关资源
    最近更新 更多