Modifier and Type | Method and Description |
---|---|
Time |
RestartStrategies.FixedDelayRestartStrategyConfiguration.getDelayBetweenAttemptsInterval() |
Time |
RestartStrategies.FailureRateRestartStrategyConfiguration.getDelayBetweenAttemptsInterval() |
Time |
RestartStrategies.FailureRateRestartStrategyConfiguration.getFailureInterval() |
Modifier and Type | Method and Description |
---|---|
static RestartStrategies.FailureRateRestartStrategyConfiguration |
RestartStrategies.failureRateRestart(int failureRate,
Time failureInterval,
Time delayInterval)
Generates a FailureRateRestartStrategyConfiguration.
|
static RestartStrategies.RestartStrategyConfiguration |
RestartStrategies.fixedDelayRestart(int restartAttempts,
Time delayInterval)
Generates a FixedDelayRestartStrategyConfiguration.
|
Constructor and Description |
---|
FailureRateRestartStrategyConfiguration(int maxFailureRate,
Time failureInterval,
Time delayBetweenAttemptsInterval) |
Modifier and Type | Method and Description |
---|---|
static Time |
Time.days(long days)
Creates a new
Time that represents the given number of days. |
static Time |
Time.hours(long hours)
Creates a new
Time that represents the given number of hours. |
static Time |
Time.milliseconds(long milliseconds)
Creates a new
Time that represents the given number of milliseconds. |
static Time |
Time.minutes(long minutes)
Creates a new
Time that represents the given number of minutes. |
static Time |
Time.of(long size,
TimeUnit unit)
|
static Time |
Time.seconds(long seconds)
Creates a new
Time that represents the given number of seconds. |
Constructor and Description |
---|
DefaultQuarantineHandler(Time timeout,
int exitCode,
org.slf4j.Logger log) |
Modifier and Type | Method and Description |
---|---|
static List<MasterState> |
MasterHooks.triggerMasterHooks(Collection<MasterTriggerRestoreHook<?>> hooks,
long checkpointId,
long timestamp,
Executor executor,
Time timeout)
Triggers all given master hooks and returns state objects for each hook that
produced a state.
|
Modifier and Type | Method and Description |
---|---|
static CompletableFuture<InetSocketAddress> |
JobClient.retrieveBlobServerAddress(JobManagerGateway jobManagerGateway,
Time timeout)
Utility method to retrieve the BlobServer address from the given JobManager gateway.
|
static ClassLoader |
JobClient.retrieveClassLoader(JobID jobID,
JobManagerGateway jobManager,
Configuration config,
HighAvailabilityServices highAvailabilityServices,
Time timeout)
Reconstructs the class loader by first requesting information about it at the JobManager
and then downloading missing jar files.
|
static void |
JobClient.submitJobDetached(JobManagerGateway jobManagerGateway,
Configuration config,
JobGraph jobGraph,
Time timeout,
ClassLoader classLoader)
Submits a job in detached mode.
|
Modifier and Type | Method and Description |
---|---|
static WebMonitor |
BootstrapTools.startWebMonitorIfConfigured(Configuration config,
HighAvailabilityServices highAvailabilityServices,
LeaderGatewayRetriever<JobManagerGateway> jobManagerRetriever,
MetricQueryServiceRetriever queryServiceRetriever,
Time timeout,
ScheduledExecutor scheduledExecutor,
org.slf4j.Logger logger)
Starts the web frontend.
|
Modifier and Type | Method and Description |
---|---|
static Time |
FutureUtils.toTime(scala.concurrent.duration.FiniteDuration finiteDuration)
Converts
FiniteDuration into Flink time. |
Modifier and Type | Method and Description |
---|---|
static <T> CompletableFuture<T> |
FutureUtils.retrySuccesfulWithDelay(java.util.function.Supplier<CompletableFuture<T>> operation,
Time retryDelay,
scala.concurrent.duration.Deadline deadline,
java.util.function.Predicate<T> acceptancePredicate,
ScheduledExecutor scheduledExecutor)
Retry the given operation with the given delay in between successful completions where the
result does not match a given predicate.
|
static <T> CompletableFuture<T> |
FutureUtils.retryWithDelay(java.util.function.Supplier<CompletableFuture<T>> operation,
int retries,
Time retryDelay,
ScheduledExecutor scheduledExecutor)
Retry the given operation with the given delay in between failures.
|
static scala.concurrent.duration.FiniteDuration |
FutureUtils.toFiniteDuration(Time time)
Converts Flink time into a
FiniteDuration . |
Modifier and Type | Method and Description |
---|---|
static ExecutionGraph |
ExecutionGraphBuilder.buildGraph(ExecutionGraph prior,
JobGraph jobGraph,
Configuration jobManagerConfig,
ScheduledExecutorService futureExecutor,
Executor ioExecutor,
SlotProvider slotProvider,
ClassLoader classLoader,
CheckpointRecoveryFactory recoveryFactory,
Time timeout,
RestartStrategy restartStrategy,
MetricGroup metrics,
int parallelismForAutoMax,
BlobWriter blobWriter,
org.slf4j.Logger log)
Builds the ExecutionGraph from the JobGraph.
|
CompletableFuture<StackTraceSampleResponse> |
Execution.requestStackTraceSample(int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStrackTraceDepth,
Time timeout)
Request a stack trace sample from the task of this execution.
|
Constructor and Description |
---|
Execution(Executor executor,
ExecutionVertex vertex,
int attemptNumber,
long globalModVersion,
long startTimestamp,
Time timeout)
Creates a new Execution attempt.
|
ExecutionGraph(JobInformation jobInformation,
ScheduledExecutorService futureExecutor,
Executor ioExecutor,
Time timeout,
RestartStrategy restartStrategy,
FailoverStrategy.Factory failoverStrategyFactory,
SlotProvider slotProvider,
ClassLoader userClassLoader,
BlobWriter blobWriter) |
ExecutionJobVertex(ExecutionGraph graph,
JobVertex jobVertex,
int defaultParallelism,
Time timeout,
long initialGlobalModVersion,
long createTimestamp) |
ExecutionVertex(ExecutionJobVertex jobVertex,
int subTaskIndex,
IntermediateResult[] producedDataSets,
Time timeout,
long initialGlobalModVersion,
long createTimestamp,
int maxPriorExecutionHistoryLength) |
Constructor and Description |
---|
FailureRateRestartStrategy(int maxFailuresPerInterval,
Time failuresInterval,
Time delayInterval) |
FailureRateRestartStrategyFactory(int maxFailuresPerInterval,
Time failuresInterval,
Time delayInterval) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<SimpleSlot> |
SlotPoolGateway.allocateSlot(ScheduledUnit task,
ResourceProfile resources,
Iterable<TaskManagerLocation> locationPreferences,
Time timeout) |
CompletableFuture<SimpleSlot> |
SlotPool.allocateSlot(ScheduledUnit task,
ResourceProfile resources,
Iterable<TaskManagerLocation> locationPreferences,
Time timeout) |
Constructor and Description |
---|
SlotPool(RpcService rpcService,
JobID jobId,
Clock clock,
Time slotRequestTimeout,
Time resourceManagerAllocationTimeout,
Time resourceManagerRequestTimeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
CompletableFuture<StackTrace> |
ActorTaskManagerGateway.requestStackTrace(Time timeout) |
CompletableFuture<StackTrace> |
TaskManagerGateway.requestStackTrace(Time timeout)
Request the stack trace from the task manager.
|
CompletableFuture<StackTraceSampleResponse> |
ActorTaskManagerGateway.requestStackTraceSample(ExecutionAttemptID executionAttemptID,
int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth,
Time timeout) |
CompletableFuture<StackTraceSampleResponse> |
TaskManagerGateway.requestStackTraceSample(ExecutionAttemptID executionAttemptID,
int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth,
Time timeout)
Request a stack trace sample from the given task.
|
CompletableFuture<TransientBlobKey> |
ActorTaskManagerGateway.requestTaskManagerLog(Time timeout) |
CompletableFuture<TransientBlobKey> |
TaskManagerGateway.requestTaskManagerLog(Time timeout)
Request the task manager log from the task manager.
|
CompletableFuture<TransientBlobKey> |
ActorTaskManagerGateway.requestTaskManagerStdout(Time timeout) |
CompletableFuture<TransientBlobKey> |
TaskManagerGateway.requestTaskManagerStdout(Time timeout)
Request the task manager stdout from the task manager.
|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Stop the given task.
|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout)
Submit a task to the task manager.
|
CompletableFuture<Acknowledge> |
ActorTaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
Modifier and Type | Field and Description |
---|---|
Time |
JobManagerServices.rpcAskTimeout |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
JobMaster.cancel(Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.cancel(Time timeout)
Cancels the currently executed job.
|
CompletableFuture<Acknowledge> |
JobManagerGateway.cancelJob(JobID jobId,
Time timeout)
Cancels the given job.
|
CompletableFuture<String> |
JobManagerGateway.cancelJobWithSavepoint(JobID jobId,
String savepointPath,
Time timeout)
Cancels the given job after taking a savepoint and returning its path.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Collection<SlotOffer>> |
JobMaster.offerSlots(ResourceID taskManagerId,
Iterable<SlotOffer> slots,
Time timeout) |
CompletableFuture<Collection<SlotOffer>> |
JobMasterGateway.offerSlots(ResourceID taskManagerId,
Iterable<SlotOffer> slots,
Time timeout)
Offers the given slots to the job manager.
|
CompletableFuture<RegistrationResponse> |
JobMaster.registerTaskManager(String taskManagerRpcAddress,
TaskManagerLocation taskManagerLocation,
Time timeout) |
CompletableFuture<RegistrationResponse> |
JobMasterGateway.registerTaskManager(String taskManagerRpcAddress,
TaskManagerLocation taskManagerLocation,
Time timeout)
Registers the task manager at the job manager.
|
CompletableFuture<AccessExecutionGraph> |
JobMaster.requestArchivedExecutionGraph(Time timeout) |
CompletableFuture<AccessExecutionGraph> |
JobMasterGateway.requestArchivedExecutionGraph(Time timeout)
Request the
ArchivedExecutionGraph of the currently executed job. |
CompletableFuture<Integer> |
JobManagerGateway.requestBlobServerPort(Time timeout)
Requests the BlobServer port.
|
CompletableFuture<Optional<org.apache.flink.runtime.messages.JobManagerMessages.ClassloadingProps>> |
JobManagerGateway.requestClassloadingProps(JobID jobId,
Time timeout)
Requests the class loading properties for the given JobID.
|
CompletableFuture<JobDetails> |
JobMaster.requestJobDetails(Time timeout) |
CompletableFuture<JobDetails> |
JobMasterGateway.requestJobDetails(Time timeout)
Request the details of the executed job.
|
CompletableFuture<JobsWithIDsOverview> |
JobManagerGateway.requestJobsOverview(Time timeout)
Requests the job overview from the JobManager.
|
CompletableFuture<JobStatus> |
JobMaster.requestJobStatus(Time timeout) |
CompletableFuture<JobStatus> |
JobMasterGateway.requestJobStatus(Time timeout)
Requests the current job status.
|
CompletableFuture<StackTrace> |
RpcTaskManagerGateway.requestStackTrace(Time timeout) |
CompletableFuture<StackTraceSampleResponse> |
RpcTaskManagerGateway.requestStackTraceSample(ExecutionAttemptID executionAttemptID,
int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth,
Time timeout) |
CompletableFuture<Optional<Instance>> |
JobManagerGateway.requestTaskManagerInstance(ResourceID resourceId,
Time timeout)
Requests the TaskManager instance registered under the given instanceId from the JobManager.
|
CompletableFuture<Collection<Instance>> |
JobManagerGateway.requestTaskManagerInstances(Time timeout)
Requests all currently registered TaskManager instances from the JobManager.
|
CompletableFuture<TransientBlobKey> |
RpcTaskManagerGateway.requestTaskManagerLog(Time timeout) |
CompletableFuture<TransientBlobKey> |
RpcTaskManagerGateway.requestTaskManagerStdout(Time timeout) |
CompletableFuture<Acknowledge> |
JobMaster.scheduleOrUpdateConsumers(ResultPartitionID partitionID,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.scheduleOrUpdateConsumers(ResultPartitionID partitionID,
Time timeout)
Notifies the JobManager about available data for a produced partition.
|
CompletableFuture<Acknowledge> |
JobMaster.start(JobMasterId newJobMasterId,
Time timeout)
Start the rpc service and begin to run the job.
|
CompletableFuture<Acknowledge> |
JobMaster.stop(Time timeout) |
CompletableFuture<Acknowledge> |
JobMasterGateway.stop(Time timeout)
Cancel the currently executed job.
|
CompletableFuture<Acknowledge> |
JobManagerGateway.stopJob(JobID jobId,
Time timeout)
Stops the given job.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
JobManagerGateway.submitJob(JobGraph jobGraph,
ListeningBehaviour listeningBehaviour,
Time timeout)
Submits a job to the JobManager.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout) |
CompletableFuture<Acknowledge> |
JobMaster.suspend(Throwable cause,
Time timeout)
Suspending job, all the running tasks will be cancelled, and communication with other components
will be disposed.
|
CompletableFuture<Acknowledge> |
RpcTaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
Constructor and Description |
---|
JobManagerServices(ScheduledExecutorService executorService,
BlobServer blobServer,
BlobLibraryCacheManager libraryCacheManager,
RestartStrategyFactory restartStrategyFactory,
Time rpcAskTimeout) |
JobMaster(RpcService rpcService,
ResourceID resourceId,
JobGraph jobGraph,
Configuration configuration,
HighAvailabilityServices highAvailabilityService,
HeartbeatServices heartbeatServices,
ScheduledExecutorService executor,
BlobServer blobServer,
BlobLibraryCacheManager libraryCacheManager,
RestartStrategyFactory restartStrategyFactory,
Time rpcAskTimeout,
JobManagerMetricGroup jobManagerMetricGroup,
OnCompletionActions jobCompletionActions,
FatalErrorHandler errorHandler,
ClassLoader userCodeLoader) |
Modifier and Type | Method and Description |
---|---|
Time |
MiniClusterConfiguration.getRpcTimeout() |
Modifier and Type | Method and Description |
---|---|
protected RpcService |
MiniCluster.createRpcService(Configuration configuration,
Time askTimeout,
boolean remoteEnabled,
String bindAddress)
Factory method to instantiate the RPC service.
|
Modifier and Type | Method and Description |
---|---|
Time |
ResourceManagerConfiguration.getHeartbeatInterval() |
Time |
ResourceManagerRuntimeServicesConfiguration.getJobTimeout() |
Time |
ResourceManagerConfiguration.getTimeout() |
Modifier and Type | Method and Description |
---|---|
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,
SlotReport slotReport,
Time timeout)
Register a
TaskExecutor at the resource manager. |
CompletableFuture<RegistrationResponse> |
ResourceManager.registerTaskExecutor(String taskExecutorAddress,
ResourceID taskExecutorResourceId,
SlotReport slotReport,
Time timeout) |
CompletableFuture<ResourceOverview> |
ResourceManagerGateway.requestResourceOverview(Time timeout)
Requests the resource overview.
|
CompletableFuture<ResourceOverview> |
ResourceManager.requestResourceOverview(Time timeout) |
CompletableFuture<Acknowledge> |
ResourceManagerGateway.requestSlot(JobMasterId jobMasterId,
SlotRequest slotRequest,
Time timeout)
Requests a slot from the resource manager.
|
CompletableFuture<Acknowledge> |
ResourceManager.requestSlot(JobMasterId jobMasterId,
SlotRequest slotRequest,
Time timeout) |
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) |
Constructor and Description |
---|
JobLeaderIdService(HighAvailabilityServices highAvailabilityServices,
ScheduledExecutor scheduledExecutor,
Time jobTimeout) |
ResourceManagerConfiguration(Time timeout,
Time heartbeatInterval) |
ResourceManagerRuntimeServicesConfiguration(Time jobTimeout,
SlotManagerConfiguration slotManagerConfiguration) |
Modifier and Type | Method and Description |
---|---|
Time |
SlotManagerConfiguration.getSlotRequestTimeout() |
Time |
SlotManagerConfiguration.getTaskManagerRequestTimeout() |
Time |
SlotManagerConfiguration.getTaskManagerTimeout() |
Constructor and Description |
---|
SlotManager(ScheduledExecutor scheduledExecutor,
Time taskManagerRequestTimeout,
Time slotRequestTimeout,
Time taskManagerTimeout) |
SlotManagerConfiguration(Time taskManagerRequestTimeout,
Time slotRequestTimeout,
Time taskManagerTimeout) |
Modifier and Type | Method and Description |
---|---|
void |
RestClient.shutdown(Time timeout) |
void |
RestServerEndpoint.shutdown(Time timeout)
Stops this REST server endpoint.
|
Modifier and Type | Field and Description |
---|---|
protected Time |
RedirectHandler.timeout |
Modifier and Type | Method and Description |
---|---|
Time |
RestHandlerConfiguration.getTimeout() |
Constructor and Description |
---|
AbstractRestHandler(CompletableFuture<String> localRestAddress,
GatewayRetriever<? extends T> leaderRetriever,
Time timeout,
MessageHeaders<R,P,M> messageHeaders) |
LegacyRestHandlerAdapter(CompletableFuture<String> localRestAddress,
GatewayRetriever<T> leaderRetriever,
Time timeout,
MessageHeaders<EmptyRequestBody,R,M> messageHeaders,
LegacyRestHandler<T,R,M> legacyRestHandler) |
RedirectHandler(CompletableFuture<String> localAddressFuture,
GatewayRetriever<? extends T> leaderRetriever,
Time timeout) |
RestHandlerConfiguration(long refreshInterval,
int maxCheckpointStatisticCacheEntries,
Time timeout,
File tmpDir) |
Constructor and Description |
---|
ClusterOverviewHandler(Executor executor,
Time timeout) |
CurrentJobIdsHandler(Executor executor,
Time timeout) |
CurrentJobsOverviewHandler(Executor executor,
Time timeout,
boolean includeRunningJobs,
boolean includeFinishedJobs) |
ExecutionGraphCache(Time timeout,
Time timeToLive) |
JobCancellationHandler(Executor executor,
Time timeout) |
JobStoppingHandler(Executor executor,
Time timeout) |
TaskManagerLogHandler(GatewayRetriever<JobManagerGateway> retriever,
Executor executor,
CompletableFuture<String> localJobManagerAddressPromise,
Time timeout,
TaskManagerLogHandler.FileMode fileMode,
Configuration config) |
TaskManagersHandler(Executor executor,
Time timeout,
MetricFetcher fetcher) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<StackTraceSample> |
StackTraceSampleCoordinator.triggerStackTraceSample(ExecutionVertex[] tasksToSample,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth)
Triggers a stack trace sample to all tasks.
|
Constructor and Description |
---|
BackPressureStatsTracker(StackTraceSampleCoordinator coordinator,
int cleanUpInterval,
int numSamples,
Time delayBetweenSamples)
Creates a back pressure statistics tracker.
|
Constructor and Description |
---|
StaticFileServerHandler(GatewayRetriever<T> retriever,
CompletableFuture<String> localJobManagerAddressFuture,
Time timeout,
File rootPath) |
Constructor and Description |
---|
MetricFetcher(GatewayRetriever<T> retriever,
MetricQueryServiceRetriever queryServiceRetriever,
Executor executor,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
static CompletableFuture<Optional<String>> |
HandlerRedirectUtils.getRedirectAddress(String localRestAddress,
RestfulGateway restfulGateway,
Time timeout) |
Modifier and Type | Field and Description |
---|---|
static Time |
RpcUtils.INF_TIMEOUT |
Modifier and Type | Method and Description |
---|---|
protected <V> CompletableFuture<V> |
RpcEndpoint.callAsync(Callable<V> callable,
Time timeout)
Execute the callable in the main thread of the underlying RPC service, returning a future for
the result of the callable.
|
<V> CompletableFuture<V> |
MainThreadExecutable.callAsync(Callable<V> callable,
Time callTimeout)
Execute the callable in the main thread of the underlying RPC endpoint and return a future for
the callable result.
|
<V> CompletableFuture<V> |
FencedMainThreadExecutable.callAsyncWithoutFencing(Callable<V> callable,
Time timeout)
Run the given callable in the main thread without attaching a fencing token.
|
protected <V> CompletableFuture<V> |
FencedRpcEndpoint.callAsyncWithoutFencing(Callable<V> callable,
Time timeout)
Run the given callable in the main thread of the RpcEndpoint without checking the fencing
token.
|
protected void |
RpcEndpoint.scheduleRunAsync(Runnable runnable,
Time delay)
Execute the runnable in the main thread of the underlying RPC endpoint, with
a delay of the given number of milliseconds.
|
static void |
RpcUtils.terminateRpcEndpoint(RpcEndpoint rpcEndpoint,
Time timeout)
Shuts the given
RpcEndpoint down and awaits its termination. |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<?> |
FencedAkkaInvocationHandler.ask(Object message,
Time timeout) |
<V> CompletableFuture<V> |
FencedAkkaInvocationHandler.callAsyncWithoutFencing(Callable<V> callable,
Time timeout) |
Constructor and Description |
---|
AkkaRpcService(akka.actor.ActorSystem actorSystem,
Time timeout) |
FencedAkkaInvocationHandler(String address,
String hostname,
akka.actor.ActorRef rpcEndpoint,
Time timeout,
long maximumFramesize,
CompletableFuture<Boolean> terminationFuture,
CompletableFuture<Void> internalTerminationFuture,
java.util.function.Supplier<F> fencingTokenSupplier) |
Modifier and Type | Method and Description |
---|---|
Time |
TaskManagerConfiguration.getInitialRegistrationPause() |
Time |
TaskManagerConfiguration.getMaxRegistrationDuration() |
Time |
TaskManagerConfiguration.getMaxRegistrationPause() |
Time |
TaskManagerConfiguration.getRefusedRegistrationPause() |
Time |
TaskManagerConfiguration.getTimeout() |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<Acknowledge> |
TaskExecutorGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutor.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
String targetAddress,
ResourceManagerId resourceManagerId,
Time timeout)
Requests a slot from the TaskManager
|
CompletableFuture<Acknowledge> |
TaskExecutor.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
String targetAddress,
ResourceManagerId resourceManagerId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Stop the given task.
|
CompletableFuture<Acknowledge> |
TaskExecutor.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.submitTask(TaskDeploymentDescriptor tdd,
JobMasterId jobMasterId,
Time timeout)
Submit a
Task to the TaskExecutor . |
CompletableFuture<Acknowledge> |
TaskExecutor.submitTask(TaskDeploymentDescriptor tdd,
JobMasterId jobMasterId,
Time timeout) |
CompletableFuture<Acknowledge> |
TaskExecutorGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
CompletableFuture<Acknowledge> |
TaskExecutor.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
Constructor and Description |
---|
TaskManagerConfiguration(int numberSlots,
String[] tmpDirectories,
Time timeout,
Time maxRegistrationDuration,
Time initialRegistrationPause,
Time maxRegistrationPause,
Time refusedRegistrationPause,
long cleanupInterval,
Configuration configuration,
boolean exitJvmOnOutOfMemory,
FlinkUserCodeClassLoaders.ResolveOrder classLoaderResolveOrder,
String[] alwaysParentFirstLoaderPatterns) |
Constructor and Description |
---|
RpcInputSplitProvider(JobMasterGateway jobMasterGateway,
JobVertexID jobVertexID,
ExecutionAttemptID executionAttemptID,
Time timeout) |
RpcResultPartitionConsumableNotifier(JobMasterGateway jobMasterGateway,
Executor executor,
Time timeout) |
Modifier and Type | Method and Description |
---|---|
boolean |
TaskSlotTable.allocateSlot(int index,
JobID jobId,
AllocationID allocationId,
Time slotTimeout)
Allocate the slot with the given index for the given job and allocation id.
|
boolean |
TaskSlotTable.markSlotInactive(AllocationID allocationId,
Time slotTimeout)
Marks the slot under the given allocation id as inactive.
|
Modifier and Type | Method and Description |
---|---|
static InetAddress |
LeaderRetrievalUtils.findConnectingAddress(LeaderRetrievalService leaderRetrievalService,
Time timeout) |
Modifier and Type | Field and Description |
---|---|
static Time |
WebRuntimeMonitor.DEFAULT_REQUEST_TIMEOUT
By default, all requests to the JobManager have a timeout of 10 seconds.
|
Modifier and Type | Method and Description |
---|---|
CompletableFuture<ClusterOverview> |
RestfulGateway.requestClusterOverview(Time timeout)
Requests the cluster status overview.
|
CompletableFuture<AccessExecutionGraph> |
RestfulGateway.requestJob(JobID jobId,
Time timeout)
Requests the AccessExecutionGraph for the given jobId.
|
CompletableFuture<MultipleJobsDetails> |
RestfulGateway.requestJobDetails(boolean includeRunning,
boolean includeFinished,
Time timeout)
Requests job details currently being executed on the Flink cluster.
|
CompletableFuture<Collection<String>> |
RestfulGateway.requestMetricQueryServicePaths(Time timeout)
Requests the paths for the
MetricQueryService to query. |
CompletableFuture<String> |
RestfulGateway.requestRestAddress(Time timeout)
Requests the REST address of this
RpcEndpoint . |
CompletableFuture<Collection<Tuple2<ResourceID,String>>> |
RestfulGateway.requestTaskManagerMetricQueryServicePaths(Time timeout)
Requests the paths for the TaskManager's
MetricQueryService to query. |
static WebMonitor |
WebMonitorUtils.startWebRuntimeMonitor(Configuration config,
HighAvailabilityServices highAvailabilityServices,
LeaderGatewayRetriever<JobManagerGateway> jobManagerRetriever,
MetricQueryServiceRetriever queryServiceRetriever,
Time timeout,
ScheduledExecutor scheduledExecutor)
Starts the web runtime monitor.
|
static <T extends RestfulGateway> |
WebMonitorUtils.tryLoadWebContent(GatewayRetriever<T> leaderRetriever,
CompletableFuture<String> restAddressFuture,
Time timeout,
File tmpDir)
Checks whether the flink-runtime-web dependency is available and if so returns a
StaticFileServerHandler which can serve the static file contents.
|
Constructor and Description |
---|
RuntimeMonitorHandler(WebMonitorConfig cfg,
RequestHandler handler,
GatewayRetriever<JobManagerGateway> retriever,
CompletableFuture<String> localJobManagerAddressFuture,
Time timeout) |
WebRuntimeMonitor(Configuration config,
LeaderRetrievalService leaderRetrievalService,
LeaderGatewayRetriever<JobManagerGateway> jobManagerRetriever,
MetricQueryServiceRetriever queryServiceRetriever,
Time timeout,
ScheduledExecutor scheduledExecutor) |
Constructor and Description |
---|
JarRunHandler(Executor executor,
File jarDirectory,
Time timeout,
Configuration clientConfig) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<MetricDumpSerialization.MetricSerializationResult> |
MetricQueryServiceGateway.queryMetrics(Time timeout) |
Modifier and Type | Method and Description |
---|---|
CompletableFuture<MetricDumpSerialization.MetricSerializationResult> |
AkkaQueryServiceGateway.queryMetrics(Time timeout) |
Constructor and Description |
---|
AkkaJobManagerRetriever(akka.actor.ActorSystem actorSystem,
Time timeout,
int retries,
Time retryDelay) |
AkkaQueryServiceRetriever(akka.actor.ActorSystem actorSystem,
Time lookupTimeout) |
RpcGatewayRetriever(RpcService rpcService,
Class<T> gatewayType,
java.util.function.Function<UUID,F> fencingTokenMapper,
int retries,
Time retryDelay) |
Modifier and Type | Field and Description |
---|---|
static Time |
FlinkKafkaProducer011.DEFAULT_KAFKA_TRANSACTION_TIMEOUT
Default value for kafka transaction timeout.
|
Modifier and Type | Field and Description |
---|---|
static Time |
TestBaseUtils.DEFAULT_HTTP_TIMEOUT |
Modifier and Type | Method and Description |
---|---|
static String |
TestBaseUtils.getFromHTTP(String url,
Time timeout) |
Copyright © 2014–2018 The Apache Software Foundation. All rights reserved.