【问题标题】:Out of Memory Issue suspecting due to suppress feature由于抑制功能而怀疑内存不足问题
【发布时间】:2020-01-31 21:10:17
【问题描述】:

目前使用 DSL kafka streaming(2.1.1) 抑制特性来存储中间聚合结果。

应用程序获取连续流并负责日窗口聚合。

应用程序在总共 9 台服务器上运行,每台服务器都有足够的内存 (64 GB) 和磁盘空间 (500 GB),并且还明确分配 21 GB 内存用于仅聚合服务,尽管它因 OOM 问题而崩溃。

抑制主题定义:application-KTABLE-SUPPRESS-STATE-STORE-0000000004-changelog PartitionCount:100 ReplicationFactor:5 Configs:cleanup.policy=compact

我对 Suppress 运算符的理解如下

1) Suppress 没有 statestore,但它依赖于由更改日志主题支持的缓冲内存。

2) 当 Suppress 运算符发出最终结果时,每个论坛都会将墓碑发送到相应的更改日志主题,因此会被删除。

但此更改日志主题的其他手动清理策略只是紧凑,因此不太确定它是如何工作的。

应用程序在几天前在生产中推出并非常频繁地观察 OOM 问题。

下面是观察..

1) 磁盘空间正在增长,因为旧的窗口记录没有从 application-KTABLE-SUPPRESS-STATE-STORE-0000000004-changelog 中删除。

2) 一旦节点获得 OOM 并且在重新启动时缓存内存确实会很快被聚合服务填满 .. 18-20 GB,基于低容量,这是无法预料的。

3) 观察到在 changelog 主题(application-KTABLE-SUPPRESS-STATE-STORE-0000000004-changelog) 下,抑制功能默认没有保留期,即使窗口已经提前,它也会发出较旧的记录。当节点由于内存问题而崩溃并再次重新启动时观察到。想知道为什么即使窗口在一天后关闭,changelog 仍然保留旧的窗口记录?可能 clean.policy 只是紧凑的?

我使用的是 kafka 流 2.1.1 版本,发现一个在 kafka 流中注册的错误,该错误已在 2.2.1 及更高版本中修复。

OutOfMemoryError when restart my Kafka Streams appplication

Kafka Streams State Store Unrecoverable from Change Log Topic

为了重新解决问题,我计划在下面。

1) 使用删除内部主题的重置应用程序工具重置 kafka 流。 2)清理kafka statestore。 3)升级kafka流版本到2.4.0,希望稳定。

如果您对 OOM 问题有其他看法,请告诉我。

须藤代码:

KTable<Windowed<String>, JsonNode> aggregateTable =
                transactions
                        .groupByKey()
                        .windowedBy(
                                TimeWindows.of(Duration.ofSeconds(windowDuration)).grace(Duration.ofSeconds(windowGraceDuration)))
                        .aggregate(() -> new AggregationService().initialize(),
                                (key, transaction, previousStats) -> new AggregationService().buildAggregation(key, transaction, previousStats, runByUnit),
                                Materialized.<String, JsonNode, WindowStore<Bytes, byte[]>>as(statStoreName).withRetention(Duration.ofSeconds((windowDuration + windowGraceDuration + windowRetentionDuration)))
                                        .withKeySerde(Serdes.String())
                                        .withValueSerde(jsonSerde))
                        .suppress(Suppressed.untilWindowCloses(Suppressed.BufferConfig.unbounded()));

感谢您的帮助。

【问题讨论】:

标签: apache-kafka-streams


【解决方案1】:

我使用的是 kafka 流媒体 2.1.1 版本,发现了一个与抑制相关的错误,后来解决到 2.2.1 及更高版本中。

https://issues.apache.org/jira/plugins/servlet/mobile#issue/KAFKA-7895

请告诉我 2.3.o 是否可以解决这个问题?

观察到的类似问题主要在重新启动节点后多次抑制属于旧窗口的记录,考虑到大量(旧和当前窗口记录),内存不足问题与此有关

【讨论】:

    猜你喜欢
    • 2015-10-30
    • 2011-12-28
    • 2018-10-27
    • 1970-01-01
    • 2020-11-21
    • 1970-01-01
    • 2011-10-19
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多