如何忽略控制流的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
相关产品推荐
相关产品推荐

