【问题标题】:Flink watermark is negativeFlink 水印为负数
【发布时间】:2019-12-03 21:17:39
【问题描述】:

我试图通过在函数内部实现AssignerWithPeriodicWatermarks 来为流分配时间戳和水印,它实现了:

override def getCurrentWatermark: Watermark = {
    // this guarantees that the watermark never goes backwards.
    val potentialWM = currentMaxTimestamp - maxOutOfOrderness
    if (potentialWM >= lastEmittedWatermark) lastEmittedWatermark = potentialWM

    new Watermark(lastEmittedWatermark)
  }

  override def extractTimestamp(element: T, previousElementTimestamp: Long): Long = {
    val timestamp = element.streamTime // something exists in the stream
    if (timestamp > currentMaxTimestamp) currentMaxTimestamp = timestamp
    timestamp
  }

但是,我仍然得到默认值-9223372036854775808的水印,当我尝试在这两个函数中添加打印时,我发现extractTimestamp中只打印了println,也就是说getCurrentWatermark的函数是从来没有打电话。

实现似乎是正确的,因为相同的代码能够在另一个脚本上运行(一些代码不是我写的)。

PS:我不是第一次遇到负水印了,我发现过了一段时间后水印会变正,但是我还是很困惑一开始发生了什么。

【问题讨论】:

    标签: scala apache-flink flink-streaming


    【解决方案1】:

    问题是您使用的是AssignerWithPeriodicWatermark,它不会在每个事件中生成水印,而是在间隔内生成水印。每当您使用AssingerWithPeriodicWatermark 时,您应该在执行环境上设置调用setTheAutowatermarkInterval。您提供的值将是调用 getCurrentWatermark 的时间间隔。 如果您没有设置它,那么该方法将永远不会被调用,因此您永远不会更改水印。 对于测试和学习,您可以考虑使用AssignerWithPunctuatedWatermark,因为这只会为每个事件发出水印。

    编辑: 正如这个答案下面提到的,autowatermarkInterval 的默认值实际上是 200 毫秒。此外,使用AssignerWithPunctuatedWatermark 并不意味着您需要为每个事件发出水印,但会为每个事件调用发出它们的方法。如果您不想发出水印,则该方法应简单地返回 null

    【讨论】:

    • 这不太正确。 autoWatermarkInterval 的默认值为 200 毫秒 - 因此当前水印将保持为负值,直到作业运行 200 毫秒,此时将调用 getCurrentWatermark。
    • 哦,那么对不起。我确信没有默认值。您能否分享确认这一点的文档的链接?非常感谢。
    • 此外,虽然使用 AssignerWithPunctuatedWatermark 您可以为每个事件发出水印,但不必以这种方式使用。只要您不想发出水印,就可以返回 null。
    • 是的,我的意思是实际上每个事件都会调用它;)也许,我并不清楚。我会更新帖子。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-07-12
    • 1970-01-01
    相关资源
    最近更新 更多