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

使用partitionCustom和keyBy的Flink代码Watermark出现MIN_VALUE的原因

问题:Flink中KeyedProcessFunction获取到的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。


解决方案与优化建议

  1. 接受初始状态:MIN_VALUE是Watermark未到达时的正常初始值,如果你需要基于Watermark做计算,可以考虑在onTimer方法中处理——该方法会在Watermark到达指定时间时触发,此时currentWatermark是正确的。
  2. 优化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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.30 21:18:49