You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

关于Flink中BoundedOutOfOrder Watermarks的原理与迟滞数据问询

事件时序说明

11:00   11:01      11:02      11:03     11:04    11:05      11:06   11:07   11:08      <=== 处理时间(Processing Time)
|   
|
|--------1---------2----------3---------4--------5---------------------------------      <== 输入事件(INPUT EVENTS)
|
|                                                      (水印)
|--------1------------2----------3-----------------5------W!-----------------4      <==== 事件时间(Event Time)
        11:01       11:02      11:03             11:05  11:05              11:04

使用的水印策略配置

WatermarkStrategy
    .<Tuple2<Long, String>>forBoundedOutOfOrderness(Duration.ofSeconds(5))
    .withTimestampAssigner((event, timestamp) -> event.f0);

问题1:Flink会每5秒向管道中注入一个Watermark事件?

回答:错误

forBoundedOutOfOrderness水印策略并非固定间隔注入水印,而是基于输入事件的事件时间戳动态生成:

  • 水印的时间戳计算公式为:当前观察到的最大事件时间戳 - 允许的乱序时长(此处为5秒)
  • 只有当新事件的事件时间戳推动了当前最大事件时间戳,且计算出的新水印值大于上一次发送的水印值时,才会向下游发送新水印
  • 若没有新事件输入,不会生成任何水印

问题2:event-4在Watermark关闭窗口后到达,是否会被视为迟滞数据,除非配置allowedLateness否则将被丢弃?

回答:是的

结合你的场景:

  1. 目标窗口为[11:00, 11:05)(左闭右开),当水印时间戳≥窗口结束时间(11:05)时,窗口会被触发计算并关闭
  2. event-4的事件时间为11:04,属于该窗口,但它在窗口关闭后才到达,属于迟到数据
  3. 默认情况下,未配置allowedLateness时,这类迟到数据会直接被丢弃;若配置了允许的迟到时长(如allowedLateness(Duration.ofSeconds(2))),则在窗口关闭后的指定时长内到达的迟到数据,仍会被纳入窗口并触发窗口重新计算

内容的提问来源于stack exchange,提问作者Mandar K

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 23:35:17