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

zhuzh 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 a4470966295 [FLINK-33985][runtime] Support obtain all partitions 
existing in cluster through ShuffleMaster.
a4470966295 is described below

commit a44709662956b306fe686623d00358a6b076f637
Author: JunRuiLee <[email protected]>
AuthorDate: Wed Mar 13 16:40:53 2024 +0800

    [FLINK-33985][runtime] Support obtain all partitions existing in cluster 
through ShuffleMaster.
    
    This closes #24553.
---
 .../io/network/NettyShuffleEnvironment.java        |  7 ++
 .../io/network/partition/ResultPartition.java      |  4 ++
 .../network/partition/ResultPartitionManager.java  | 16 +++++
 .../partition/TaskExecutorPartitionTracker.java    |  8 +++
 .../TaskExecutorPartitionTrackerImpl.java          | 15 +++++
 .../apache/flink/runtime/jobmaster/JobMaster.java  | 24 +++++++
 .../flink/runtime/jobmaster/JobMasterGateway.java  |  7 ++
 ...ntext.java => DefaultPartitionWithMetrics.java} | 35 +++++-----
 ...ffleContext.java => DefaultShuffleMetrics.java} | 29 ++++----
 .../flink/runtime/shuffle/JobShuffleContext.java   |  6 ++
 .../runtime/shuffle/JobShuffleContextImpl.java     |  6 ++
 .../flink/runtime/shuffle/NettyShuffleMaster.java  | 22 ++++++
 ...uffleContext.java => PartitionWithMetrics.java} | 24 ++-----
 .../flink/runtime/shuffle/ShuffleEnvironment.java  | 13 ++++
 .../flink/runtime/shuffle/ShuffleMaster.java       | 13 ++++
 ...{JobShuffleContext.java => ShuffleMetrics.java} | 23 ++-----
 .../flink/runtime/taskexecutor/TaskExecutor.java   | 23 +++++++
 .../runtime/taskexecutor/TaskExecutorGateway.java  | 13 ++++
 .../taskexecutor/partition/PartitionTable.java     |  5 ++
 .../partition/ResultPartitionManagerTest.java      | 20 ++++++
 .../TaskExecutorPartitionTrackerImplTest.java      | 23 +++++++
 .../TestingTaskExecutorPartitionTracker.java       |  5 ++
 .../flink/runtime/jobmaster/JobMasterTest.java     | 78 ++++++++++++++++++++++
 .../taskexecutor/TaskExecutorSubmissionTest.java   | 73 ++++++++++++++++++++
 .../taskexecutor/TestingTaskExecutorGateway.java   | 14 ++++
 .../TestingTaskExecutorGatewayBuilder.java         | 13 ++++
 .../taskexecutor/partition/PartitionTableTest.java |  9 +++
 27 files changed, 458 insertions(+), 70 deletions(-)

diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/NettyShuffleEnvironment.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/NettyShuffleEnvironment.java
index f3ae8ff7d2b..c04678c3d4f 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/NettyShuffleEnvironment.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/NettyShuffleEnvironment.java
@@ -45,6 +45,7 @@ import 
org.apache.flink.runtime.shuffle.NettyShuffleDescriptor;
 import org.apache.flink.runtime.shuffle.ShuffleDescriptor;
 import org.apache.flink.runtime.shuffle.ShuffleEnvironment;
 import org.apache.flink.runtime.shuffle.ShuffleIOOwnerContext;
+import org.apache.flink.runtime.shuffle.ShuffleMetrics;
 import 
org.apache.flink.runtime.taskmanager.NettyShuffleEnvironmentConfiguration;
 import org.apache.flink.util.Preconditions;
 
@@ -198,6 +199,12 @@ public class NettyShuffleEnvironment
         return resultPartitionManager.getUnreleasedPartitions();
     }
 
+    @Override
+    public Optional<ShuffleMetrics> 
getMetricsIfPartitionOccupyingLocalResource(
+            ResultPartitionID partitionId) {
+        return resultPartitionManager.getMetricsOfPartition(partitionId);
+    }
+
     // 
--------------------------------------------------------------------------------------------
     //  Create Output Writers and Input Readers
     // 
--------------------------------------------------------------------------------------------
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartition.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartition.java
index 6361e07ce3a..6cbcfc0c598 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartition.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartition.java
@@ -215,6 +215,10 @@ public abstract class ResultPartition implements 
ResultPartitionWriter {
         return partitionType;
     }
 
+    public ResultPartitionBytesCounter getResultPartitionBytes() {
+        return resultPartitionBytes;
+    }
+
     // ------------------------------------------------------------------------
 
     @Override
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManager.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManager.java
index 5013cd7075f..7dd952c3286 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManager.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManager.java
@@ -19,6 +19,8 @@
 package org.apache.flink.runtime.io.network.partition;
 
 import org.apache.flink.annotation.VisibleForTesting;
+import org.apache.flink.runtime.shuffle.DefaultShuffleMetrics;
+import org.apache.flink.runtime.shuffle.ShuffleMetrics;
 import org.apache.flink.util.CollectionUtil;
 import org.apache.flink.util.concurrent.ScheduledExecutor;
 
@@ -293,4 +295,18 @@ public class ResultPartitionManager implements 
ResultPartitionProvider {
             return registeredPartitions.keySet();
         }
     }
+
+    public Optional<ShuffleMetrics> getMetricsOfPartition(ResultPartitionID 
partitionId) {
+        synchronized (registeredPartitions) {
+            final ResultPartition partition = 
registeredPartitions.get(partitionId);
+
+            if (partition == null) {
+                return Optional.empty();
+            }
+
+            return Optional.of(
+                    new DefaultShuffleMetrics(
+                            
partition.getResultPartitionBytes().createSnapshot()));
+        }
+    }
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTracker.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTracker.java
index 59d16cb51de..bb2246ab88c 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTracker.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTracker.java
@@ -44,6 +44,14 @@ public interface TaskExecutorPartitionTracker
      */
     void stopTrackingAndReleaseJobPartitionsFor(JobID producingJobId);
 
