Package | Description |
---|---|
org.apache.flink.runtime.broadcast | |
org.apache.flink.runtime.checkpoint | |
org.apache.flink.runtime.checkpoint.savepoint | |
org.apache.flink.runtime.checkpoint.stats | |
org.apache.flink.runtime.execution | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.instance | |
org.apache.flink.runtime.jobgraph | |
org.apache.flink.runtime.jobgraph.tasks | |
org.apache.flink.runtime.jobmanager.scheduler | |
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.metrics.groups | |
org.apache.flink.runtime.taskmanager |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
BroadcastVariableKey.getVertexId() |
Constructor and Description |
---|
BroadcastVariableKey(JobVertexID vertexId,
String name,
int superstep) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
TaskState.getJobVertexID() |
Modifier and Type | Method and Description |
---|---|
Map<JobVertexID,TaskState> |
PendingCheckpoint.getTaskStates() |
Map<JobVertexID,TaskState> |
CompletedCheckpoint.getTaskStates() |
Modifier and Type | Method and Description |
---|---|
TaskState |
CompletedCheckpoint.getTaskState(JobVertexID jobVertexID) |
Modifier and Type | Method and Description |
---|---|
boolean |
CheckpointCoordinator.restoreLatestCheckpointedState(Map<JobVertexID,ExecutionJobVertex> tasks,
boolean errorIfNoCheckpoint,
boolean allOrNothingState) |
Constructor and Description |
---|
TaskState(JobVertexID jobVertexID,
int parallelism) |
Constructor and Description |
---|
CompletedCheckpoint(JobID job,
long checkpointID,
long timestamp,
long completionTimestamp,
Map<JobVertexID,TaskState> taskStates) |
Modifier and Type | Method and Description |
---|---|
void |
SavepointCoordinator.restoreSavepoint(Map<JobVertexID,ExecutionJobVertex> tasks,
String savepointPath,
boolean allowNonRestoredState)
Resets the state of
Execution instances back to the state of a savepoint. |
Modifier and Type | Method and Description |
---|---|
scala.Option<OperatorCheckpointStats> |
SimpleCheckpointStatsTracker.getOperatorStats(JobVertexID operatorId) |
scala.Option<OperatorCheckpointStats> |
DisabledCheckpointStatsTracker.getOperatorStats(JobVertexID operatorId) |
scala.Option<OperatorCheckpointStats> |
CheckpointStatsTracker.getOperatorStats(JobVertexID operatorId)
Returns a snapshot of the checkpoint statistics for an operator.
|
Modifier and Type | Method and Description |
---|---|
JobVertexID |
Environment.getJobVertexId()
Gets the ID of the JobVertex for which this task executes a parallel subtask.
|
Modifier and Type | Method and Description |
---|---|
JobVertexID |
ExecutionVertex.getJobvertexId() |
JobVertexID |
TaskInformation.getJobVertexId() |
JobVertexID |
ExecutionJobVertex.getJobVertexId() |
Modifier and Type | Method and Description |
---|---|
Map<JobVertexID,ExecutionJobVertex> |
ExecutionGraph.getAllVertices() |
Modifier and Type | Method and Description |
---|---|
ExecutionJobVertex |
ExecutionGraph.getJobVertex(JobVertexID id) |
Constructor and Description |
---|
TaskInformation(JobVertexID jobVertexId,
String taskName,
int parallelism,
String invokableClassName,
Configuration taskConfiguration) |
Modifier and Type | Method and Description |
---|---|
SimpleSlot |
SlotSharingGroupAssignment.addSharedSlotAndAllocateSubSlot(SharedSlot sharedSlot,
Locality locality,
JobVertexID groupId) |
Modifier and Type | Method and Description |
---|---|
static JobVertexID |
JobVertexID.fromHexString(String hexString) |
JobVertexID |
JobVertex.getID()
Returns the ID of this job vertex.
|
Modifier and Type | Method and Description |
---|---|
JobVertex |
JobGraph.findVertexByID(JobVertexID id)
Searches for a vertex with a matching ID and returns it.
|
Constructor and Description |
---|
InputFormatVertex(String name,
JobVertexID id) |
JobVertex(String name,
JobVertexID id)
Constructs a new job vertex and assigns it with the given name.
|
Modifier and Type | Method and Description |
---|---|
List<JobVertexID> |
JobSnapshottingSettings.getVerticesToAcknowledge() |
List<JobVertexID> |
JobSnapshottingSettings.getVerticesToConfirm() |
List<JobVertexID> |
JobSnapshottingSettings.getVerticesToTrigger() |
Constructor and Description |
---|
JobSnapshottingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints) |
JobSnapshottingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints) |
JobSnapshottingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
ScheduledUnit.getJobVertexId() |
Modifier and Type | Method and Description |
---|---|
Set<JobVertexID> |
SlotSharingGroup.getJobVertexIds() |
Modifier and Type | Method and Description |
---|---|
void |
SlotSharingGroup.addVertexToGroup(JobVertexID id) |
void |
SlotSharingGroup.removeVertexFromGroup(JobVertexID id) |
Constructor and Description |
---|
SlotSharingGroup(JobVertexID... sharedVertices) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
ExecutionGraphMessages.ExecutionStateChanged.vertexID() |
JobVertexID |
JobManagerMessages.RequestNextInputSplit.vertexID() |
Constructor and Description |
---|
ExecutionStateChanged(JobID jobID,
JobVertexID vertexID,
String taskName,
int totalNumberOfSubTasks,
int subtaskIndex,
ExecutionAttemptID executionID,
ExecutionState newExecutionState,
long timestamp,
String optionalMessage) |
RequestNextInputSplit(JobID jobID,
JobVertexID vertexID,
ExecutionAttemptID executionAttempt) |
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 |
---|---|
JobVertexID |
Task.getJobVertexId() |
JobVertexID |
RuntimeEnvironment.getJobVertexId() |
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) |
TaskInputSplitProvider(ActorGateway jobManager,
JobID jobId,
JobVertexID vertexId,
ExecutionAttemptID executionID,
ClassLoader userCodeClassLoader,
scala.concurrent.duration.FiniteDuration timeout) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.