[FLINK-3024] Fix TimestampExtractor.getCurrentWatermark() Behaviour
Previously the internal currentWatermark would be updated even if the value returned from getCurrentWatermark was lower than the current watermark. This can lead to problems with chaining because the watermark is directly forwarded without going through the watermark logic that ensures correct behaviour (monotonically increasing). This adds a test that verifies that the timestamp extractor does not emit decreasing watermarks.
Showing
想要评论请 注册 或 登录