【发布时间】:2019-06-10 20:47:17
【问题描述】:
我在 Kafka 流应用程序中编写了这段代码:
KGroupedStream<String, foo> groupedStream = stream.groupByKey();
groupedStream.windowedBy(
SessionWindows.with(Duration.ofSeconds(3)).grace(Duration.ofSeconds(3)))
.aggregate(() -> {...})
.suppress(Suppressed.untilWindowCloses(unbounded()))
.toStream()...
应该(如果我理解正确的话)在窗口关闭后为每个键发出记录。 不知何故,行为如下:
流不会发出第一条记录,即使使用不同的 Key,也只会在第二条记录之后转发它,然后第二条记录仅在第三条之后发出,依此类推..
我已经尝试了多个带有“exactly_once”的 StreamConfig,并且无论有没有缓存,这种行为仍然存在。
提前感谢您的帮助!
【问题讨论】:
-
如果你希望你的数据按时间段而不是“会话”聚合,我想你需要使用
TimeWindows而不是SessionWindows。 -
这对我不起作用。有一个定时窗口,但在为同一个键添加新事件之前,它仍然无法完成对旧窗口的抑制效果。非常令人沮丧和违反直觉!
标签: apache-kafka apache-kafka-streams window-functions suppress