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 |
---|---|
static CompletedCheckpoint |
Checkpoints.loadAndValidateCheckpoint(JobID jobId,
Map<JobVertexID,ExecutionJobVertex> tasks,
CompletedCheckpointStorageLocation location,
ClassLoader classLoader,
boolean allowNonRestoredState) |
boolean |
CheckpointCoordinator.restoreLatestCheckpointedState(Map<JobVertexID,ExecutionJobVertex> tasks,
boolean errorIfNoCheckpoint,
boolean allowNonRestoredState)
Deprecated.
|
boolean |
CheckpointCoordinator.restoreSavepoint(String savepointPointer,
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.
|
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.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<OperatorBackPressureStatsResponse> |
Dispatcher.requestOperatorBackPressureStats(JobID jobId,
JobVertexID jobVertexId) |
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 |
ArchivedExecutionJobVertex.getJobVertexId() |
JobVertexID |
AccessExecutionJobVertex.getJobVertexId()
Returns the
JobVertexID for this job vertex. |
JobVertexID |
TaskInformation.getJobVertexId() |
JobVertexID |
ExecutionJobVertex.getJobVertexId() |
Modifier and Type | Method and Description |
---|---|
Map<JobVertexID,? extends AccessExecutionJobVertex> |
AccessExecutionGraph.getAllVertices()
Returns a map containing all job vertices for this execution graph.
|
Map<JobVertexID,ExecutionJobVertex> |
ExecutionGraph.getAllVertices() |
Map<JobVertexID,AccessExecutionJobVertex> |
ArchivedExecutionGraph.getAllVertices() |
static Map<JobVertexID,ExecutionJobVertex> |
ExecutionJobVertex.includeLegacyJobVertexIDs(Map<JobVertexID,ExecutionJobVertex> tasks) |
Modifier and Type | Method and Description |
---|---|
AccessExecutionJobVertex |
AccessExecutionGraph.getJobVertex(JobVertexID id)
Returns the job vertex for the given
JobVertexID . |
ExecutionJobVertex |
ExecutionGraph.getJobVertex(JobVertexID id) |
ArchivedExecutionJobVertex |
ArchivedExecutionGraph.getJobVertex(JobVertexID id) |
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,
ResourceProfile resourceProfile,
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,
ErrorInfo failureCause,
String jsonPlan,
StringifiedAccumulatorResult[] archivedUserAccumulators,
Map<String,SerializedValue<OptionalFailure<Object>>> serializedUserAccumulators,
ArchivedExecutionConfig executionConfig,
boolean isStoppable,
CheckpointCoordinatorConfiguration jobCheckpointingConfiguration,
CheckpointStatsSnapshot checkpointStatsSnapshot,
String stateBackendName) |
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 |
---|
InputOutputFormatVertex(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 |
---|
InputOutputFormatVertex(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() |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
DefaultLogicalVertex.getId() |
Modifier and Type | Method and Description |
---|---|
Set<JobVertexID> |
LogicalPipelinedRegion.getVertexIDs() |
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,
ResourceSpec resource) |
void |
SlotSharingGroup.removeVertexFromGroup(JobVertexID id,
ResourceSpec resource) |
Constructor and Description |
---|
ScheduledUnit(Execution task,
JobVertexID jobVertexId,
SlotSharingGroupId slotSharingGroupId,
CoLocationConstraint coLocationConstraint) |
ScheduledUnit(JobVertexID jobVertexId,
SlotSharingGroupId slotSharingGroupId,
CoLocationConstraint coLocationConstraint) |
Modifier and Type | Field and Description |
---|---|
protected JobVertexID |
TaskMetricGroup.vertexId |
Modifier and Type | Method and Description |
---|---|
TaskMetricGroup |
UnregisteredMetricGroups.UnregisteredTaskManagerJobMetricGroup.addTask(JobVertexID jobVertexId,
ExecutionAttemptID executionAttemptID,
String taskName,
int subtaskIndex,
int attemptNumber) |
TaskMetricGroup |
TaskManagerJobMetricGroup.addTask(JobVertexID jobVertexId,
ExecutionAttemptID executionAttemptID,
String taskName,
int subtaskIndex,
int attemptNumber) |
TaskMetricGroup |
UnregisteredMetricGroups.UnregisteredTaskManagerMetricGroup.addTaskForJob(JobID jobId,
String jobName,
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) |
Constructor and Description |
---|
TaskMetricGroup(MetricRegistry registry,
TaskManagerJobMetricGroup parent,
JobVertexID vertexId,
AbstractID executionId,
String taskName,
int subtaskIndex,
int attemptNumber) |
Modifier and Type | Method and Description |
---|---|
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 |
KvStateLocationRegistry.notifyKvStateRegistered(JobVertexID jobVertexId,
KeyGroupRange keyGroupRange,
String registrationName,
KvStateID kvStateId,
InetSocketAddress 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 |
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.
|
Constructor and Description |
---|
KvStateLocationRegistry(JobID jobId,
Map<JobVertexID,ExecutionJobVertex> jobVertices)
Creates the registry for the job.
|
Modifier and Type | Method and Description |
---|---|
protected JobVertexID |
JobVertexIdPathParameter.convertFromString(String value) |
Modifier and Type | Method and Description |
---|---|
protected String |
JobVertexIdPathParameter.convertToString(JobVertexID value) |
Constructor and Description |
---|
JobVertexDetailsInfo(JobVertexID id,
String name,
int parallelism,
long now,
List<SubtaskExecutionAttemptDetailsInfo> subtasks) |
JobVertexTaskManagersInfo(JobVertexID jobVertexID,
String name,
long now,
Collection<JobVertexTaskManagersInfo.TaskManagersInfo> taskManagerInfos) |
Modifier and Type | Method and Description |
---|---|
Map<JobVertexID,TaskCheckpointStatistics> |
CheckpointStatistics.getCheckpointStatisticsPerTask() |
Constructor and Description |
---|
CompletedCheckpointStatistics(long id,
CheckpointStatsStatus status,
boolean savepoint,
long triggerTimestamp,
long latestAckTimestamp,
long stateSize,
long duration,
long alignmentBuffered,
int numSubtasks,
int numAckSubtasks,
Map<JobVertexID,TaskCheckpointStatistics> checkpointingStatisticsPerTask,
String externalPath,
boolean discarded) |
FailedCheckpointStatistics(long id,
CheckpointStatsStatus status,
boolean savepoint,
long triggerTimestamp,
long latestAckTimestamp,
long stateSize,
long duration,
long alignmentBuffered,
int numSubtasks,
int numAckSubtasks,
Map<JobVertexID,TaskCheckpointStatistics> checkpointingStatisticsPerTask,
long failureTimestamp,
String failureMessage) |
PendingCheckpointStatistics(long id,
CheckpointStatsStatus status,
boolean savepoint,
long triggerTimestamp,
long latestAckTimestamp,
long stateSize,
long duration,
long alignmentBuffered,
int numSubtasks,
int numAckSubtasks,
Map<JobVertexID,TaskCheckpointStatistics> checkpointingStatisticsPerTask) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
JobDetailsInfo.JobVertexDetailsInfo.getJobVertexID() |
Modifier and Type | Method and Description |
---|---|
static SubtaskExecutionAttemptDetailsInfo |
SubtaskExecutionAttemptDetailsInfo.create(AccessExecution execution,
MetricFetcher metricFetcher,
JobID jobID,
JobVertexID jobVertexID) |
Constructor and Description |
---|
JobVertexDetailsInfo(JobVertexID jobVertexID,
String name,
int parallelism,
ExecutionState executionState,
long startTime,
long endTime,
long duration,
Map<ExecutionState,Integer> tasksPerState,
IOMetricsInfo jobVertexMetrics) |
SubtasksAllAccumulatorsInfo(JobVertexID jobVertexId,
int parallelism,
Collection<SubtasksAllAccumulatorsInfo.SubtaskAccumulatorsInfo> subtaskAccumulatorsInfos) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
JobVertexIDDeserializer.deserialize(org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonParser p,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.DeserializationContext ctxt) |
Modifier and Type | Method and Description |
---|---|
void |
JobVertexIDKeySerializer.serialize(JobVertexID value,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonGenerator gen,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.SerializerProvider provider) |
void |
JobVertexIDSerializer.serialize(JobVertexID value,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.core.JsonGenerator gen,
org.apache.flink.shaded.jackson2.com.fasterxml.jackson.databind.SerializerProvider provider) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
ExecutionVertexID.getJobVertexId() |
Constructor and Description |
---|
ExecutionVertexID(JobVertexID jobVertexId,
int subtaskIndex) |
Modifier and Type | Method and Description |
---|---|
TaskLocalStateStore |
TaskExecutorLocalStateStoresManager.localStateStoreForSubtask(JobID jobId,
AllocationID allocationID,
JobVertexID jobVertexID,
int subtaskIndex) |
Constructor and Description |
---|
LocalRecoveryDirectoryProviderImpl(File[] allocationBaseDirs,
JobID jobID,
JobVertexID jobVertexID,
int subtaskIndex) |
LocalRecoveryDirectoryProviderImpl(File allocationBaseDir,
JobID jobID,
JobVertexID jobVertexID,
int subtaskIndex) |
TaskLocalStateStoreImpl(JobID jobID,
AllocationID allocationID,
JobVertexID jobVertexID,
int subtaskIndex,
LocalRecoveryConfig localRecoveryConfig,
Executor discardExecutor) |
Modifier and Type | Method and Description |
---|---|
TaskKvStateRegistry |
KvStateService.createKvStateTaskRegistry(JobID jobId,
JobVertexID jobVertexId) |
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(JobMasterGateway jobMasterGateway,
JobVertexID jobVertexID,
ExecutionAttemptID executionAttemptID,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
JobVertexID |
RuntimeEnvironment.getJobVertexId() |
JobVertexID |
Task.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,
TaskStateManager taskStateManager,
GlobalAggregateManager aggregateManager,
AccumulatorRegistry accumulatorRegistry,
TaskKvStateRegistry kvStateRegistry,
InputSplitProvider splitProvider,
Map<String,Future<Path>> distCacheEntries,
ResultPartitionWriter[] writers,
InputGate[] inputGates,
TaskEventDispatcher taskEventDispatcher,
CheckpointResponder checkpointResponder,
TaskManagerRuntimeInfo taskManagerInfo,
TaskMetricGroup metrics,
Task containingTask) |
Modifier and Type | Method and Description |
---|---|
default CompletableFuture<OperatorBackPressureStatsResponse> |
RestfulGateway.requestOperatorBackPressureStats(JobID jobId,
JobVertexID jobVertexId)
Requests the statistics on operator back pressure.
|
Modifier and Type | Method and Description |
---|---|
JobVertexID |
SavepointEnvironment.getJobVertexId() |
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.