Apache Flink接入广播流后Watermark不推进如何解决?
解决方案
可以实现,核心原理是Flink多输入算子的Watermark默认取所有输入通道当前Watermark的最小值,只要让广播流的Watermark永远不会拖慢非广播流即可,最常用的实现方案如下:
方案1:给广播流配置固定最大值Watermark(最推荐,实现最简单)
在S3流map操作之后、broadcast操作之前,手动分配一个永远返回最大值的Watermark策略,代码示例:
val enrichWithElements: BroadcastStream[EnrichWithElement] = flinkEnv.readFile(textInputFormat, "s3 path ...", FileProcessingMode.PROCESS_CONTINUOUSLY, 30) .map(s3Element => { EnrichWithElement(.....) }) // 给广播流配置最大Watermark,避免拖慢整体水位推进 .assignTimestampsAndWatermarks( WatermarkStrategy.forBoundedOutOfOrderness[EnrichWithElement](Duration.ZERO) .withTimestampAssigner(new SerializableTimestampAssigner[EnrichWithElement] { override def extractTimestamp(element: EnrichWithElement, recordTimestamp: Long): Long = { // 广播流不需要有效业务时间戳,赋值为最大值即可 Long.MaxValue } }) ) .broadcast(new MapStateDescriptor......)
这个方案的运行逻辑是:两个输入流的Watermark分别为非广播流正常推进的业务水位、广播流固定为Long.MaxValue,算子取两者最小值作为整体水位,最终水位推进完全和非广播流同步,完全符合你的需求。而且适配你当前S3文件低频次更新、30秒轮询的场景,不会对补全逻辑产生任何影响。
方案2:自定义多输入水位合并逻辑
如果使用Flink 1.17及以上版本,也可以通过自定义WatermarkGeneratorSupplier直接忽略广播流的水位输入,不过该方案实现复杂度更高,普通场景优先选择方案1即可。
内容的提问来源于stack exchange,提问作者Gaurav Kumar
相关产品推荐
相关产品推荐

