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. |
Modifier and Type | Method and Description |
---|---|
Time |
AkkaUtils$.getDefaultTimeout() |
static Time |
AkkaUtils.getDefaultTimeout() |
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 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,
org.slf4j.Logger log)
Builds the ExecutionGraph from the JobGraph.
|
Future<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(ScheduledExecutorService futureExecutor,
Executor ioExecutor,
JobID jobId,
String jobName,
Configuration jobConfig,
SerializedValue<ExecutionConfig> serializedConfig,
Time timeout,
RestartStrategy restartStrategy,
FailoverStrategy.Factory failoverStrategyFactory,
List<BlobKey> requiredJarFiles,
List<URL> requiredClasspaths,
SlotProvider slotProvider,
ClassLoader userClassLoader) |
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 |
---|---|
Future<SimpleSlot> |
SlotPoolGateway.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 |
---|---|
Future<Acknowledge> |
TaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
Future<Acknowledge> |
ActorTaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
Future<StackTrace> |
TaskManagerGateway.requestStackTrace(Time timeout)
Request the stack trace from the task manager.
|
Future<StackTrace> |
ActorTaskManagerGateway.requestStackTrace(Time timeout) |
Future<StackTraceSampleResponse> |
TaskManagerGateway.requestStackTraceSample(ExecutionAttemptID executionAttemptID,
int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth,
Time timeout)
Request a stack trace sample from the given task.
|
Future<StackTraceSampleResponse> |
ActorTaskManagerGateway.requestStackTraceSample(ExecutionAttemptID executionAttemptID,
int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth,
Time timeout) |
Future<BlobKey> |
TaskManagerGateway.requestTaskManagerLog(Time timeout)
Request the task manager log from the task manager.
|
Future<BlobKey> |
ActorTaskManagerGateway.requestTaskManagerLog(Time timeout) |
Future<BlobKey> |
TaskManagerGateway.requestTaskManagerStdout(Time timeout)
Request the task manager stdout from the task manager.
|
Future<BlobKey> |
ActorTaskManagerGateway.requestTaskManagerStdout(Time timeout) |
Future<Acknowledge> |
TaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Stop the given task.
|
Future<Acknowledge> |
ActorTaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
Future<Acknowledge> |
TaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout)
Submit a task to the task manager.
|
Future<Acknowledge> |
ActorTaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout) |
Future<Acknowledge> |
TaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
Future<Acknowledge> |
ActorTaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
Modifier and Type | Field and Description |
---|---|
Time |
JobManagerServices.rpcAskTimeout |
Modifier and Type | Method and Description |
---|---|
Future<Acknowledge> |
RpcTaskManagerGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
Future<Iterable<SlotOffer>> |
JobMasterGateway.offerSlots(ResourceID taskManagerId,
Iterable<SlotOffer> slots,
UUID leaderId,
Time timeout)
Offer the given slots to the job manager.
|
Future<RegistrationResponse> |
JobMasterGateway.registerTaskManager(String taskManagerRpcAddress,
TaskManagerLocation taskManagerLocation,
UUID leaderId,
Time timeout)
Register the task manager at the job manager.
|
Future<StackTrace> |
RpcTaskManagerGateway.requestStackTrace(Time timeout) |
Future<StackTraceSampleResponse> |
RpcTaskManagerGateway.requestStackTraceSample(ExecutionAttemptID executionAttemptID,
int sampleId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth,
Time timeout) |
Future<BlobKey> |
RpcTaskManagerGateway.requestTaskManagerLog(Time timeout) |
Future<BlobKey> |
RpcTaskManagerGateway.requestTaskManagerStdout(Time timeout) |
Future<Acknowledge> |
JobMasterGateway.scheduleOrUpdateConsumers(UUID leaderSessionID,
ResultPartitionID partitionID,
Time timeout)
Notifies the JobManager about available data for a produced partition.
|
Future<Acknowledge> |
RpcTaskManagerGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout) |
Future<Acknowledge> |
RpcTaskManagerGateway.submitTask(TaskDeploymentDescriptor tdd,
Time timeout) |
Future<Acknowledge> |
RpcTaskManagerGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout) |
Constructor and Description |
---|
JobManagerServices(ScheduledExecutorService executorService,
BlobLibraryCacheManager libraryCacheManager,
RestartStrategyFactory restartStrategyFactory,
Time rpcAskTimeout) |
JobMaster(RpcService rpcService,
ResourceID resourceId,
JobGraph jobGraph,
Configuration configuration,
HighAvailabilityServices highAvailabilityService,
HeartbeatServices heartbeatServices,
ScheduledExecutorService executor,
BlobLibraryCacheManager libraryCacheManager,
RestartStrategyFactory restartStrategyFactory,
Time rpcAskTimeout,
JobManagerMetricGroup jobManagerMetricGroup,
OnCompletionActions jobCompletionActions,
FatalErrorHandler errorHandler,
ClassLoader userCodeLoader) |
Modifier and Type | Method and Description |
---|---|
Time |
StackTraceSampleMessages.TriggerStackTraceSample.delayBetweenSamples() |
Time |
StackTraceSampleMessages.SampleTaskStackTrace.delayBetweenSamples() |
Constructor and Description |
---|
SampleTaskStackTrace(int sampleId,
ExecutionAttemptID executionId,
Time delayBetweenSamples,
int maxStackTraceDepth,
int numRemainingSamples,
List<StackTraceElement[]> currentTraces,
akka.actor.ActorRef sender) |
TriggerStackTraceSample(int sampleId,
ExecutionAttemptID executionId,
int numSamples,
Time delayBetweenSamples,
int maxStackTraceDepth) |
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 |
---|---|
Future<RegistrationResponse> |
ResourceManagerGateway.registerJobManager(UUID resourceManagerLeaderId,
UUID jobMasterLeaderId,
ResourceID jobMasterResourceId,
String jobMasterAddress,
JobID jobID,
Time timeout)
Register a
JobMaster at the resource manager. |
Future<RegistrationResponse> |
ResourceManagerGateway.registerTaskExecutor(UUID resourceManagerLeaderId,
String taskExecutorAddress,
ResourceID resourceID,
SlotReport slotReport,
Time timeout)
Register a
TaskExecutor at the resource manager. |
Future<Acknowledge> |
ResourceManagerGateway.requestSlot(UUID resourceManagerLeaderID,
UUID jobMasterLeaderID,
SlotRequest slotRequest,
Time timeout)
Requests a slot from the resource manager.
|
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 |
---|---|
protected <V> Future<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> Future<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.
|
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.
|
Constructor and Description |
---|
AkkaRpcService(akka.actor.ActorSystem actorSystem,
Time timeout) |
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 |
---|---|
Future<Acknowledge> |
TaskExecutorGateway.cancelTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Cancel the given task.
|
Future<Acknowledge> |
TaskExecutorGateway.requestSlot(SlotID slotId,
JobID jobId,
AllocationID allocationId,
String targetAddress,
UUID resourceManagerLeaderId,
Time timeout)
Requests a slot from the TaskManager
|
Future<Acknowledge> |
TaskExecutorGateway.stopTask(ExecutionAttemptID executionAttemptID,
Time timeout)
Stop the given task.
|
Future<Acknowledge> |
TaskExecutorGateway.submitTask(TaskDeploymentDescriptor tdd,
UUID leaderId,
Time timeout)
Submit a
Task to the TaskExecutor . |
Future<Acknowledge> |
TaskExecutorGateway.updatePartitions(ExecutionAttemptID executionAttemptID,
Iterable<PartitionInfo> partitionInfos,
Time timeout)
Update the task where the given partitions can be found.
|
Constructor and Description |
---|
TaskManagerConfiguration(int numberSlots,
String[] tmpDirectories,
Time timeout,
Time maxRegistrationDuration,
Time initialRegistrationPause,
Time maxRegistrationPause,
Time refusedRegistrationPause,
long cleanupInterval,
Configuration configuration,
boolean exitJvmOnOutOfMemory) |
Constructor and Description |
---|
RpcInputSplitProvider(UUID jobMasterLeaderId,
JobMasterGateway jobMasterGateway,
JobID jobID,
JobVertexID jobVertexID,
ExecutionAttemptID executionAttemptID,
Time timeout) |
RpcResultPartitionConsumableNotifier(UUID jobMasterLeaderId,
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 | Method and Description |
---|---|
Future<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.
|
Modifier and Type | Method and Description |
---|---|
StreamQueryConfig |
StreamQueryConfig.withIdleStateRetentionTime(Time time)
Specifies the time interval for how long idle state, i.e., state which was not updated, will
be retained.
|
StreamQueryConfig |
StreamQueryConfig.withIdleStateRetentionTime(Time minTime,
Time maxTime)
Specifies a minimum and a maximum time interval for how long idle state, i.e., state which
was not updated, will be retained.
|
Copyright © 2014–2018 The Apache Software Foundation. All rights reserved.