Package | Description |
---|---|
org.apache.flink.runtime.accumulators | |
org.apache.flink.runtime.checkpoint | |
org.apache.flink.runtime.deployment | |
org.apache.flink.runtime.execution | |
org.apache.flink.runtime.execution.librarycache | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.io.network.netty | |
org.apache.flink.runtime.io.network.partition | |
org.apache.flink.runtime.io.network.partition.consumer | |
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.messages.checkpoint |
This package contains the messages that are sent between
JobManager
and TaskManager to coordinate the checkpoint snapshots of the
distributed dataflow. |
org.apache.flink.runtime.metrics.groups | |
org.apache.flink.runtime.taskmanager | |
org.apache.flink.runtime.testingUtils | |
org.apache.flink.runtime.webmonitor |
Modifier and Type | Field and Description |
---|---|
protected ExecutionAttemptID |
AccumulatorRegistry.taskID |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
AccumulatorSnapshot.getExecutionAttemptID() |
Constructor and Description |
---|
AccumulatorRegistry(JobID jobID,
ExecutionAttemptID taskID) |
AccumulatorSnapshot(JobID jobID,
ExecutionAttemptID executionAttemptID,
Map<AccumulatorRegistry.Metric,Accumulator<?,?>> flinkAccumulators,
Map<String,Accumulator<?,?>> userAccumulators) |
Modifier and Type | Method and Description |
---|---|
PendingCheckpoint.TaskAcknowledgeResult |
PendingCheckpoint.acknowledgeTask(ExecutionAttemptID executionAttemptId,
SerializedValue<StateHandle<?>> state,
long stateSize,
Map<Integer,SerializedValue<StateHandle<?>>> kvState) |
Constructor and Description |
---|
PendingCheckpoint(JobID jobId,
long checkpointId,
long checkpointTimestamp,
Map<ExecutionAttemptID,ExecutionVertex> verticesToConfirm,
Executor executor) |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
TaskDeploymentDescriptor.getExecutionAttemptId() |
Constructor and Description |
---|
TaskDeploymentDescriptor(SerializedValue<JobInformation> serializedJobInformation,
SerializedValue<TaskInformation> serializedTaskInformation,
ExecutionAttemptID executionAttemptId,
int subtaskIndex,
int attemptNumber,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors) |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
Environment.getExecutionId()
Gets the ID of the task execution attempt.
|
Modifier and Type | Method and Description |
---|---|
void |
LibraryCacheManager.registerTask(JobID id,
ExecutionAttemptID execution,
Collection<BlobKey> requiredJarFiles,
Collection<URL> requiredClasspaths)
Registers a job task execution with its required jar files and classpaths.
|
void |
FallbackLibraryCacheManager.registerTask(JobID id,
ExecutionAttemptID execution,
Collection<BlobKey> requiredJarFiles,
Collection<URL> requiredClasspaths) |
void |
BlobLibraryCacheManager.registerTask(JobID jobId,
ExecutionAttemptID task,
Collection<BlobKey> requiredJarFiles,
Collection<URL> requiredClasspaths) |
void |
LibraryCacheManager.unregisterTask(JobID id,
ExecutionAttemptID execution)
Unregisters a job from the library cache manager.
|
void |
FallbackLibraryCacheManager.unregisterTask(JobID id,
ExecutionAttemptID execution) |
void |
BlobLibraryCacheManager.unregisterTask(JobID jobId,
ExecutionAttemptID task) |
Modifier and Type | Method and Description |
---|---|
static ExecutionAttemptID |
ExecutionAttemptID.fromByteBuf(io.netty.buffer.ByteBuf buf) |
ExecutionAttemptID |
Execution.getAttemptId() |
Modifier and Type | Method and Description |
---|---|
Map<ExecutionAttemptID,Map<AccumulatorRegistry.Metric,Accumulator<?,?>>> |
ExecutionGraph.getFlinkAccumulators()
Gets the internal flink accumulator map of maps which contains some metrics.
|
Map<ExecutionAttemptID,Execution> |
ExecutionGraph.getRegisteredExecutions() |
Modifier and Type | Method and Description |
---|---|
boolean |
ExecutionVertex.sendMessageToCurrentExecution(Serializable message,
ExecutionAttemptID attemptID) |
boolean |
ExecutionVertex.sendMessageToCurrentExecution(Serializable message,
ExecutionAttemptID attemptID,
ActorGateway sender) |
Modifier and Type | Method and Description |
---|---|
void |
PartitionProducerStateChecker.requestPartitionProducerState(JobID jobId,
ExecutionAttemptID receiverExecutionId,
IntermediateDataSetID intermediateDataSetId,
ResultPartitionID resultPartitionId)
Requests the execution state of the execution producing a result partition.
|
Modifier and Type | Field and Description |
---|---|
com.google.common.collect.Table<ExecutionAttemptID,IntermediateResultPartitionID,ResultPartition> |
ResultPartitionManager.registeredPartitions |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
ResultPartitionID.getProducerId() |
Modifier and Type | Method and Description |
---|---|
void |
ResultPartitionManager.releasePartitionsProducedBy(ExecutionAttemptID executionId) |
void |
ResultPartitionManager.releasePartitionsProducedBy(ExecutionAttemptID executionId,
Throwable cause) |
Constructor and Description |
---|
ResultPartitionID(IntermediateResultPartitionID partitionId,
ExecutionAttemptID producerId) |
Modifier and Type | Method and Description |
---|---|
static SingleInputGate |
SingleInputGate.create(String owningTaskName,
JobID jobId,
ExecutionAttemptID executionId,
InputGateDeploymentDescriptor igdd,
NetworkEnvironment networkEnvironment,
IOMetricGroup metrics)
Creates an input gate and all of its input channels.
|
Constructor and Description |
---|
SingleInputGate(String owningTaskName,
JobID jobId,
ExecutionAttemptID executionId,
IntermediateDataSetID consumedResultId,
int consumedSubpartitionIndex,
int numberOfInputChannels,
PartitionProducerStateChecker partitionStateChecker,
IOMetricGroup metrics) |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
TaskMessages.CancelTask.attemptID() |
ExecutionAttemptID |
TaskMessages.StopTask.attemptID() |
ExecutionAttemptID |
JobManagerMessages.RequestNextInputSplit.executionAttempt() |
ExecutionAttemptID |
StackTraceSampleMessages.TriggerStackTraceSample.executionId() |
ExecutionAttemptID |
StackTraceSampleMessages.ResponseStackTraceSampleSuccess.executionId() |
ExecutionAttemptID |
StackTraceSampleMessages.ResponseStackTraceSampleFailure.executionId() |
ExecutionAttemptID |
StackTraceSampleMessages.SampleTaskStackTrace.executionId() |
ExecutionAttemptID |
ExecutionGraphMessages.ExecutionStateChanged.executionID() |
ExecutionAttemptID |
TaskMessages.FailTask.executionID() |
ExecutionAttemptID |
TaskMessages.TaskInFinalState.executionID() |
ExecutionAttemptID |
TaskMessages.UpdateTaskSinglePartitionInfo.executionID() |
ExecutionAttemptID |
TaskMessages.UpdateTaskMultiplePartitionInfos.executionID() |
ExecutionAttemptID |
TaskMessages.FailIntermediateResultPartitions.executionID() |
ExecutionAttemptID |
TaskMessages.TaskOperationResult.executionID() |
abstract ExecutionAttemptID |
TaskMessages.UpdatePartitionInfo.executionID() |
ExecutionAttemptID |
TaskMessages.PartitionProducerState.receiverExecutionId() |
ExecutionAttemptID |
JobManagerMessages.RequestPartitionProducerState.receiverExecutionId() |
Modifier and Type | Method and Description |
---|---|
static TaskMessages.UpdateTaskMultiplePartitionInfos |
TaskMessages.createUpdateTaskMultiplePartitionInfos(ExecutionAttemptID executionID,
List<IntermediateDataSetID> resultIDs,
List<InputChannelDeploymentDescriptor> partitionInfos) |
TaskMessages.UpdateTaskMultiplePartitionInfos |
TaskMessages$.createUpdateTaskMultiplePartitionInfos(ExecutionAttemptID executionID,
List<IntermediateDataSetID> resultIDs,
List<InputChannelDeploymentDescriptor> partitionInfos) |
Constructor and Description |
---|
CancelTask(ExecutionAttemptID attemptID) |
ExecutionStateChanged(JobID jobID,
JobVertexID vertexID,
String taskName,
int totalNumberOfSubTasks,
int subtaskIndex,
ExecutionAttemptID executionID,
ExecutionState newExecutionState,
long timestamp,
String optionalMessage) |
FailIntermediateResultPartitions(ExecutionAttemptID executionID) |
FailTask(ExecutionAttemptID executionID,
Throwable cause) |
PartitionProducerState(ExecutionAttemptID receiverExecutionId,
scala.util.Either<scala.Tuple3<IntermediateDataSetID,ResultPartitionID,ExecutionState>,Exception> result) |
RequestNextInputSplit(JobID jobID,
JobVertexID vertexID,
ExecutionAttemptID executionAttempt) |
RequestPartitionProducerState(JobID jobId,
ExecutionAttemptID receiverExecutionId,
IntermediateDataSetID intermediateDataSetId,
ResultPartitionID resultPartitionId) |
ResponseStackTraceSampleFailure(int sampleId,
ExecutionAttemptID executionId,
Exception cause) |
ResponseStackTraceSampleSuccess(int sampleId,
ExecutionAttemptID executionId,
List<StackTraceElement[]> samples) |
SampleTaskStackTrace(int sampleId,
ExecutionAttemptID executionId,
scala.concurrent.duration.FiniteDuration delayBetweenSamples,
int maxStackTraceDepth,
int numRemainingSamples,
List<StackTraceElement[]> currentTraces,
akka.actor.ActorRef sender) |
StopTask(ExecutionAttemptID attemptID) |
TaskInFinalState(ExecutionAttemptID executionID) |
TaskOperationResult(ExecutionAttemptID executionID,
boolean success) |
TaskOperationResult(ExecutionAttemptID executionID,
boolean success,
String description) |
TriggerStackTraceSample(int sampleId,
ExecutionAttemptID executionId,
int numSamples,
scala.concurrent.duration.FiniteDuration delayBetweenSamples,
int maxStackTraceDepth) |
UpdateTaskMultiplePartitionInfos(ExecutionAttemptID executionID,
scala.collection.Seq<scala.Tuple2<IntermediateDataSetID,InputChannelDeploymentDescriptor>> partitionInfos) |
UpdateTaskSinglePartitionInfo(ExecutionAttemptID executionID,
IntermediateDataSetID resultId,
InputChannelDeploymentDescriptor partitionInfo) |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
AbstractCheckpointMessage.getTaskExecutionId() |
Constructor and Description |
---|
AbstractCheckpointMessage(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId) |
AcknowledgeCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId) |
AcknowledgeCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId,
SerializedValue<StateHandle<?>> state,
long stateSize) |
DeclineCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId) |
DeclineCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId,
Throwable reason) |
NotifyCheckpointComplete(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId,
long timestamp) |
TriggerCheckpoint(JobID job,
ExecutionAttemptID taskExecutionId,
long checkpointId,
long timestamp) |
Modifier and Type | Method and Description |
---|---|
TaskMetricGroup |
TaskManagerJobMetricGroup.addTask(JobVertexID jobVertexId,
ExecutionAttemptID executionAttemptID,
String taskName,
int subtaskIndex,
int attemptNumber) |
TaskMetricGroup |
TaskManagerMetricGroup.addTaskForJob(JobID jobId,
String jobName,
JobVertexID jobVertexId,
ExecutionAttemptID executionAttemptId,
String taskName,
int subtaskIndex,
int attemptNumber) |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
Task.getExecutionId() |
ExecutionAttemptID |
RuntimeEnvironment.getExecutionId() |
ExecutionAttemptID |
TaskExecutionState.getID()
Returns the ID of the task this result belongs to
|
Modifier and Type | Method and Description |
---|---|
protected HashMap<ExecutionAttemptID,Task> |
TaskManager.runningTasks()
Registry of all tasks currently executed by this TaskManager
|
Constructor and Description |
---|
RuntimeEnvironment(JobID jobId,
JobVertexID jobVertexId,
ExecutionAttemptID executionId,
ExecutionConfig executionConfig,
TaskInfo taskInfo,
Configuration jobConfiguration,
Configuration taskConfiguration,
ClassLoader userCodeClassLoader,
MemoryManager memManager,
IOManager ioManager,
BroadcastVariableManager bcVarManager,
AccumulatorRegistry accumulatorRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
ActorGateway jobManager,
TaskManagerRuntimeInfo taskManagerInfo,
TaskMetricGroup metrics,
Task containingTask) |
Task(JobInformation jobInformation,
TaskInformation taskInformation,
ExecutionAttemptID executionAttemptID,
int subtaskIndex,
int attemptNumber,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
MemoryManager memManager,
IOManager ioManager,
NetworkEnvironment networkEnvironment,
BroadcastVariableManager bcVarManager,
ActorGateway taskManagerActor,
ActorGateway jobManagerActor,
scala.concurrent.duration.FiniteDuration actorAskTimeout,
LibraryCacheManager libraryCache,
FileCache fileCache,
TaskManagerRuntimeInfo taskManagerConfig,
TaskMetricGroup metricGroup)
IMPORTANT: This constructor may not start any work that would need to
be undone in the case of a failing task deployment.
|
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)
Creates a new task execution state update, with an attached exception.
|
TaskInputSplitProvider(ActorGateway jobManager,
JobID jobId,
JobVertexID vertexId,
ExecutionAttemptID executionID,
ClassLoader userCodeClassLoader,
scala.concurrent.duration.FiniteDuration timeout) |
Modifier and Type | Method and Description |
---|---|
ExecutionAttemptID |
TestingTaskManagerMessages.NotifyWhenTaskRemoved.executionID() |
ExecutionAttemptID |
TestingTaskManagerMessages.NotifyWhenTaskIsRunning.executionID() |
Modifier and Type | Method and Description |
---|---|
Map<ExecutionAttemptID,Task> |
TestingTaskManagerMessages.ResponseRunningTasks.asJava() |
Map<ExecutionAttemptID,Map<AccumulatorRegistry.Metric,Accumulator<?,?>>> |
TestingJobManagerMessages.UpdatedAccumulators.flinkAccumulators() |
scala.collection.immutable.Map<ExecutionAttemptID,Task> |
TestingTaskManagerMessages.ResponseRunningTasks.tasks() |
scala.collection.mutable.HashSet<ExecutionAttemptID> |
TestingTaskManagerLike.unregisteredTasks() |
scala.collection.mutable.HashMap<ExecutionAttemptID,scala.collection.immutable.Set<akka.actor.ActorRef>> |
TestingTaskManagerLike.waitForRemoval() |
scala.collection.mutable.HashMap<ExecutionAttemptID,scala.collection.immutable.Set<akka.actor.ActorRef>> |
TestingTaskManagerLike.waitForRunning() |
Constructor and Description |
---|
NotifyWhenTaskIsRunning(ExecutionAttemptID executionID) |
NotifyWhenTaskRemoved(ExecutionAttemptID executionID) |
Constructor and Description |
---|
ResponseRunningTasks(scala.collection.immutable.Map<ExecutionAttemptID,Task> tasks) |
UpdatedAccumulators(JobID jobID,
Map<ExecutionAttemptID,Map<AccumulatorRegistry.Metric,Accumulator<?,?>>> flinkAccumulators,
Map<String,Accumulator<?,?>> userAccumulators) |
Modifier and Type | Method and Description |
---|---|
Map<ExecutionAttemptID,List<StackTraceElement[]>> |
StackTraceSample.getStackTraces()
Returns the a map of stack traces by execution ID.
|
Modifier and Type | Method and Description |
---|---|
void |
StackTraceSampleCoordinator.collectStackTraces(int sampleId,
ExecutionAttemptID executionId,
List<StackTraceElement[]> stackTraces)
Collects stack traces of a task.
|
Constructor and Description |
---|
StackTraceSample(int sampleId,
long startTime,
long endTime,
Map<ExecutionAttemptID,List<StackTraceElement[]>> stackTracesByTask)
Creates a stack trace sample
|
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.