public class RpcTaskManagerGateway extends Object implements TaskManagerGateway
TaskManagerGateway
for Flink's RPC system.Constructor and Description |
---|
RpcTaskManagerGateway(TaskExecutorGateway taskExecutorGateway,
JobMasterId jobMasterId) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
void |
failPartition(ExecutionAttemptID executionAttemptID)
Fail all intermediate result partitions of the given task.
|
CompletableFuture<Acknowledge> |
freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout)
Frees the slot with the given allocation ID.
|
String |
getAddress()
Return the address of the task manager with which the gateway is associated.
|
void |
notifyCheckpointComplete(ExecutionAttemptID executionAttemptID,
JobID jobId,
long checkpointId,
long timestamp)
Notify the given task about a completed checkpoint.
|
CompletableFuture<StackTraceSampleResponse> |
requestStackTraceSample(ExecutionAttemptID executionAttemptID,
int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth,
Time timeout)
Request a stack trace sample from the given task.
|
CompletableFuture<Acknowledge> |
stopTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Stop the given task.
|
CompletableFuture<Acknowledge> |
submitTask(TaskDeploymentDescriptor tdd,
Time timeout)
Submit a task to the task manager.
|
void |
triggerCheckpoint(ExecutionAttemptID executionAttemptID,
JobID jobId,
long checkpointId,
long timestamp,
CheckpointOptions checkpointOptions)
Trigger for the given task a checkpoint.
|
CompletableFuture<Acknowledge> |
updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
public RpcTaskManagerGateway(TaskExecutorGateway taskExecutorGateway, JobMasterId jobMasterId)
public String getAddress()
TaskManagerGateway
getAddress
in interface TaskManagerGateway
public CompletableFuture<StackTraceSampleResponse> requestStackTraceSample(ExecutionAttemptID executionAttemptID, int sampleId, int numSamples, Time delayBetweenSamples, int maxStackTraceDepth, Time timeout)
TaskManagerGateway
requestStackTraceSample
in interface TaskManagerGateway
executionAttemptID
- identifying the task to samplesampleId
- of the samplenumSamples
- to take from the given taskdelayBetweenSamples
- to wait formaxStackTraceDepth
- of the returned sampletimeout
- of the requestpublic CompletableFuture<Acknowledge> submitTask(TaskDeploymentDescriptor tdd, Time timeout)
TaskManagerGateway
submitTask
in interface TaskManagerGateway
tdd
- describing the task to submittimeout
- of the submit operationpublic CompletableFuture<Acknowledge> stopTask(ExecutionAttemptID executionAttemptID, Time timeout)
TaskManagerGateway
stopTask
in interface TaskManagerGateway
executionAttemptID
- identifying the tasktimeout
- of the submit operationpublic CompletableFuture<Acknowledge> cancelTask(ExecutionAttemptID executionAttemptID, Time timeout)
TaskManagerGateway
cancelTask
in interface TaskManagerGateway
executionAttemptID
- identifying the tasktimeout
- of the submit operationpublic CompletableFuture<Acknowledge> updatePartitions(ExecutionAttemptID executionAttemptID, Iterable<PartitionInfo> partitionInfos, Time timeout)
TaskManagerGateway
updatePartitions
in interface TaskManagerGateway
executionAttemptID
- identifying the taskpartitionInfos
- telling where the partition can be retrieved fromtimeout
- of the submit operationpublic void failPartition(ExecutionAttemptID executionAttemptID)
TaskManagerGateway
failPartition
in interface TaskManagerGateway
executionAttemptID
- identifying the taskpublic void notifyCheckpointComplete(ExecutionAttemptID executionAttemptID, JobID jobId, long checkpointId, long timestamp)
TaskManagerGateway
notifyCheckpointComplete
in interface TaskManagerGateway
executionAttemptID
- identifying the taskjobId
- identifying the job to which the task belongscheckpointId
- of the completed checkpointtimestamp
- of the completed checkpointpublic void triggerCheckpoint(ExecutionAttemptID executionAttemptID, JobID jobId, long checkpointId, long timestamp, CheckpointOptions checkpointOptions)
TaskManagerGateway
triggerCheckpoint
in interface TaskManagerGateway
executionAttemptID
- identifying the taskjobId
- identifying the job to which the task belongscheckpointId
- of the checkpoint to triggertimestamp
- of the checkpoint to triggercheckpointOptions
- of the checkpoint to triggerpublic CompletableFuture<Acknowledge> freeSlot(AllocationID allocationId, Throwable cause, Time timeout)
TaskManagerGateway
freeSlot
in interface TaskManagerGateway
allocationId
- identifying the slot to freecause
- of the freeing operationtimeout
- for the operationCopyright © 2014–2020 The Apache Software Foundation. All rights reserved.