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());

Reply via email to