【问题标题】:How does kafka streams compute watermarks?kafka流如何计算水印?
【发布时间】:2019-02-08 01:47:16
【问题描述】:

Kafka Streams 是否在内部计算水印?是否可以(仅)在窗口完成时(即水印通过窗口结束时)观察窗口的结果?

【问题讨论】:

    标签: apache-kafka apache-kafka-streams


    【解决方案1】:

    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()))
    

    【讨论】:

    • 那么各个处理器使用事件时间戳(加上滞后)来决定窗口何时完成并发出一次结果?如果下游还有另一个窗口运算符,它将使用哪些时间戳来确定该窗口是否完整?
    • 是的,一次。并且“如果下游有另一个窗口操作符,它将使用哪些时间戳来确定该窗口是否完整?”——我不确定,我从未对窗口化流进行窗口化......将不得不考虑这一点。
    【解决方案2】:

    Kafka Streams 是否在内部计算水印?

    没有。 Kafka Streams 遵循不需要水印的持续更新 处理模型。您可以在网上找到更多详细信息:

    是否可以(仅)在窗口完成时(即水印通过窗口结束时)观察窗口的结果?

    您可以在任何时间点观察窗口的结果。通过例如KTable#toStream()#foreach()(即基于推送的方法)或通过Interactive Queries 订阅结果更改日志流,让您主动查询结果窗口(即基于拉取的方法)。

    正如@Dmitry 所说,对于基于推送的方法,如果您只对窗口的最终结果感兴趣,也可以使用suppress() 运算符。

    【讨论】:

    • 感谢您提供有趣的见解。连续表更新模型看起来很有趣。对于基于推送的模型,它使用简单的事件时间戳(加上延迟)来确定何时将窗口视为已完成,因此看起来像是完整水印模型的子集。
    • 对于基于推送的模型,您会收到关于窗口的每次更新的通知,即使尚未达到窗口结束时间。您可以使用suppress() 获得与水印类似的行为。
    • @MatthiasJ.Sax:如果将流重新键入中间流,如何计算消息的时间戳?例如:即使“源”流分区具有单调时间戳,中间流分区也可能会乱序接收事件
    • 不重新计算时间戳,但保留原始输入记录时间戳。没错,它可能会导致数据乱序,即使原始输入主题不包含乱序数据。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-08-07
    • 1970-01-01
    • 1970-01-01
    • 2014-01-06
    • 2015-01-12
    相关资源
    最近更新 更多