+    /**
+     * Get all partitions tracked for the given job.
+     *
+     * @param producingJobId the job id
+     * @return the tracked partitions
+     */
+    Collection<TaskExecutorPartitionInfo> getTrackedPartitionsFor(JobID 
producingJobId);
+
     /** Promotes the given partitions. */
     void promoteJobPartitions(Collection<ResultPartitionID> 
partitionsToPromote);
 
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImpl.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImpl.java
index 17f526268e0..8bd44b9801d 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImpl.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImpl.java
@@ -35,6 +35,8 @@ import java.util.Map;
 import java.util.Set;
 import java.util.stream.Collectors;
 
+import static java.util.stream.Collectors.toList;
+
 /**
  * Utility for tracking partitions and issuing release calls to task executors 
and shuffle masters.
  */
@@ -83,6 +85,19 @@ public class TaskExecutorPartitionTrackerImpl
         shuffleEnvironment.releasePartitionsLocally(partitionsForJob);
     }
 
+    @Override
+    public Collection<TaskExecutorPartitionInfo> getTrackedPartitionsFor(JobID 
producingJobId) {
+        return partitionTable.getTrackedPartitions(producingJobId).stream()
+                .map(
+                        partitionId -> {
+                            final PartitionInfo<JobID, 
TaskExecutorPartitionInfo> partitionInfo =
+                                    partitionInfos.get(partitionId);
+                            Preconditions.checkNotNull(partitionInfo);
+                            return partitionInfo.getMetaInfo();
+                        })
+                .collect(toList());
+    }
+
     @Override
     public void promoteJobPartitions(Collection<ResultPartitionID> 
partitionsToPromote) {
         LOG.debug("Promoting Job Partitions {}", partitionsToPromote);
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMaster.java 
b/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMaster.java
index 5cb409e2b9f..6c30fe35b63 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMaster.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMaster.java
@@ -90,6 +90,7 @@ import org.apache.flink.runtime.scheduler.ExecutionGraphInfo;
 import org.apache.flink.runtime.scheduler.SchedulerNG;
 import org.apache.flink.runtime.shuffle.JobShuffleContext;
 import org.apache.flink.runtime.shuffle.JobShuffleContextImpl;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.runtime.shuffle.ShuffleMaster;
 import org.apache.flink.runtime.slots.ResourceRequirement;
 import org.apache.flink.runtime.state.KeyGroupRange;
@@ -115,9 +116,11 @@ import javax.annotation.Nullable;
 
 import java.io.IOException;
 import java.net.InetSocketAddress;
+import java.util.ArrayList;
 import java.util.Collection;
 import java.util.HashMap;
 import java.util.HashSet;
+import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Optional;
@@ -957,6 +960,27 @@ public class JobMaster extends 
FencedRpcEndpoint<JobMasterId>
         return future;
     }
 
+    @Override
+    public CompletableFuture<Collection<PartitionWithMetrics>>
+            getAllPartitionWithMetricsOnTaskManagers() {
+        final List<CompletableFuture<Collection<PartitionWithMetrics>>> 
allFutures =
+                new ArrayList<>();
+        registeredTaskManagers
+                .values()
+                .forEach(
+                        taskManager ->
+                                allFutures.add(
+                                        taskManager
+                                                .getTaskExecutorGateway()
+                                                
.getPartitionWithMetrics(jobGraph.getJobID())));
+        return FutureUtils.combineAll(allFutures)
+                .thenApply(
+                        partitions ->
+                                partitions.stream()
+                                        .flatMap(Collection::stream)
+                                        .collect(Collectors.toList()));
+    }
+
     @Override
     public CompletableFuture<Acknowledge> 
