变更日志max.compaction.lag.ms设为Infinity的含义及是否会被压缩
Kafka Streams窗口聚合变更日志配置解答
你使用以下代码实现6小时窗口的聚合:
windowedBy(TimeWindows.of(Duration.ofHours(6))) .aggregate(aggregator, aggregator, Materialized.as("my-agg"))
对应的变更日志配置如下:
cleanup.policy = compact,delete retention.ms = 2d max.compaction.lag.ms = Infinity
1. max.compaction.lag.ms = Infinity的含义
max.compaction.lag.ms用于限制日志段从最后一次写入到允许被压缩的最大等待时长。设置为Infinity表示取消这个延迟限制:只要满足压缩触发条件(比如符合cleanup.policy的压缩规则、日志段达到指定大小等),对应的日志段就可以立即被压缩,无需等待固定时长。
对于窗口聚合的变更日志而言,这个配置能让已经过期的窗口数据尽快通过压缩被清理,避免额外的存储占用。
2. 该变更日志是否会被压缩?
会被压缩。原因如下:
- 配置中
cleanup.policy = compact,delete同时启用了**压缩(compact)和删除(delete)**两种清理策略:- 压缩策略会保留每个key的最新消息,剔除旧的冗余消息;
- 删除策略会在消息达到
retention.ms = 2d(2天)的留存期限后,清理过期消息。
- 再加上
max.compaction.lag.ms = Infinity取消了压缩的延迟限制,只要满足压缩触发条件,压缩操作就会执行。
内容的提问来源于stack exchange,提问作者Abe
相关产品推荐
相关产品推荐

