Package | Description |
---|---|
org.apache.flink.client.program | |
org.apache.flink.client.program.rest | |
org.apache.flink.runtime.akka | |
org.apache.flink.runtime.dispatcher | |
org.apache.flink.runtime.jobmanager.scheduler | |
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 actors, like the
JobManager and
TaskManager to coordinate the distributed operations. |
org.apache.flink.runtime.minicluster | |
org.apache.flink.runtime.resourcemanager | |
org.apache.flink.runtime.resourcemanager.slotmanager | |
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> |
ClusterClient.disposeSavepoint(String savepointPath) |
CompletableFuture<Acknowledge> |
MiniClusterClient.disposeSavepoint(String savepointPath) |
CompletableFuture<Acknowledge> |
ClusterClient.rescaleJob(JobID jobId,
int newParallelism)
Rescales the specified job such that it will have the new parallelism.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
RestClusterClient.disposeSavepoint(String savepointPath) |
CompletableFuture<Acknowledge> |
RestClusterClient.rescaleJob(JobID jobId,
int newParallelism) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
AkkaJobManagerGateway.cancelJob(JobID jobId,
Time timeout) |
CompletableFuture<Acknowledge> |
AkkaJobManagerGateway.stopJob(JobID jobId,
Time timeout) |
CompletableFuture<Acknowledge> |
AkkaJobManagerGateway.submitJob(JobGraph jobGraph,
ListeningBehaviour listeningBehaviour,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
Scheduler.cancelSlotRequest(SlotRequestId slotRequestId,
SlotSharingGroupId slotSharingGroupId,
Throwable cause) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout)
Frees the slot with the given allocation ID.
|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Stop the given task.
|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout)
Submit a task to the task manager.
|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
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> |
JobMaster.cancel(Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.cancel(Time timeout)
Cancels the currently executed job.
|
CompletableFuture<Acknowledge> |
JobManagerGateway.cancelJob(JobID jobId,
Time timeout)
Cancels the given job.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMaster.disconnectTaskManager(ResourceID resourceID,
Exception cause) |
CompletableFuture<Acknowledge> |
JobMasterGateway.disconnectTaskManager(ResourceID resourceID,
Exception cause)
Disconnects the given
TaskExecutor from the
JobMaster . |
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout) |
CompletableFuture<Acknowledge> |
KvStateRegistryGateway.notifyKvStateRegistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
InetSocketAddress kvStateServerAddress)
Notifies that queryable state has been registered.
|
CompletableFuture<Acknowledge> |
JobMaster.notifyKvStateRegistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
InetSocketAddress kvStateServerAddress) |
CompletableFuture<Acknowledge> |
KvStateRegistryGateway.notifyKvStateUnregistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName)
Notifies that queryable state has been unregistered.
|
CompletableFuture<Acknowledge> |
JobMaster.notifyKvStateUnregistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName) |
CompletableFuture<Acknowledge> |
JobMaster.rescaleJob(int newParallelism,
RescalingBehaviour rescalingBehaviour,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.rescaleJob(int newParallelism,
RescalingBehaviour rescalingBehaviour,
Time timeout)
Triggers rescaling of the executed job.
|
CompletableFuture<Acknowledge> |
JobMaster.rescaleOperators(Collection<JobVertexID> operators,
int newParallelism,
RescalingBehaviour rescalingBehaviour,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.rescaleOperators(Collection<JobVertexID> operators,
int newParallelism,
RescalingBehaviour rescalingBehaviour,
Time timeout)
Triggers rescaling of the given set of operators.
|
CompletableFuture<Acknowledge> |
JobMaster.scheduleOrUpdateConsumers(ResultPartitionID partitionID,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.scheduleOrUpdateConsumers(ResultPartitionID partitionID,
Time timeout)
Notifies the JobManager about available data for a produced partition.
|
CompletableFuture<Acknowledge> |
JobMaster.start(JobMasterId newJobMasterId,
Time timeout)
Start the rpc service and begin to run the job.
|
CompletableFuture<Acknowledge> |
JobMaster.stop(Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.stop(Time timeout)
Cancel the currently executed job.
|
CompletableFuture<Acknowledge> |
JobManagerGateway.stopJob(JobID jobId,
Time timeout)
Stops the given job.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
JobManagerGateway.submitJob(JobGraph jobGraph,
ListeningBehaviour listeningBehaviour,
Time timeout)
Submits a job to the JobManager.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMaster.suspend(Exception cause,
Time timeout)
Suspending job, all the running tasks will be cancelled, and communication with other components
will be disposed.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMaster.updateTaskExecutionState(TaskExecutionState taskExecutionState)
Updates the task execution state for a given task.
|
CompletableFuture<Acknowledge> |
JobMasterGateway.updateTaskExecutionState(TaskExecutionState taskExecutionState)
Updates the task execution state for a given task.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
SlotPool.ProviderAndOwner.cancelSlotRequest(SlotRequestId slotRequestId,
SlotSharingGroupId slotSharingGroupId,
Throwable cause) |
CompletableFuture<Acknowledge> |
SlotProvider.cancelSlotRequest(SlotRequestId slotRequestId,
SlotSharingGroupId slotSharingGroupId,
Throwable cause)
Cancels the slot request with the given
SlotRequestId and SlotSharingGroupId . |
CompletableFuture<Acknowledge> |
SlotPoolGateway.registerTaskManager(ResourceID resourceID)
Registers a TaskExecutor with the given
ResourceID at SlotPool . |
CompletableFuture<Acknowledge> |
SlotPool.registerTaskManager(ResourceID resourceID)
Register TaskManager to this pool, only those slots come from registered TaskManager will be considered valid.
|
CompletableFuture<Acknowledge> |
SlotPool.releaseSlot(SlotRequestId slotRequestId,
SlotSharingGroupId slotSharingGroupId,
Throwable cause) |
CompletableFuture<Acknowledge> |
AllocatedSlotActions.releaseSlot(SlotRequestId slotRequestId,
SlotSharingGroupId slotSharingGroupId,
Throwable cause)
Releases the slot with the given
SlotRequestId . |
CompletableFuture<Acknowledge> |
SlotPoolGateway.releaseTaskManager(ResourceID resourceId,
Exception cause)
Releases a TaskExecutor with the given
ResourceID from the SlotPool . |
CompletableFuture<Acknowledge> |
SlotPool.releaseTaskManager(ResourceID resourceId,
Exception cause)
Unregister TaskManager from this pool, all the related slots will be released and tasks be canceled.
|
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) |
CompletableFuture<Acknowledge> |
MiniCluster.stopJob(JobID jobId) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
ResourceManagerGateway.deregisterApplication(ApplicationStatus finalStatus,
String diagnostics)
Deregister Flink from the underlying resource management system.
|
CompletableFuture<Acknowledge> |
ResourceManager.deregisterApplication(ApplicationStatus finalStatus,
String diagnostics)
Cleanup application and shut down cluster.
|
CompletableFuture<Acknowledge> |
ResourceManagerGateway.requestSlot(JobMasterId jobMasterId,
SlotRequest slotRequest,
Time timeout)
Requests a slot from the resource manager.
|
CompletableFuture<Acknowledge> |
ResourceManager.requestSlot(JobMasterId jobMasterId,
SlotRequest slotRequest,
Time timeout) |
CompletableFuture<Acknowledge> |
ResourceManagerGateway.sendSlotReport(ResourceID taskManagerResourceId,
InstanceID taskManagerRegistrationId,
SlotReport slotReport,
Time timeout)
Sends the given
SlotReport to the ResourceManager. |
CompletableFuture<Acknowledge> |
ResourceManager.sendSlotReport(ResourceID taskManagerResourceId,
InstanceID taskManagerRegistrationId,
SlotReport slotReport,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
PendingSlotRequest.getRequestFuture() |
Modifier and Type | Method and Description |
---|---|
void |
PendingSlotRequest.setRequestFuture(CompletableFuture<Acknowledge> requestFuture) |
Modifier and Type | Method and Description |
---|---|
protected CompletableFuture<Acknowledge> |
RescalingHandlers.RescalingTriggerHandler.triggerOperation(HandlerRequest<EmptyRequestBody,RescalingTriggerMessageParameters> 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,EmptyMessageParameters> request,
RestfulGateway gateway) |
Modifier and Type | Method and Description |
---|---|
protected AsynchronousOperationInfo |
SavepointDisposalHandlers.SavepointDisposalStatusHandler.operationResultResponse(Acknowledge operationResult) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskExecutorGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutor.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.confirmCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointId,
long checkpointTimestamp)
Confirm a checkpoint for the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutor.confirmCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointId,
long checkpointTimestamp) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout)
Frees the slot with the given allocation ID.
|
CompletableFuture<Acknowledge> |
TaskExecutor.freeSlot(AllocationID allocationId,
Throwable cause,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
String targetAddress,
ResourceManagerId resourceManagerId,
Time timeout)
Requests a slot from the TaskManager.
|
CompletableFuture<Acknowledge> |
TaskExecutor.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
String targetAddress,
ResourceManagerId resourceManagerId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Stop the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutor.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.submitTask(TaskDeploymentDescriptor tdd,
JobMasterId jobMasterId,
Time timeout)
Submit a
Task to the TaskExecutor . |
CompletableFuture<Acknowledge> |
TaskExecutor.submitTask(TaskDeploymentDescriptor tdd,
JobMasterId jobMasterId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.triggerCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointID,
long checkpointTimestamp,
CheckpointOptions checkpointOptions)
Trigger the checkpoint for the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutor.triggerCheckpoint(ExecutionAttemptID executionAttemptID,
long checkpointId,
long checkpointTimestamp,
CheckpointOptions checkpointOptions) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
CompletableFuture<Acknowledge> |
TaskExecutor.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
RestfulGateway.cancelJob(JobID jobId,
Time timeout)
Cancel the given job.
|
default CompletableFuture<Acknowledge> |
RestfulGateway.disposeSavepoint(String savepointPath,
Time timeout)
Dispose the given savepoint.
|
default CompletableFuture<Acknowledge> |
RestfulGateway.rescaleJob(JobID jobId,
int newParallelism,
RescalingBehaviour rescalingBehaviour,
Time timeout)
Trigger rescaling of the given job.
|
default CompletableFuture<Acknowledge> |
RestfulGateway.shutDownCluster() |
CompletableFuture<Acknowledge> |
RestfulGateway.stopJob(JobID jobId,
Time timeout)
Stop the given job.
|
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.