Modifier and Type | Method and Description |
---|---|
static JobID |
JobID.fromByteArray(byte[] bytes)
Creates a new JobID from the given byte sequence.
|
static JobID |
JobID.fromByteBuffer(ByteBuffer buf) |
static JobID |
JobID.fromHexString(String hexString) |
static JobID |
JobID.generate()
Creates a new (statistically) random JobID.
|
JobID |
Plan.getJobId()
Gets the ID of the job that the dataflow plan belongs to.
|
JobID |
JobSubmissionResult.getJobID()
Returns the JobID assigned to the job by the Flink runtime.
|
Modifier and Type | Method and Description |
---|---|
abstract void |
PlanExecutor.endSession(JobID jobID)
Ends the job session, identified by the given JobID.
|
void |
Plan.setJobId(JobID jobId)
Sets the ID of the job that the dataflow plan belongs to.
|
Constructor and Description |
---|
JobExecutionResult(JobID jobID,
long netRuntime,
Map<String,Object> accumulators)
Creates a new JobExecutionResult.
|
JobSubmissionResult(JobID jobID) |
Modifier and Type | Field and Description |
---|---|
protected JobID |
ExecutionEnvironment.jobID
The ID of the session, defined by this execution environment.
|
Modifier and Type | Method and Description |
---|---|
JobID |
ExecutionEnvironment.getId()
Gets the JobID by which this environment is identified.
|
Modifier and Type | Method and Description |
---|---|
JobID |
ExecutionEnvironment.getId()
Gets the UUID by which this environment is identified.
|
Modifier and Type | Method and Description |
---|---|
void |
RemoteExecutor.endSession(JobID jobID) |
void |
LocalExecutor.endSession(JobID jobID) |
Modifier and Type | Method and Description |
---|---|
JobID |
DetachedEnvironment.DetachedJobExecutionResult.getJobID() |
Modifier and Type | Method and Description |
---|---|
void |
ClusterClient.cancel(JobID jobId)
Cancels a job identified by the job id.
|
void |
ClusterClient.endSession(JobID jobId)
Tells the JobManager to finish the session (job) defined by the given ID.
|
Map<String,Object> |
ClusterClient.getAccumulators(JobID jobID)
Requests and returns the accumulators for the given job identifier.
|
Map<String,Object> |
ClusterClient.getAccumulators(JobID jobID,
ClassLoader loader)
Requests and returns the accumulators for the given job identifier.
|
void |
ClusterClient.stop(JobID jobId)
Stops a program on Flink cluster whose job-manager is configured in this client's configuration.
|
Modifier and Type | Method and Description |
---|---|
void |
ClusterClient.endSessions(List<JobID> jobIds)
Tells the JobManager to finish the sessions (jobs) defined by the given IDs.
|
Modifier and Type | Method and Description |
---|---|
JobGraph |
JobGraphGenerator.compileJobGraph(OptimizedPlan program,
JobID jobId) |
Modifier and Type | Field and Description |
---|---|
protected JobID |
AccumulatorRegistry.jobID |
Modifier and Type | Method and Description |
---|---|
JobID |
AccumulatorSnapshot.getJobID() |
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 |
---|---|
void |
BlobClient.delete(JobID jobId,
String key)
Deletes the BLOB identified by the given job ID and key from the BLOB server.
|
void |
BlobClient.deleteAll(JobID jobId)
Deletes all BLOBs belonging to the job with the given ID from the BLOB server
|
InputStream |
BlobClient.get(JobID jobID,
String key)
Downloads the BLOB identified by the given job ID and key from the BLOB server.
|
void |
BlobClient.put(JobID jobId,
String key,
byte[] value)
Uploads the data of the given byte array to the BLOB server and stores it under the given job ID and key.
|
void |
BlobClient.put(JobID jobId,
String key,
byte[] value,
int offset,
int len)
Uploads data from the given byte array to the BLOB server and stores it under the given job ID and key.
|
void |
BlobClient.put(JobID jobId,
String key,
InputStream inputStream)
Uploads data from the given input stream to the BLOB server and stores it under the given job ID and key.
|
Modifier and Type | Method and Description |
---|---|
JobID |
PendingCheckpoint.getJobId() |
JobID |
CompletedCheckpoint.getJobId() |
Modifier and Type | Method and Description |
---|---|
CheckpointIDCounter |
ZooKeeperCheckpointRecoveryFactory.createCheckpointIDCounter(JobID jobID) |
CheckpointIDCounter |
StandaloneCheckpointRecoveryFactory.createCheckpointIDCounter(JobID ignored) |
CheckpointIDCounter |
CheckpointRecoveryFactory.createCheckpointIDCounter(JobID jobId)
Creates a
CheckpointIDCounter instance for a job. |
CompletedCheckpointStore |
ZooKeeperCheckpointRecoveryFactory.createCheckpointStore(JobID jobId,
ClassLoader userClassLoader) |
CompletedCheckpointStore |
StandaloneCheckpointRecoveryFactory.createCheckpointStore(JobID jobId,
ClassLoader userClassLoader) |
CompletedCheckpointStore |
CheckpointRecoveryFactory.createCheckpointStore(JobID jobId,
ClassLoader userClassLoader)
Creates a
CompletedCheckpointStore instance for a job. |
Constructor and Description |
---|
CheckpointCoordinator(JobID job,
long baseInterval,
long checkpointTimeout,
int numberKeyGroups,
ExecutionVertex[] tasksToTrigger,
ExecutionVertex[] tasksToWaitFor,
ExecutionVertex[] tasksToCommitTo,
ClassLoader userClassLoader,
CheckpointIDCounter checkpointIDCounter,
CompletedCheckpointStore completedCheckpointStore,
RecoveryMode recoveryMode,
Executor executor) |
CheckpointCoordinator(JobID job,
long baseInterval,
long checkpointTimeout,
long minPauseBetweenCheckpoints,
int maxConcurrentCheckpointAttempts,
int numberKeyGroups,
ExecutionVertex[] tasksToTrigger,
ExecutionVertex[] tasksToWaitFor,
ExecutionVertex[] tasksToCommitTo,
ClassLoader userClassLoader,
CheckpointIDCounter checkpointIDCounter,
CompletedCheckpointStore completedCheckpointStore,
RecoveryMode recoveryMode,
CheckpointStatsTracker statsTracker,
Executor executor) |
CompletedCheckpoint(JobID job,
long checkpointID,
long timestamp,
long completionTimestamp,
Map<JobVertexID,TaskState> taskStates) |
PendingCheckpoint(JobID jobId,
long checkpointId,
long checkpointTimestamp,
Map<ExecutionAttemptID,ExecutionVertex> verticesToConfirm,
Executor executor) |
Constructor and Description |
---|
SavepointCoordinator(JobID jobId,
long baseInterval,
long checkpointTimeout,
int numberKeyGroups,
ExecutionVertex[] tasksToTrigger,
ExecutionVertex[] tasksToWaitFor,
ExecutionVertex[] tasksToCommitTo,
ClassLoader userClassLoader,
CheckpointIDCounter checkpointIDCounter,
SavepointStore savepointStore,
CheckpointStatsTracker statsTracker,
Executor executor) |
Modifier and Type | Method and Description |
---|---|
JobID |
SerializedJobExecutionResult.getJobId() |
JobID |
JobStatusMessage.getJobId() |
JobID |
JobExecutionException.getJobID() |
Constructor and Description |
---|
JobCancellationException(JobID jobID,
String msg,
Throwable cause) |
JobExecutionException(JobID jobID,
String msg) |
JobExecutionException(JobID jobID,
String msg,
Throwable cause)
Constructs a new job execution exception.
|
JobExecutionException(JobID jobID,
Throwable cause) |
JobStatusMessage(JobID jobId,
String jobName,
JobStatus jobState,
long startTime) |
JobSubmissionException(JobID jobID,
String msg) |
JobSubmissionException(JobID jobID,
String msg,
Throwable cause) |
JobTimeoutException(JobID jobID,
String msg,
Throwable cause) |
SerializedJobExecutionResult(JobID jobID,
long netRuntime,
Map<String,SerializedValue<Object>> accumulators)
Creates a new SerializedJobExecutionResult.
|
Modifier and Type | Method and Description |
---|---|
JobID |
ShutdownClusterAfterJob.jobId() |
Constructor and Description |
---|
ShutdownClusterAfterJob(JobID jobId) |
Modifier and Type | Method and Description |
---|---|
JobID |
Environment.getJobID()
Returns the ID of the job that the task belongs to.
|
Modifier and Type | Method and Description |
---|---|
ClassLoader |
LibraryCacheManager.getClassLoader(JobID id)
Returns the user code class loader associated with id.
|
ClassLoader |
FallbackLibraryCacheManager.getClassLoader(JobID id) |
ClassLoader |
BlobLibraryCacheManager.getClassLoader(JobID id) |
int |
BlobLibraryCacheManager.getNumberOfReferenceHolders(JobID jobId) |
void |
LibraryCacheManager.registerJob(JobID id,
Collection<BlobKey> requiredJarFiles,
Collection<URL> requiredClasspaths)
Registers a job with its required jar files and classpaths.
|
void |
FallbackLibraryCacheManager.registerJob(JobID id,
Collection<BlobKey> requiredJarFiles,
Collection<URL> requiredClasspaths) |
void |
BlobLibraryCacheManager.registerJob(JobID id,
Collection<BlobKey> requiredJarFiles,
Collection<URL> requiredClasspaths) |
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.unregisterJob(JobID id)
Unregisters a job from the library cache manager.
|
void |
FallbackLibraryCacheManager.unregisterJob(JobID id) |
void |
BlobLibraryCacheManager.unregisterJob(JobID id) |
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 |
---|---|
JobID |
JobInformation.getJobId() |
JobID |
ExecutionVertex.getJobId() |
JobID |
ExecutionJobVertex.getJobId() |
JobID |
ExecutionGraph.getJobID() |
Constructor and Description |
---|
ExecutionGraph(Executor futureExecutor,
Executor ioExecutor,
JobID jobId,
String jobName,
Configuration jobConfig,
SerializedValue<ExecutionConfig> serializedConfig,
scala.concurrent.duration.FiniteDuration timeout,
RestartStrategy restartStrategy,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
Scheduler scheduler,
ClassLoader userClassLoader,
MetricGroup metricGroup) |
JobInformation(JobID jobId,
String jobName,
SerializedValue<ExecutionConfig> serializedExecutionConfig,
Configuration jobConfiguration,
Collection<BlobKey> requiredJarFileBlobKeys,
Collection<URL> requiredClasspathURLs) |
Modifier and Type | Method and Description |
---|---|
Future<Path> |
FileCache.createTmpFile(String name,
DistributedCache.DistributedCacheEntry entry,
JobID jobID)
If the file doesn't exists locally, it will copy the file to the temp directory.
|
void |
FileCache.deleteTmpFile(String name,
JobID jobID)
Deletes the local file after a 5 second delay.
|
Modifier and Type | Method and Description |
---|---|
JobID |
Slot.getJobID()
Returns the ID of the job this allocated slot belongs to.
|
Modifier and Type | Method and Description |
---|---|
SharedSlot |
Instance.allocateSharedSlot(JobID jobID,
SlotSharingGroupAssignment sharingGroupAssignment)
Allocates a shared slot on this TaskManager instance.
|
SimpleSlot |
Instance.allocateSimpleSlot(JobID jobID)
Allocates a simple slot on this TaskManager instance.
|
Constructor and Description |
---|
SharedSlot(JobID jobID,
Instance instance,
int slotNumber,
SlotSharingGroupAssignment assignmentGroup)
Creates a new shared slot that has no parent (is a root slot) and does not belong to any task group.
|
SharedSlot(JobID jobID,
Instance instance,
int slotNumber,
SlotSharingGroupAssignment assignmentGroup,
SharedSlot parent,
AbstractID groupId)
Creates a new shared slot that has is a sub-slot of the given parent shared slot, and that belongs
to the given task group.
|
SimpleSlot(JobID jobID,
Instance instance,
int slotNumber)
Creates a new simple slot that stands alone and does not belong to shared slot.
|
SimpleSlot(JobID jobID,
Instance instance,
int slotNumber,
SharedSlot parent,
AbstractID groupID)
Creates a new simple slot that belongs to the given shared slot and
is identified by the given ID..
|
Slot(JobID jobID,
Instance instance,
int slotNumber,
SharedSlot parent,
AbstractID groupID)
Base constructor for slots.
|
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 | Method and Description |
---|---|
JobID |
ResultPartition.getJobId() |
Modifier and Type | Method and Description |
---|---|
void |
ResultPartitionConsumableNotifier.notifyPartitionConsumable(JobID jobId,
ResultPartitionID partitionId) |
Constructor and Description |
---|
ResultPartition(String owningTaskName,
JobID jobId,
ResultPartitionID partitionId,
ResultPartitionType partitionType,
int numberOfSubpartitions,
ResultPartitionManager partitionManager,
ResultPartitionConsumableNotifier partitionConsumableNotifier,
IOManager ioManager,
boolean sendScheduleOrUpdateConsumersMessage) |
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 |
---|---|
JobID |
JobGraph.getJobID()
Returns the ID of the job.
|
Constructor and Description |
---|
JobGraph(JobID jobId,
String jobName)
Constructs a new job graph with the given job ID (or a random ID, if
null is passed),
the given name and the given execution configuration (see ExecutionConfig ). |
JobGraph(JobID jobId,
String jobName,
JobVertex... vertices)
Constructs a new job graph with the given name, the given
ExecutionConfig ,
the given jobId or a random one if null supplied, and the given job vertices. |
Modifier and Type | Method and Description |
---|---|
JobID |
SubmittedJobGraph.getJobId()
|
static JobID |
ZooKeeperSubmittedJobGraphStore.jobIdfromPath(String path)
Returns the JobID from the given path in ZooKeeper.
|
Modifier and Type | Method and Description |
---|---|
protected scala.collection.mutable.HashMap<JobID,scala.Tuple2<ExecutionGraph,JobInfo>> |
JobManager.currentJobs()
Either running or not yet archived jobs (session hasn't been ended).
|
Collection<JobID> |
ZooKeeperSubmittedJobGraphStore.getJobIds() |
Collection<JobID> |
SubmittedJobGraphStore.getJobIds()
Get all job ids of submitted job graphs to the submitted job graph store.
|
Collection<JobID> |
StandaloneSubmittedJobGraphStore.getJobIds() |
protected scala.collection.mutable.LinkedHashMap<JobID,ExecutionGraph> |
MemoryArchivist.graphs() |
Modifier and Type | Method and Description |
---|---|
static String |
ZooKeeperSubmittedJobGraphStore.getPathForJob(JobID jobId)
Returns the JobID as a String (with leading slash).
|
void |
SubmittedJobGraphStore.SubmittedJobGraphListener.onAddedJobGraph(JobID jobId)
Callback for
SubmittedJobGraph instances added by a different SubmittedJobGraphStore instance. |
void |
JobManager.onAddedJobGraph(JobID jobId) |
void |
SubmittedJobGraphStore.SubmittedJobGraphListener.onRemovedJobGraph(JobID jobId)
Callback for
SubmittedJobGraph instances removed by a different SubmittedJobGraphStore instance. |
void |
JobManager.onRemovedJobGraph(JobID jobId) |
SubmittedJobGraph |
ZooKeeperSubmittedJobGraphStore.recoverJobGraph(JobID jobId) |
SubmittedJobGraph |
SubmittedJobGraphStore.recoverJobGraph(JobID jobId)
Returns the
SubmittedJobGraph with the given JobID . |
SubmittedJobGraph |
StandaloneSubmittedJobGraphStore.recoverJobGraph(JobID jobId) |
void |
ZooKeeperSubmittedJobGraphStore.removeJobGraph(JobID jobId) |
void |
SubmittedJobGraphStore.removeJobGraph(JobID jobId)
Removes the
SubmittedJobGraph with the given JobID if it exists. |
void |
StandaloneSubmittedJobGraphStore.removeJobGraph(JobID jobId) |
Modifier and Type | Method and Description |
---|---|
JobID |
JobManagerMessages.RecoverJob.jobId() |
JobID |
JobManagerMessages.RequestPartitionProducerState.jobId() |
JobID |
JobManagerMessages.ScheduleOrUpdateConsumers.jobId() |
JobID |
JobManagerMessages.JobSubmitSuccess.jobId() |
JobID |
JobManagerMessages.TriggerSavepoint.jobId() |
JobID |
JobManagerMessages.TriggerSavepointSuccess.jobId() |
JobID |
JobManagerMessages.TriggerSavepointFailure.jobId() |
JobID |
ExecutionGraphMessages.ExecutionStateChanged.jobID() |
JobID |
ExecutionGraphMessages.JobStatusChanged.jobID() |
JobID |
ArchiveMessages.ArchiveExecutionGraph.jobID() |
JobID |
ArchiveMessages.RequestArchivedJob.jobID() |
JobID |
JobManagerMessages.CancelJob.jobID() |
JobID |
JobManagerMessages.StopJob.jobID() |
JobID |
JobManagerMessages.RequestNextInputSplit.jobID() |
JobID |
JobManagerMessages.RequestJobStatus.jobID() |
JobID |
JobManagerMessages.CurrentJobStatus.jobID() |
JobID |
JobManagerMessages.CancellationSuccess.jobID() |
JobID |
JobManagerMessages.CancellationFailure.jobID() |
JobID |
JobManagerMessages.StoppingSuccess.jobID() |
JobID |
JobManagerMessages.StoppingFailure.jobID() |
JobID |
JobManagerMessages.RequestJob.jobID() |
JobID |
JobManagerMessages.JobFound.jobID() |
JobID |
JobManagerMessages.JobNotFound.jobID() |
JobID |
JobManagerMessages.RemoveJob.jobID() |
JobID |
JobManagerMessages.RemoveCachedJob.jobID() |
JobID |
JobManagerMessages.JobStatusResponse.jobID() |
JobID |
JobManagerMessages.CancellationResponse.jobID() |
JobID |
JobManagerMessages.StoppingResponse.jobID() |
JobID |
JobManagerMessages.JobResponse.jobID() |
Modifier and Type | Method and Description |
---|---|
Object |
JobManagerMessages$.getRequestJobStatus(JobID jobId) |
static Object |
JobManagerMessages.getRequestJobStatus(JobID jobId) |
Modifier and Type | Method and Description |
---|---|
JobID |
RequestAccumulatorResultsStringified.jobID() |
JobID |
AccumulatorMessage.jobID()
ID of the job that the accumulator belongs to
|
JobID |
RequestAccumulatorResults.jobID() |
JobID |
AccumulatorResultsFound.jobID() |
JobID |
AccumulatorResultsErroneous.jobID() |
JobID |
AccumulatorResultsNotFound.jobID() |
JobID |
AccumulatorResultStringsFound.jobID() |
Constructor and Description |
---|
AccumulatorResultsErroneous(JobID jobID,
Exception cause) |
AccumulatorResultsFound(JobID jobID,
Map<String,SerializedValue<Object>> result) |
AccumulatorResultsNotFound(JobID jobID) |
AccumulatorResultStringsFound(JobID jobID,
StringifiedAccumulatorResult[] result) |
RequestAccumulatorResults(JobID jobID) |
RequestAccumulatorResultsStringified(JobID jobID) |
Modifier and Type | Method and Description |
---|---|
JobID |
AbstractCheckpointMessage.getJob() |
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 |
---|---|
JobID |
JobDetails.getJobId() |
Modifier and Type | Method and Description |
---|---|
List<JobID> |
JobsWithIDsOverview.getJobsCancelled() |
List<JobID> |
JobsWithIDsOverview.getJobsFailed() |
List<JobID> |
JobsWithIDsOverview.getJobsFinished() |
List<JobID> |
JobsWithIDsOverview.getJobsRunningOrPending() |
Constructor and Description |
---|
JobDetails(JobID jobId,
String jobName,
long startTime,
long endTime,
JobStatus status,
long lastUpdateTime,
int[] numVerticesPerExecutionState,
int numTasks) |
Constructor and Description |
---|
JobsWithIDsOverview(List<JobID> jobsRunningOrPending,
List<JobID> jobsFinished,
List<JobID> jobsCancelled,
List<JobID> jobsFailed) |
JobsWithIDsOverview(List<JobID> jobsRunningOrPending,
List<JobID> jobsFinished,
List<JobID> jobsCancelled,
List<JobID> jobsFailed) |
JobsWithIDsOverview(List<JobID> jobsRunningOrPending,
List<JobID> jobsFinished,
List<JobID> jobsCancelled,
List<JobID> jobsFailed) |
JobsWithIDsOverview(List<JobID> jobsRunningOrPending,
List<JobID> jobsFinished,
List<JobID> jobsCancelled,
List<JobID> jobsFailed) |
Modifier and Type | Field and Description |
---|---|
protected JobID |
JobMetricGroup.jobId
The ID of the job represented by this metrics group
|
Modifier and Type | Method and Description |
---|---|
JobID |
JobMetricGroup.jobId() |
Modifier and Type | Method and Description |
---|---|
TaskMetricGroup |
TaskManagerMetricGroup.addTaskForJob(JobID jobId,
String jobName,
JobVertexID jobVertexId,
ExecutionAttemptID executionAttemptId,
String taskName,
int subtaskIndex,
int attemptNumber) |
void |
JobManagerMetricGroup.removeJob(JobID jobId) |
void |
TaskManagerMetricGroup.removeJobMetricsGroup(JobID jobId,
TaskManagerJobMetricGroup group) |
Constructor and Description |
---|
JobManagerJobMetricGroup(MetricRegistry registry,
JobManagerMetricGroup parent,
JobID jobId,
String jobName) |
JobManagerJobMetricGroup(MetricRegistry registry,
JobManagerMetricGroup parent,
JobManagerJobScopeFormat scopeFormat,
JobID jobId,
String jobName) |
JobMetricGroup(MetricRegistry registry,
JobID jobId,
String jobName,
String[] scope) |
TaskManagerJobMetricGroup(MetricRegistry registry,
TaskManagerMetricGroup parent,
JobID jobId,
String jobName) |
TaskManagerJobMetricGroup(MetricRegistry registry,
TaskManagerMetricGroup parent,
TaskManagerJobScopeFormat scopeFormat,
JobID jobId,
String jobName) |
Modifier and Type | Method and Description |
---|---|
String[] |
JobManagerJobScopeFormat.formatScope(JobManagerMetricGroup parent,
JobID jid,
String jobName) |
String[] |
TaskManagerJobScopeFormat.formatScope(TaskManagerMetricGroup parent,
JobID jid,
String jobName) |
Modifier and Type | Method and Description |
---|---|
scala.collection.Iterable<JobID> |
LocalFlinkMiniCluster.currentlyRunningJobs() |
List<JobID> |
LocalFlinkMiniCluster.getCurrentlyRunningJobsJava() |
Modifier and Type | Method and Description |
---|---|
akka.actor.ActorSystem |
FlinkMiniCluster.startJobClientActorSystem(JobID jobID) |
void |
LocalFlinkMiniCluster.stopJob(JobID id) |
Modifier and Type | Method and Description |
---|---|
JobID |
TaskExecutionState.getJobID()
The ID of the job the task belongs to
|
JobID |
Task.getJobID() |
JobID |
RuntimeEnvironment.getJobID() |
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) |
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 |
---|---|
JobID |
TestingTaskManagerMessages.RegisterSubmitTaskListener.jobId() |
JobID |
TestingTaskManagerMessages.AccumulatorsChanged.jobID() |
JobID |
TestingJobManagerMessages.RequestExecutionGraph.jobID() |
JobID |
TestingJobManagerMessages.ExecutionGraphFound.jobID() |
JobID |
TestingJobManagerMessages.ExecutionGraphNotFound.jobID() |
JobID |
TestingJobManagerMessages.WaitForAllVerticesToBeRunning.jobID() |
JobID |
TestingJobManagerMessages.WaitForAllVerticesToBeRunningOrFinished.jobID() |
JobID |
TestingJobManagerMessages.AllVerticesRunning.jobID() |
JobID |
TestingJobManagerMessages.NotifyWhenJobRemoved.jobID() |
JobID |
TestingJobManagerMessages.RequestWorkingTaskManager.jobID() |
JobID |
TestingJobManagerMessages.NotifyWhenJobStatus.jobID() |
JobID |
TestingJobManagerMessages.JobStatusIs.jobID() |
JobID |
TestingJobManagerMessages.NotifyWhenAccumulatorChange.jobID() |
JobID |
TestingJobManagerMessages.UpdatedAccumulators.jobID() |
JobID |
TestingJobManagerMessages.ResponseExecutionGraph.jobID() |
JobID |
TestingMessages.CheckIfJobRemoved.jobID() |
Modifier and Type | Method and Description |
---|---|
scala.collection.mutable.HashMap<JobID,akka.actor.ActorRef> |
TestingTaskManagerLike.registeredSubmitTaskListeners()
Map of registered task submit listeners
|
scala.collection.mutable.HashMap<JobID,scala.Tuple2<Object,scala.collection.immutable.Set<akka.actor.ActorRef>>> |
TestingJobManagerLike.waitForAccumulatorUpdate() |
scala.collection.mutable.HashMap<JobID,scala.collection.immutable.Set<akka.actor.ActorRef>> |
TestingJobManagerLike.waitForAllVerticesToBeRunning() |
scala.collection.mutable.HashMap<JobID,scala.collection.immutable.Set<akka.actor.ActorRef>> |
TestingJobManagerLike.waitForAllVerticesToBeRunningOrFinished() |
scala.collection.mutable.HashMap<JobID,scala.collection.mutable.HashMap<JobStatus,scala.collection.immutable.Set<akka.actor.ActorRef>>> |
TestingJobManagerLike.waitForJobStatus() |
Modifier and Type | Method and Description |
---|---|
boolean |
TestingJobManagerLike.checkIfAllVerticesRunning(JobID jobID) |
boolean |
TestingJobManagerLike.checkIfAllVerticesRunningOrFinished(JobID jobID) |
void |
TestingJobManagerLike.notifyListeners(JobID jobID) |
Constructor and Description |
---|
AccumulatorsChanged(JobID jobID) |
AllVerticesRunning(JobID jobID) |
CheckIfJobRemoved(JobID jobID) |
ExecutionGraphFound(JobID jobID,
ExecutionGraph executionGraph) |
ExecutionGraphNotFound(JobID jobID) |
JobStatusIs(JobID jobID,
JobStatus state) |
NotifyWhenAccumulatorChange(JobID jobID) |
NotifyWhenJobRemoved(JobID jobID) |
NotifyWhenJobStatus(JobID jobID,
JobStatus state) |
RegisterSubmitTaskListener(JobID jobId) |
RequestExecutionGraph(JobID jobID) |
RequestWorkingTaskManager(JobID jobID) |
UpdatedAccumulators(JobID jobID,
Map<ExecutionAttemptID,Map<AccumulatorRegistry.Metric,Accumulator<?,?>>> flinkAccumulators,
Map<String,Accumulator<?,?>> userAccumulators) |
WaitForAllVerticesToBeRunning(JobID jobID) |
WaitForAllVerticesToBeRunningOrFinished(JobID jobID) |
Modifier and Type | Method and Description |
---|---|
static ZooKeeperCheckpointIDCounter |
ZooKeeperUtils.createCheckpointIDCounter(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
JobID jobId)
Creates a
ZooKeeperCheckpointIDCounter instance. |
static CompletedCheckpointStore |
ZooKeeperUtils.createCompletedCheckpoints(org.apache.curator.framework.CuratorFramework client,
Configuration configuration,
JobID jobId,
int maxNumberOfCheckpointsToRetain,
ClassLoader userClassLoader,
Executor executor)
Creates a
ZooKeeperCompletedCheckpointStore instance. |
Modifier and Type | Method and Description |
---|---|
ExecutionGraph |
ExecutionGraphHolder.getExecutionGraph(JobID jid,
ActorGateway jobManager)
Retrieves the execution graph with
JobID jid or null if it cannot be found. |
Modifier and Type | Method and Description |
---|---|
static String |
StreamIterationHead.createBrokerIdString(JobID jid,
String iterationID,
int subtaskIndex)
Creates the identification string with which head and tail task find the shared blocking
queue for the back channel.
|
Modifier and Type | Method and Description |
---|---|
JobID |
YarnJobManager.stopWhenJobFinished() |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.