Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
InputGateDeploymentDescriptor.getConsumedResultId() |
IntermediateDataSetID |
ResultPartitionDeploymentDescriptor.getResultId() |
Constructor and Description |
---|
InputGateDeploymentDescriptor(IntermediateDataSetID consumedResultId,
ResultPartitionType consumedPartitionType,
int consumedSubpartitionIndex,
ShuffleDescriptor[] inputChannels) |
InputGateDeploymentDescriptor(IntermediateDataSetID consumedResultId,
ResultPartitionType consumedPartitionType,
SubpartitionIndexRange consumedSubpartitionIndexRange,
TaskDeploymentDescriptor.MaybeOffloaded<ShuffleDescriptor[]> serializedInputChannels) |
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
IntermediateResult.getId() |
IntermediateDataSetID |
PartitionInfo.getIntermediateDataSetID() |
Modifier and Type | Method and Description |
---|---|
Map<IntermediateDataSetID,IntermediateResult> |
DefaultExecutionGraph.getAllIntermediateResults() |
Map<IntermediateDataSetID,IntermediateResult> |
ExecutionGraph.getAllIntermediateResults() |
Modifier and Type | Method and Description |
---|---|
List<ShuffleDescriptor> |
DefaultExecutionGraph.getClusterPartitionShuffleDescriptors(IntermediateDataSetID intermediateDataSetID) |
List<ShuffleDescriptor> |
InternalExecutionGraphAccessor.getClusterPartitionShuffleDescriptors(IntermediateDataSetID intermediateResultPartition)
Get the shuffle descriptors of the cluster partitions ordered by partition number.
|
Modifier and Type | Method and Description |
---|---|
void |
ExecutionJobVertex.connectToPredecessors(Map<IntermediateDataSetID,IntermediateResult> intermediateDataSets) |
Constructor and Description |
---|
PartitionInfo(IntermediateDataSetID intermediateResultPartitionID,
ShuffleDescriptor shuffleDescriptor) |
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
TaskExecutorPartitionInfo.getIntermediateDataSetId() |
Modifier and Type | Method and Description |
---|---|
Map<IntermediateDataSetID,DataSetMetaInfo> |
ResourceManagerPartitionTrackerImpl.listDataSets() |
CompletableFuture<Map<IntermediateDataSetID,DataSetMetaInfo>> |
ClusterPartitionManager.listDataSets()
Returns all datasets for which partitions are being tracked.
|
Map<IntermediateDataSetID,DataSetMetaInfo> |
ResourceManagerPartitionTracker.listDataSets()
Returns all data sets for which partitions are being tracked.
|
Modifier and Type | Method and Description |
---|---|
List<ShuffleDescriptor> |
JobMasterPartitionTracker.getClusterPartitionShuffleDescriptors(IntermediateDataSetID intermediateDataSetID)
Get the shuffle descriptors of the cluster partitions ordered by partition number.
|
List<ShuffleDescriptor> |
ResourceManagerPartitionTrackerImpl.getClusterPartitionShuffleDescriptors(IntermediateDataSetID dataSetID) |
List<ShuffleDescriptor> |
ResourceManagerPartitionTracker.getClusterPartitionShuffleDescriptors(IntermediateDataSetID dataSetID)
Returns all the shuffle descriptors of cluster partitions for the intermediate dataset.
|
List<ShuffleDescriptor> |
JobMasterPartitionTrackerImpl.getClusterPartitionShuffleDescriptors(IntermediateDataSetID intermediateDataSetID) |
CompletableFuture<List<ShuffleDescriptor>> |
ClusterPartitionManager.getClusterPartitionsShuffleDescriptors(IntermediateDataSetID intermediateDataSetID)
Get the shuffle descriptors of the cluster partitions ordered by partition number.
|
CompletableFuture<Void> |
ResourceManagerPartitionTrackerImpl.releaseClusterPartitions(IntermediateDataSetID dataSetId) |
CompletableFuture<Void> |
ClusterPartitionManager.releaseClusterPartitions(IntermediateDataSetID dataSetToRelease)
Releases all partitions associated with the given dataset.
|
CompletableFuture<Void> |
ResourceManagerPartitionTracker.releaseClusterPartitions(IntermediateDataSetID dataSetId)
Issues a release calls to all task executors that are hosting partitions of the given data
set.
|
void |
PartitionProducerStateProvider.requestPartitionProducerState(IntermediateDataSetID intermediateDataSetId,
ResultPartitionID resultPartitionId,
java.util.function.Consumer<? super PartitionProducerStateProvider.ResponseHandle> responseConsumer)
Trigger the producer execution state request.
|
Modifier and Type | Method and Description |
---|---|
void |
TaskExecutorClusterPartitionReleaser.releaseClusterPartitions(ResourceID taskExecutorId,
Set<IntermediateDataSetID> dataSetsToRelease) |
void |
TaskExecutorPartitionTracker.stopTrackingAndReleaseClusterPartitions(Collection<IntermediateDataSetID> dataSetsToRelease)
Releases partitions associated with the given datasets and stops tracking of partitions that
were released.
|
void |
TaskExecutorPartitionTrackerImpl.stopTrackingAndReleaseClusterPartitions(Collection<IntermediateDataSetID> dataSetsToRelease) |
Constructor and Description |
---|
TaskExecutorPartitionInfo(ShuffleDescriptor shuffleDescriptor,
IntermediateDataSetID intermediateDataSetId,
int numberOfPartitions) |
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
InputGateID.getConsumedResultID() |
Constructor and Description |
---|
InputGateID(IntermediateDataSetID consumedResultID,
ExecutionAttemptID consumerID) |
SingleInputGate(String owningTaskName,
int gateIndex,
IntermediateDataSetID consumedResultId,
ResultPartitionType consumedPartitionType,
SubpartitionIndexRange subpartitionIndexRange,
int numberOfInputChannels,
PartitionProducerStateProvider partitionProducerStateProvider,
SupplierWithException<BufferPool,IOException> bufferPoolFactory,
BufferDecompressor bufferDecompressor,
MemorySegmentProvider memorySegmentProvider,
int segmentSize,
ThroughputCalculator throughputCalculator,
BufferDebloater bufferDebloater) |
Modifier and Type | Method and Description |
---|---|
static IntermediateDataSetID |
IntermediateDataSetID.fromByteBuf(org.apache.flink.shaded.netty4.io.netty.buffer.ByteBuf buf) |
IntermediateDataSetID |
IntermediateDataSet.getId() |
IntermediateDataSetID |
IntermediateResultPartitionID.getIntermediateDataSetID() |
IntermediateDataSetID |
JobEdge.getSourceId()
Gets the ID of the consumed data set.
|
Modifier and Type | Method and Description |
---|---|
List<IntermediateDataSetID> |
JobVertex.getIntermediateDataSetIdsToConsume() |
Modifier and Type | Method and Description |
---|---|
void |
JobVertex.addIntermediateDataSetIdToConsume(IntermediateDataSetID intermediateDataSetId) |
JobEdge |
JobVertex.connectNewDataSetAsInput(JobVertex input,
DistributionPattern distPattern,
ResultPartitionType partitionType,
IntermediateDataSetID intermediateDataSetId,
boolean isBroadcast) |
IntermediateDataSet |
JobVertex.getOrCreateResultDataSet(IntermediateDataSetID id,
ResultPartitionType partitionType) |
Constructor and Description |
---|
IntermediateDataSet(IntermediateDataSetID id,
ResultPartitionType resultType,
JobVertex producer) |
IntermediateResultPartitionID(IntermediateDataSetID intermediateDataSetID,
int partitionNum)
Creates an new intermediate result partition ID with
IntermediateDataSetID and the
partitionNum. |
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
DefaultLogicalResult.getId() |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<ExecutionState> |
JobMaster.requestPartitionState(IntermediateDataSetID intermediateResultId,
ResultPartitionID resultPartitionId) |
CompletableFuture<ExecutionState> |
JobMasterGateway.requestPartitionState(IntermediateDataSetID intermediateResultId,
ResultPartitionID partitionId)
Requests the current state of the partition.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Map<IntermediateDataSetID,DataSetMetaInfo>> |
ResourceManager.listDataSets() |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<List<ShuffleDescriptor>> |
ResourceManager.getClusterPartitionsShuffleDescriptors(IntermediateDataSetID intermediateDataSetID) |
CompletableFuture<Void> |
ResourceManager.releaseClusterPartitions(IntermediateDataSetID dataSetId) |
Modifier and Type | Method and Description |
---|---|
protected IntermediateDataSetID |
ClusterDataSetIdPathParameter.convertFromString(String value) |
Modifier and Type | Method and Description |
---|---|
protected String |
ClusterDataSetIdPathParameter.convertToString(IntermediateDataSetID value) |
Modifier and Type | Method and Description |
---|---|
static ClusterDataSetListResponseBody |
ClusterDataSetListResponseBody.from(Map<IntermediateDataSetID,DataSetMetaInfo> dataSets) |
Modifier and Type | Method and Description |
---|---|
List<IntermediateDataSetID> |
ClusterDatasetCorruptedException.getCorruptedClusterDatasetIds() |
Modifier and Type | Method and Description |
---|---|
ExecutionState |
SchedulerNG.requestPartitionState(IntermediateDataSetID intermediateResultId,
ResultPartitionID resultPartitionId) |
ExecutionState |
ExecutionGraphHandler.requestPartitionState(IntermediateDataSetID intermediateResultId,
ResultPartitionID resultPartitionId) |
ExecutionState |
SchedulerBase.requestPartitionState(IntermediateDataSetID intermediateResultId,
ResultPartitionID resultPartitionId) |
Constructor and Description |
---|
ClusterDatasetCorruptedException(Throwable cause,
List<IntermediateDataSetID> corruptedClusterDatasetIds) |
Modifier and Type | Method and Description |
---|---|
ExecutionState |
AdaptiveScheduler.requestPartitionState(IntermediateDataSetID intermediateResultId,
ResultPartitionID resultPartitionId) |
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
ConsumedPartitionGroup.getIntermediateDataSetID()
Get the ID of IntermediateDataSet this ConsumedPartitionGroup belongs to.
|
IntermediateDataSetID |
SchedulingResultPartition.getResultId()
Gets id of the intermediate result.
|
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
PartitionDescriptor.getResultId() |
Modifier and Type | Method and Description |
---|---|
Map<IntermediateDataSetID,Integer> |
TaskInputsOutputsDescriptor.getInputChannelNums() |
Map<IntermediateDataSetID,ResultPartitionType> |
TaskInputsOutputsDescriptor.getPartitionTypes() |
Map<IntermediateDataSetID,Integer> |
TaskInputsOutputsDescriptor.getSubpartitionNums() |
Modifier and Type | Method and Description |
---|---|
static int |
NettyShuffleUtils.computeNetworkBuffersForAnnouncing(int numBuffersPerChannel,
int numFloatingBuffersPerGate,
int sortShuffleMinParallelism,
int sortShuffleMinBuffers,
int numTotalInputChannels,
int numTotalInputGates,
Map<IntermediateDataSetID,Integer> subpartitionNums,
Map<IntermediateDataSetID,ResultPartitionType> partitionTypes) |
static int |
NettyShuffleUtils.computeNetworkBuffersForAnnouncing(int numBuffersPerChannel,
int numFloatingBuffersPerGate,
int sortShuffleMinParallelism,
int sortShuffleMinBuffers,
int numTotalInputChannels,
int numTotalInputGates,
Map<IntermediateDataSetID,Integer> subpartitionNums,
Map<IntermediateDataSetID,ResultPartitionType> partitionTypes) |
static TaskInputsOutputsDescriptor |
TaskInputsOutputsDescriptor.from(int inputGateNums,
Map<IntermediateDataSetID,Integer> inputChannelNums,
Map<IntermediateDataSetID,Integer> subpartitionNums,
Map<IntermediateDataSetID,ResultPartitionType> partitionTypes) |
static TaskInputsOutputsDescriptor |
TaskInputsOutputsDescriptor.from(int inputGateNums,
Map<IntermediateDataSetID,Integer> inputChannelNums,
Map<IntermediateDataSetID,Integer> subpartitionNums,
Map<IntermediateDataSetID,ResultPartitionType> partitionTypes) |
static TaskInputsOutputsDescriptor |
TaskInputsOutputsDescriptor.from(int inputGateNums,
Map<IntermediateDataSetID,Integer> inputChannelNums,
Map<IntermediateDataSetID,Integer> subpartitionNums,
Map<IntermediateDataSetID,ResultPartitionType> partitionTypes) |
Constructor and Description |
---|
PartitionDescriptor(IntermediateDataSetID resultId,
int totalNumberOfPartitions,
IntermediateResultPartitionID partitionId,
ResultPartitionType partitionType,
int numberOfSubpartitions,
int connectionIndex,
boolean isBroadcast,
boolean isAllToAllDistribution) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<ExecutionState> |
PartitionProducerStateChecker.requestPartitionProducerState(JobID jobId,
IntermediateDataSetID intermediateDataSetId,
ResultPartitionID resultPartitionId)
Requests the execution state of the execution producing a result partition.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.releaseClusterPartitions(Collection<IntermediateDataSetID> dataSetsToRelease,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutor.releaseClusterPartitions(Collection<IntermediateDataSetID> dataSetsToRelease,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.releaseClusterPartitions(Collection<IntermediateDataSetID> dataSetsToRelease,
Time timeout)
Releases all cluster partitions belong to any of the given data sets.
|
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
ClusterPartitionReport.ClusterPartitionReportEntry.getDataSetId() |
Constructor and Description |
---|
ClusterPartitionReportEntry(IntermediateDataSetID dataSetId,
int numTotalPartitions,
Map<ResultPartitionID,ShuffleDescriptor> shuffleDescriptors) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<ExecutionState> |
RpcPartitionStateChecker.requestPartitionProducerState(JobID jobId,
IntermediateDataSetID resultId,
ResultPartitionID partitionId) |
Modifier and Type | Method and Description |
---|---|
void |
Task.requestPartitionProducerState(IntermediateDataSetID intermediateDataSetId,
ResultPartitionID resultPartitionId,
java.util.function.Consumer<? super PartitionProducerStateProvider.ResponseHandle> responseConsumer) |
Modifier and Type | Method and Description |
---|---|
IntermediateDataSetID |
StreamNode.getConsumeClusterDatasetId() |
IntermediateDataSetID |
NonChainedOutput.getDataSetId() |
IntermediateDataSetID |
StreamEdge.getIntermediateDatasetIdToProduce() |
IntermediateDataSetID |
NonChainedOutput.getPersistentDataSetId() |
Modifier and Type | Method and Description |
---|---|
void |
StreamGraph.addEdge(Integer upStreamVertexID,
Integer downStreamVertexID,
int typeNumber,
IntermediateDataSetID intermediateDataSetId) |
void |
StreamNode.setConsumeClusterDatasetId(IntermediateDataSetID consumeClusterDatasetId) |
Constructor and Description |
---|
NonChainedOutput(boolean supportsUnalignedCheckpoints,
int sourceNodeId,
int consumerParallelism,
int consumerMaxParallelism,
long bufferTimeout,
boolean isPersistentDataSet,
IntermediateDataSetID dataSetId,
OutputTag<?> outputTag,
StreamPartitioner<?> partitioner,
ResultPartitionType partitionType) |
StreamEdge(StreamNode sourceVertex,
StreamNode targetVertex,
int typeNumber,
long bufferTimeout,
StreamPartitioner<?> outputPartitioner,
OutputTag outputTag,
StreamExchangeMode exchangeMode,
int uniqueId,
IntermediateDataSetID intermediateDatasetId) |
StreamEdge(StreamNode sourceVertex,
StreamNode targetVertex,
int typeNumber,
StreamPartitioner<?> outputPartitioner,
OutputTag outputTag,
StreamExchangeMode exchangeMode,
int uniqueId,
IntermediateDataSetID intermediateDatasetId) |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.