【发布时间】: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