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);
+ }
}