From ca757a08127ac69d25f7ae5da346671659fafc43 Mon Sep 17 00:00:00 2001 From: Martijn Visser <2989614+MartijnVisser@users.noreply.github.com> Date: Fri, 28 Aug 2026 14:33:08 +0200 Subject: [PATCH] [FLINK-40503][core] Document effective idleness detection timing The withIdleness javadoc suggested a partition is considered idle once no records have arrived for idleTimeout, but WatermarksWithIdleness anchors its countdown at the first quiet periodic probe (up to two watermark intervals after the last record) and only switches to idle at the first probe strictly exceeding the timeout, so detection can take up to idleTimeout plus three watermark intervals. Document this actual behavior in the withIdleness and WatermarksWithIdleness javadocs and in the idle-sources documentation instead of changing the long-standing timing, and add a test pinning the documented contract. The Chinese documentation page translates this section and is left for a follow-up translation update. Generated-by: Claude Code (Fable 5) --- .../event-time/generating_watermarks.md | 8 ++ .../common/eventtime/WatermarkStrategy.java | 8 ++ .../eventtime/WatermarksWithIdleness.java | 7 ++ .../WatermarksWithIdlenessTimeoutTest.java | 79 +++++++++++++++++++ 4 files changed, 102 insertions(+) create mode 100644 flink-core/src/test/java/org/apache/flink/api/common/eventtime/WatermarksWithIdlenessTimeoutTest.java 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(); + } +}