Package | Description |
---|---|
org.apache.flink.runtime.execution | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.io.network.netty | |
org.apache.flink.runtime.jobmaster | |
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.taskexecutor.rpc | |
org.apache.flink.runtime.taskmanager |
Modifier and Type | Method and Description |
---|---|
static ExecutionState |
ExecutionState.valueOf(String name)
Returns the enum constant of this type with the specified name.
|
static ExecutionState[] |
ExecutionState.values()
Returns an array containing the constants of this enum type, in
the order they are declared.
|
Modifier and Type | Method and Description |
---|---|
static ExecutionState |
ExecutionJobVertex.getAggregateJobVertexState(int[] verticesPerState,
int parallelism)
A utility function that computes an "aggregated" state for the vertex.
|
ExecutionState |
ExecutionJobVertex.getAggregateState() |
ExecutionState |
ArchivedExecutionJobVertex.getAggregateState() |
ExecutionState |
AccessExecutionJobVertex.getAggregateState()
Returns the aggregated
ExecutionState for this job vertex. |
ExecutionState |
ExecutionVertex.getExecutionState() |
ExecutionState |
ArchivedExecutionVertex.getExecutionState() |
ExecutionState |
AccessExecutionVertex.getExecutionState()
Returns the current
ExecutionState for this execution vertex. |
ExecutionState |
Execution.getState() |
ExecutionState |
ArchivedExecution.getState() |
ExecutionState |
AccessExecution.getState()
Returns the current
ExecutionState for this execution. |
Modifier and Type | Method and Description |
---|---|
Future<ExecutionState> |
ExecutionVertex.cancel() |
Future<ExecutionState> |
Execution.getTerminationFuture()
Gets a future that completes once the task execution reaches a terminal state.
|
Modifier and Type | Method and Description |
---|---|
void |
StatusListenerMessenger.executionStatusChanged(JobID jobID,
JobVertexID vertexID,
String taskName,
int taskParallelism,
int subtaskIndex,
ExecutionAttemptID executionID,
ExecutionState newExecutionState,
long timestamp,
String optionalMessage) |
void |
ExecutionStatusListener.executionStatusChanged(JobID jobID,
JobVertexID vertexID,
String taskName,
int totalNumberOfSubTasks,
int subtaskIndex,
ExecutionAttemptID executionID,
ExecutionState newExecutionState,
long timestamp,
String optionalMessage)
Called whenever the execution status of a task changes.
|
long |
ExecutionVertex.getStateTimestamp(ExecutionState state) |
long |
Execution.getStateTimestamp(ExecutionState state) |
long |
ArchivedExecutionVertex.getStateTimestamp(ExecutionState state) |
long |
ArchivedExecution.getStateTimestamp(ExecutionState state) |
long |
AccessExecutionVertex.getStateTimestamp(ExecutionState state)
Returns the timestamp for the given
ExecutionState . |
long |
AccessExecution.getStateTimestamp(ExecutionState state)
Returns the timestamp for the given
ExecutionState . |
Constructor and Description |
---|
ArchivedExecution(StringifiedAccumulatorResult[] userAccumulators,
IOMetrics ioMetrics,
ExecutionAttemptID attemptId,
int attemptNumber,
ExecutionState state,
String failureCause,
TaskManagerLocation assignedResourceLocation,
int parallelSubtaskIndex,
long[] stateTimestamps) |
IllegalExecutionStateException(Execution execution,
ExecutionState expected,
ExecutionState actual)
Creates a new IllegalExecutionStateException with the error message indicating
the expected and actual state.
|
IllegalExecutionStateException(ExecutionState expected,
ExecutionState actual)
Creates a new IllegalExecutionStateException with the error message indicating
the expected and actual state.
|
Modifier and Type | Method and Description |
---|---|
Future<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 |
---|---|
ExecutionState |
JobMaster.requestPartitionState(UUID leaderSessionID,
IntermediateDataSetID intermediateResultId,
ResultPartitionID resultPartitionId) |
Modifier and Type | Method and Description |
---|---|
Future<ExecutionState> |
JobMasterGateway.requestPartitionState(UUID leaderSessionID,
IntermediateDataSetID intermediateResultId,
ResultPartitionID partitionId)
Requests the current state of the partition.
|
Modifier and Type | Method and Description |
---|---|
ExecutionState |
ExecutionGraphMessages.ExecutionStateChanged.newExecutionState() |
Constructor and Description |
---|
ExecutionStateChanged(JobID jobID,
JobVertexID vertexID,
String taskName,
int totalNumberOfSubTasks,
int subtaskIndex,
ExecutionAttemptID executionID,
ExecutionState newExecutionState,
long timestamp,
String optionalMessage) |
Modifier and Type | Method and Description |
---|---|
Future<ExecutionState> |
RpcPartitionStateChecker.requestPartitionProducerState(JobID jobId,
IntermediateDataSetID resultId,
ResultPartitionID partitionId) |
Modifier and Type | Method and Description |
---|---|
ExecutionState |
TaskExecutionState.getExecutionState()
Returns the new execution state of the task.
|
ExecutionState |
Task.getExecutionState()
Returns the current execution state of the task.
|
Modifier and Type | Method and Description |
---|---|
Future<ExecutionState> |
ActorGatewayPartitionProducerStateChecker.requestPartitionProducerState(JobID jobId,
IntermediateDataSetID intermediateDataSetId,
ResultPartitionID resultPartitionId) |
Constructor and Description |
---|
TaskExecutionState(JobID jobID,
ExecutionAttemptID executionId,
ExecutionState executionState)
Creates a new task execution state update, with no attached exception and no accumulators.
|
TaskExecutionState(JobID jobID,
ExecutionAttemptID executionId,
ExecutionState executionState,
Throwable error)
Creates a new task execution state update, with an attached exception but no accumulators.
|
TaskExecutionState(JobID jobID,
ExecutionAttemptID executionId,
ExecutionState executionState,
Throwable error,
AccumulatorSnapshot accumulators,
IOMetrics ioMetrics)
Creates a new task execution state update, with an attached exception.
|
Copyright © 2014–2018 The Apache Software Foundation. All rights reserved.