public class YarnJobManager extends JobManager
JobManager
with additional messages
to start/administer/stop the Yarn session.
Constructor and Description |
---|
YarnJobManager(Configuration flinkConfiguration,
Executor futureExecutor,
Executor ioExecutor,
InstanceManager instanceManager,
Scheduler scheduler,
BlobLibraryCacheManager libraryCacheManager,
akka.actor.ActorRef archive,
RestartStrategyFactory restartStrategyFactory,
scala.concurrent.duration.FiniteDuration timeout,
LeaderElectionService leaderElectionService,
SubmittedJobGraphStore submittedJobGraphs,
CheckpointRecoveryFactory checkpointRecoveryFactory,
SavepointStore savepointStore,
scala.concurrent.duration.FiniteDuration jobRecoveryTimeout,
scala.Option<MetricRegistry> metricsRegistry) |
Modifier and Type | Method and Description |
---|---|
scala.concurrent.duration.FiniteDuration |
DEFAULT_YARN_HEARTBEAT_DELAY() |
scala.PartialFunction<Object,scala.runtime.BoxedUnit> |
handleMessage()
Central work method of the JobManager actor.
|
scala.PartialFunction<Object,scala.runtime.BoxedUnit> |
handleYarnMessage() |
JobID |
stopWhenJobFinished() |
scala.concurrent.duration.FiniteDuration |
YARN_HEARTBEAT_DELAY() |
ARCHIVE_NAME, archive, checkpointRecoveryFactory, createJobManagerComponents, currentJobs, currentResourceManager, executionContext, flinkConfiguration, futureExecutor, futuresToComplete, getAddress, getJobManagerActorRef, getJobManagerActorRef, getJobManagerActorRef, getJobManagerActorRefFuture, getJobManagerAkkaURL, getLocalJobManagerAkkaURL, getRemoteJobManagerAkkaURL, getRemoteJobManagerAkkaURL, grantLeadership, handleError, instanceManager, ioExecutor, JOB_MANAGER_NAME, jobManagerMetricGroup, jobRecoveryTimeout, leaderElectionService, leaderSessionID, libraryCacheManager, log, LOG, main, metricsRegistry, onAddedJobGraph, onRemovedJobGraph, parseArgs, postStop, preStart, recoveryMode, restartStrategyFactory, retryOnBindException, revokeLeadership, runJobManager, runJobManager, RUNTIME_FAILURE_RETURN_CODE, savepointStore, scheduler, shutdown, startActorSystemAndJobManagerActors, startJobManagerActors, startJobManagerActors, STARTUP_FAILURE_RETURN_CODE, submittedJobGraphs, timeout, unhandled, webMonitorPort
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
decorateMessage, handleDiscardedMessage, receive
receive
public YarnJobManager(Configuration flinkConfiguration, Executor futureExecutor, Executor ioExecutor, InstanceManager instanceManager, Scheduler scheduler, BlobLibraryCacheManager libraryCacheManager, akka.actor.ActorRef archive, RestartStrategyFactory restartStrategyFactory, scala.concurrent.duration.FiniteDuration timeout, LeaderElectionService leaderElectionService, SubmittedJobGraphStore submittedJobGraphs, CheckpointRecoveryFactory checkpointRecoveryFactory, SavepointStore savepointStore, scala.concurrent.duration.FiniteDuration jobRecoveryTimeout, scala.Option<MetricRegistry> metricsRegistry)
public scala.concurrent.duration.FiniteDuration DEFAULT_YARN_HEARTBEAT_DELAY()
public scala.concurrent.duration.FiniteDuration YARN_HEARTBEAT_DELAY()
public JobID stopWhenJobFinished()
public scala.PartialFunction<Object,scala.runtime.BoxedUnit> handleMessage()
JobManager
handleMessage
in interface FlinkActor
handleMessage
in class JobManager
public scala.PartialFunction<Object,scala.runtime.BoxedUnit> handleYarnMessage()
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.