This is an automated email from the ASF dual-hosted git repository.

davidradl pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/flink.git


The following commit(s) were added to refs/heads/master by this push:
     new 0a6d74100d4 [FLINK-35321][runtime] Ensure pendingCommitables metric is 
not re-registered while copying CommittableCollector (#27598)
0a6d74100d4 is described below

commit 0a6d74100d4f0a84373bbb2aafbf7f88c0d6c190
Author: Piotr Rudnicki <[email protected]>
AuthorDate: Thu Sep 10 12:36:03 2026 +0200

    [FLINK-35321][runtime] Ensure pendingCommitables metric is not 
re-registered while copying CommittableCollector (#27598)
    
    * [FLINK-35321] Add flag for metric registration
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Add unit test and refactor constructor method
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Rename local param
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Refactor test to not use Mockito
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Add license agreement
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Rename plus code format
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Revert flag in serialization path
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Custom test metric group implementation
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Get rid on flag and move metric registration to static 
constructor
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Correct tests
    
    Signed-off-by: deamondev <[email protected]>
    
    * [FLINK-35321] Restore static class for testing
    
    Signed-off-by: Piotr Rudnicki <[email protected]>
    
    ---------
    
    Signed-off-by: deamondev <[email protected]>
    Signed-off-by: Piotr Rudnicki <[email protected]>
---
 .../sink/committables/CommittableCollector.java    |  5 +-
 .../metrics/groups/MetricsGroupTestUtils.java      | 92 ++++++++++++++++++++++
 .../committables/CommittableCollectorTest.java     | 23 +++++-
 3 files changed, 115 insertions(+), 5 deletions(-)

diff --git 
a/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollector.java
 
b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollector.java
index 96585a632d1..b63f1a0903e 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollector.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollector.java
@@ -64,7 +64,6 @@ public class CommittableCollector<CommT> {
             SinkCommitterMetricGroup metricGroup) {
         this.checkpointCommittables = new 
TreeMap<>(checkNotNull(checkpointCommittables));
         this.metricGroup = metricGroup;
-        
this.metricGroup.setCurrentPendingCommittablesGauge(this::getNumPending);
     }
 
     private int getNumPending() {
@@ -82,7 +81,9 @@ public class CommittableCollector<CommT> {
      * @return {@link CommittableCollector}
      */
     public static <CommT> CommittableCollector<CommT> 
of(SinkCommitterMetricGroup metricGroup) {
-        return new CommittableCollector<>(metricGroup);
+        CommittableCollector<CommT> collector = new 
CommittableCollector<>(metricGroup);
+        
metricGroup.setCurrentPendingCommittablesGauge(collector::getNumPending);
+        return collector;
     }
 
     /**
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/MetricsGroupTestUtils.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/MetricsGroupTestUtils.java
index 1308bde0b18..4e9e055d8af 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/MetricsGroupTestUtils.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/metrics/groups/MetricsGroupTestUtils.java
@@ -18,9 +18,17 @@
 
 package org.apache.flink.runtime.metrics.groups;
 
+import org.apache.flink.annotation.VisibleForTesting;
+import org.apache.flink.metrics.Counter;
+import org.apache.flink.metrics.Gauge;
 import org.apache.flink.metrics.MetricGroup;
 import org.apache.flink.metrics.groups.OperatorIOMetricGroup;
+import org.apache.flink.metrics.groups.OperatorMetricGroup;
+import org.apache.flink.metrics.groups.SinkCommitterMetricGroup;
 import org.apache.flink.metrics.groups.UnregisteredMetricsGroup;
+import org.apache.flink.runtime.metrics.MetricNames;
+
+import java.util.concurrent.atomic.AtomicInteger;
 
 /** Util class to create metric groups for SinkV2 tests. */
 public class MetricsGroupTestUtils {
@@ -46,4 +54,88 @@ public class MetricsGroupTestUtils {
                 new UnregisteredMetricsGroup(),
                 UnregisteredMetricsGroup.createOperatorIOMetricGroup());
     }
+
+    public static TrackableCommitterMetricGroup 
mockTrackableCommitterMetricGroup() {
+        return new TrackableCommitterMetricGroup(
+                new UnregisteredMetricsGroup(),
+                UnregisteredMetricsGroup.createOperatorIOMetricGroup());
+    }
+
+    public static class TrackableCommitterMetricGroup extends 
ProxyMetricGroup<MetricGroup>
+            implements SinkCommitterMetricGroup {
+
+        private final AtomicInteger gaugeCallCount = new AtomicInteger(0);
+
+        private final Counter numCommittablesTotal;
+        private final Counter numCommittablesFailure;
+        private final Counter numCommittablesRetry;
+        private final Counter numCommitatblesSuccess;
+        private final Counter numCommitatblesAlreadyCommitted;
+        private final OperatorIOMetricGroup operatorIOMetricGroup;
+
+        @VisibleForTesting
+        public TrackableCommitterMetricGroup(
+                MetricGroup parentMetricGroup, OperatorIOMetricGroup 
operatorIOMetricGroup) {
+            super(parentMetricGroup);
+            numCommittablesTotal = 
parentMetricGroup.counter(MetricNames.TOTAL_COMMITTABLES);
+            numCommittablesFailure = 
parentMetricGroup.counter(MetricNames.FAILED_COMMITTABLES);
+            numCommittablesRetry = 
parentMetricGroup.counter(MetricNames.RETRIED_COMMITTABLES);
+            numCommitatblesSuccess = 
parentMetricGroup.counter(MetricNames.SUCCESSFUL_COMMITTABLES);
+            numCommitatblesAlreadyCommitted =
+                    
parentMetricGroup.counter(MetricNames.ALREADY_COMMITTED_COMMITTABLES);
+
+            this.operatorIOMetricGroup = operatorIOMetricGroup;
+        }
+
+        public static TrackableCommitterMetricGroup wrap(OperatorMetricGroup 
operatorMetricGroup) {
+            return new TrackableCommitterMetricGroup(
+                    operatorMetricGroup, 
operatorMetricGroup.getIOMetricGroup());
+        }
+
+        @Override
+        public OperatorIOMetricGroup getIOMetricGroup() {
+            return operatorIOMetricGroup;
+        }
+
+        @Override
+        public Counter getNumCommittablesTotalCounter() {
+            return numCommittablesTotal;
+        }
+
+        @Override
+        public Counter getNumCommittablesFailureCounter() {
+            return numCommittablesFailure;
+        }
+
+        @Override
+        public Counter getNumCommittablesRetryCounter() {
+            return numCommittablesRetry;
+        }
+
+        @Override
+        public Counter getNumCommittablesSuccessCounter() {
+            return numCommitatblesSuccess;
+        }
+
+        @Override
+        public Counter getNumCommittablesAlreadyCommittedCounter() {
+            return numCommitatblesAlreadyCommitted;
+        }
+
+        @Override
+        public void setCurrentPendingCommittablesGauge(
+                Gauge<Integer> currentPendingCommittablesGauge) {
+            gaugeCallCount.incrementAndGet();
+            parentMetricGroup.gauge(
+                    MetricNames.PENDING_COMMITTABLES, 
currentPendingCommittablesGauge);
+        }
+
+        public int getGaugeCallCount() {
+            return gaugeCallCount.get();
+        }
+
+        public void resetGaugeCallCount() {
+            gaugeCallCount.set(0);
+        }
+    }
 }
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorTest.java
 
b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorTest.java
index 892b3785e25..a697ee94e3d 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorTest.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/streaming/runtime/operators/sink/committables/CommittableCollectorTest.java
@@ -18,17 +18,22 @@
 
 package org.apache.flink.streaming.runtime.operators.sink.committables;
 
-import org.apache.flink.metrics.groups.SinkCommitterMetricGroup;
 import org.apache.flink.runtime.metrics.groups.MetricsGroupTestUtils;
 import org.apache.flink.streaming.api.connector.sink2.CommittableSummary;
 
+import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 
 import static org.assertj.core.api.Assertions.assertThat;
 
 class CommittableCollectorTest {
-    private static final SinkCommitterMetricGroup METRIC_GROUP =
-            MetricsGroupTestUtils.mockCommitterMetricGroup();
+    private static final MetricsGroupTestUtils.TrackableCommitterMetricGroup 
METRIC_GROUP =
+            MetricsGroupTestUtils.mockTrackableCommitterMetricGroup();
+
+    @BeforeEach
+    public void setUp() {
+        METRIC_GROUP.resetGaugeCallCount();
+    }
 
     @Test
     void testGetCheckpointCommittablesUpTo() {
@@ -42,4 +47,16 @@ class CommittableCollectorTest {
 
         
assertThat(committableCollector.getCheckpointCommittablesUpTo(2)).hasSize(2);
     }
+
+    @Test
+    void testSetPendingGaugeNotCalledOnCopy() {
+        final CommittableCollector<Integer> committableCollector =
+                CommittableCollector.of(METRIC_GROUP);
+
+        assertThat(METRIC_GROUP.getGaugeCallCount()).isEqualTo(1);
+
+        committableCollector.copy();
+
+        assertThat(METRIC_GROUP.getGaugeCallCount()).isEqualTo(1);
+    }
 }

Reply via email to