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 98bd2659f5b53600c7486aea4b0e5ed44bab0e9c
Author: Piotr Nowojski <[email protected]>
AuthorDate: Tue Apr 2 14:59:35 2024 +0200

    [FLINK-35065][metrics] Add numFiredTimers and numFiredTimersPerSecond
---
 docs/content.zh/docs/ops/metrics.md                | 12 ++++-
 docs/content/docs/ops/metrics.md                   | 12 ++++-
 .../apache/flink/runtime/metrics/MetricNames.java  |  3 ++
 .../runtime/metrics/groups/TaskIOMetricGroup.java  | 10 +++++
 .../api/operators/InternalTimeServiceManager.java  |  2 +
 .../operators/InternalTimeServiceManagerImpl.java  |  8 +++-
 .../api/operators/InternalTimerServiceImpl.java    |  6 +++
 .../operators/StreamTaskStateInitializerImpl.java  |  1 +
 .../BatchExecutionInternalTimeServiceManager.java  |  2 +
 .../operators/InternalTimerServiceImplTest.java    | 52 +++++++++++++++++-----
 .../StateInitializationContextImplTest.java        |  2 +
 .../StreamTaskStateInitializerImplTest.java        |  2 +
 .../BatchExecutionInternalTimeServiceTest.java     | 15 +++++++
 .../util/AbstractStreamOperatorTestHarness.java    |  3 ++
 .../restore/StreamOperatorSnapshotRestoreTest.java |  2 +
 15 files changed, 119 insertions(+), 13 deletions(-)

diff --git a/docs/content.zh/docs/ops/metrics.md 
b/docs/content.zh/docs/ops/metrics.md
index 39ddfabcd99..fd16d9af824 100644
--- a/docs/content.zh/docs/ops/metrics.md
+++ b/docs/content.zh/docs/ops/metrics.md
@@ -1632,7 +1632,7 @@ Note that the metrics are only available via reporters.
       <td>Histogram</td>
     </tr>
     <tr>
-      <th rowspan="25"><strong>Task</strong></th>
+      <th rowspan="27"><strong>Task</strong></th>
       <td>numBytesInLocal</td>
       <td><span class="label label-danger">Attention:</span> deprecated, use 
<a href="{{< ref "docs/ops/metrics" >}}#default-shuffle-service">Default 
shuffle service metrics</a>.</td>
       <td>Counter</td>
@@ -1692,6 +1692,16 @@ Note that the metrics are only available via reporters.
       <td>The number of network buffers this task emits per second.</td>
       <td>Meter</td>
     </tr>
+    <tr>
+      <td>numFiredTimers</td>
+      <td>The total number of timers this task has fired.</td>
+      <td>Counter</td>
+    </tr>
+    <tr>
+      <td>numFiredTimersPerSecond</td>
+      <td>The number of timers this task fires per second.</td>
+      <td>Meter</td>
+    </tr>
     <tr>
       <td>isBackPressured</td>
       <td>Whether the task is back-pressured.</td>
diff --git a/docs/content/docs/ops/metrics.md b/docs/content/docs/ops/metrics.md
index 91a654d7d89..14c8887aa2a 100644
--- a/docs/content/docs/ops/metrics.md
+++ b/docs/content/docs/ops/metrics.md
@@ -1622,7 +1622,7 @@ Note that the metrics are only available via reporters.
       <td>Histogram</td>
     </tr>
     <tr>
-      <th rowspan="25"><strong>Task</strong></th>
+      <th rowspan="27"><strong>Task</strong></th>
       <td>numBytesInLocal</td>
       <td><span class="label label-danger">Attention:</span> deprecated, use 
<a href="{{< ref "docs/ops/metrics" >}}#default-shuffle-service">Default 
shuffle service metrics</a>.</td>
       <td>Counter</td>
@@ -1682,6 +1682,16 @@ Note that the metrics are only available via reporters.
       <td>The number of network buffers this task emits per second.</td>
       <td>Meter</td>
     </tr>
