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

如何忽略控制流的Watermark?解决多输入算子水印停滞问题

Flink多输入算子Watermark停滞问题与解决方案

在某些场景下,拥有多个输入的算子的整体(最小)Watermark可能会因速度极慢、近乎静态的控制流(或多个此类控制流)而停滞。

目前已有多种被使用或推荐的解决方案:

  • 定义返回MAX_WATERMARK的Watermark生成器
  • 或者类似地,使用withIdleness并设置极短的空闲周期
  • 使用Operator API包装用户自定义函数,示例如下:
public class SingleWatermarkKeyedCoProcessFunction<K, IN1, IN2, OUT>
extends KeyedCoProcessOperator<K, IN1, IN2, OUT> {
  public SingleWatermarkKeyedCoProcessFunction(KeyedCoProcessFunction<K, IN1, IN2, OUT> flatMapper) {
    super(flatMapper);
  }

  @Override
  public void processWatermark1(Watermark mark) throws Exception {
    super.processWatermark(mark);
  }

  @Override
  public void processWatermark2(Watermark mark) { }
}

在我看来,这更像是流本身的属性,因此我更倾向于方案1或2——不过我个人在一些项目中使用过方案3。

或许可以新增一种Watermark生成器来覆盖此类场景?比如针对方案1的实现:

static <A> DocumentWatermarkStrategy<A> ignoreWatermarks() {
    return (ctx) -> new IgnoreWatermarksGenerator<>();
}

其中:

public class IgnoreWatermarksGenerator<T> implements WatermarkGenerator<T> {
    @Override
    public void onEvent(T event, long eventTimestamp, WatermarkOutput output) { /* Do nothing */ }

    @Override
    public void onPeriodicEmit(WatermarkOutput output) {
        output.emitWatermark(Watermark.MAX_WATERMARK);
    }
}

不知各位对此有何看法?或许文档应更明确地介绍这一常见问题与陷阱?

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 10:03:20