【发布时间】:2017-06-07 19:02:01
【问题描述】:
我们正在使用 Kafka 流的 SessionWindows 来聚合相关事件的到达。除了聚合,我们还使用until() API 指定窗口的保留时间。
直播信息:
会话窗口(不活动时间)为 1 分钟,传递给until() 的保留时间为 2 分钟。
我们正在使用自定义的TimestampExtractor 来映射事件的时间。
示例:
事件:e1;活动时间:上午 10:00:00;到达时间:下午 2 点(同一天)
事件:e2;活动时间:上午 10:00:30;到达时间下午 2:10(同一天)
第二个事件的到达时间是 e1 到达后的 10 分钟,超过了保留时间 + 不活动时间。但较旧的事件 e1 仍然是聚合的一部分,尽管保留时间为 2 分钟。
问题:
1) kafka 流如何使用until() API 清理状态存储?由于指定为参数的保留值是“窗口将保持多长时间的下限”。究竟何时清除窗口?
2) 是否有后台线程定期清理状态存储?如果是,那么有没有办法确定清除窗口的实际时间。
3) 任何会在保留时间后清除窗口数据的流配置。
【问题讨论】:
标签: apache-kafka apache-kafka-streams