+    <tr>
+      <td>numFiredTimers</td>
+      <td>The total number of timers this task has fired.</td>
+      <td>Counter</td>
+    </tr>
+    <tr>
+      <td>numFiredTimersPerSecond</td>
+      <td>The number of timers this task fires per second.</td>
+      <td>Meter</td>
+    </tr>
     <tr>
       <td>isBackPressured</td>
       <td>Whether the task is back-pressured.</td>
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/MetricNames.java 
b/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/MetricNames.java
index 392ad847117..b0987786745 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/MetricNames.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/MetricNames.java
@@ -44,6 +44,9 @@ public class MetricNames {
     public static final String IO_CURRENT_INPUT_WATERMARK_PATERN = 
"currentInput%dWatermark";
     public static final String IO_CURRENT_OUTPUT_WATERMARK = 
"currentOutputWatermark";
 
+    public static final String NUM_FIRED_TIMERS = "numFiredTimers";
+    public static final String NUM_FIRED_TIMERS_RATE = "numFiredTimers" + 
SUFFIX_RATE;
+
     public static final String NUM_RUNNING_JOBS = "numRunningJobs";
     public static final String TASK_SLOTS_AVAILABLE = "taskSlotsAvailable";
     public static final String TASK_SLOTS_TOTAL = "taskSlotsTotal";
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/TaskIOMetricGroup.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/TaskIOMetricGroup.java
index a571b078c2e..69ead952442 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/TaskIOMetricGroup.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/metrics/groups/TaskIOMetricGroup.java
@@ -55,6 +55,8 @@ public class TaskIOMetricGroup extends 
ProxyMetricGroup<TaskMetricGroup> {
     private final SumCounter numRecordsIn;
     private final SumCounter numRecordsOut;
     private final Counter numBuffersOut;
+    private final Counter numFiredTimers;
+    private final MeterView numFiredTimersRate;
     private final Counter numMailsProcessed;
 
     private final Meter numBytesInRate;
@@ -140,6 +142,10 @@ public class TaskIOMetricGroup extends 
ProxyMetricGroup<TaskMetricGroup> {
         this.accumulatedIdleTime =
                 gauge(MetricNames.ACC_TASK_IDLE_TIME, 
idleTimePerSecond::getAccumulatedCount);
 
+        this.numFiredTimers = counter(MetricNames.NUM_FIRED_TIMERS, new 
SimpleCounter());
+        this.numFiredTimersRate =
+                meter(MetricNames.NUM_FIRED_TIMERS_RATE, new 
MeterView(numFiredTimers));
+
         this.numMailsProcessed = new SimpleCounter();
         this.mailboxThroughput =
                 meter(MetricNames.MAILBOX_THROUGHPUT, new 
MeterView(numMailsProcessed));
@@ -207,6 +213,10 @@ public class TaskIOMetricGroup extends 
ProxyMetricGroup<TaskMetricGroup> {
         return numBuffersOut;
     }
 
+    public Counter getNumFiredTimers() {
+        return numFiredTimers;
+    }
+
     public Counter getNumMailsProcessedCounter() {
         return numMailsProcessed;
     }
diff --git 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManager.java
 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManager.java
index 439789c3709..b51ea67bc6e 100644
--- 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManager.java
+++ 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManager.java
@@ -20,6 +20,7 @@ package org.apache.flink.streaming.api.operators;
 
 import org.apache.flink.annotation.Internal;
 import org.apache.flink.api.common.typeutils.TypeSerializer;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.state.CheckpointableKeyedStateBackend;
 import org.apache.flink.runtime.state.KeyGroupStatePartitionStreamProvider;
 import org.apache.flink.runtime.state.KeyedStateCheckpointOutputStream;
@@ -73,6 +74,7 @@ public interface InternalTimeServiceManager<K> {
     @FunctionalInterface
     interface Provider extends Serializable {
         <K> InternalTimeServiceManager<K> create(
+                TaskIOMetricGroup taskIOMetricGroup,
                 CheckpointableKeyedStateBackend<K> keyedStatedBackend,
                 ClassLoader userClassloader,
                 KeyContext keyContext,
diff --git 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManagerImpl.java
 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManagerImpl.java
index 51a280bdda2..0a0a76cd530 100644
--- 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManagerImpl.java
+++ 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimeServiceManagerImpl.java
@@ -23,6 +23,7 @@ import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.api.common.typeutils.TypeSerializer;
 import org.apache.flink.core.memory.DataOutputView;
 import org.apache.flink.core.memory.DataOutputViewStreamWrapper;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.state.CheckpointableKeyedStateBackend;
 import org.apache.flink.runtime.state.KeyGroupRange;
 import org.apache.flink.runtime.state.KeyGroupStatePartitionStreamProvider;
@@ -66,6 +67,7 @@ public class InternalTimeServiceManagerImpl<K> implements 
InternalTimeServiceMan
 
     @VisibleForTesting static final String EVENT_TIMER_PREFIX = 
TIMER_STATE_PREFIX + "/event_";
 
+    private final TaskIOMetricGroup taskIOMetricGroup;
     private final KeyGroupRange localKeyGroupRange;
     private final KeyContext keyContext;
     private final PriorityQueueSetFactory priorityQueueSetFactory;
@@ -75,12 +77,13 @@ public class InternalTimeServiceManagerImpl<K> implements 
InternalTimeServiceMan
     private final Map<String, InternalTimerServiceImpl<K, ?>> timerServices;
 
     private InternalTimeServiceManagerImpl(
+            TaskIOMetricGroup taskIOMetricGroup,
             KeyGroupRange localKeyGroupRange,
             KeyContext keyContext,
             PriorityQueueSetFactory priorityQueueSetFactory,
             ProcessingTimeService processingTimeService,
             StreamTaskCancellationContext cancellationContext) {
-
+        this.taskIOMetricGroup = taskIOMetricGroup;
         this.localKeyGroupRange = 
Preconditions.checkNotNull(localKeyGroupRange);
         this.priorityQueueSetFactory = 
Preconditions.checkNotNull(priorityQueueSetFactory);
         this.keyContext = Preconditions.checkNotNull(keyContext);
@@ -96,6 +99,7 @@ public class InternalTimeServiceManagerImpl<K> implements 
InternalTimeServiceMan
      * <p><b>IMPORTANT:</b> Keep in sync with {@link 
InternalTimeServiceManager.Provider}.
      */
     public static <K> InternalTimeServiceManagerImpl<K> create(
+            TaskIOMetricGroup taskIOMetricGroup,
             CheckpointableKeyedStateBackend<K> keyedStateBackend,
             ClassLoader userClassloader,
             KeyContext keyContext,
@@ -107,6 +111,7 @@ public class InternalTimeServiceManagerImpl<K> implements 
InternalTimeServiceMan
 
         final InternalTimeServiceManagerImpl<K> timeServiceManager =
                 new InternalTimeServiceManagerImpl<>(
+                        taskIOMetricGroup,
                         keyGroupRange,
                         keyContext,
                         keyedStateBackend,
@@ -160,6 +165,7 @@ public class InternalTimeServiceManagerImpl<K> implements 
InternalTimeServiceMan
 
             timerService =
                     new InternalTimerServiceImpl<>(
+                            taskIOMetricGroup,
                             localKeyGroupRange,
                             keyContext,
                             processingTimeService,
diff --git 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImpl.java
 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImpl.java
index 1345916132a..198fcb58d76 100644
--- 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImpl.java
+++ 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImpl.java
@@ -21,6 +21,7 @@ package org.apache.flink.streaming.api.operators;
 import org.apache.flink.annotation.VisibleForTesting;
 import org.apache.flink.api.common.typeutils.TypeSerializer;
 import org.apache.flink.api.common.typeutils.TypeSerializerSchemaCompatibility;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.state.InternalPriorityQueue;
 import org.apache.flink.runtime.state.KeyGroupRange;
 import org.apache.flink.runtime.state.KeyGroupedInternalPriorityQueue;
@@ -45,6 +46,7 @@ public class InternalTimerServiceImpl<K, N> implements 
InternalTimerService<N> {
 
     private final ProcessingTimeService processingTimeService;
 
+    private final TaskIOMetricGroup taskIOMetricGroup;
     private final KeyContext keyContext;
 
     /** Processing time timers that are currently in-flight. */
@@ -93,12 +95,14 @@ public class InternalTimerServiceImpl<K, N> implements 
InternalTimerService<N> {
     private InternalTimersSnapshot<K, N> restoredTimersSnapshot;
 
     InternalTimerServiceImpl(
+            TaskIOMetricGroup taskIOMetricGroup,
             KeyGroupRange localKeyGroupRange,
             KeyContext keyContext,
             ProcessingTimeService processingTimeService,
             KeyGroupedInternalPriorityQueue<TimerHeapInternalTimer<K, N>> 
processingTimeTimersQueue,
             KeyGroupedInternalPriorityQueue<TimerHeapInternalTimer<K, N>> 
eventTimeTimersQueue,
             StreamTaskCancellationContext cancellationContext) {
+        this.taskIOMetricGroup = taskIOMetricGroup;
 
         this.keyContext = checkNotNull(keyContext);
         this.processingTimeService = checkNotNull(processingTimeService);
@@ -292,6 +296,7 @@ public class InternalTimerServiceImpl<K, N> implements 
InternalTimerService<N> {
             keyContext.setCurrentKey(timer.getKey());
             processingTimeTimersQueue.poll();
             triggerTarget.onProcessingTime(timer);
+            taskIOMetricGroup.getNumFiredTimers().inc();
         }
 
         if (timer != null && nextTimer == null) {
@@ -312,6 +317,7 @@ public class InternalTimerServiceImpl<K, N> implements 
InternalTimerService<N> {
             keyContext.setCurrentKey(timer.getKey());
             eventTimeTimersQueue.poll();
             triggerTarget.onEventTime(timer);
+            taskIOMetricGroup.getNumFiredTimers().inc();
         }
     }
 
diff --git 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImpl.java
 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImpl.java
index 9cd11a09a81..50b3c4e0867 100644
--- 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImpl.java
+++ 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImpl.java
@@ -228,6 +228,7 @@ public class StreamTaskStateInitializerImpl implements 
StreamTaskStateInitialize
 
                 timeServiceManager =
                         timeServiceManagerProvider.create(
+                                
environment.getMetricGroup().getIOMetricGroup(),
                                 keyedStatedBackend,
                                 
environment.getUserCodeClassLoader().asClassLoader(),
                                 keyContext,
diff --git 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceManager.java
 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceManager.java
index c1152725e05..9e0ea4a2749 100644
--- 
a/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceManager.java
+++ 
b/flink-streaming-java/src/main/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceManager.java
@@ -19,6 +19,7 @@
 package org.apache.flink.streaming.api.operators.sorted.state;
 
 import org.apache.flink.api.common.typeutils.TypeSerializer;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.state.CheckpointableKeyedStateBackend;
 import org.apache.flink.runtime.state.KeyGroupStatePartitionStreamProvider;
 import org.apache.flink.runtime.state.KeyedStateBackend;
@@ -85,6 +86,7 @@ public class BatchExecutionInternalTimeServiceManager<K>
     }
 
     public static <K> InternalTimeServiceManager<K> create(
+            TaskIOMetricGroup taskIOMetricGroup,
             CheckpointableKeyedStateBackend<K> keyedStatedBackend,
             ClassLoader userClassloader,
             KeyContext keyContext, // the operator
diff --git 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImplTest.java
 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImplTest.java
index 5caa0aedae4..544cc9128a5 100644
--- 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImplTest.java
+++ 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/InternalTimerServiceImplTest.java
@@ -24,6 +24,8 @@ import 
org.apache.flink.api.common.typeutils.base.StringSerializer;
 import org.apache.flink.api.java.tuple.Tuple3;
 import org.apache.flink.core.memory.DataInputViewStreamWrapper;
 import org.apache.flink.core.memory.DataOutputViewStreamWrapper;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
+import org.apache.flink.runtime.metrics.groups.UnregisteredMetricGroups;
 import org.apache.flink.runtime.state.KeyGroupRange;
 import org.apache.flink.runtime.state.KeyGroupRangeAssignment;
 import org.apache.flink.runtime.state.KeyGroupedInternalPriorityQueue;
@@ -89,6 +91,8 @@ class InternalTimerServiceImplTest {
 
         InternalTimerServiceImpl<Integer, String> service =
                 createInternalTimerService(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         testKeyGroupList,
                         keyContext,
                         processingTimeService,
@@ -119,6 +123,8 @@ class InternalTimerServiceImplTest {
 
         InternalTimerServiceImpl<Integer, String> timerService =
                 createInternalTimerService(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyGroupRange,
                         keyContext,
                         new TestProcessingTimeService(),
@@ -401,13 +407,22 @@ class InternalTimerServiceImplTest {
 
         TestKeyContext keyContext = new TestKeyContext();
         TestProcessingTimeService processingTimeService = new 
TestProcessingTimeService();
-        InternalTimerServiceImpl<Integer, String> timerService =
-                createAndStartInternalTimerService(
-                        mockTriggerable,
+        PriorityQueueSetFactory priorityQueueSetFactory = createQueueFactory();
+        TaskIOMetricGroup taskIOMetricGroup =
+                
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup().getIOMetricGroup();
+        InternalTimerServiceImpl<Integer, String> service =
+                createInternalTimerService(
+                        taskIOMetricGroup,
+                        testKeyGroupRange,
                         keyContext,
                         processingTimeService,
-                        testKeyGroupRange,
-                        createQueueFactory());
+                        IntSerializer.INSTANCE,
+                        StringSerializer.INSTANCE,
+                        priorityQueueSetFactory);
+
+        service.startTimerService(
+                IntSerializer.INSTANCE, StringSerializer.INSTANCE, 
mockTriggerable);
+        InternalTimerServiceImpl<Integer, String> timerService = service;
 
         // get two different keys
         int key1 = getKeyInKeyGroupRange(testKeyGroupRange, maxParallelism);
@@ -431,6 +446,7 @@ class InternalTimerServiceImplTest {
         assertThat(timerService.numEventTimeTimers("ciao")).isEqualTo(2);
 
         timerService.advanceWatermark(10);
+        
assertThat(taskIOMetricGroup.getNumFiredTimers().getCount()).isEqualTo(4);
 
         verify(mockTriggerable, times(4)).onEventTime(anyInternalTimer());
         verify(mockTriggerable, times(1))
@@ -453,13 +469,22 @@ class InternalTimerServiceImplTest {
 
         TestKeyContext keyContext = new TestKeyContext();
         TestProcessingTimeService processingTimeService = new 
TestProcessingTimeService();
-        InternalTimerServiceImpl<Integer, String> timerService =
-                createAndStartInternalTimerService(
-                        mockTriggerable,
+        PriorityQueueSetFactory priorityQueueSetFactory = createQueueFactory();
+        TaskIOMetricGroup taskIOMetricGroup =
+                
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup().getIOMetricGroup();
+        InternalTimerServiceImpl<Integer, String> service =
+                createInternalTimerService(
+                        taskIOMetricGroup,
+                        testKeyGroupRange,
                         keyContext,
                         processingTimeService,
-                        testKeyGroupRange,
-                        createQueueFactory());
+                        IntSerializer.INSTANCE,
+                        StringSerializer.INSTANCE,
+                        priorityQueueSetFactory);
+
+        service.startTimerService(
+                IntSerializer.INSTANCE, StringSerializer.INSTANCE, 
mockTriggerable);
+        InternalTimerServiceImpl<Integer, String> timerService = service;
 
         // get two different keys
         int key1 = getKeyInKeyGroupRange(testKeyGroupRange, maxParallelism);
@@ -483,6 +508,7 @@ class InternalTimerServiceImplTest {
         assertThat(timerService.numProcessingTimeTimers("ciao")).isEqualTo(2);
 
         processingTimeService.setCurrentTime(10);
+        
assertThat(taskIOMetricGroup.getNumFiredTimers().getCount()).isEqualTo(4);
 
         verify(mockTriggerable, times(4)).onProcessingTime(anyInternalTimer());
         verify(mockTriggerable, times(1))
@@ -1002,6 +1028,8 @@ class InternalTimerServiceImplTest {
             PriorityQueueSetFactory priorityQueueSetFactory) {
         InternalTimerServiceImpl<Integer, String> service =
                 createInternalTimerService(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyGroupList,
                         keyContext,
                         processingTimeService,
@@ -1026,6 +1054,8 @@ class InternalTimerServiceImplTest {
         // create an empty service
         InternalTimerServiceImpl<Integer, String> service =
                 createInternalTimerService(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyGroupsList,
                         keyContext,
                         processingTimeService,
@@ -1083,6 +1113,7 @@ class InternalTimerServiceImplTest {
     }
 
     private static <K, N> InternalTimerServiceImpl<K, N> 
createInternalTimerService(
+            TaskIOMetricGroup taskIOMetricGroup,
             KeyGroupRange keyGroupsList,
             KeyContext keyContext,
             ProcessingTimeService processingTimeService,
@@ -1094,6 +1125,7 @@ class InternalTimerServiceImplTest {
                 new TimerSerializer<>(keySerializer, namespaceSerializer);
 
         return new InternalTimerServiceImpl<>(
+                taskIOMetricGroup,
                 keyGroupsList,
                 keyContext,
                 processingTimeService,
diff --git 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StateInitializationContextImplTest.java
 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StateInitializationContextImplTest.java
index f1cf2ab3135..28498824479 100644
--- 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StateInitializationContextImplTest.java
+++ 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StateInitializationContextImplTest.java
@@ -35,6 +35,7 @@ import 
org.apache.flink.runtime.checkpoint.StateObjectCollection;
 import org.apache.flink.runtime.checkpoint.SubTaskInitializationMetricsBuilder;
 import org.apache.flink.runtime.checkpoint.TaskStateSnapshot;
 import org.apache.flink.runtime.jobgraph.OperatorID;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.operators.testutils.DummyEnvironment;
 import org.apache.flink.runtime.state.CheckpointableKeyedStateBackend;
 import org.apache.flink.runtime.state.DefaultOperatorStateBackend;
@@ -194,6 +195,7 @@ class StateInitializationContextImplTest {
                         new InternalTimeServiceManager.Provider() {
                             @Override
                             public <K> InternalTimeServiceManager<K> create(
+                                    TaskIOMetricGroup taskIOMetricGroup,
                                     CheckpointableKeyedStateBackend<K> 
keyedStatedBackend,
                                     ClassLoader userClassloader,
                                     KeyContext keyContext,
diff --git 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImplTest.java
 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImplTest.java
index cc99c87908e..8e85ba875c4 100644
--- 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImplTest.java
+++ 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/StreamTaskStateInitializerImplTest.java
@@ -32,6 +32,7 @@ import 
org.apache.flink.runtime.checkpoint.metadata.CheckpointTestUtils;
 import org.apache.flink.runtime.executiongraph.ExecutionAttemptID;
 import org.apache.flink.runtime.jobgraph.OperatorID;
 import org.apache.flink.runtime.metrics.MetricNames;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.operators.testutils.DummyEnvironment;
 import org.apache.flink.runtime.state.AbstractKeyedStateBackend;
 import org.apache.flink.runtime.state.CheckpointableKeyedStateBackend;
@@ -343,6 +344,7 @@ class StreamTaskStateInitializerImplTest {
                     new InternalTimeServiceManager.Provider() {
                         @Override
                         public <K> InternalTimeServiceManager<K> create(
+                                TaskIOMetricGroup taskIOMetricGroup,
                                 CheckpointableKeyedStateBackend<K> 
keyedStatedBackend,
                                 ClassLoader userClassloader,
                                 KeyContext keyContext,
diff --git 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceTest.java
 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceTest.java
index 72ce7e4b6aa..36a64c8319c 100644
--- 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceTest.java
+++ 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/api/operators/sorted/state/BatchExecutionInternalTimeServiceTest.java
@@ -23,6 +23,7 @@ import org.apache.flink.api.common.JobID;
 import org.apache.flink.api.common.typeutils.base.IntSerializer;
 import org.apache.flink.core.fs.CloseableRegistry;
 import org.apache.flink.metrics.groups.UnregisteredMetricsGroup;
+import org.apache.flink.runtime.metrics.groups.UnregisteredMetricGroups;
 import org.apache.flink.runtime.operators.testutils.MockEnvironment;
 import org.apache.flink.runtime.query.TaskKvStateRegistry;
 import org.apache.flink.runtime.state.AbstractKeyedStateBackend;
@@ -87,6 +88,8 @@ class BatchExecutionInternalTimeServiceTest {
         assertThatThrownBy(
                         () ->
                                 
BatchExecutionInternalTimeServiceManager.create(
+                                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                                .getIOMetricGroup(),
                                         stateBackend,
                                         this.getClass().getClassLoader(),
                                         new DummyKeyContext(),
@@ -141,6 +144,8 @@ class BatchExecutionInternalTimeServiceTest {
                         KEY_SERIALIZER, new KeyGroupRange(0, 1), new 
ExecutionConfig());
         InternalTimeServiceManager<Integer> timeServiceManager =
                 BatchExecutionInternalTimeServiceManager.create(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyedStatedBackend,
                         this.getClass().getClassLoader(),
                         new DummyKeyContext(),
@@ -177,6 +182,8 @@ class BatchExecutionInternalTimeServiceTest {
                         KEY_SERIALIZER, new KeyGroupRange(0, 1), new 
ExecutionConfig());
         InternalTimeServiceManager<Integer> timeServiceManager =
                 BatchExecutionInternalTimeServiceManager.create(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyedStatedBackend,
                         this.getClass().getClassLoader(),
                         new DummyKeyContext(),
@@ -206,6 +213,8 @@ class BatchExecutionInternalTimeServiceTest {
                         KEY_SERIALIZER, new KeyGroupRange(0, 1), new 
ExecutionConfig());
         InternalTimeServiceManager<Integer> timeServiceManager =
                 BatchExecutionInternalTimeServiceManager.create(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyedStatedBackend,
                         this.getClass().getClassLoader(),
                         new DummyKeyContext(),
@@ -253,6 +262,8 @@ class BatchExecutionInternalTimeServiceTest {
         TestProcessingTimeService processingTimeService = new 
TestProcessingTimeService();
         InternalTimeServiceManager<Integer> timeServiceManager =
                 BatchExecutionInternalTimeServiceManager.create(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyedStatedBackend,
                         this.getClass().getClassLoader(),
                         new DummyKeyContext(),
@@ -288,6 +299,8 @@ class BatchExecutionInternalTimeServiceTest {
         TestProcessingTimeService processingTimeService = new 
TestProcessingTimeService();
         InternalTimeServiceManager<Integer> timeServiceManager =
                 BatchExecutionInternalTimeServiceManager.create(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyedStatedBackend,
                         this.getClass().getClassLoader(),
                         new DummyKeyContext(),
@@ -328,6 +341,8 @@ class BatchExecutionInternalTimeServiceTest {
         TestProcessingTimeService processingTimeService = new 
TestProcessingTimeService();
         InternalTimeServiceManager<Integer> timeServiceManager =
                 BatchExecutionInternalTimeServiceManager.create(
+                        
UnregisteredMetricGroups.createUnregisteredTaskMetricGroup()
+                                .getIOMetricGroup(),
                         keyedStatedBackend,
                         this.getClass().getClassLoader(),
                         new DummyKeyContext(),
diff --git 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/util/AbstractStreamOperatorTestHarness.java
 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/util/AbstractStreamOperatorTestHarness.java
index fc5a71cd61a..f31dee84ad8 100644
--- 
a/flink-streaming-java/src/test/java/org/apache/flink/streaming/util/AbstractStreamOperatorTestHarness.java
+++ 
b/flink-streaming-java/src/test/java/org/apache/flink/streaming/util/AbstractStreamOperatorTestHarness.java
@@ -37,6 +37,7 @@ import 
org.apache.flink.runtime.checkpoint.SubTaskInitializationMetricsBuilder;
 import org.apache.flink.runtime.checkpoint.TaskStateSnapshot;
 import org.apache.flink.runtime.execution.Environment;
 import org.apache.flink.runtime.jobgraph.OperatorID;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.operators.testutils.MockEnvironment;
 import org.apache.flink.runtime.operators.testutils.MockEnvironmentBuilder;
 import org.apache.flink.runtime.operators.testutils.MockInputSplitProvider;
@@ -149,6 +150,7 @@ public class AbstractStreamOperatorTestHarness<OUT> 
implements AutoCloseable {
             new InternalTimeServiceManager.Provider() {
                 @Override
                 public <K> InternalTimeServiceManager<K> create(
+                        TaskIOMetricGroup taskIOMetricGroup,
                         CheckpointableKeyedStateBackend<K> keyedStatedBackend,
                         ClassLoader userClassloader,
                         KeyContext keyContext,
@@ -158,6 +160,7 @@ public class AbstractStreamOperatorTestHarness<OUT> 
implements AutoCloseable {
                         throws Exception {
                     InternalTimeServiceManagerImpl<K> typedTimeServiceManager =
                             InternalTimeServiceManagerImpl.create(
+                                    taskIOMetricGroup,
                                     keyedStatedBackend,
                                     userClassloader,
                                     keyContext,
diff --git 
a/flink-tests/src/test/java/org/apache/flink/test/state/operator/restore/StreamOperatorSnapshotRestoreTest.java
 
b/flink-tests/src/test/java/org/apache/flink/test/state/operator/restore/StreamOperatorSnapshotRestoreTest.java
index a2e0c7e5c71..344e401dc02 100644
--- 
a/flink-tests/src/test/java/org/apache/flink/test/state/operator/restore/StreamOperatorSnapshotRestoreTest.java
+++ 
b/flink-tests/src/test/java/org/apache/flink/test/state/operator/restore/StreamOperatorSnapshotRestoreTest.java
@@ -33,6 +33,7 @@ import org.apache.flink.core.memory.DataOutputView;
 import org.apache.flink.core.memory.DataOutputViewStreamWrapper;
 import org.apache.flink.runtime.checkpoint.OperatorSubtaskState;
 import org.apache.flink.runtime.jobgraph.JobVertexID;
+import org.apache.flink.runtime.metrics.groups.TaskIOMetricGroup;
 import org.apache.flink.runtime.operators.testutils.MockEnvironment;
 import org.apache.flink.runtime.operators.testutils.MockEnvironmentBuilder;
 import org.apache.flink.runtime.operators.testutils.MockInputSplitProvider;
@@ -236,6 +237,7 @@ public class StreamOperatorSnapshotRestoreTest extends 
TestLogger {
                 new InternalTimeServiceManager.Provider() {
                     @Override
                     public <K> InternalTimeServiceManager<K> create(
+                            TaskIOMetricGroup taskIOMetricGroup,
                             CheckpointableKeyedStateBackend<K> 
keyedStatedBackend,
                             ClassLoader userClassloader,
                             KeyContext keyContext,


Reply via email to