使用partitionCustom和keyBy的Flink代码Watermark出现MIN_VALUE的原因
我编写了一段测试Flink Watermark的代码,期望通过partitionCustom将标识为A、B的数据路由到不同分区,从而让它们使用各自独立的Watermark。但运行代码后,发现多个key对应的currentWatermark显示为MIN_VALUE(-9223372036854775808),无法理解该现象的原因。
测试代码
import org.apache.flink.api.common.functions.Partitioner import org.apache.flink.configuration.Configuration import org.apache.flink.streaming.api.TimeCharacteristic import org.apache.flink.streaming.api.functions.{AssignerWithPunctuatedWatermarks, KeyedProcessFunction} import org.apache.flink.streaming.api.scala.{StreamExecutionEnvironment, _} import org.apache.flink.streaming.api.watermark.Watermark import org.apache.flink.util.Collector import java.text.SimpleDateFormat import java.util.Date object Test{ def to_milli(str: String) = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").parse(str).getTime def to_char(milli: Long) = { val date = if (milli <= 0) new Date(0) else new Date(milli) new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(date) } val data = Seq( ("A", "2020-08-30 10:50:11"), ("B", "2020-08-30 10:50:13"), ("B", "2020-08-30 10:50:04"), ("A", "2020-08-30 10:50:08") ) def main(args: Array[String]): Unit = { val env = StreamExecutionEnvironment.createLocalEnvironmentWithWebUI(new Configuration()) env.setStreamTimeCharacteristic(TimeCharacteristic.EventTime) env.setParallelism(4) env.fromCollection(data).setParallelism(1).partitionCustom(new Partitioner[String] { override def partition(key: String, numPartitions: Int): Int = key.hashCode() % numPartitions }, e => e._1).assignTimestampsAndWatermarks(new AssignerWithPunctuatedWatermarks[(String, String)]() { var maxSeen = Long.MinValue override def checkAndGetNextWatermark(lastElement: (String, String), extractedTimestamp: Long): Watermark = { val eventTime = to_milli(lastElement._2) if (eventTime > maxSeen) { maxSeen = eventTime } println(s"key: ${lastElement._1}, event time: ${to_char(eventTime)}, watermark: ${to_char(maxSeen - 4 * 1000)} ") new Watermark(maxSeen - 4 * 1000) } override def extractTimestamp(element: (String, String), previousElementTimestamp: Long): Long = to_milli(element._2) }).setParallelism(2).keyBy(_._1).process(new KeyedProcessFunction[String, (String, String), (String, String)] { override def processElement(value: (String, String), ctx: KeyedProcessFunction[String, (String, String), (String, String)]#Context, out: Collector[(String, String)]): Unit = { val watermark = ctx.timerService().currentWatermark() val key = ctx.getCurrentKey val s = if (watermark == Long.MinValue) s"MIN_VALUE: $watermark" else to_char(watermark) System.out.println(s"watermark:$key is: $s ") val eventTime = value._2 if (eventTime > to_char(watermark)) out.collect(value) } }).setParallelism(2).print() env.execute() } }
运行输出
key: B, event time: 2020-08-30 10:50:13, watermark: 2020-08-30 10:50:09 key: A, event time: 2020-08-30 10:50:11, watermark: 2020-08-30 10:50:07 key: B, event time: 2020-08-30 10:50:04, watermark: 2020-08-30 10:50:09 key: A, event time: 2020-08-30 10:50:08, watermark: 2020-08-30 10:50:07 watermark:A is: MIN_VALUE: -9223372036854775808 watermark:B is: MIN_VALUE: -9223372036854775808 watermark:A is: MIN_VALUE: -9223372036854775808 watermark:B is: 2020-08-30 10:50:09 3> (B,2020-08-30 10:50:13) 4> (A,2020-08-30 10:50:08) 3> (A,2020-08-30 10:50:11)
我无法理解为何会出现以下输出:
watermark:B is: MIN_VALUE: -9223372036854775808
watermark:A is: MIN_VALUE: -9223372036854775808
watermark:A is: MIN_VALUE: -9223372036854775808
原因分析
1. Watermark与元素的传递顺序
Flink中,元素会先被发送到下游算子,而Watermark的传递是异步滞后于元素的。当processElement第一次处理某个key的元素时,上游生成的Watermark还没来得及传递到下游算子,此时下游算子的currentWatermark会保持初始值MIN_VALUE。
2. Punctuated Watermark的特性
你使用的AssignerWithPunctuatedWatermarks会为每个元素生成一个Watermark,但Watermark的发送时机晚于元素本身。所以前几个被处理的元素,对应的下游算子还未收到任何Watermark,自然显示MIN_VALUE。从输出可以看到,最后一个B的元素处理时,Watermark已经到达,所以显示了正确的值。
3. 并行度与分区的影响
你的assignTimestampsAndWatermarks和后续的process算子都设置了并行度2,A、B的数据被分到不同的上游并行实例后,每个上游实例的Watermark会独立传递到下游。下游的每个并行实例在收到第一个Watermark之前,所有元素处理时的currentWatermark都是MIN_VALUE。
解决方案与优化建议
- 接受初始状态:MIN_VALUE是Watermark未到达时的正常初始值,如果你需要基于Watermark做计算,可以考虑在
onTimer方法中处理——该方法会在Watermark到达指定时间时触发,此时currentWatermark是正确的。 - 优化Watermark生成逻辑:
checkAndGetNextWatermark方法的extractedTimestamp参数已经是你提取的事件时间,不需要重复调用to_milli(lastElement._2),可以直接使用该参数减少重复解析:
override def checkAndGetNextWatermark(lastElement: (String, String), extractedTimestamp: Long): Watermark = { if (extractedTimestamp > maxSeen) { maxSeen = extractedTimestamp } println(s"key: ${lastElement._1}, event time: ${to_char(extractedTimestamp)}, watermark: ${to_char(maxSeen - 4 * 1000)} ") new Watermark(maxSeen - 4 * 1000) }
内容的提问来源于stack exchange,提问作者Tom

