【发布时间】:2018-06-27 02:29:43
【问题描述】:
我们使用以下代码在会话窗口中进行聚合:
.windowedBy(SessionWindows.with(...))
.aggregate(..., ..., ...)
为我们自动创建的状态存储由带有cleanup.policy=compact 的更改日志主题支持。
在重新部署拓扑时,我们发现恢复状态存储所用的时间比预期的要长得多(10 多分钟)。解释似乎是即使会话已关闭,它仍然存在于更改日志主题中。
我们注意到会话窗口的默认维护持续时间为 1 天,但即使在超过不活动 + 维护持续时间之后,看起来消息也不会从变更日志主题中删除。
a) 我们是否需要手动删除“旧”(根据我们的定义)消息来控制变更日志主题的大小? (这可能是 [1] 所暗示的情况。)
b) 是否有可能以某种方式使用cleanup.policy=compact,delete 创建更改日志主题,这是否有意义?
[1] 会话存储似乎是由 Kafka Stream 的 UnwindowedChangelogTopicConfig(而不是 WindowedChangelogTopicConfig)在内部创建的,这可能会使来自 Kafka Streams - reducing the memory footprint for large state stores 的评论相关:“对于非窗口存储,没有保留策略. 基础主题仅被压缩。因此,如果您知道不再需要记录,则需要通过墓碑将其删除。但实现起来有点棘手...... – Matthias J. Sax Jun 2017 年 20 月 27 日 22:07"
【问题讨论】:
标签: apache-kafka apache-kafka-streams