Modifier and Type | Method and Description |
---|---|
protected ResourceManager<?> |
StandaloneJobClusterEntryPoint.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Modifier and Type | Method and Description |
---|---|
protected ResourceManager<?> |
MesosJobClusterEntrypoint.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
protected ResourceManager<?> |
MesosSessionClusterEntrypoint.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
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.
|
protected RegisteredMesosWorkerNode |
MesosResourceManager.workerStarted(ResourceID resourceID)
Callback when a worker was started.
|
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 |
---|
MesosResourceManager(RpcService rpcService,
String resourceManagerEndpointId,
ResourceID resourceId,
ResourceManagerConfiguration resourceManagerConfiguration,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
SlotManager slotManager,
MetricRegistry metricRegistry,
JobLeaderIdService jobLeaderIdService,
ClusterInformation clusterInformation,
FatalErrorHandler fatalErrorHandler,
Configuration flinkConfig,
MesosServices mesosServices,
MesosConfiguration mesosConfig,
MesosTaskManagerParameters taskManagerParameters,
ContainerSpecification taskManagerContainerSpec,
String webUiUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Collection<Tuple2<ResourceID,String>>> |
AkkaJobManagerGateway.requestTaskManagerMetricQueryServicePaths(Time timeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Optional<Instance>> |
AkkaJobManagerGateway.requestTaskManagerInstance(ResourceID resourceId,
Time timeout) |
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()
Gets the ResourceID of the object.
|
ResourceID |
ResourceID.getResourceID()
A ResourceID can always retrieve a ResourceID.
|
ResourceID |
SlotID.getResourceID() |
Constructor and Description |
---|
SlotID(ResourceID resourceId,
int slotNumber) |
Modifier and Type | Method and Description |
---|---|
static InputChannelDeploymentDescriptor[] |
InputChannelDeploymentDescriptor.fromEdges(ExecutionEdge[] edges,
ResourceID consumerResourceId,
boolean allowLazyDeployment)
Creates an input channel deployment descriptor for each partition.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Collection<Tuple2<ResourceID,String>>> |
Dispatcher.requestTaskManagerMetricQueryServicePaths(Time timeout) |
Modifier and Type | Method and Description |
---|---|
JobManagerRunner |
Dispatcher.JobManagerRunnerFactory.createJobManagerRunner(ResourceID resourceId,
JobGraph jobGraph,
Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
BlobServer blobServer,
JobManagerSharedServices jobManagerServices,
JobManagerJobMetricGroupFactory jobManagerJobMetricGroupFactory,
FatalErrorHandler fatalErrorHandler) |
JobManagerRunner |
Dispatcher.DefaultJobManagerRunnerFactory.createJobManagerRunner(ResourceID resourceId,
JobGraph jobGraph,
Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
BlobServer blobServer,
JobManagerSharedServices jobManagerServices,
JobManagerJobMetricGroupFactory jobManagerJobMetricGroupFactory,
FatalErrorHandler fatalErrorHandler) |
Modifier and Type | Method and Description |
---|---|
protected abstract ResourceManager<?> |
ClusterEntrypoint.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
protected ResourceManager<?> |
StandaloneSessionClusterEntrypoint.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Modifier and Type | Method and Description |
---|---|
<I,O> HeartbeatManager<I,O> |
HeartbeatServices.createHeartbeatManager(ResourceID resourceId,
HeartbeatListener<I,O> heartbeatListener,
ScheduledExecutor scheduledExecutor,
org.slf4j.Logger log)
Creates a heartbeat manager which does not actively send heartbeats.
|
<I,O> HeartbeatManager<I,O> |
HeartbeatServices.createHeartbeatManagerSender(ResourceID resourceId,
HeartbeatListener<I,O> heartbeatListener,
ScheduledExecutor scheduledExecutor,
org.slf4j.Logger log)
Creates a heartbeat manager which actively sends heartbeats to monitoring targets.
|
long |
HeartbeatManager.getLastHeartbeatFrom(ResourceID resourceId)
Returns the last received heartbeat from the given target.
|
long |
HeartbeatManagerImpl.getLastHeartbeatFrom(ResourceID resourceId) |
void |
HeartbeatManager.monitorTarget(ResourceID resourceID,
HeartbeatTarget<O> heartbeatTarget)
Start monitoring a
HeartbeatTarget . |
void |
HeartbeatManagerImpl.monitorTarget(ResourceID resourceID,
HeartbeatTarget<O> heartbeatTarget) |
void |
HeartbeatListener.notifyHeartbeatTimeout(ResourceID resourceID)
Callback which is called if a heartbeat for the machine identified by the given resource
ID times out.
|
void |
HeartbeatManagerImpl.receiveHeartbeat(ResourceID heartbeatOrigin,
I heartbeatPayload) |
void |
HeartbeatTarget.receiveHeartbeat(ResourceID heartbeatOrigin,
I heartbeatPayload)
Sends a heartbeat response to the target.
|
void |
HeartbeatListener.reportPayload(ResourceID resourceID,
I payload)
Callback which is called whenever a heartbeat with an associated payload is received.
|
void |
HeartbeatManagerImpl.requestHeartbeat(ResourceID requestOrigin,
I heartbeatPayload) |
void |
HeartbeatTarget.requestHeartbeat(ResourceID requestOrigin,
I heartbeatPayload)
Requests a heartbeat from the target.
|
CompletableFuture<O> |
HeartbeatListener.retrievePayload(ResourceID resourceID)
Retrieves the payload value for the next heartbeat message.
|
void |
HeartbeatManager.unmonitorTarget(ResourceID resourceID)
Stops monitoring the heartbeat target with the associated resource ID.
|
void |
HeartbeatManagerImpl.unmonitorTarget(ResourceID resourceID) |
Constructor and Description |
---|
HeartbeatManagerImpl(long heartbeatTimeoutIntervalMs,
ResourceID ownResourceID,
HeartbeatListener<I,O> heartbeatListener,
Executor executor,
ScheduledExecutor scheduledExecutor,
org.slf4j.Logger log) |
HeartbeatManagerSenderImpl(long heartbeatPeriod,
long heartbeatTimeout,
ResourceID ownResourceID,
HeartbeatListener<I,O> heartbeatListener,
Executor executor,
ScheduledExecutor scheduledExecutor,
org.slf4j.Logger log) |
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 |
---|---|
Instance |
Scheduler.getInstance(ResourceID resourceId) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
JMTMRegistrationSuccess.getResourceID() |
ResourceID |
JobMasterRegistrationSuccess.getResourceManagerResourceId() |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
JobMaster.disconnectTaskManager(ResourceID resourceID,
Exception cause) |
CompletableFuture<Acknowledge> |
JobMasterGateway.disconnectTaskManager(ResourceID resourceID,
Exception cause)
Disconnects the given
TaskExecutor from the
JobMaster . |
void |
JobMaster.failSlot(ResourceID taskManagerId,
AllocationID allocationId,
Exception cause) |
void |
JobMasterGateway.failSlot(ResourceID taskManagerId,
AllocationID allocationId,
Exception cause)
Fails the slot with the given allocation id and cause.
|
void |
JobMaster.heartbeatFromResourceManager(ResourceID resourceID) |
void |
JobMasterGateway.heartbeatFromResourceManager(ResourceID resourceID)
Sends heartbeat request from the resource manager.
|
void |
JobMaster.heartbeatFromTaskManager(ResourceID resourceID,
AccumulatorReport accumulatorReport) |
void |
JobMasterGateway.heartbeatFromTaskManager(ResourceID resourceID,
AccumulatorReport accumulatorReport)
Sends the heartbeat to job manager from task manager.
|
CompletableFuture<Collection<SlotOffer>> |
JobMaster.offerSlots(ResourceID taskManagerId,
Collection<SlotOffer> slots,
Time timeout) |
CompletableFuture<Collection<SlotOffer>> |
JobMasterGateway.offerSlots(ResourceID taskManagerId,
Collection<SlotOffer> slots,
Time timeout)
Offers the given slots to the job manager.
|
CompletableFuture<Optional<Instance>> |
JobManagerGateway.requestTaskManagerInstance(ResourceID resourceId,
Time timeout)
Requests the TaskManager instance registered under the given instanceId from the JobManager.
|
Constructor and Description |
---|
JMTMRegistrationSuccess(ResourceID resourceID) |
JobManagerRunner(ResourceID resourceId,
JobGraph jobGraph,
Configuration configuration,
RpcService rpcService,
HighAvailabilityServices haServices,
HeartbeatServices heartbeatServices,
BlobServer blobServer,
JobManagerSharedServices jobManagerSharedServices,
JobManagerJobMetricGroupFactory jobManagerJobMetricGroupFactory,
FatalErrorHandler fatalErrorHandler)
Exceptions that occur while creating the JobManager or JobManagerRunner are directly
thrown and not reported to the given
FatalErrorHandler . |
JobMaster(RpcService rpcService,
JobMasterConfiguration jobMasterConfiguration,
ResourceID resourceId,
JobGraph jobGraph,
HighAvailabilityServices highAvailabilityService,
SlotPoolFactory slotPoolFactory,
JobManagerSharedServices jobManagerSharedServices,
HeartbeatServices heartbeatServices,
BlobServer blobServer,
JobManagerJobMetricGroupFactory jobMetricGroupFactory,
OnCompletionActions jobCompletionActions,
FatalErrorHandler fatalErrorHandler,
ClassLoader userCodeLoader) |
JobMasterRegistrationSuccess(long heartbeatInterval,
ResourceManagerId resourceManagerId,
ResourceID resourceManagerResourceId) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
SlotPoolGateway.registerTaskManager(ResourceID resourceID)
Registers a TaskExecutor with the given
ResourceID at SlotPool . |
CompletableFuture<Acknowledge> |
SlotPool.registerTaskManager(ResourceID resourceID)
Register TaskManager to this pool, only those slots come from registered TaskManager will be considered valid.
|
CompletableFuture<Acknowledge> |
SlotPoolGateway.releaseTaskManager(ResourceID resourceId,
Exception cause)
Releases a TaskExecutor with the given
ResourceID from the SlotPool . |
CompletableFuture<Acknowledge> |
SlotPool.releaseTaskManager(ResourceID resourceId,
Exception cause)
Unregister TaskManager from this pool, all the related slots will be released and tasks be canceled.
|
Modifier and Type | Method and Description |
---|---|
void |
MetricRegistryImpl.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,
long maximumFramesize)
Starts the MetricQueryService actor in the given actor system.
|
Modifier and Type | Method and Description |
---|---|
protected ResourceID |
StandaloneResourceManager.workerStarted(ResourceID resourceID) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Collection<Tuple2<ResourceID,String>>> |
ResourceManagerGateway.requestTaskManagerMetricQueryServicePaths(Time timeout)
Requests the paths for the TaskManager's
MetricQueryService to query. |
CompletableFuture<Collection<Tuple2<ResourceID,String>>> |
ResourceManager.requestTaskManagerMetricQueryServicePaths(Time timeout) |
Modifier and Type | Method and Description |
---|---|
protected void |
ResourceManager.closeTaskManagerConnection(ResourceID resourceID,
Exception cause)
This method should be called by the framework once it detects that a currently registered
task executor has failed.
|
void |
ResourceManagerGateway.disconnectTaskManager(ResourceID resourceID,
Exception cause)
Disconnects a TaskManager specified by the given resourceID from the
ResourceManager . |
void |
ResourceManager.disconnectTaskManager(ResourceID resourceId,
Exception cause) |
void |
ResourceManagerGateway.heartbeatFromJobManager(ResourceID heartbeatOrigin)
Sends the heartbeat to resource manager from job manager
|
void |
ResourceManager.heartbeatFromJobManager(ResourceID resourceID) |
void |
ResourceManagerGateway.heartbeatFromTaskManager(ResourceID heartbeatOrigin,
SlotReport slotReport)
Sends the heartbeat to resource manager from task manager
|
void |
ResourceManager.heartbeatFromTaskManager(ResourceID resourceID,
SlotReport slotReport) |
CompletableFuture<RegistrationResponse> |
ResourceManagerGateway.registerJobManager(JobMasterId jobMasterId,
ResourceID jobMasterResourceId,
String jobMasterAddress,
JobID jobId,
Time timeout)
Register a
JobMaster at the resource manager. |
CompletableFuture<RegistrationResponse> |
ResourceManager.registerJobManager(JobMasterId jobMasterId,
ResourceID jobManagerResourceId,
String jobManagerAddress,
JobID jobId,
Time timeout) |
CompletableFuture<RegistrationResponse> |
ResourceManagerGateway.registerTaskExecutor(String taskExecutorAddress,
ResourceID resourceId,
int dataPort,
HardwareDescription hardwareDescription,
Time timeout)
Register a
TaskExecutor at the resource manager. |
CompletableFuture<RegistrationResponse> |
ResourceManager.registerTaskExecutor(String taskExecutorAddress,
ResourceID taskExecutorResourceId,
int dataPort,
HardwareDescription hardwareDescription,
Time timeout) |
CompletableFuture<TransientBlobKey> |
ResourceManagerGateway.requestTaskManagerFileUpload(ResourceID taskManagerId,
FileType fileType,
Time timeout)
Request the file upload from the given
TaskExecutor to the cluster's BlobServer . |
CompletableFuture<TransientBlobKey> |
ResourceManager.requestTaskManagerFileUpload(ResourceID taskManagerId,
FileType fileType,
Time timeout) |
CompletableFuture<TaskManagerInfo> |
ResourceManagerGateway.requestTaskManagerInfo(ResourceID taskManagerId,
Time timeout)
Requests information about the given
TaskExecutor . |
CompletableFuture<TaskManagerInfo> |
ResourceManager.requestTaskManagerInfo(ResourceID resourceId,
Time timeout) |
CompletableFuture<Acknowledge> |
ResourceManagerGateway.sendSlotReport(ResourceID taskManagerResourceId,
InstanceID taskManagerRegistrationId,
SlotReport slotReport,
Time timeout)
Sends the given
SlotReport to the ResourceManager. |
CompletableFuture<Acknowledge> |
ResourceManager.sendSlotReport(ResourceID taskManagerResourceId,
InstanceID taskManagerRegistrationId,
SlotReport slotReport,
Time timeout) |
boolean |
StandaloneResourceManager.stopWorker(ResourceID resourceID) |
protected abstract WorkerType |
ResourceManager.workerStarted(ResourceID resourceID)
Callback when a worker was started.
|
protected ResourceID |
StandaloneResourceManager.workerStarted(ResourceID resourceID) |
Constructor and Description |
---|
ResourceManager(RpcService rpcService,
String resourceManagerEndpointId,
ResourceID resourceId,
ResourceManagerConfiguration resourceManagerConfiguration,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
SlotManager slotManager,
MetricRegistry metricRegistry,
JobLeaderIdService jobLeaderIdService,
ClusterInformation clusterInformation,
FatalErrorHandler fatalErrorHandler,
JobManagerMetricGroup jobManagerMetricGroup) |
ResourceManagerRunner(ResourceID resourceId,
String resourceManagerEndpointId,
Configuration configuration,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
ClusterInformation clusterInformation,
JobManagerMetricGroup jobManagerMetricGroup) |
StandaloneResourceManager(RpcService rpcService,
String resourceManagerEndpointId,
ResourceID resourceId,
ResourceManagerConfiguration resourceManagerConfiguration,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
SlotManager slotManager,
MetricRegistry metricRegistry,
JobLeaderIdService jobLeaderIdService,
ClusterInformation clusterInformation,
FatalErrorHandler fatalErrorHandler,
JobManagerMetricGroup jobManagerMetricGroup) |
Constructor and Description |
---|
UnknownTaskExecutorException(ResourceID taskExecutorId) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
JobManagerRegistration.getJobManagerResourceID() |
ResourceID |
TaskExecutorConnection.getResourceID() |
Constructor and Description |
---|
JobManagerRegistration(JobID jobID,
ResourceID jobManagerResourceID,
JobMasterGateway jobManagerGateway) |
TaskExecutorConnection(ResourceID resourceID,
TaskExecutorGateway taskExecutorGateway) |
Modifier and Type | Method and Description |
---|---|
protected abstract CompletableFuture<TransientBlobKey> |
AbstractTaskManagerFileHandler.requestFileUpload(ResourceManagerGateway resourceManagerGateway,
ResourceID taskManagerResourceId) |
protected CompletableFuture<TransientBlobKey> |
TaskManagerLogFileHandler.requestFileUpload(ResourceManagerGateway resourceManagerGateway,
ResourceID taskManagerResourceId) |
protected CompletableFuture<TransientBlobKey> |
TaskManagerStdoutFileHandler.requestFileUpload(ResourceManagerGateway resourceManagerGateway,
ResourceID taskManagerResourceId) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
TaskManagersFilterQueryParameter.convertStringToValue(String value) |
Modifier and Type | Method and Description |
---|---|
String |
TaskManagersFilterQueryParameter.convertValueToString(ResourceID value) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
ResourceIDDeserializer.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 |
ResourceIDSerializer.serialize(ResourceID 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 |
---|---|
protected ResourceID |
TaskManagerIdPathParameter.convertFromString(String value) |
ResourceID |
TaskManagerInfo.getResourceId() |
Modifier and Type | Method and Description |
---|---|
protected String |
TaskManagerIdPathParameter.convertToString(ResourceID value) |
Constructor and Description |
---|
TaskManagerDetailsInfo(ResourceID resourceId,
String address,
int dataPort,
long lastHeartbeat,
int numberSlots,
int numberAvailableSlots,
HardwareDescription hardwareDescription,
TaskManagerMetricsInfo taskManagerMetrics) |
TaskManagerInfo(ResourceID resourceId,
String address,
int dataPort,
long lastHeartbeat,
int numberSlots,
int numberAvailableSlots,
HardwareDescription hardwareDescription) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
TaskExecutor.getResourceID() |
ResourceID |
JobManagerConnection.getResourceID() |
ResourceID |
TaskExecutorRegistrationSuccess.getResourceManagerId()
Gets the unique ID that identifies the ResourceManager.
|
Modifier and Type | Method and Description |
---|---|
static TaskManagerServices |
TaskManagerServices.fromConfiguration(TaskManagerServicesConfiguration taskManagerServicesConfiguration,
ResourceID resourceID,
Executor taskIOExecutor,
long freeHeapMemoryWithDefrag,
long maxJvmHeapMemory)
Creates and returns the task manager services.
|
void |
TaskExecutorGateway.heartbeatFromJobManager(ResourceID heartbeatOrigin)
Heartbeat request from the job manager.
|
void |
TaskExecutor.heartbeatFromJobManager(ResourceID resourceID) |
void |
TaskExecutorGateway.heartbeatFromResourceManager(ResourceID heartbeatOrigin)
Heartbeat request from the resource manager.
|
void |
TaskExecutor.heartbeatFromResourceManager(ResourceID resourceID) |
static void |
TaskManagerRunner.runTaskManager(Configuration configuration,
ResourceID resourceId) |
static TaskExecutor |
TaskManagerRunner.startTaskManager(Configuration configuration,
ResourceID resourceID,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
BlobCacheService blobCacheService,
boolean localCommunicationOnly,
FatalErrorHandler fatalErrorHandler) |
Constructor and Description |
---|
JobManagerConnection(JobID jobID,
ResourceID resourceID,
JobMasterGateway jobMasterGateway,
TaskManagerActions taskManagerActions,
CheckpointResponder checkpointResponder,
LibraryCacheManager libraryCacheManager,
ResultPartitionConsumableNotifier resultPartitionConsumableNotifier,
PartitionProducerStateChecker partitionStateChecker) |
TaskExecutorRegistrationSuccess(InstanceID registrationId,
ResourceID resourceManagerResourceId,
long heartbeatInterval,
ClusterInformation clusterInformation)
Create a new
TaskExecutorRegistrationSuccess message. |
TaskExecutorToResourceManagerConnection(org.slf4j.Logger log,
RpcService rpcService,
String taskManagerAddress,
ResourceID taskManagerResourceId,
int dataPort,
HardwareDescription hardwareDescription,
String resourceManagerAddress,
ResourceManagerId resourceManagerId,
Executor executor,
RegistrationConnectionListener<TaskExecutorToResourceManagerConnection,TaskExecutorRegistrationSuccess> registrationListener) |
TaskManagerRunner(Configuration configuration,
ResourceID resourceId) |
Modifier and Type | Method and Description |
---|---|
SlotReport |
TaskSlotTable.createSlotReport(ResourceID resourceId) |
Modifier and Type | Method and Description |
---|---|
ResourceID |
TaskManagerLocation.getResourceID()
Gets the ID of the resource in which the TaskManager is started.
|
Constructor and Description |
---|
TaskManagerLocation(ResourceID resourceID,
InetAddress inetAddress,
int dataPort)
Constructs a new instance connection info object.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Collection<Tuple2<ResourceID,String>>> |
RestfulGateway.requestTaskManagerMetricQueryServicePaths(Time timeout)
Requests the paths for the TaskManager's
MetricQueryService to query. |
Modifier and Type | Method and Description |
---|---|
ResourceID |
YarnContainerInLaunch.getResourceID() |
ResourceID |
YarnWorkerNode.getResourceID() |
ResourceID |
RegisteredYarnWorkerNode.getResourceID() |
Modifier and Type | Method and Description |
---|---|
protected void |
YarnFlinkResourceManager.releasePendingWorker(ResourceID id) |
protected YarnWorkerNode |
YarnResourceManager.workerStarted(ResourceID resourceID) |
protected RegisteredYarnWorkerNode |
YarnFlinkResourceManager.workerStarted(ResourceID resourceID) |
Modifier and Type | Method and Description |
---|---|
protected Collection<RegisteredYarnWorkerNode> |
YarnFlinkResourceManager.reacceptRegisteredWorkers(Collection<ResourceID> toConsolidate) |
Constructor and Description |
---|
YarnResourceManager(RpcService rpcService,
String resourceManagerEndpointId,
ResourceID resourceId,
Configuration flinkConfig,
Map<String,String> env,
ResourceManagerConfiguration resourceManagerConfiguration,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
SlotManager slotManager,
MetricRegistry metricRegistry,
JobLeaderIdService jobLeaderIdService,
ClusterInformation clusterInformation,
FatalErrorHandler fatalErrorHandler,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Modifier and Type | Method and Description |
---|---|
protected ResourceManager<?> |
YarnJobClusterEntrypoint.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
protected ResourceManager<?> |
YarnSessionClusterEntrypoint.createResourceManager(Configuration configuration,
ResourceID resourceId,
RpcService rpcService,
HighAvailabilityServices highAvailabilityServices,
HeartbeatServices heartbeatServices,
MetricRegistry metricRegistry,
FatalErrorHandler fatalErrorHandler,
ClusterInformation clusterInformation,
String webInterfaceUrl,
JobManagerMetricGroup jobManagerMetricGroup) |
Copyright © 2014–2020 The Apache Software Foundation. All rights reserved.