Package | Description |
---|---|
org.apache.flink.mesos.runtime.clusterframework | |
org.apache.flink.runtime.clusterframework |
This package contains the cluster resource management functionality.
|
org.apache.flink.runtime.clusterframework.messages |
This package contains the actor messages that are sent between the
cluster resource framework and the JobManager, as well as the generic
messages sent between the cluster resource framework and the client.
|
org.apache.flink.runtime.clusterframework.standalone | |
org.apache.flink.runtime.clusterframework.types | |
org.apache.flink.runtime.instance | |
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 | |
org.apache.flink.runtime.metrics.dump | |
org.apache.flink.runtime.minicluster | |
org.apache.flink.runtime.taskmanager | |
org.apache.flink.yarn |
Modifier and Type | Method and Description |
---|---|
ResourceID |
RegisteredMesosWorkerNode.getResourceID() |
Modifier and Type | Method and Description |
---|---|
protected void |
MesosFlinkResourceManager.releasePendingWorker(ResourceID id)
Release the given pending worker.
|
protected RegisteredMesosWorkerNode |
MesosFlinkResourceManager.workerStarted(ResourceID resourceID)
Accept the given started worker into the internal state.
|
Modifier and Type | Method and Description |
---|---|
protected Collection<RegisteredMesosWorkerNode> |
MesosFlinkResourceManager.reacceptRegisteredWorkers(Collection<ResourceID> toConsolidate)
Accept the given registered workers into the internal state.
|
Constructor and Description |
---|
MesosTaskManager(TaskManagerConfiguration config,
ResourceID resourceID,
TaskManagerLocation taskManagerLocation,
MemoryManager memoryManager,
IOManager ioManager,
NetworkEnvironment network,
int numberOfSlots,
LeaderRetrievalService leaderRetrievalService,
MetricRegistry metricRegistry) |
Modifier and Type | Method and Description |
---|---|
boolean |
FlinkResourceManager.isStarted(ResourceID resourceId)
Gets the started worker for a given resource ID, if one is available.
|
void |
FlinkResourceManager.notifyWorkerFailed(ResourceID resourceID,
String message)
This method should be called by the framework once it detects that a currently registered
worker has failed.
|
protected abstract void |
FlinkResourceManager.releasePendingWorker(ResourceID resourceID)
Trigger a release of a pending worker.
|
protected abstract WorkerType |
FlinkResourceManager.workerStarted(ResourceID resourceID)
Callback when a worker was started.
|
Modifier and Type | Method and Description |
---|---|
protected abstract Collection<WorkerType> |
FlinkResourceManager.reacceptRegisteredWorkers(Collection<ResourceID> registered)
This method is called when the resource manager starts after a failure and reconnects to
the leader JobManager, who still has some workers registered.
|
Modifier and Type | Method and Description |
---|---|
ResourceID |
NotifyResourceStarted.getResourceID() |
ResourceID |
ResourceRemoved.resourceId()
Gets the ID under which the resource is registered (for example container ID).
|
ResourceID |
RemoveResource.resourceId()
Gets the ID under which the resource is registered (for example container ID).
|
Modifier and Type | Method and Description |
---|---|
Collection<ResourceID> |
RegisterResourceManagerSuccessful.currentlyRegisteredTaskManagers() |
Constructor and Description |
---|
NotifyResourceStarted(ResourceID resourceID) |
RemoveResource(ResourceID resourceId)
Constructor for a shutdown of a registered resource.
|
ResourceRemoved(ResourceID resourceId,
String message)
Constructor for a shutdown of a registered resource.
|
Constructor and Description |
---|
RegisterResourceManagerSuccessful(akka.actor.ActorRef jobManager,
Collection<ResourceID> currentlyRegisteredTaskManagers)
Creates a new message with a list of existing known TaskManagers.
|
Modifier and Type | Method and Description |
---|---|
protected ResourceID |
StandaloneResourceManager.workerStarted(ResourceID resourceID) |
Modifier and Type | Method and Description |
---|---|
protected Collection<ResourceID> |
StandaloneResourceManager.reacceptRegisteredWorkers(Collection<ResourceID> toConsolidate) |
Modifier and Type | Method and Description |
---|---|
protected void |
StandaloneResourceManager.releasePendingWorker(ResourceID resourceID) |
protected void |
StandaloneResourceManager.releaseStartedWorker(ResourceID resourceID) |
protected ResourceID |
StandaloneResourceManager.workerStarted(ResourceID resourceID) |
Modifier and Type | Method and Description |
---|---|
protected Collection<ResourceID> |
StandaloneResourceManager.reacceptRegisteredWorkers(Collection<ResourceID> toConsolidate) |
Modifier and Type | Method and Description |
---|---|
static ResourceID |
ResourceID.generate()
Generate a random resource id.
|
ResourceID |
ResourceIDRetrievable.getResourceID() |
ResourceID |
ResourceID.getResourceID()
A ResourceID can always retrieve a ResourceID.
|
Modifier and Type | Method and Description |
---|---|
ResourceID |
Slot.getTaskManagerID()
Gets the ID of the TaskManager that offers this slot.
|
ResourceID |
Instance.getTaskManagerID() |
Modifier and Type | Method and Description |
---|---|
Instance |
InstanceManager.getRegisteredInstance(ResourceID ref) |
boolean |
InstanceManager.isRegistered(ResourceID resourceId) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
RegistrationMessages.RegisterTaskManager.resourceId() |
Constructor and Description |
---|
RegisterTaskManager(ResourceID resourceId,
TaskManagerLocation connectionInfo,
HardwareDescription resources,
int numberOfSlots) |
Modifier and Type | Method and Description |
---|---|
void |
MetricRegistry.startQueryService(akka.actor.ActorSystem actorSystem,
ResourceID resourceID)
Initializes the MetricQueryService.
|
Modifier and Type | Method and Description |
---|---|
static akka.actor.ActorRef |
MetricQueryService.startMetricQueryService(akka.actor.ActorSystem actorSystem,
ResourceID resourceID)
Starts the MetricQueryService actor in the given actor system.
|
Modifier and Type | Method and Description |
---|---|
akka.actor.Props |
LocalFlinkMiniCluster.getTaskManagerProps(Class<? extends TaskManager> taskManagerClass,
TaskManagerConfiguration taskManagerConfig,
ResourceID resourceID,
TaskManagerLocation taskManagerLocation,
MemoryManager memoryManager,
IOManager ioManager,
NetworkEnvironment networkEnvironment,
LeaderRetrievalService leaderRetrievalService,
MetricRegistry metricsRegistry) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
TaskManagerLocation.getResourceID()
Gets the ID of the resource in which the TaskManager is started.
|
protected ResourceID |
TaskManager.resourceID() |
Modifier and Type | Method and Description |
---|---|
scala.Tuple7<TaskManagerConfiguration,TaskManagerLocation,MemoryManager,IOManager,NetworkEnvironment,LeaderRetrievalService,MetricRegistry> |
TaskManager$.createTaskManagerComponents(Configuration configuration,
ResourceID resourceID,
String taskManagerHostname,
boolean localTaskManagerCommunication,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption) |
static scala.Tuple7<TaskManagerConfiguration,TaskManagerLocation,MemoryManager,IOManager,NetworkEnvironment,LeaderRetrievalService,MetricRegistry> |
TaskManager.createTaskManagerComponents(Configuration configuration,
ResourceID resourceID,
String taskManagerHostname,
boolean localTaskManagerCommunication,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption) |
akka.actor.Props |
TaskManager$.getTaskManagerProps(Class<? extends TaskManager> taskManagerClass,
TaskManagerConfiguration taskManagerConfig,
ResourceID resourceID,
TaskManagerLocation taskManagerLocation,
MemoryManager memoryManager,
IOManager ioManager,
NetworkEnvironment networkEnvironment,
LeaderRetrievalService leaderRetrievalService,
MetricRegistry metricsRegistry) |
static akka.actor.Props |
TaskManager.getTaskManagerProps(Class<? extends TaskManager> taskManagerClass,
TaskManagerConfiguration taskManagerConfig,
ResourceID resourceID,
TaskManagerLocation taskManagerLocation,
MemoryManager memoryManager,
IOManager ioManager,
NetworkEnvironment networkEnvironment,
LeaderRetrievalService leaderRetrievalService,
MetricRegistry metricsRegistry) |
void |
TaskManager$.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration)
Starts and runs the TaskManager.
|
static void |
TaskManager.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration)
Starts and runs the TaskManager.
|
void |
TaskManager$.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
static void |
TaskManager.runTaskManager(String taskManagerHostname,
ResourceID resourceID,
int actorSystemPort,
Configuration configuration,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
void |
TaskManager$.selectNetworkInterfaceAndRunTaskManager(Configuration configuration,
ResourceID resourceID,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
static void |
TaskManager.selectNetworkInterfaceAndRunTaskManager(Configuration configuration,
ResourceID resourceID,
Class<? extends TaskManager> taskManagerClass)
Starts and runs the TaskManager.
|
akka.actor.ActorRef |
TaskManager$.startTaskManagerComponentsAndActor(Configuration configuration,
ResourceID resourceID,
akka.actor.ActorSystem actorSystem,
String taskManagerHostname,
scala.Option<String> taskManagerActorName,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption,
boolean localTaskManagerCommunication,
Class<? extends TaskManager> taskManagerClass)
Starts the task manager actor.
|
static akka.actor.ActorRef |
TaskManager.startTaskManagerComponentsAndActor(Configuration configuration,
ResourceID resourceID,
akka.actor.ActorSystem actorSystem,
String taskManagerHostname,
scala.Option<String> taskManagerActorName,
scala.Option<LeaderRetrievalService> leaderRetrievalServiceOption,
boolean localTaskManagerCommunication,
Class<? extends TaskManager> taskManagerClass)
Starts the task manager actor.
|
Constructor and Description |
---|
TaskManager(TaskManagerConfiguration config,
ResourceID resourceID,
TaskManagerLocation location,
MemoryManager memoryManager,
IOManager ioManager,
NetworkEnvironment network,
int numberOfSlots,
LeaderRetrievalService leaderRetrievalService,
MetricRegistry metricsRegistry) |
TaskManagerLocation(ResourceID resourceID,
InetAddress inetAddress,
int dataPort)
Constructs a new instance connection info object.
|
Modifier and Type | Method and Description |
---|---|
ResourceID |
YarnContainerInLaunch.getResourceID() |
ResourceID |
RegisteredYarnWorkerNode.getResourceID() |
Modifier and Type | Method and Description |
---|---|
protected void |
YarnFlinkResourceManager.releasePendingWorker(ResourceID id) |
protected RegisteredYarnWorkerNode |
YarnFlinkResourceManager.workerStarted(ResourceID resourceID) |
Modifier and Type | Method and Description |
---|---|
protected Collection<RegisteredYarnWorkerNode> |
YarnFlinkResourceManager.reacceptRegisteredWorkers(Collection<ResourceID> toConsolidate) |
Constructor and Description |
---|
YarnTaskManager(TaskManagerConfiguration config,
ResourceID resourceID,
TaskManagerLocation taskManagerLocation,
MemoryManager memoryManager,
IOManager ioManager,
NetworkEnvironment network,
int numberOfSlots,
LeaderRetrievalService leaderRetrievalService,
MetricRegistry metricRegistry) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.