【发布时间】:2019-02-08 01:47:16
【问题描述】:
Kafka Streams 是否在内部计算水印?是否可以(仅)在窗口完成时(即水印通过窗口结束时)观察窗口的结果?
【问题讨论】:
标签: apache-kafka apache-kafka-streams
Kafka Streams 是否在内部计算水印?是否可以(仅)在窗口完成时(即水印通过窗口结束时)观察窗口的结果?
【问题讨论】:
标签: apache-kafka apache-kafka-streams
Kafka Streams 内部不使用水印,但 2.1.0 中的一项新功能可让您观察窗口关闭时的结果。它被称为Suppressed,您可以在文档中了解它:Window Final Results:
KGroupedStream<UserId, Event> grouped = ...;
grouped
.windowedBy(TimeWindows.of(Duration.ofHours(1)).grace(ofMinutes(10)))
.count()
.suppress(Suppressed.untilWindowCloses(unbounded()))
【讨论】:
Kafka Streams 是否在内部计算水印?
没有。 Kafka Streams 遵循不需要水印的持续更新 处理模型。您可以在网上找到更多详细信息:
是否可以(仅)在窗口完成时(即水印通过窗口结束时)观察窗口的结果?
您可以在任何时间点观察窗口的结果。通过例如KTable#toStream()#foreach()(即基于推送的方法)或通过Interactive Queries 订阅结果更改日志流,让您主动查询结果窗口(即基于拉取的方法)。
正如@Dmitry 所说,对于基于推送的方法,如果您只对窗口的最终结果感兴趣,也可以使用suppress() 运算符。
【讨论】:
suppress() 获得与水印类似的行为。