diff --git a/docs/content/docs/dev/datastream/event-time/generating_watermarks.md b/docs/content/docs/dev/datastream/event-time/generating_watermarks.md index b6d766d5b9992f..b80c428df107ed 100644 --- a/docs/content/docs/dev/datastream/event-time/generating_watermarks.md +++ b/docs/content/docs/dev/datastream/event-time/generating_watermarks.md @@ -206,6 +206,14 @@ WatermarkStrategy \ {{< /tab >}} {{< /tabs >}} +Note that idleness is only evaluated when watermarks are emitted periodically: +the idle timeout countdown starts at the first periodic watermark emission that +saw no records since the previous one, and the input switches to idle at the +first periodic emission after strictly more than the configured timeout has +elapsed since then. An input is therefore marked idle no earlier than the idle +timeout after its last record and, in the worst case, up to the idle timeout +plus three times the `pipeline.auto-watermark-interval` after it. + ## Watermark alignment In the previous paragraph we discussed a situation when splits/partitions/shards or sources are idle diff --git a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkStrategy.java b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkStrategy.java index 388ef68c7fd51c..b02e0306c1dc1f 100644 --- a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkStrategy.java +++ b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarkStrategy.java @@ -143,6 +143,14 @@ default WatermarkStrategy withTimestampAssigner( *

Idleness can be important if some partitions have little data and might not have events * during some periods. Without idleness, these streams can stall the overall event time * progress of the application. + * + *

Idleness is only evaluated on each periodic watermark emission (see {@code + * pipeline.auto-watermark-interval}). The timeout countdown starts at the first periodic + * emission that saw no records since the previous one — which can be up to two watermark + * intervals after the last record — and the partition switches to idle at the first periodic + * emission after strictly more than {@code idleTimeout} has elapsed since then. A partition is + * therefore marked idle no earlier than {@code idleTimeout} after its last record and, in the + * worst case, up to {@code idleTimeout} plus three watermark intervals after it. */ default WatermarkStrategy withIdleness(Duration idleTimeout) { checkNotNull(idleTimeout, "idleTimeout"); diff --git a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java index 4cc651cf6e9139..ed8a02dd965e17 100644 --- a/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java +++ b/flink-core/src/main/java/org/apache/flink/api/common/eventtime/WatermarksWithIdleness.java @@ -33,6 +33,13 @@ * A WatermarkGenerator that adds idleness detection to another WatermarkGenerator. If no events * come within a certain time (timeout duration) then this generator marks the stream as idle, until * the next watermark is generated. + * + *

Idleness is only evaluated in {@link #onPeriodicEmit(WatermarkOutput)}: the timeout countdown + * starts at the first periodic emission that saw no events since the previous one, and the stream + * is marked idle at the first periodic emission after strictly more than the timeout has elapsed + * since the countdown started. The stream is therefore marked idle no earlier than the timeout + * after the last event and, in the worst case, up to the timeout plus three periodic emission + * intervals after it. */ @Public public class WatermarksWithIdleness implements WatermarkGenerator { diff --git a/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTimeoutTest.java b/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTimeoutTest.java new file mode 100644 index 00000000000000..bdc52f31f304f9 --- /dev/null +++ b/flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTimeoutTest.java @@ -0,0 +1,79 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.flink.api.common.eventtime; + +import org.apache.flink.util.clock.ManualClock; + +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.util.concurrent.TimeUnit; + +import static org.assertj.core.api.Assertions.assertThat; + +/** + * Tests (FLINK-40503) pinning the documented idleness detection timing of {@link + * WatermarksWithIdleness}: the timeout countdown is anchored at the first quiet periodic probe (not + * at the last event), and idleness fires at the first probe strictly more than the timeout after + * that anchor. With timeout T = 10ms, an event at t=4ms, and probes every 5ms, idleness is thus + * declared at t=25ms — 21ms after the last event, within the documented worst case of T plus three + * probe intervals. + */ +class WatermarksWithIdlenessTimeoutTest { + + private static final Duration TIMEOUT = Duration.ofMillis(10); + + @Test + void idleFiresAtFirstProbeStrictlyExceedingTimeoutAfterFirstQuietProbe() { + // non-zero start: IdlenessTimer treats a relative timestamp of 0 as "no timer started" + ManualClock clock = new ManualClock(1_000_000_000L); + WatermarksWithIdleness generator = + new WatermarksWithIdleness<>(new NoWatermarksGenerator<>(), TIMEOUT, clock); + TestingWatermarkOutput output = new TestingWatermarkOutput(); + + // event at t=4ms + clock.advanceTime(4, TimeUnit.MILLISECONDS); + generator.onEvent(new Object(), 1L, output); + + // probe t=5ms: sees activity -> resets, no timer yet + clock.advanceTime(1, TimeUnit.MILLISECONDS); + generator.onPeriodicEmit(output); + assertThat(output.isIdle()).isFalse(); + + // probe t=10ms: first quiet probe -> countdown starts HERE (not at the last event) + clock.advanceTime(5, TimeUnit.MILLISECONDS); + generator.onPeriodicEmit(output); + assertThat(output.isIdle()).isFalse(); + + // probe t=15ms: elapsed since countdown start = 5ms, not > 10ms + clock.advanceTime(5, TimeUnit.MILLISECONDS); + generator.onPeriodicEmit(output); + assertThat(output.isIdle()).isFalse(); + + // probe t=20ms: elapsed = exactly 10ms; strict '>' comparison -> still NOT idle + clock.advanceTime(5, TimeUnit.MILLISECONDS); + generator.onPeriodicEmit(output); + assertThat(output.isIdle()).isFalse(); + + // probe t=25ms: elapsed = 15ms > 10ms -> idle, 21ms after the last event + clock.advanceTime(5, TimeUnit.MILLISECONDS); + generator.onPeriodicEmit(output); + assertThat(output.isIdle()).isTrue(); + } +}