Package | Description |
---|---|
org.apache.flink.client.program | |
org.apache.flink.client.program.rest | |
org.apache.flink.runtime.blocklist | |
org.apache.flink.runtime.dispatcher | |
org.apache.flink.runtime.dispatcher.cleanup | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.jobmanager.slots | |
org.apache.flink.runtime.jobmaster | |
org.apache.flink.runtime.jobmaster.slotpool | |
org.apache.flink.runtime.messages |
This package contains the messages that are sent between Flink's distributed components to
coordinate the distributed operations.
|
org.apache.flink.runtime.minicluster | |
org.apache.flink.runtime.operators.coordination | |
org.apache.flink.runtime.resourcemanager | |
org.apache.flink.runtime.rest.handler.job.rescaling | |
org.apache.flink.runtime.rest.handler.job.savepoints | |
org.apache.flink.runtime.taskexecutor | |
org.apache.flink.runtime.webmonitor |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
MiniClusterClient.cancel(JobID jobId) |
CompletableFuture<Acknowledge> |
ClusterClient.cancel(JobID jobId)
Cancels a job identified by the job id.
|
CompletableFuture<Acknowledge> |
MiniClusterClient.disposeSavepoint(String savepointPath) |
CompletableFuture<Acknowledge> |
ClusterClient.disposeSavepoint(String savepointPath)
Dispose the savepoint under the given path.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
RestClusterClient.cancel(JobID jobID) |
CompletableFuture<Acknowledge> |
RestClusterClient.disposeSavepoint(String savepointPath) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
BlocklistListener.notifyNewBlockedNodes(Collection<BlockedNode> newNodes)
Notify new blocked node records.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
CheckpointResourcesCleanupRunner.cancel(Time timeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
Execution.sendOperatorEvent(OperatorID operatorId,
SerializedValue<OperatorEvent> event)
Sends the operator event to the Task on the Task Executor.
|
CompletableFuture<Acknowledge> |
Execution.triggerCheckpoint(long checkpointId,
long timestamp,
CheckpointOptions checkpointOptions)
Trigger a new checkpoint on the task of this execution.
|
CompletableFuture<Acknowledge> |
Execution.triggerSynchronousSavepoint(long checkpointId,
long timestamp,
CheckpointOptions checkpointOptions)
Trigger a new checkpoint on the task of this execution.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
CompletableFuture<Acknowledge> |
TaskManagerGateway.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout)
Frees the slot with the given allocation ID.
|
CompletableFuture<Acknowledge> |
TaskManagerGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout)
Submit a task to the task manager.
|
CompletableFuture<Acknowledge> |
TaskManagerGateway.triggerCheckpoint(ExecutionAttemptID executionAttemptID,
JobID jobId,
long checkpointId,
long timestamp,
CheckpointOptions checkpointOptions)
Trigger for the given task a checkpoint.
|
CompletableFuture<Acknowledge> |
TaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
DeclareResourceRequirementServiceConnectionManager.DeclareResourceRequirementsService.declareResourceRequirements(ResourceRequirements resourceRequirements) |
Modifier and Type | Method and Description |
---|---|
static Acknowledge |
Acknowledge.get()
Gets the singleton instance.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
MiniCluster.cancelJob(JobID jobId) |
CompletableFuture<Acknowledge> |
MiniCluster.disposeSavepoint(String savepointPath) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
OperatorCoordinator.SubtaskGateway.sendEvent(OperatorEvent evt)
Sends an event to the parallel subtask with the given subtask index.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
ResourceManager.declareRequiredResources(JobMasterId jobMasterId,
ResourceRequirements resourceRequirements,
Time timeout) |
CompletableFuture<Acknowledge> |
ResourceManagerGateway.declareRequiredResources(JobMasterId jobMasterId,
ResourceRequirements resourceRequirements,
Time timeout)
Declares the absolute resource requirements for a job.
|
CompletableFuture<Acknowledge> |
ResourceManager.deregisterApplication(ApplicationStatus finalStatus,
String diagnostics)
Cleanup application and shut down cluster.
|
CompletableFuture<Acknowledge> |
ResourceManagerGateway.deregisterApplication(ApplicationStatus finalStatus,
String diagnostics)
Deregister Flink from the underlying resource management system.
|
CompletableFuture<Acknowledge> |
ResourceManager.notifyNewBlockedNodes(Collection<BlockedNode> newNodes) |
CompletableFuture<Acknowledge> |
ResourceManager.sendSlotReport(ResourceID taskManagerResourceId,
InstanceID taskManagerRegistrationId,
SlotReport slotReport,
Time timeout) |
CompletableFuture<Acknowledge> |
ResourceManagerGateway.sendSlotReport(ResourceID taskManagerResourceId,
InstanceID taskManagerRegistrationId,
SlotReport slotReport,
Time timeout)
Sends the given
SlotReport to the ResourceManager. |
Modifier and Type | Method and Description |
---|---|
protected CompletableFuture<Acknowledge> |
RescalingHandlers.RescalingTriggerHandler.triggerOperation(HandlerRequest<EmptyRequestBody> request,
RestfulGateway gateway) |
Modifier and Type | Method and Description |
---|---|
protected AsynchronousOperationInfo |
RescalingHandlers.RescalingStatusHandler.operationResultResponse(Acknowledge operationResult) |
Modifier and Type | Method and Description |
---|---|
protected CompletableFuture<Acknowledge> |
SavepointDisposalHandlers.SavepointDisposalTriggerHandler.triggerOperation(HandlerRequest<SavepointDisposalRequest> request,
RestfulGateway gateway) |
protected CompletableFuture<Acknowledge> |
SavepointHandlers.SavepointTriggerHandler.triggerOperation(HandlerRequest<SavepointTriggerRequestBody> request,
AsynchronousJobOperationKey operationKey,
RestfulGateway gateway) |
protected CompletableFuture<Acknowledge> |
SavepointHandlers.StopWithSavepointHandler.triggerOperation(HandlerRequest<StopWithSavepointRequestBody> request,
AsynchronousJobOperationKey operationKey,
RestfulGateway gateway) |
Modifier and Type | Method and Description |
---|---|
protected AsynchronousOperationInfo |
SavepointDisposalHandlers.SavepointDisposalStatusHandler.operationResultResponse(Acknowledge operationResult) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.abortCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointId,
long latestCompletedCheckpointId,
long checkpointTimestamp) |
CompletableFuture<Acknowledge> |
TaskExecutor.abortCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointId,
long latestCompletedCheckpointId,
long checkpointTimestamp) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.abortCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointId,
long latestCompletedCheckpointId,
long checkpointTimestamp)
Abort a checkpoint for the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutor.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.confirmCheckpoint(ExecutionAttemptID executionAttemptID,
long completedCheckpointId,
long completedCheckpointTimestamp,
long lastSubsumedCheckpointId) |
CompletableFuture<Acknowledge> |
TaskExecutor.confirmCheckpoint(ExecutionAttemptID executionAttemptID,
long completedCheckpointId,
long completedCheckpointTimestamp,
long lastSubsumedCheckpointId) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.confirmCheckpoint(ExecutionAttemptID executionAttemptID,
long completedCheckpointId,
long completedCheckpointTimestamp,
long lastSubsumedCheckpointId)
Confirm a checkpoint for the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutor.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout)
Frees the slot with the given allocation ID.
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.promotePartitions(JobID jobId,
Set<ResultPartitionID> partitionIds) |
CompletableFuture<Acknowledge> |
TaskExecutor.promotePartitions(JobID jobId,
Set<ResultPartitionID> partitionIds) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.promotePartitions(JobID jobId,
Set<ResultPartitionID> partitionIds)
Batch promote intermediate result partitions.
|
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.
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
ResourceProfile resourceProfile,
String targetAddress,
ResourceManagerId resourceManagerId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutor.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
ResourceProfile resourceProfile,
String targetAddress,
ResourceManagerId resourceManagerId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
ResourceProfile resourceProfile,
String targetAddress,
ResourceManagerId resourceManagerId,
Time timeout)
Requests a slot from the TaskManager.
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
CompletableFuture<Acknowledge> |
TaskExecutor.sendOperatorEventToTask(ExecutionAttemptID executionAttemptID,
OperatorID operatorId,
SerializedValue<OperatorEvent> evt) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt) |
CompletableFuture<Acknowledge> |
TaskExecutorOperatorEventGateway.sendOperatorEventToTask(ExecutionAttemptID task,
OperatorID operator,
SerializedValue<OperatorEvent> evt)
Sends an operator event to an operator in a task executed by the Task Manager (Task
Executor).
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.submitTask(TaskDeploymentDescriptor tdd,
JobMasterId jobMasterId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutor.submitTask(TaskDeploymentDescriptor tdd,
JobMasterId jobMasterId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.submitTask(TaskDeploymentDescriptor tdd,
JobMasterId jobMasterId,
Time timeout)
Submit a
Task to the TaskExecutor . |
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.triggerCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointID,
long checkpointTimestamp,
CheckpointOptions checkpointOptions) |
CompletableFuture<Acknowledge> |
TaskExecutor.triggerCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointId,
long checkpointTimestamp,
CheckpointOptions checkpointOptions) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.triggerCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointID,
long checkpointTimestamp,
CheckpointOptions checkpointOptions)
Trigger the checkpoint for the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutorGatewayDecoratorBase.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutor.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
RestfulGateway.cancelJob(JobID jobId,
Time timeout)
Cancel the given job.
|
CompletableFuture<Acknowledge> |
NonLeaderRetrievalRestfulGateway.cancelJob(JobID jobId,
Time timeout) |
default CompletableFuture<Acknowledge> |
RestfulGateway.disposeSavepoint(String savepointPath,
Time timeout)
Dispose the given savepoint.
|
default CompletableFuture<Acknowledge> |
RestfulGateway.shutDownCluster() |
default CompletableFuture<Acknowledge> |
RestfulGateway.stopWithSavepoint(AsynchronousJobOperationKey operationKey,
String targetDirectory,
SavepointFormatType formatType,
TriggerSavepointMode savepointMode,
Time timeout)
Stops the job with a savepoint, returning a future that completes when the operation is
started.
|
default CompletableFuture<Acknowledge> |
RestfulGateway.triggerCheckpoint(AsynchronousJobOperationKey operationKey,
CheckpointType checkpointType,
Time timeout)
Triggers a checkpoint with the given savepoint directory as a target.
|
default CompletableFuture<Acknowledge> |
RestfulGateway.triggerSavepoint(AsynchronousJobOperationKey operationKey,
String targetDirectory,
SavepointFormatType formatType,
TriggerSavepointMode savepointMode,
Time timeout)
Triggers a savepoint with the given savepoint directory as a target, returning a future that
completes when the operation is started.
|
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.