Modifier and Type | Method and Description |
---|---|
protected ActorGateway |
CliFrontend.getJobManagerGateway(CommandLineOptions options)
Retrieves the
ActorGateway for the JobManager. |
Modifier and Type | Method and Description |
---|---|
ActorGateway |
ClusterClient.getJobManagerGateway()
Returns the
ActorGateway of the current job manager leader using
the LeaderRetrievalService . |
Modifier and Type | Method and Description |
---|---|
static List<BlobKey> |
BlobClient.uploadJarFiles(ActorGateway jobManager,
scala.concurrent.duration.FiniteDuration askTimeout,
List<Path> jars)
Retrieves the
BlobServer address from the JobManager and uploads
the JAR files to it. |
Modifier and Type | Method and Description |
---|---|
ActorGateway |
CheckpointCoordinator.createActivatorDeactivator(akka.actor.ActorSystem actorSystem,
UUID leaderSessionID) |
protected ActorGateway |
CheckpointCoordinator.getJobStatusListener() |
Modifier and Type | Method and Description |
---|---|
protected void |
CheckpointCoordinator.setJobStatusListener(ActorGateway jobStatusListener) |
Modifier and Type | Method and Description |
---|---|
ActorGateway |
SavepointCoordinator.createActivatorDeactivator(akka.actor.ActorSystem actorSystem,
UUID leaderSessionID) |
Modifier and Type | Method and Description |
---|---|
static void |
JobClient.submitJobDetached(ActorGateway jobManagerGateway,
JobGraph jobGraph,
scala.concurrent.duration.FiniteDuration timeout,
ClassLoader classLoader)
Submits a job in detached mode.
|
Modifier and Type | Method and Description |
---|---|
void |
ExecutionGraph.registerExecutionListener(ActorGateway listener) |
void |
ExecutionGraph.registerJobStatusListener(ActorGateway listener) |
boolean |
ExecutionVertex.sendMessageToCurrentExecution(Serializable message,
ExecutionAttemptID attemptID,
ActorGateway sender) |
Modifier and Type | Class and Description |
---|---|
class |
AkkaActorGateway
Concrete
ActorGateway implementation which uses Akka to communicate with remote actors. |
Modifier and Type | Method and Description |
---|---|
ActorGateway |
Instance.getActorGateway()
Returns the InstanceGateway of this Instance.
|
Modifier and Type | Method and Description |
---|---|
void |
AkkaActorGateway.forward(Object message,
ActorGateway sender)
Forwards a message.
|
void |
ActorGateway.forward(Object message,
ActorGateway sender)
Forwards a message.
|
void |
AkkaActorGateway.tell(Object message,
ActorGateway sender)
Sends a message asynchronously without a result with sender being the sender.
|
void |
ActorGateway.tell(Object message,
ActorGateway sender)
Sends a message asynchronously without a result with sender being the sender.
|
Constructor and Description |
---|
Instance(ActorGateway actorGateway,
InstanceConnectionInfo connectionInfo,
ResourceID resourceId,
InstanceID id,
HardwareDescription resources,
int numberOfSlots)
Constructs an instance reflecting a registered TaskManager.
|
Modifier and Type | Method and Description |
---|---|
void |
NetworkEnvironment.associateWithTaskManagerAndJobManager(ActorGateway jobManagerGateway,
ActorGateway taskManagerGateway)
This associates the network environment with a TaskManager and JobManager.
|
Modifier and Type | Method and Description |
---|---|
void |
JobGraph.uploadUserJars(ActorGateway jobManager,
scala.concurrent.duration.FiniteDuration askTimeout)
Uploads the previously added user JAR files to the job manager through
the job manager's BLOB server.
|
Modifier and Type | Method and Description |
---|---|
ActorGateway |
FlinkMiniCluster.getLeaderGateway(scala.concurrent.duration.FiniteDuration timeout) |
Modifier and Type | Method and Description |
---|---|
scala.concurrent.Future<ActorGateway> |
FlinkMiniCluster.getLeaderGatewayFuture() |
scala.concurrent.Promise<ActorGateway> |
FlinkMiniCluster.leaderGateway()
Future to the
ActorGateway of the current leader |
Modifier and Type | Method and Description |
---|---|
void |
Task.registerExecutionListener(ActorGateway listener) |
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) |
Task(JobInformation jobInformation,
TaskInformation taskInformation,
ExecutionAttemptID executionAttemptID,
int subtaskIndex,
int attemptNumber,
Collection<ResultPartitionDeploymentDescriptor> resultPartitionDeploymentDescriptors,
Collection<InputGateDeploymentDescriptor> inputGateDeploymentDescriptors,
int targetSlotNumber,
SerializedValue<StateHandle<?>> operatorState,
MemoryManager memManager,
IOManager ioManager,
NetworkEnvironment networkEnvironment,
BroadcastVariableManager bcVarManager,
ActorGateway taskManagerActor,
ActorGateway jobManagerActor,
scala.concurrent.duration.FiniteDuration actorAskTimeout,
LibraryCacheManager libraryCache,
FileCache fileCache,
TaskManagerRuntimeInfo taskManagerConfig,
TaskMetricGroup metricGroup)
IMPORTANT: This constructor may not start any work that would need to
be undone in the case of a failing task deployment.
|
TaskInputSplitProvider(ActorGateway jobManager,
JobID jobId,
JobVertexID vertexId,
ExecutionAttemptID executionID,
ClassLoader userCodeClassLoader,
scala.concurrent.duration.FiniteDuration timeout) |
Modifier and Type | Method and Description |
---|---|
scala.Option<ActorGateway> |
TestingJobManagerMessages.WorkingTaskManager.gatewayOption() |
Constructor and Description |
---|
WorkingTaskManager(scala.Option<ActorGateway> gatewayOption) |
Modifier and Type | Method and Description |
---|---|
static ActorGateway |
LeaderRetrievalUtils.retrieveLeaderGateway(LeaderRetrievalService leaderRetrievalService,
akka.actor.ActorSystem actorSystem,
scala.concurrent.duration.FiniteDuration timeout)
Retrieves the current leader gateway using the given
LeaderRetrievalService . |
Modifier and Type | Method and Description |
---|---|
scala.concurrent.Future<ActorGateway> |
LeaderRetrievalUtils.LeaderGatewayListener.getActorGatewayFuture() |
Modifier and Type | Method and Description |
---|---|
scala.Tuple2<ActorGateway,Integer> |
JobManagerRetriever.awaitJobManagerGatewayAndWebPort()
Awaits the leading job manager gateway and its web monitor port.
|
scala.Option<scala.Tuple2<ActorGateway,Integer>> |
JobManagerRetriever.getJobManagerGatewayAndWebPort()
Returns the currently known leading job manager gateway and its web monitor port.
|
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. |
protected abstract void |
RuntimeMonitorHandlerBase.respondAsLeader(io.netty.channel.ChannelHandlerContext ctx,
io.netty.handler.codec.http.router.Routed routed,
ActorGateway jobManager) |
protected void |
RuntimeMonitorHandler.respondAsLeader(io.netty.channel.ChannelHandlerContext ctx,
io.netty.handler.codec.http.router.Routed routed,
ActorGateway jobManager) |
Modifier and Type | Method and Description |
---|---|
String |
TaskManagersHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
RequestHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager)
Core method that handles the request and generates the response.
|
String |
JobStoppingHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JobManagerConfigHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JobCancellationHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JarUploadHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JarRunHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JarPlanHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JarListHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JarDeleteHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
JarAccessDeniedHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
DashboardConfigHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
CurrentJobsOverviewHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
CurrentJobIdsHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
ClusterOverviewHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
String |
AbstractExecutionGraphRequestHandler.handleRequest(Map<String,String> pathParams,
Map<String,String> queryParams,
ActorGateway jobManager) |
protected void |
TaskManagerLogHandler.respondAsLeader(io.netty.channel.ChannelHandlerContext ctx,
io.netty.handler.codec.http.router.Routed routed,
ActorGateway jobManager)
Response when running with leading JobManager.
|
Modifier and Type | Method and Description |
---|---|
static String |
HandlerRedirectUtils.getRedirectAddress(String localJobManagerAddress,
scala.Tuple2<ActorGateway,Integer> leader) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.