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<>();