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,