notifyNewBlockedNodes(Collection<BlockedNode> newNodes) {
         blocklistHandler.addNewBlockedNodes(newNodes);
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMasterGateway.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMasterGateway.java
index 02c3c7d501a..994492807ff 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMasterGateway.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/jobmaster/JobMasterGateway.java
@@ -46,6 +46,7 @@ import 
org.apache.flink.runtime.resourcemanager.ResourceManagerId;
 import org.apache.flink.runtime.rpc.FencedRpcGateway;
 import org.apache.flink.runtime.rpc.RpcTimeout;
 import org.apache.flink.runtime.scheduler.ExecutionGraphInfo;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.runtime.slots.ResourceRequirement;
 import 
org.apache.flink.runtime.taskexecutor.TaskExecutorToJobManagerHeartbeatPayload;
 import org.apache.flink.runtime.taskexecutor.slot.SlotOffer;
@@ -55,6 +56,7 @@ import org.apache.flink.util.SerializedValue;
 import javax.annotation.Nullable;
 
 import java.util.Collection;
+import java.util.Collections;
 import java.util.concurrent.CompletableFuture;
 
 /** {@link JobMaster} rpc gateway interface. */
@@ -301,6 +303,11 @@ public interface JobMasterGateway
     CompletableFuture<?> stopTrackingAndReleasePartitions(
             Collection<ResultPartitionID> partitionIds);
 
+    default CompletableFuture<Collection<PartitionWithMetrics>>
+            getAllPartitionWithMetricsOnTaskManagers() {
+        return CompletableFuture.completedFuture(Collections.emptyList());
+    }
+
     /**
      * Read current {@link JobResourceRequirements job resource requirements}.
      *
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/DefaultPartitionWithMetrics.java
similarity index 51%
copy from 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
copy to 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/DefaultPartitionWithMetrics.java
index 3cfa996f0c1..d50a2b5acb1 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/DefaultPartitionWithMetrics.java
@@ -18,25 +18,26 @@
 
 package org.apache.flink.runtime.shuffle;
 
-import org.apache.flink.api.common.JobID;
-import org.apache.flink.runtime.io.network.partition.ResultPartitionID;
+import static org.apache.flink.util.Preconditions.checkNotNull;
 
-import java.util.Collection;
-import java.util.concurrent.CompletableFuture;
+/** Default {@link PartitionWithMetrics} implementation. */
+public class DefaultPartitionWithMetrics implements PartitionWithMetrics {
+    private final ShuffleDescriptor shuffleDescriptor;
+    private final ShuffleMetrics partitionMetrics;
 
-/**
- * Job level shuffle context which can offer some job information like job ID 
and through it, the
- * shuffle plugin notify the job to stop tracking the lost result partitions.
- */
-public interface JobShuffleContext {
+    public DefaultPartitionWithMetrics(
+            ShuffleDescriptor shuffleDescriptor, ShuffleMetrics 
partitionMetrics) {
+        this.shuffleDescriptor = checkNotNull(shuffleDescriptor);
+        this.partitionMetrics = checkNotNull(partitionMetrics);
+    }
 
-    /** @return the corresponding {@link JobID}. */
-    JobID getJobId();
+    @Override
+    public ShuffleMetrics getPartitionMetrics() {
+        return partitionMetrics;
+    }
 
-    /**
-     * Notifies the job to stop tracking and release the target result 
partitions, which means these
-     * partitions will be removed and will be reproduced if used afterwards.
-     */
-    CompletableFuture<?> stopTrackingAndReleasePartitions(
-            Collection<ResultPartitionID> partitionIds);
+    @Override
+    public ShuffleDescriptor getPartition() {
+        return shuffleDescriptor;
+    }
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/DefaultShuffleMetrics.java
similarity index 51%
copy from 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
copy to 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/DefaultShuffleMetrics.java
index 3cfa996f0c1..8d3aeb3bc8e 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/DefaultShuffleMetrics.java
@@ -18,25 +18,20 @@
 
 package org.apache.flink.runtime.shuffle;
 
-import org.apache.flink.api.common.JobID;
-import org.apache.flink.runtime.io.network.partition.ResultPartitionID;
+import org.apache.flink.runtime.executiongraph.ResultPartitionBytes;
 
-import java.util.Collection;
-import java.util.concurrent.CompletableFuture;
+import static org.apache.flink.util.Preconditions.checkNotNull;
 
-/**
- * Job level shuffle context which can offer some job information like job ID 
and through it, the
- * shuffle plugin notify the job to stop tracking the lost result partitions.
- */
-public interface JobShuffleContext {
+/** Default {@link ShuffleMetrics} implementation. */
+public class DefaultShuffleMetrics implements ShuffleMetrics {
+    private final ResultPartitionBytes partitionBytes;
 
-    /** @return the corresponding {@link JobID}. */
-    JobID getJobId();
+    public DefaultShuffleMetrics(ResultPartitionBytes partitionBytes) {
+        this.partitionBytes = checkNotNull(partitionBytes);
+    }
 
-    /**
-     * Notifies the job to stop tracking and release the target result 
partitions, which means these
-     * partitions will be removed and will be reproduced if used afterwards.
-     */
-    CompletableFuture<?> stopTrackingAndReleasePartitions(
-            Collection<ResultPartitionID> partitionIds);
+    @Override
+    public ResultPartitionBytes getPartitionBytes() {
+        return partitionBytes;
+    }
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
index 3cfa996f0c1..585b73806d9 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
@@ -39,4 +39,10 @@ public interface JobShuffleContext {
      */
     CompletableFuture<?> stopTrackingAndReleasePartitions(
             Collection<ResultPartitionID> partitionIds);
+
+    /**
+     * Retrieves a collection containing descriptions and metrics of existing 
result partitions from
+     * all TaskManagers.
+     */
+    CompletableFuture<Collection<PartitionWithMetrics>> 
getAllPartitionWithMetricsOnTaskManagers();
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContextImpl.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContextImpl.java
index d05689706d2..546cc3025b4 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContextImpl.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContextImpl.java
@@ -49,4 +49,10 @@ public class JobShuffleContextImpl implements 
JobShuffleContext {
             Collection<ResultPartitionID> partitionIds) {
         return jobMasterGateway.stopTrackingAndReleasePartitions(partitionIds);
     }
+
+    @Override
+    public CompletableFuture<Collection<PartitionWithMetrics>>
+            getAllPartitionWithMetricsOnTaskManagers() {
+        return jobMasterGateway.getAllPartitionWithMetricsOnTaskManagers();
+    }
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/NettyShuffleMaster.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/NettyShuffleMaster.java
index 6dbea9e16cb..9cce16cf495 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/NettyShuffleMaster.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/NettyShuffleMaster.java
@@ -31,6 +31,9 @@ import org.apache.flink.runtime.util.ConfigurationParserUtils;
 
 import javax.annotation.Nullable;
 
+import java.util.Collection;
+import java.util.HashMap;
+import java.util.Map;
 import java.util.Optional;
 import java.util.concurrent.CompletableFuture;
 
@@ -58,6 +61,8 @@ public class NettyShuffleMaster implements 
ShuffleMaster<NettyShuffleDescriptor>
 
     @Nullable private final TieredInternalShuffleMaster 
tieredInternalShuffleMaster;
 
+    private final Map<JobID, JobShuffleContext> jobShuffleContexts = new 
HashMap<>();
+
     public NettyShuffleMaster(Configuration conf) {
         checkNotNull(conf);
         buffersPerInputChannel =
@@ -165,4 +170,21 @@ public class NettyShuffleMaster implements 
ShuffleMaster<NettyShuffleDescriptor>
                         || conf.get(BATCH_SHUFFLE_MODE) == 
ALL_EXCHANGES_HYBRID_SELECTIVE)
                 && conf.get(NETWORK_HYBRID_SHUFFLE_ENABLE_NEW_MODE);
     }
+
+    @Override
+    public CompletableFuture<Collection<PartitionWithMetrics>> 
getAllPartitionWithMetrics(
+            JobID jobId) {
+        return checkNotNull(jobShuffleContexts.get(jobId))
+                .getAllPartitionWithMetricsOnTaskManagers();
+    }
+
+    @Override
+    public void registerJob(JobShuffleContext context) {
+        jobShuffleContexts.put(context.getJobId(), context);
+    }
+
+    @Override
+    public void unregisterJob(JobID jobId) {
+        jobShuffleContexts.remove(jobId);
+    }
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/PartitionWithMetrics.java
similarity index 51%
copy from 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
copy to 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/PartitionWithMetrics.java
index 3cfa996f0c1..872f092c6e9 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/PartitionWithMetrics.java
@@ -18,25 +18,11 @@
 
 package org.apache.flink.runtime.shuffle;
 
-import org.apache.flink.api.common.JobID;
-import org.apache.flink.runtime.io.network.partition.ResultPartitionID;
+import java.io.Serializable;
 
-import java.util.Collection;
-import java.util.concurrent.CompletableFuture;
+/** Interface representing the description and metrics of a result partition. 
*/
+public interface PartitionWithMetrics extends Serializable {
+    ShuffleMetrics getPartitionMetrics();
 
-/**
- * Job level shuffle context which can offer some job information like job ID 
and through it, the
- * shuffle plugin notify the job to stop tracking the lost result partitions.
- */
-public interface JobShuffleContext {
-
-    /** @return the corresponding {@link JobID}. */
-    JobID getJobId();
-
-    /**
-     * Notifies the job to stop tracking and release the target result 
partitions, which means these
-     * partitions will be removed and will be reproduced if used afterwards.
-     */
-    CompletableFuture<?> stopTrackingAndReleasePartitions(
-            Collection<ResultPartitionID> partitionIds);
+    ShuffleDescriptor getPartition();
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleEnvironment.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleEnvironment.java
index 4914a011447..7811f7f75fe 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleEnvironment.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleEnvironment.java
@@ -33,6 +33,7 @@ import 
org.apache.flink.runtime.io.network.partition.consumer.InputGate;
 import java.io.IOException;
 import java.util.Collection;
 import java.util.List;
+import java.util.Optional;
 
 /**
  * Interface for the implementation of shuffle service local environment.
@@ -162,6 +163,18 @@ public interface ShuffleEnvironment<P extends 
ResultPartitionWriter, G extends I
      */
     Collection<ResultPartitionID> getPartitionsOccupyingLocalResources();
 
+    /**
+     * Get metrics of the partition if it still occupies some resources 
locally and have not been
+     * released yet.
+     *
+     * @param partitionId the partition id
+     * @return An Optional of {@link ShuffleMetrics}, if found, of the given 
partition
+     */
+    default Optional<ShuffleMetrics> 
getMetricsIfPartitionOccupyingLocalResource(
+            ResultPartitionID partitionId) {
+        return Optional.empty();
+    }
+
     /**
      * Factory method for the {@link InputGate InputGates} to consume result 
partitions.
      *
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleMaster.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleMaster.java
index 82aa4530d57..e459571201a 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleMaster.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleMaster.java
@@ -22,6 +22,7 @@ import org.apache.flink.api.common.JobID;
 import org.apache.flink.configuration.MemorySize;
 
 import java.util.Collection;
+import java.util.Collections;
 import java.util.concurrent.CompletableFuture;
 
 /**
@@ -108,4 +109,16 @@ public interface ShuffleMaster<T extends 
ShuffleDescriptor> extends AutoCloseabl
             TaskInputsOutputsDescriptor taskInputsOutputsDescriptor) {
         return MemorySize.ZERO;
     }
+
+    /**
+     * Get all partitions and their metrics, the metrics include sizes of 
sub-partitions in a result
+     * partition.
+     *
+     * @param jobId ID of the target job
+     * @return All partitions belong to the target job and their metrics
+     */
+    default CompletableFuture<Collection<PartitionWithMetrics>> 
getAllPartitionWithMetrics(
+            JobID jobId) {
+        return CompletableFuture.completedFuture(Collections.emptyList());
+    }
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleMetrics.java
similarity index 52%
copy from 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
copy to 
flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleMetrics.java
index 3cfa996f0c1..23e4089fda1 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/JobShuffleContext.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/shuffle/ShuffleMetrics.java
@@ -18,25 +18,14 @@
 
 package org.apache.flink.runtime.shuffle;
 
-import org.apache.flink.api.common.JobID;
-import org.apache.flink.runtime.io.network.partition.ResultPartitionID;
+import org.apache.flink.runtime.executiongraph.ResultPartitionBytes;
 
-import java.util.Collection;
-import java.util.concurrent.CompletableFuture;
+import java.io.Serializable;
 
 /**
- * Job level shuffle context which can offer some job information like job ID 
and through it, the
- * shuffle plugin notify the job to stop tracking the lost result partitions.
+ * Interface provides access to the shuffle metrics which includes the meta 
information of
+ * partition(partition bytes, etc).
  */
-public interface JobShuffleContext {
-
-    /** @return the corresponding {@link JobID}. */
-    JobID getJobId();
-
-    /**
-     * Notifies the job to stop tracking and release the target result 
partitions, which means these
-     * partitions will be removed and will be reproduced if used afterwards.
-     */
-    CompletableFuture<?> stopTrackingAndReleasePartitions(
-            Collection<ResultPartitionID> partitionIds);
+public interface ShuffleMetrics extends Serializable {
+    ResultPartitionBytes getPartitionBytes();
 }
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java
index 8be2e04d24a..99b4ca7d370 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutor.java
@@ -99,6 +99,8 @@ import org.apache.flink.runtime.rpc.RpcEndpoint;
 import org.apache.flink.runtime.rpc.RpcService;
 import org.apache.flink.runtime.rpc.RpcServiceUtils;
 import 
org.apache.flink.runtime.security.token.DelegationTokenReceiverRepository;
+import org.apache.flink.runtime.shuffle.DefaultPartitionWithMetrics;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.runtime.shuffle.ShuffleDescriptor;
 import org.apache.flink.runtime.shuffle.ShuffleEnvironment;
 import 
org.apache.flink.runtime.state.TaskExecutorChannelStateExecutorFactoryManager;
@@ -1441,6 +1443,27 @@ public class TaskExecutor extends RpcEndpoint implements 
TaskExecutorGateway {
         }
     }
 
+    @Override
+    public CompletableFuture<Collection<PartitionWithMetrics>> 
getPartitionWithMetrics(
+            JobID jobId) {
+        Collection<TaskExecutorPartitionInfo> partitionInfoList =
+                partitionTracker.getTrackedPartitionsFor(jobId);
+        List<PartitionWithMetrics> partitionWithMetrics = new ArrayList<>();
+        partitionInfoList.forEach(
+                info -> {
+                    ResultPartitionID partitionId = 
info.getResultPartitionId();
+                    shuffleEnvironment
+                            
.getMetricsIfPartitionOccupyingLocalResource(partitionId)
+                            .ifPresent(
+                                    metrics ->
+                                            partitionWithMetrics.add(
+                                                    new 
DefaultPartitionWithMetrics(
+                                                            
info.getShuffleDescriptor(), metrics)));
+                });
+
+        return CompletableFuture.completedFuture(partitionWithMetrics);
+    }
+
     @Override
     public CompletableFuture<ProfilingInfo> requestProfiling(
             int duration, ProfilingInfo.ProfilingMode mode, Duration timeout) {
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorGateway.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorGateway.java
index d3e9dd119a5..f5686238650 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorGateway.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/TaskExecutorGateway.java
@@ -43,6 +43,7 @@ import org.apache.flink.runtime.rest.messages.ProfilingInfo;
 import org.apache.flink.runtime.rest.messages.ThreadDumpInfo;
 import org.apache.flink.runtime.rpc.RpcGateway;
 import org.apache.flink.runtime.rpc.RpcTimeout;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.runtime.taskmanager.Task;
 import org.apache.flink.types.SerializableOptional;
 import org.apache.flink.util.SerializedValue;
@@ -316,6 +317,18 @@ public interface TaskExecutorGateway
     CompletableFuture<Acknowledge> updateDelegationTokens(
             ResourceManagerId resourceManagerId, byte[] tokens);
 
+    /**
+     * Get all partitions and their metrics located on this task executor, the 
metrics mainly
+     * includes the meta information of partition(partition bytes, etc).
+     *
+     * @param jobId ID of the target job
+     * @return All partitions belong to the target job and their metrics
+     */
+    default CompletableFuture<Collection<PartitionWithMetrics>> 
getPartitionWithMetrics(
+            JobID jobId) {
+        throw new UnsupportedOperationException();
+    }
+
     /**
      * Requests the profiling from this TaskManager.
      *
diff --git 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTable.java
 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTable.java
index b1db983f695..d8d9a176c9e 100644
--- 
a/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTable.java
+++ 
b/flink-runtime/src/main/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTable.java
@@ -41,6 +41,11 @@ public class PartitionTable<K> {
         return trackedPartitionsPerKey.containsKey(key);
     }
 
+    public Collection<ResultPartitionID> getTrackedPartitions(K key) {
+        Preconditions.checkNotNull(key);
+        return trackedPartitionsPerKey.getOrDefault(key, 
Collections.emptySet());
+    }
+
     /** Starts the tracking of the given partition for the given key. */
     public void startTrackingPartitions(K key, Collection<ResultPartitionID> 
newPartitionIds) {
         Preconditions.checkNotNull(key);
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManagerTest.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManagerTest.java
index 69923490ace..6e9c16ea807 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManagerTest.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/ResultPartitionManagerTest.java
@@ -20,10 +20,13 @@ package org.apache.flink.runtime.io.network.partition;
 
 import org.apache.flink.runtime.io.network.netty.NettyPartitionRequestListener;
 import org.apache.flink.runtime.io.network.partition.consumer.InputChannelID;
+import org.apache.flink.runtime.shuffle.ShuffleMetrics;
 import org.apache.flink.util.concurrent.ManuallyTriggeredScheduledExecutor;
 
 import org.junit.jupiter.api.Test;
 
+import java.util.Arrays;
+import java.util.Optional;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.TimeUnit;
 
@@ -158,6 +161,23 @@ class ResultPartitionManagerTest {
         verifyCreateSubpartitionViewThrowsException(partitionManager, 
partition.getPartitionId());
     }
 
+    @Test
+    void testGetMetricsOfPartition() throws Exception {
+        final ResultPartitionManager partitionManager = new 
ResultPartitionManager();
+        final ResultPartition partition = createPartition();
+        partition.resultPartitionBytes.incAll(100);
+
+        partitionManager.registerResultPartition(partition);
+        Optional<ShuffleMetrics> metricsOfPartition =
+                partitionManager.getMetricsOfPartition(partition.partitionId);
+        assertThat(metricsOfPartition)
+                .hasValueSatisfying(
+                        metrics ->
+                                Arrays.equals(
+                                        
metrics.getPartitionBytes().getSubpartitionBytes(),
+                                        new long[] {100}));
+    }
+
     /** Test notifier timeout in {@link ResultPartitionManager}. */
     @Test
     void testCreateViewReaderForNotifierTimeout() throws Exception {
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImplTest.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImplTest.java
index 225062e221e..13749d99352 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImplTest.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TaskExecutorPartitionTrackerImplTest.java
@@ -153,6 +153,29 @@ class TaskExecutorPartitionTrackerImplTest {
                 .contains(resultPartitionId1);
     }
 
+    @Test
+    void testGetTrackedPartitionsFor() {
+        final TestingShuffleEnvironment testingShuffleEnvironment = new 
TestingShuffleEnvironment();
+
+        final JobID jobId = new JobID();
+        final ResultPartitionID resultPartitionId = new ResultPartitionID();
+
+        final TaskExecutorPartitionTracker partitionTracker =
+                new 
TaskExecutorPartitionTrackerImpl(testingShuffleEnvironment);
+        TaskExecutorPartitionInfo partitionInfo =
+                new TaskExecutorPartitionInfo(
+                        new TestingShuffleDescriptor(resultPartitionId),
+                        new IntermediateDataSetID(),
+                        1);
+
+        partitionTracker.startTrackingPartition(jobId, partitionInfo);
+        Collection<TaskExecutorPartitionInfo> partitions =
+                partitionTracker.getTrackedPartitionsFor(jobId);
+
+        assertThat(partitions).hasSize(1);
+        assertThat(partitions.iterator().next()).isEqualTo(partitionInfo);
+    }
+
     @Test
     void promoteJobPartitions() throws Exception {
         final TestingShuffleEnvironment testingShuffleEnvironment = new 
TestingShuffleEnvironment();
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TestingTaskExecutorPartitionTracker.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TestingTaskExecutorPartitionTracker.java
index cfd0fd86e82..8816125f3c7 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TestingTaskExecutorPartitionTracker.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/io/network/partition/TestingTaskExecutorPartitionTracker.java
@@ -115,6 +115,11 @@ public class TestingTaskExecutorPartitionTracker 
implements TaskExecutorPartitio
         stopTrackingAndReleaseAllPartitionsConsumer.accept(producingJobId);
     }
 
+    @Override
+    public Collection<TaskExecutorPartitionInfo> getTrackedPartitionsFor(JobID 
producingJobId) {
+        return Collections.emptyList();
+    }
+
     @Override
     public void promoteJobPartitions(Collection<ResultPartitionID> 
partitionsToPromote) {
         promotePartitionsConsumer.accept(partitionsToPromote);
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTest.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTest.java
index e8e63652283..bc03aa604f7 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTest.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/jobmaster/JobMasterTest.java
@@ -59,6 +59,7 @@ import 
org.apache.flink.runtime.executiongraph.AccessExecution;
 import org.apache.flink.runtime.executiongraph.AccessExecutionVertex;
 import org.apache.flink.runtime.executiongraph.ArchivedExecutionGraph;
 import org.apache.flink.runtime.executiongraph.ExecutionAttemptID;
+import org.apache.flink.runtime.executiongraph.ResultPartitionBytes;
 import 
org.apache.flink.runtime.executiongraph.failover.FailoverStrategyFactoryLoader;
 import org.apache.flink.runtime.heartbeat.HeartbeatServices;
 import org.apache.flink.runtime.heartbeat.HeartbeatServicesImpl;
@@ -103,6 +104,10 @@ import 
org.apache.flink.runtime.scheduler.ExecutionGraphInfo;
 import org.apache.flink.runtime.scheduler.SchedulerTestingUtils;
 import org.apache.flink.runtime.scheduler.TestingSchedulerNG;
 import org.apache.flink.runtime.scheduler.TestingSchedulerNGFactory;
+import org.apache.flink.runtime.shuffle.DefaultPartitionWithMetrics;
+import org.apache.flink.runtime.shuffle.DefaultShuffleMetrics;
+import org.apache.flink.runtime.shuffle.NettyShuffleDescriptor;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.runtime.state.CompletedCheckpointStorageLocation;
 import org.apache.flink.runtime.state.StreamStateHandle;
 import org.apache.flink.runtime.taskexecutor.TaskExecutorGateway;
@@ -115,6 +120,7 @@ import 
org.apache.flink.runtime.taskmanager.TaskManagerLocation;
 import org.apache.flink.runtime.taskmanager.UnresolvedTaskManagerLocation;
 import org.apache.flink.runtime.testtasks.NoOpInvokable;
 import org.apache.flink.runtime.testutils.CommonTestUtils;
+import org.apache.flink.runtime.util.NettyShuffleDescriptorBuilder;
 import org.apache.flink.runtime.util.TestingFatalErrorHandler;
 import org.apache.flink.testutils.TestingUtils;
 import org.apache.flink.util.FlinkException;
@@ -1902,6 +1908,78 @@ class JobMasterTest {
         }
     }
 
+    @Test
+    void testGetAllPartitionWithMetrics() throws Exception {
+        JobVertex jobVertex = new JobVertex("jobVertex");
+        jobVertex.setInvokableClass(NoOpInvokable.class);
+        jobVertex.setParallelism(1);
+        final JobGraph jobGraph = JobGraphTestUtils.batchJobGraph(jobVertex);
+
+        try (final JobMaster jobMaster =
+                new JobMasterBuilder(jobGraph, rpcService)
+                        .withConfiguration(configuration)
+                        .withHighAvailabilityServices(haServices)
+                        .withHeartbeatServices(heartbeatServices)
+                        .withBlocklistHandlerFactory(
+                                new 
DefaultBlocklistHandler.Factory(Duration.ofMillis((100L))))
+                        .createJobMaster()) {
+
+            jobMaster.start();
+
+            final JobMasterGateway jobMasterGateway =
+                    jobMaster.getSelfGateway(JobMasterGateway.class);
+
+            NettyShuffleDescriptor shuffleDescriptor =
+                    NettyShuffleDescriptorBuilder.newBuilder().buildLocal();
+            DefaultShuffleMetrics shuffleMetrics =
+                    new DefaultShuffleMetrics(new ResultPartitionBytes(new 
long[] {1, 2, 3}));
+            Collection<PartitionWithMetrics> defaultPartitionWithMetrics =
+                    Collections.singletonList(
+                            new DefaultPartitionWithMetrics(shuffleDescriptor, 
shuffleMetrics));
+            final TestingTaskExecutorGateway taskExecutorGateway =
+                    new TestingTaskExecutorGatewayBuilder()
+                            .setRequestPartitionWithMetricsFunction(
+                                    ignored ->
+                                            CompletableFuture.completedFuture(
+                                                    
defaultPartitionWithMetrics))
+                            .createTestingTaskExecutorGateway();
+
+            final LocalUnresolvedTaskManagerLocation 
taskManagerUnresolvedLocation =
+                    new LocalUnresolvedTaskManagerLocation();
+            final Collection<SlotOffer> slotOffers =
+                    registerSlotsAtJobMaster(
+                            1,
+                            jobMasterGateway,
+                            jobGraph.getJobID(),
+                            taskExecutorGateway,
+                            taskManagerUnresolvedLocation);
+            assertThat(slotOffers).hasSize(1);
+
+            waitUntilAllExecutionsAreScheduledOrDeployed(jobMasterGateway);
+
+            PartitionWithMetrics metrics =
+                    jobMasterGateway
+                            .getAllPartitionWithMetricsOnTaskManagers()
+                            .get()
+                            .iterator()
+                            .next();
+            PartitionWithMetrics expectedMetrics = 
defaultPartitionWithMetrics.iterator().next();
+
+            
assertThat(metrics.getPartitionMetrics().getPartitionBytes().getSubpartitionBytes())
+                    .isEqualTo(
+                            expectedMetrics
+                                    .getPartitionMetrics()
+                                    .getPartitionBytes()
+                                    .getSubpartitionBytes());
+            assertThat(metrics.getPartition().getResultPartitionID())
+                    
.isEqualTo(expectedMetrics.getPartition().getResultPartitionID());
+            assertThat(metrics.getPartition().isUnknown())
+                    .isEqualTo(expectedMetrics.getPartition().isUnknown());
+            assertThat(metrics.getPartition().storesLocalResourcesOn())
+                    
.isEqualTo(expectedMetrics.getPartition().storesLocalResourcesOn());
+        }
+    }
+
     @Test
     void testBlockResourcesWillTriggerReleaseFreeSlots() throws Exception {
         JobVertex jobVertex = new JobVertex("jobVertex");
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java
index e4a3b889bfe..8b7070175dd 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TaskExecutorSubmissionTest.java
@@ -50,6 +50,7 @@ import org.apache.flink.runtime.messages.Acknowledge;
 import org.apache.flink.runtime.shuffle.NettyShuffleDescriptor;
 import org.apache.flink.runtime.shuffle.PartitionDescriptor;
 import org.apache.flink.runtime.shuffle.PartitionDescriptorBuilder;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.runtime.shuffle.ShuffleEnvironment;
 import org.apache.flink.runtime.taskexecutor.slot.TaskSlotTable;
 import org.apache.flink.runtime.taskmanager.Task;
@@ -327,6 +328,78 @@ class TaskExecutorSubmissionTest {
         }
     }
 
+    @Test
+    void testGetPartitionWithMetrics() throws Exception {
+        ResourceID producerLocation = ResourceID.generate();
+        NettyShuffleDescriptor sdd =
+                createRemoteWithIdAndLocation(
+                        new IntermediateResultPartitionID(), producerLocation);
+
+        PartitionDescriptor partitionDescriptor =
+                PartitionDescriptorBuilder.newBuilder()
+                        
.setPartitionId(sdd.getResultPartitionID().getPartitionId())
+                        .setPartitionType(ResultPartitionType.BLOCKING)
+                        .build();
+
+        ResultPartitionDeploymentDescriptor 
resultPartitionDeploymentDescriptor =
+                new ResultPartitionDeploymentDescriptor(partitionDescriptor, 
sdd, 1);
+        TaskDeploymentDescriptor tdd =
+                createTestTaskDeploymentDescriptor(
+                        "task",
+                        sdd.getResultPartitionID().getProducerId(),
+                        TestingAbstractInvokables.Sender.class,
+                        1,
+                        
Collections.singletonList(resultPartitionDeploymentDescriptor),
+                        Collections.emptyList());
+
+        ExecutionAttemptID eid = tdd.getExecutionAttemptId();
+
+        final CompletableFuture<Void> taskFinishedFuture = new 
CompletableFuture<>();
+
+        final JobMasterId jobMasterId = JobMasterId.generate();
+        TestingJobMasterGateway testingJobMasterGateway =
+                new TestingJobMasterGatewayBuilder()
+                        .setFencingTokenSupplier(() -> jobMasterId)
+                        .build();
+
+        try (TaskSubmissionTestEnvironment env =
+                new TaskSubmissionTestEnvironment.Builder(jobId)
+                        .setResourceID(producerLocation)
+                        .setSlotSize(1)
+                        .addTaskManagerActionListener(
+                                eid, ExecutionState.FINISHED, 
taskFinishedFuture)
+                        .setJobMasterId(jobMasterId)
+                        .setJobMasterGateway(testingJobMasterGateway)
+                        .useRealNonMockShuffleEnvironment()
+                        .build(EXECUTOR_EXTENSION.getExecutor())) {
+            TaskExecutorGateway tmGateway = env.getTaskExecutorGateway();
+            TaskSlotTable<Task> taskSlotTable = env.getTaskSlotTable();
+
+            taskSlotTable.allocateSlot(0, jobId, tdd.getAllocationId(), 
Duration.ofSeconds(60));
+            tmGateway.submitTask(tdd, jobMasterId, timeout).get();
+
+            taskFinishedFuture.get();
+
+            Collection<PartitionWithMetrics> partitionWithMetricsCollection =
+                    tmGateway.getPartitionWithMetrics(jobId).get();
+            assertThat(partitionWithMetricsCollection.size()).isOne();
+            PartitionWithMetrics partitionWithMetrics =
+                    partitionWithMetricsCollection.iterator().next();
+            
assertThat(partitionWithMetrics.getPartition().getResultPartitionID())
+                    .isEqualTo(sdd.getResultPartitionID());
+            
assertThat(partitionWithMetrics.getPartition().isUnknown()).isEqualTo(sdd.isUnknown());
+            
assertThat(partitionWithMetrics.getPartition().storesLocalResourcesOn())
+                    .isEqualTo(sdd.storesLocalResourcesOn());
+            assertThat(
+                            partitionWithMetrics
+                                    .getPartitionMetrics()
+                                    .getPartitionBytes()
+                                    .getSubpartitionBytes())
+                    // the sender task will send two int value, the expected 
bytes is 16.
+                    .isEqualTo(new long[] {16});
+        }
+    }
+
     /**
      * This tests creates two tasks. The sender sends data but fails to send 
the state update back
      * to the job manager. the second one blocks to be canceled
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
index 765206b6728..7795be6f52c 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGateway.java
@@ -42,6 +42,7 @@ import 
org.apache.flink.runtime.resourcemanager.ResourceManagerId;
 import org.apache.flink.runtime.rest.messages.LogInfo;
 import org.apache.flink.runtime.rest.messages.ProfilingInfo;
 import org.apache.flink.runtime.rest.messages.ThreadDumpInfo;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.runtime.webmonitor.threadinfo.ThreadInfoSamplesRequest;
 import org.apache.flink.types.SerializableOptional;
 import org.apache.flink.util.Preconditions;
@@ -80,6 +81,9 @@ public class TestingTaskExecutorGateway implements 
TaskExecutorGateway {
                     CompletableFuture<Acknowledge>>
             requestSlotFunction;
 
+    private final Function<JobID, 
CompletableFuture<Collection<PartitionWithMetrics>>>
+            requestPartitionWithMetricsFunction;
+
     private final BiFunction<AllocationID, Throwable, 
CompletableFuture<Acknowledge>>
             freeSlotFunction;
 
@@ -140,6 +144,8 @@ public class TestingTaskExecutorGateway implements 
TaskExecutorGateway {
                                     ResourceManagerId>,
                             CompletableFuture<Acknowledge>>
                     requestSlotFunction,
+            Function<JobID, 
CompletableFuture<Collection<PartitionWithMetrics>>>
+                    requestPartitionWithMetricsFunction,
             BiFunction<AllocationID, Throwable, 
CompletableFuture<Acknowledge>> freeSlotFunction,
             Consumer<JobID> freeInactiveSlotsConsumer,
             Function<ResourceID, CompletableFuture<Void>> 
heartbeatResourceManagerFunction,
@@ -175,6 +181,8 @@ public class TestingTaskExecutorGateway implements 
TaskExecutorGateway {
                 Preconditions.checkNotNull(disconnectJobManagerConsumer);
         this.submitTaskConsumer = 
Preconditions.checkNotNull(submitTaskConsumer);
         this.requestSlotFunction = 
Preconditions.checkNotNull(requestSlotFunction);
+        this.requestPartitionWithMetricsFunction =
+                
Preconditions.checkNotNull(requestPartitionWithMetricsFunction);
         this.freeSlotFunction = Preconditions.checkNotNull(freeSlotFunction);
         this.freeInactiveSlotsConsumer = 
Preconditions.checkNotNull(freeInactiveSlotsConsumer);
         this.heartbeatResourceManagerFunction = 
heartbeatResourceManagerFunction;
@@ -211,6 +219,12 @@ public class TestingTaskExecutorGateway implements 
TaskExecutorGateway {
                         resourceManagerId));
     }
 
+    @Override
+    public CompletableFuture<Collection<PartitionWithMetrics>> 
getPartitionWithMetrics(
+            JobID jobId) {
+        return requestPartitionWithMetricsFunction.apply(jobId);
+    }
+
     @Override
     public CompletableFuture<Acknowledge> submitTask(
             TaskDeploymentDescriptor tdd, JobMasterId jobMasterId, Time 
timeout) {
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGatewayBuilder.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGatewayBuilder.java
index 31f5df12337..6f21873af39 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGatewayBuilder.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/TestingTaskExecutorGatewayBuilder.java
@@ -38,6 +38,7 @@ import 
org.apache.flink.runtime.operators.coordination.OperatorEvent;
 import org.apache.flink.runtime.resourcemanager.ResourceManagerId;
 import org.apache.flink.runtime.rest.messages.ProfilingInfo;
 import org.apache.flink.runtime.rest.messages.ThreadDumpInfo;
+import org.apache.flink.runtime.shuffle.PartitionWithMetrics;
 import org.apache.flink.util.SerializedValue;
 import org.apache.flink.util.concurrent.FutureUtils;
 import org.apache.flink.util.function.QuadFunction;
@@ -45,6 +46,7 @@ import org.apache.flink.util.function.TriConsumer;
 import org.apache.flink.util.function.TriFunction;
 
 import java.util.Collection;
+import java.util.Collections;
 import java.util.Set;
 import java.util.concurrent.CompletableFuture;
 import java.util.function.BiConsumer;
@@ -129,6 +131,9 @@ public class TestingTaskExecutorGatewayBuilder {
                     Tuple6<SlotID, JobID, AllocationID, ResourceProfile, 
String, ResourceManagerId>,
                     CompletableFuture<Acknowledge>>
             requestSlotFunction = NOOP_REQUEST_SLOT_FUNCTION;
+    private Function<JobID, 
CompletableFuture<Collection<PartitionWithMetrics>>>
+            requestPartitionWithMetricsFunction =
+                    ignored -> 
CompletableFuture.completedFuture(Collections.emptyList());
     private BiFunction<AllocationID, Throwable, 
CompletableFuture<Acknowledge>> freeSlotFunction =
             NOOP_FREE_SLOT_FUNCTION;
     private Consumer<JobID> freeInactiveSlotsConsumer = 
NOOP_FREE_INACTIVE_SLOTS_CONSUMER;
@@ -217,6 +222,13 @@ public class TestingTaskExecutorGatewayBuilder {
         return this;
     }
 
+    public TestingTaskExecutorGatewayBuilder 
setRequestPartitionWithMetricsFunction(
+            Function<JobID, 
CompletableFuture<Collection<PartitionWithMetrics>>>
+                    requestPartitionWithMetricsFunction) {
+        this.requestPartitionWithMetricsFunction = 
requestPartitionWithMetricsFunction;
+        return this;
+    }
+
     public TestingTaskExecutorGatewayBuilder setFreeSlotFunction(
             BiFunction<AllocationID, Throwable, 
CompletableFuture<Acknowledge>> freeSlotFunction) {
         this.freeSlotFunction = freeSlotFunction;
@@ -325,6 +337,7 @@ public class TestingTaskExecutorGatewayBuilder {
                 disconnectJobManagerConsumer,
                 submitTaskConsumer,
                 requestSlotFunction,
+                requestPartitionWithMetricsFunction,
                 freeSlotFunction,
                 freeInactiveSlotsConsumer,
                 heartbeatResourceManagerFunction,
diff --git 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTableTest.java
 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTableTest.java
index 53262ef14b4..283890d7797 100644
--- 
a/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTableTest.java
+++ 
b/flink-runtime/src/test/java/org/apache/flink/runtime/taskexecutor/partition/PartitionTableTest.java
@@ -54,6 +54,15 @@ class PartitionTableTest {
         assertThat(table.hasTrackedPartitions(JOB_ID)).isTrue();
     }
 
+    @Test
+    void testGetTrackedPartitions() {
+        final PartitionTable<JobID> table = new PartitionTable<>();
+
+        table.startTrackingPartitions(JOB_ID, 
Collections.singletonList(PARTITION_ID));
+
+        
assertThat(table.getTrackedPartitions(JOB_ID)).containsExactly(PARTITION_ID);
+    }
+
     @Test
     void testStartTrackingZeroPartitionDoesNotMutateState() {
         final PartitionTable<JobID> table = new PartitionTable<>();

Reply via email to