Package | Description |
---|---|
org.apache.flink.migration.runtime.checkpoint | |
org.apache.flink.runtime.broadcast | |
org.apache.flink.runtime.checkpoint | |
org.apache.flink.runtime.checkpoint.savepoint | |
org.apache.flink.runtime.execution | |
org.apache.flink.runtime.executiongraph | |
org.apache.flink.runtime.instance | |
org.apache.flink.runtime.io.network | |
org.apache.flink.runtime.jobgraph | |
org.apache.flink.runtime.jobgraph.tasks | |
org.apache.flink.runtime.jobmanager.scheduler | |
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.metrics.groups | |
org.apache.flink.runtime.query |
This package contains all KvState query related classes.
|
org.apache.flink.runtime.taskexecutor.rpc | |
org.apache.flink.runtime.taskmanager | |
org.apache.flink.runtime.webmonitor.handlers |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
TaskState.getJobVertexID()
Deprecated.
|
Constructor and Description |
---|
TaskState(JobVertexID jobVertexID,
int parallelism)
Deprecated.
|
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 |
TaskStateStats.getJobVertexId()
Returns the ID of the operator the statistics belong to.
|
JobVertexID |
TaskState.getJobVertexID()
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
TaskStateStats |
AbstractCheckpointStats.getTaskStateStats(JobVertexID jobVertexId)
Returns the task state stats for the given job vertex ID or
null if no task with such an ID is available. |
Modifier and Type | Method and Description |
---|---|
boolean |
CheckpointCoordinator.restoreLatestCheckpointedState(Map<JobVertexID,ExecutionJobVertex> tasks,
boolean errorIfNoCheckpoint,
boolean allowNonRestoredState)
Restores the latest checkpointed state.
|
boolean |
CheckpointCoordinator.restoreSavepoint(String savepointPath,
boolean allowNonRestored,
Map<JobVertexID,ExecutionJobVertex> tasks,
ClassLoader userClassLoader)
Restore the state with given savepoint
|
Constructor and Description |
---|
TaskState(JobVertexID jobVertexID,
int parallelism,
int maxParallelism,
int chainLength)
Deprecated.
|
Constructor and Description |
---|
StateAssignmentOperation(Map<JobVertexID,ExecutionJobVertex> tasks,
Map<OperatorID,OperatorState> operatorStates,
boolean allowNonRestoredState) |
Modifier and Type | Method and Description |
---|---|
static Savepoint |
SavepointV2.convertToOperatorStateSavepointV2(Map<JobVertexID,ExecutionJobVertex> tasks,
Savepoint savepoint)
Deprecated.
Only kept for backwards-compatibility with versions < 1.3. Will be removed in the future.
|
static CompletedCheckpoint |
SavepointLoader.loadAndValidateSavepoint(JobID jobId,
Map<JobVertexID,ExecutionJobVertex> tasks,
String savepointPath,
ClassLoader classLoader,
boolean allowNonRestoredState)
Loads a savepoint back as a
CompletedCheckpoint . |
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() |
JobVertexID |
ArchivedExecutionJobVertex.getJobVertexId() |
JobVertexID |
AccessExecutionJobVertex.getJobVertexId()
Returns the
JobVertexID for this job vertex. |
Modifier and Type | Method and Description |
---|---|
Map<JobVertexID,ExecutionJobVertex> |
ExecutionGraph.getAllVertices() |
Map<JobVertexID,AccessExecutionJobVertex> |
ArchivedExecutionGraph.getAllVertices() |
Map<JobVertexID,? extends AccessExecutionJobVertex> |
AccessExecutionGraph.getAllVertices()
Returns a map containing all job vertices for this execution graph.
|
static Map<JobVertexID,ExecutionJobVertex> |
ExecutionJobVertex.includeLegacyJobVertexIDs(Map<JobVertexID,ExecutionJobVertex> tasks) |
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.
|
ExecutionJobVertex |
ExecutionGraph.getJobVertex(JobVertexID id) |
ArchivedExecutionJobVertex |
ArchivedExecutionGraph.getJobVertex(JobVertexID id) |
AccessExecutionJobVertex |
AccessExecutionGraph.getJobVertex(JobVertexID id)
Returns the job vertex for the given
JobVertexID . |
Modifier and Type | Method and Description |
---|---|
static Map<JobVertexID,ExecutionJobVertex> |
ExecutionJobVertex.includeLegacyJobVertexIDs(Map<JobVertexID,ExecutionJobVertex> tasks) |
Constructor and Description |
---|
ArchivedExecutionJobVertex(ArchivedExecutionVertex[] taskVertices,
JobVertexID id,
String name,
int parallelism,
int maxParallelism,
StringifiedAccumulatorResult[] archivedUserAccumulators) |
TaskInformation(JobVertexID jobVertexId,
String taskName,
int numberOfSubtasks,
int maxNumberOfSubtaks,
String invokableClassName,
Configuration taskConfiguration) |
Constructor and Description |
---|
ArchivedExecutionGraph(JobID jobID,
String jobName,
Map<JobVertexID,ArchivedExecutionJobVertex> tasks,
List<ArchivedExecutionJobVertex> verticesInCreationOrder,
long[] stateTimestamps,
JobStatus state,
String failureCause,
String jsonPlan,
StringifiedAccumulatorResult[] archivedUserAccumulators,
Map<String,SerializedValue<Object>> serializedUserAccumulators,
ArchivedExecutionConfig executionConfig,
boolean isStoppable,
JobCheckpointingSettings jobCheckpointingSettings,
CheckpointStatsSnapshot checkpointStatsSnapshot) |
Modifier and Type | Method and Description |
---|---|
SimpleSlot |
SlotSharingGroupAssignment.addSharedSlotAndAllocateSubSlot(SharedSlot sharedSlot,
Locality locality,
JobVertexID groupId) |
Modifier and Type | Method and Description |
---|---|
TaskKvStateRegistry |
NetworkEnvironment.createKvStateTaskRegistry(JobID jobId,
JobVertexID jobVertexId) |
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 |
---|---|
List<JobVertexID> |
JobVertex.getIdAlternatives()
Returns a list of all alternative IDs 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.
|
static OperatorID |
OperatorID.fromJobVertexID(JobVertexID id) |
Constructor and Description |
---|
InputFormatVertex(String name,
JobVertexID id) |
InputFormatVertex(String name,
JobVertexID id,
List<JobVertexID> alternativeIds,
List<OperatorID> operatorIds,
List<OperatorID> alternativeOperatorIds) |
JobVertex(String name,
JobVertexID id)
Constructs a new job vertex and assigns it with the given name.
|
JobVertex(String name,
JobVertexID primaryId,
List<JobVertexID> alternativeIds,
List<OperatorID> operatorIds,
List<OperatorID> alternativeOperatorIds)
Constructs a new job vertex and assigns it with the given name.
|
Constructor and Description |
---|
InputFormatVertex(String name,
JobVertexID id,
List<JobVertexID> alternativeIds,
List<OperatorID> operatorIds,
List<OperatorID> alternativeOperatorIds) |
JobVertex(String name,
JobVertexID primaryId,
List<JobVertexID> alternativeIds,
List<OperatorID> operatorIds,
List<OperatorID> alternativeOperatorIds)
Constructs a new job vertex and assigns it with the given name.
|
Modifier and Type | Method and Description |
---|---|
List<JobVertexID> |
JobCheckpointingSettings.getVerticesToAcknowledge() |
List<JobVertexID> |
JobCheckpointingSettings.getVerticesToConfirm() |
List<JobVertexID> |
JobCheckpointingSettings.getVerticesToTrigger() |
Constructor and Description |
---|
JobCheckpointingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints,
ExternalizedCheckpointSettings externalizedCheckpointSettings,
SerializedValue<StateBackend> defaultStateBackend,
boolean isExactlyOnce) |
JobCheckpointingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints,
ExternalizedCheckpointSettings externalizedCheckpointSettings,
SerializedValue<StateBackend> defaultStateBackend,
boolean isExactlyOnce) |
JobCheckpointingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints,
ExternalizedCheckpointSettings externalizedCheckpointSettings,
SerializedValue<StateBackend> defaultStateBackend,
boolean isExactlyOnce) |
JobCheckpointingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints,
ExternalizedCheckpointSettings externalizedCheckpointSettings,
SerializedValue<StateBackend> defaultStateBackend,
SerializedValue<MasterTriggerRestoreHook.Factory[]> masterHooks,
boolean isExactlyOnce) |
JobCheckpointingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints,
ExternalizedCheckpointSettings externalizedCheckpointSettings,
SerializedValue<StateBackend> defaultStateBackend,
SerializedValue<MasterTriggerRestoreHook.Factory[]> masterHooks,
boolean isExactlyOnce) |
JobCheckpointingSettings(List<JobVertexID> verticesToTrigger,
List<JobVertexID> verticesToAcknowledge,
List<JobVertexID> verticesToConfirm,
long checkpointInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpoints,
ExternalizedCheckpointSettings externalizedCheckpointSettings,
SerializedValue<StateBackend> defaultStateBackend,
SerializedValue<MasterTriggerRestoreHook.Factory[]> masterHooks,
boolean isExactlyOnce) |
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 |
---|---|
void |
JobMasterGateway.notifyKvStateRegistered(JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
KvStateServerAddress kvStateServerAddress) |
void |
JobMaster.notifyKvStateRegistered(JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
KvStateServerAddress kvStateServerAddress) |
void |
JobMasterGateway.notifyKvStateUnregistered(JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName) |
void |
JobMaster.notifyKvStateUnregistered(JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName) |
Future<SerializedInputSplit> |
JobMasterGateway.requestNextInputSplit(UUID leaderSessionID,
JobVertexID vertexID,
ExecutionAttemptID executionAttempt)
Requesting next input split for the
ExecutionJobVertex . |
SerializedInputSplit |
JobMaster.requestNextInputSplit(UUID leaderSessionID,
JobVertexID vertexID,
ExecutionAttemptID executionAttempt) |
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 |
KvStateMessage.NotifyKvStateRegistered.getJobVertexId()
Returns the JobVertexID the KvState instance belongs to
|
JobVertexID |
KvStateMessage.NotifyKvStateUnregistered.getJobVertexId()
Returns the JobVertexID the KvState instance belongs to
|
JobVertexID |
KvStateLocation.getJobVertexId()
Returns the JobVertexID the KvState instances belong to.
|
Modifier and Type | Method and Description |
---|---|
TaskKvStateRegistry |
KvStateRegistry.createTaskRegistry(JobID jobId,
JobVertexID jobVertexId)
Creates a
TaskKvStateRegistry facade for the Task
identified by the given JobID and JobVertexID instance. |
void |
KvStateRegistryListener.notifyKvStateRegistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId)
Notifies the listener about a registered KvState instance.
|
void |
KvStateRegistryGateway.notifyKvStateRegistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
KvStateServerAddress kvStateServerAddress)
Notifies the listener about a registered KvState instance.
|
void |
KvStateLocationRegistry.notifyKvStateRegistered(JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
KvStateServerAddress kvStateServerAddress)
Notifies the registry about a registered KvState instance.
|
void |
KvStateRegistryListener.notifyKvStateUnregistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName)
Notifies the listener about an unregistered KvState instance.
|
void |
KvStateRegistryGateway.notifyKvStateUnregistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName)
Notifies the listener about an unregistered KvState instance.
|
void |
KvStateLocationRegistry.notifyKvStateUnregistered(JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName)
Notifies the registry about an unregistered KvState instance.
|
KvStateID |
KvStateRegistry.registerKvState(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
InternalKvState<?> kvState)
Registers the KvState instance and returns the assigned ID.
|
void |
KvStateRegistry.unregisterKvState(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId)
Unregisters the KvState instance identified by the given KvStateID.
|
Constructor and Description |
---|
KvStateLocation(JobID jobId,
JobVertexID jobVertexId,
int numKeyGroups,
String registrationName)
Creates the location information
|
NotifyKvStateRegistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
KvStateServerAddress kvStateServerAddress)
Notifies the JobManager about a registered
InternalKvState instance. |
NotifyKvStateUnregistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName)
Notifies the JobManager about an unregistered
InternalKvState instance. |
Constructor and Description |
---|
KvStateLocationRegistry(JobID jobId,
Map<JobVertexID,ExecutionJobVertex> jobVertices)
Creates the registry for the job.
|
Modifier and Type | Method and Description |
---|---|
void |
RpcKvStateRegistryListener.notifyKvStateRegistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId) |
void |
RpcKvStateRegistryListener.notifyKvStateUnregistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName) |
Constructor and Description |
---|
RpcInputSplitProvider(UUID jobMasterLeaderId,
JobMasterGateway jobMasterGateway,
JobID jobID,
JobVertexID jobVertexID,
ExecutionAttemptID executionAttemptID,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
Task.getJobVertexId() |
JobVertexID |
RuntimeEnvironment.getJobVertexId() |
Modifier and Type | Method and Description |
---|---|
void |
ActorGatewayKvStateRegistryListener.notifyKvStateRegistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId) |
void |
ActorGatewayKvStateRegistryListener.notifyKvStateUnregistered(JobID jobId,
JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName) |
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,
TaskKvStateRegistry kvStateRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
CheckpointResponder checkpointResponder,
TaskManagerRuntimeInfo taskManagerInfo,
TaskMetricGroup metrics,
Task containingTask) |
TaskInputSplitProvider(ActorGateway jobManager,
JobID jobID,
JobVertexID vertexID,
ExecutionAttemptID executionID,
scala.concurrent.duration.FiniteDuration timeout) |
Modifier and Type | Method and Description |
---|---|
static JobVertexID |
AbstractJobVertexRequestHandler.parseJobVertexId(Map<String,String> params)
Returns the job vertex ID parsed from the provided parameters.
|
Copyright © 2014–2018 The Apache Software Foundation. All rights reserved.