This is an automated email from the ASF dual-hosted git repository. pnowojski pushed a commit to branch master in repository https://gitbox.apache.org/repos/asf/flink.git
commit 7fcc6b8d11aa58f95c14cfe5efd01925e6eee893 Author: Efrat Levitan <[email protected]> AuthorDate: Thu Sep 3 19:28:13 2026 +0300 [FLINK-40549][metrics] Close splitMetricGroup along with the child splitWatermarkMetricGroup Previousely we incorrectly closed the child splitWatermarkMetricGroup on split finished, leaving the parent metricGroup unclosed --- .../groups/InternalSourceSplitMetricGroup.java | 19 +++++++++++++------ .../groups/InternalSourceSplitMetricGroupTest.java | 12 ++++++++++++ 2 files changed, 25 insertions(+), 6 deletions(-) diff --git a/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroup.java b/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroup.java index a379a434764..1426e72a158 100644 --- a/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroup.java +++ b/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroup.java @@ -47,6 +47,7 @@ public class InternalSourceSplitMetricGroup extends ProxyMetricGroup<MetricGroup private static final String WATERMARK = "watermark"; private static final long SPLIT_NOT_STARTED = -1L; private long splitStartTime = SPLIT_NOT_STARTED; + private final MetricGroup splitMetricGroup; private final MetricGroup splitWatermarkMetricGroup; private final String splitId; @@ -58,7 +59,8 @@ public class InternalSourceSplitMetricGroup extends ProxyMetricGroup<MetricGroup super(parentMetricGroup); this.clock = clock; this.splitId = splitId; - splitWatermarkMetricGroup = parentMetricGroup.addGroup(SPLIT, splitId).addGroup(WATERMARK); + splitMetricGroup = parentMetricGroup.addGroup(SPLIT, splitId); + splitWatermarkMetricGroup = splitMetricGroup.addGroup(WATERMARK); pausedTimePerSecond = splitWatermarkMetricGroup.gauge( MetricNames.SPLIT_PAUSED_TIME, new TimerGauge(clock)); @@ -187,17 +189,22 @@ public class InternalSourceSplitMetricGroup extends ProxyMetricGroup<MetricGroup } public void onSplitFinished() { - if (splitWatermarkMetricGroup instanceof AbstractMetricGroup) { - ((AbstractMetricGroup) splitWatermarkMetricGroup).close(); + if (splitMetricGroup instanceof AbstractMetricGroup) { + ((AbstractMetricGroup<?>) splitMetricGroup).close(); } else { - if (splitWatermarkMetricGroup != null) { + if (splitMetricGroup != null) { LOG.warn( - "Split watermark metric group can not be closed, expecting an instance of AbstractMetricGroup but got: ", - splitWatermarkMetricGroup.getClass().getName()); + "Split metric group can not be closed, expecting an instance of AbstractMetricGroup but got: {}", + splitMetricGroup.getClass().getName()); } } } + @VisibleForTesting + public MetricGroup getSplitMetricGroup() { + return splitMetricGroup; + } + @VisibleForTesting public MetricGroup getSplitWatermarkMetricGroup() { return splitWatermarkMetricGroup; diff --git a/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroupTest.java b/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroupTest.java index 5dd63dcdaa5..b48c5e16359 100644 --- a/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroupTest.java +++ b/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/InternalSourceSplitMetricGroupTest.java @@ -40,6 +40,18 @@ class InternalSourceSplitMetricGroupTest { null); } + @Test + void testOnSplitFinishedClosesSplitMetricGroup() { + InternalSourceSplitMetricGroup metricGroup = getMetricGroupWithClock(new ManualClock()); + + metricGroup.onSplitFinished(); + + assertThat(((AbstractMetricGroup<?>) metricGroup.getSplitMetricGroup()).isClosed()) + .isTrue(); + assertThat(((AbstractMetricGroup<?>) metricGroup.getSplitWatermarkMetricGroup()).isClosed()) + .isTrue(); + } + @Test void testClocksStartTickingAfterSplitStarted() throws InterruptedException { ManualClock clock = new ManualClock(System.currentTimeMillis());
