Class SingleThreadEventExecutor
- java.lang.Object
-
- java.util.concurrent.AbstractExecutorService
-
- io.netty.util.concurrent.AbstractEventExecutor
-
- io.netty.util.concurrent.AbstractScheduledEventExecutor
-
- io.netty.util.concurrent.SingleThreadEventExecutor
-
- All Implemented Interfaces:
EventExecutor,EventExecutorGroup,OrderedEventExecutor,ThreadAwareExecutor,java.lang.Iterable<EventExecutor>,java.util.concurrent.Executor,java.util.concurrent.ExecutorService,java.util.concurrent.ScheduledExecutorService
- Direct Known Subclasses:
DefaultEventExecutor,SingleThreadEventLoop
public abstract class SingleThreadEventExecutor extends AbstractScheduledEventExecutor implements OrderedEventExecutor
Abstract base class forOrderedEventExecutor's that execute all its submitted tasks in a single thread.
-
-
Nested Class Summary
Nested Classes Modifier and Type Class Description private static classSingleThreadEventExecutor.DefaultThreadPropertiesprotected static interfaceSingleThreadEventExecutor.NonWakeupRunnableDeprecated.overridewakesUpForTask(java.lang.Runnable)to re-create this behaviour-
Nested classes/interfaces inherited from class io.netty.util.concurrent.AbstractEventExecutor
AbstractEventExecutor.LazyRunnable
-
-
Field Summary
Fields Modifier and Type Field Description private static java.util.concurrent.atomic.AtomicLongFieldUpdater<SingleThreadEventExecutor>ACCUMULATED_ACTIVE_TIME_NANOS_UPDATERprivate longaccumulatedActiveTimeNanosprivate booleanaddTaskWakesUpprivate static java.util.concurrent.atomic.AtomicIntegerFieldUpdater<SingleThreadEventExecutor>CONSECUTIVE_BUSY_CYCLES_UPDATERprivate static java.util.concurrent.atomic.AtomicIntegerFieldUpdater<SingleThreadEventExecutor>CONSECUTIVE_IDLE_CYCLES_UPDATERprivate intconsecutiveBusyCyclesTracks the number of consecutive monitor cycles this executor's utilization has been above the scale-up threshold.private intconsecutiveIdleCyclesTracks the number of consecutive monitor cycles this executor's utilization has been below the scale-down threshold.(package private) static intDEFAULT_MAX_PENDING_EXECUTOR_TASKSprivate java.util.concurrent.Executorexecutorprivate longgracefulShutdownQuietPeriodprivate longgracefulShutdownStartTimeprivate longgracefulShutdownTimeoutprivate booleaninterruptedprivate longlastActivityTimeNanosprivate longlastExecutionTimeprivate static InternalLoggerloggerprivate intmaxPendingTasksprivate static java.lang.RunnableNOOP_TASKprivate java.util.concurrent.locks.LockprocessingLockprivate static java.util.concurrent.atomic.AtomicReferenceFieldUpdater<SingleThreadEventExecutor,ThreadProperties>PROPERTIES_UPDATERprivate RejectedExecutionHandlerrejectedExecutionHandlerprivate static longSCHEDULE_PURGE_INTERVALprivate java.util.Set<java.lang.Runnable>shutdownHooksprivate static intST_NOT_STARTEDprivate static intST_SHUTDOWNprivate static intST_SHUTTING_DOWNprivate static intST_STARTEDprivate static intST_SUSPENDEDprivate static intST_SUSPENDINGprivate static intST_TERMINATEDprivate intstateprivate static java.util.concurrent.atomic.AtomicIntegerFieldUpdater<SingleThreadEventExecutor>STATE_UPDATERprivate booleansupportSuspensionprivate java.util.Queue<java.lang.Runnable>taskQueueprivate Promise<?>terminationFutureprivate java.lang.Threadthreadprivate java.util.concurrent.CountDownLatchthreadLockprivate ThreadPropertiesthreadProperties-
Fields inherited from class io.netty.util.concurrent.AbstractScheduledEventExecutor
nextTaskId, scheduledTaskQueue, WAKEUP_TASK
-
Fields inherited from class io.netty.util.concurrent.AbstractEventExecutor
DEFAULT_SHUTDOWN_QUIET_PERIOD, DEFAULT_SHUTDOWN_TIMEOUT
-
-
Constructor Summary
Constructors Modifier Constructor Description protectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp)Create a new instanceprotectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, boolean supportSuspension, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)Create a new instanceprotectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, boolean supportSuspension, java.util.Queue<java.lang.Runnable> taskQueue, RejectedExecutionHandler rejectedHandler)protectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)Create a new instanceprotectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, java.util.Queue<java.lang.Runnable> taskQueue, RejectedExecutionHandler rejectedHandler)protectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.ThreadFactory threadFactory, boolean addTaskWakesUp)Create a new instanceprotectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.ThreadFactory threadFactory, boolean addTaskWakesUp, boolean supportSuspension, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)Create a new instanceprotectedSingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.ThreadFactory threadFactory, boolean addTaskWakesUp, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)Create a new instance
-
Method Summary
All Methods Static Methods Instance Methods Abstract Methods Concrete Methods Deprecated Methods Modifier and Type Method Description voidaddShutdownHook(java.lang.Runnable task)Add aRunnablewhich will be executed on shutdown of this instanceprotected voidaddTask(java.lang.Runnable task)Add a task to the task queue, or throws aRejectedExecutionExceptionif this instance was shutdown before.protected voidafterRunningAllTasks()Invoked before returning fromrunAllTasks()andrunAllTasks(long).booleanawaitTermination(long timeout, java.util.concurrent.TimeUnit unit)protected booleancanSuspend()protected booleancanSuspend(int state)protected voidcleanup()Do nothing, sub-classes may overrideprotected booleanconfirmShutdown()Confirm that the shutdown if the instance should be done now!protected longdeadlineNanos()Returns the absolute point in time (relative toAbstractScheduledEventExecutor.getCurrentTimeNanos()) at which the next closest scheduled task should run.protected longdelayNanos(long currentTimeNanos)Returns the amount of time left until the scheduled task with the closest dead line is executed.private voiddoStartThread()(package private) intdrainTasks()private booleanensureThreadStarted(int oldState)voidexecute(java.lang.Runnable task)private voidexecute(java.lang.Runnable task, boolean immediate)private voidexecute0(java.lang.Runnable task)private booleanexecuteExpiredScheduledTasks()private booleanfetchFromScheduledTaskQueue()protected intgetAndIncrementBusyCycles()Atomically increments the counter for consecutive monitor cycles where utilization was above the scale-up threshold.protected intgetAndIncrementIdleCycles()Atomically increments the counter for consecutive monitor cycles where utilization was below the scale-down threshold.protected longgetAndResetAccumulatedActiveTimeNanos()Returns the accumulated active time since the last call and resets the counter.protected longgetLastActivityTimeNanos()Returns the timestamp of the last known activity (tasks + I/O).protected intgetNumOfRegisteredChannels()Returns the number of registered channels for auto-scaling related decisions.protected booleanhasTasks()booleaninEventLoop(java.lang.Thread thread)Returntrueif the givenThreadis executed in the event loop,falseotherwise.protected voidinterruptThread()Interrupt the current runningThread.<T> java.util.List<java.util.concurrent.Future<T>>invokeAll(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks)<T> java.util.List<java.util.concurrent.Future<T>>invokeAll(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks, long timeout, java.util.concurrent.TimeUnit unit)<T> TinvokeAny(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks)<T> TinvokeAny(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks, long timeout, java.util.concurrent.TimeUnit unit)booleanisShutdown()booleanisShuttingDown()Returnstrueif and only if allEventExecutors managed by thisEventExecutorGroupare being shut down gracefully or was shut down.booleanisSuspended()Returnstrueif theEventExecutoris considered suspended.protected booleanisSuspensionSupported()Returnstrueif thisSingleThreadEventExecutorsupports suspension.booleanisTerminated()voidlazyExecute(java.lang.Runnable task)LikeExecutor.execute(Runnable)but does not guarantee the task will be run until either a non-lazy task is executed or the executor is shut down.private voidlazyExecute0(java.lang.Runnable task)protected java.util.Queue<java.lang.Runnable>newTaskQueue()Deprecated.Please use and overridenewTaskQueue(int).protected java.util.Queue<java.lang.Runnable>newTaskQueue(int maxPendingTasks)Create a newQueuewhich will holds the tasks to execute.(package private) booleanofferTask(java.lang.Runnable task)protected java.lang.RunnablepeekTask()intpendingTasks()Return the number of tasks that are pending for processing.protected java.lang.RunnablepollTask()protected static java.lang.RunnablepollTaskFrom(java.util.Queue<java.lang.Runnable> taskQueue)protected static voidreject()protected voidreject(java.lang.Runnable task)Offers the task to the associatedRejectedExecutionHandler.voidremoveShutdownHook(java.lang.Runnable task)Remove a previous addedRunnableas a shutdown hookprotected booleanremoveTask(java.lang.Runnable task)protected voidreportActiveIoTime(long nanos)Adds the given duration to the total active time for the current measurement window.protected voidresetBusyCycles()Resets the counter for consecutive busy cycles to zero.protected voidresetIdleCycles()Resets the counter for consecutive idle cycles to zero.protected abstract voidrun()Runs the task-processing loop untilconfirmShutdown()returnstrue.protected booleanrunAllTasks()Poll all tasks from the task queue and run them viaRunnable.run()method.protected booleanrunAllTasks(long timeoutNanos)Poll all tasks from the task queue and run them viaRunnable.run()method.protected booleanrunAllTasksFrom(java.util.Queue<java.lang.Runnable> taskQueue)Runs all tasks from the passedtaskQueue.private booleanrunExistingTasksFrom(java.util.Queue<java.lang.Runnable> taskQueue)What ever tasks are present intaskQueuewhen this method is invoked will beRunnable.run().protected booleanrunScheduledAndExecutorTasks(int maxDrainAttempts)Execute all expired scheduled tasks and all current tasks in the executor queue until both queues are empty, ormaxDrainAttemptshas been exceeded.private booleanrunShutdownHooks()(package private) voidscheduleRemoveScheduled(ScheduledFutureTask<?> task)voidshutdown()Deprecated.private voidshutdown0(long quietPeriod, long timeout, int shutdownState)Future<?>shutdownGracefully(long quietPeriod, long timeout, java.util.concurrent.TimeUnit unit)Signals this executor that the caller wants the executor to be shut down.private voidstartThread()protected java.lang.RunnabletakeTask()Take the nextRunnablefrom the task queue and so will block if no task is currently present.Future<?>terminationFuture()Returns theFuturewhich is notified when allEventExecutors managed by thisEventExecutorGrouphave been terminated.ThreadPropertiesthreadProperties()private voidthrowIfInEventLoop(java.lang.String method)booleantrySuspend()Try to suspend thisEventExecutorand returntrueif suspension was successful.protected voidupdateLastExecutionTime()Updates the internal timestamp that tells when a submitted task was executed most recently.protected booleanwakesUpForTask(java.lang.Runnable task)Can be overridden to control which tasks require waking theEventExecutorthread if it is waiting so that they can be run immediately.protected voidwakeup(boolean inEventLoop)-
Methods inherited from class io.netty.util.concurrent.AbstractScheduledEventExecutor
afterScheduledTaskSubmitted, beforeScheduledTaskSubmitted, cancelScheduledTasks, deadlineNanos, deadlineToDelayNanos, defaultCurrentTimeNanos, delayNanos, fetchFromScheduledTaskQueue, getCurrentTimeNanos, hasScheduledTasks, initialNanoTime, nanoTime, nextScheduledTaskDeadlineNanos, nextScheduledTaskNano, peekScheduledTask, pollScheduledTask, pollScheduledTask, removeScheduled, schedule, schedule, scheduleAtFixedRate, scheduledTaskQueue, scheduleFromEventLoop, scheduleWithFixedDelay, ticker, validateScheduled
-
Methods inherited from class io.netty.util.concurrent.AbstractEventExecutor
iterator, newTaskFor, newTaskFor, next, parent, runTask, safeExecute, shutdownGracefully, shutdownNow, submit, submit, submit
-
Methods inherited from class java.lang.Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
-
Methods inherited from interface io.netty.util.concurrent.EventExecutor
inEventLoop, isExecutorThread, newFailedFuture, newProgressivePromise, newPromise, newSucceededFuture, parent
-
Methods inherited from interface io.netty.util.concurrent.EventExecutorGroup
iterator, next, schedule, schedule, scheduleAtFixedRate, scheduleWithFixedDelay, shutdownGracefully, shutdownNow, submit, submit, submit, ticker
-
-
-
-
Field Detail
-
DEFAULT_MAX_PENDING_EXECUTOR_TASKS
static final int DEFAULT_MAX_PENDING_EXECUTOR_TASKS
-
logger
private static final InternalLogger logger
-
ST_NOT_STARTED
private static final int ST_NOT_STARTED
- See Also:
- Constant Field Values
-
ST_SUSPENDING
private static final int ST_SUSPENDING
- See Also:
- Constant Field Values
-
ST_SUSPENDED
private static final int ST_SUSPENDED
- See Also:
- Constant Field Values
-
ST_STARTED
private static final int ST_STARTED
- See Also:
- Constant Field Values
-
ST_SHUTTING_DOWN
private static final int ST_SHUTTING_DOWN
- See Also:
- Constant Field Values
-
ST_SHUTDOWN
private static final int ST_SHUTDOWN
- See Also:
- Constant Field Values
-
ST_TERMINATED
private static final int ST_TERMINATED
- See Also:
- Constant Field Values
-
NOOP_TASK
private static final java.lang.Runnable NOOP_TASK
-
STATE_UPDATER
private static final java.util.concurrent.atomic.AtomicIntegerFieldUpdater<SingleThreadEventExecutor> STATE_UPDATER
-
PROPERTIES_UPDATER
private static final java.util.concurrent.atomic.AtomicReferenceFieldUpdater<SingleThreadEventExecutor,ThreadProperties> PROPERTIES_UPDATER
-
ACCUMULATED_ACTIVE_TIME_NANOS_UPDATER
private static final java.util.concurrent.atomic.AtomicLongFieldUpdater<SingleThreadEventExecutor> ACCUMULATED_ACTIVE_TIME_NANOS_UPDATER
-
CONSECUTIVE_IDLE_CYCLES_UPDATER
private static final java.util.concurrent.atomic.AtomicIntegerFieldUpdater<SingleThreadEventExecutor> CONSECUTIVE_IDLE_CYCLES_UPDATER
-
CONSECUTIVE_BUSY_CYCLES_UPDATER
private static final java.util.concurrent.atomic.AtomicIntegerFieldUpdater<SingleThreadEventExecutor> CONSECUTIVE_BUSY_CYCLES_UPDATER
-
taskQueue
private final java.util.Queue<java.lang.Runnable> taskQueue
-
thread
private volatile java.lang.Thread thread
-
threadProperties
private volatile ThreadProperties threadProperties
-
executor
private final java.util.concurrent.Executor executor
-
interrupted
private volatile boolean interrupted
-
processingLock
private final java.util.concurrent.locks.Lock processingLock
-
threadLock
private final java.util.concurrent.CountDownLatch threadLock
-
shutdownHooks
private final java.util.Set<java.lang.Runnable> shutdownHooks
-
addTaskWakesUp
private final boolean addTaskWakesUp
-
maxPendingTasks
private final int maxPendingTasks
-
rejectedExecutionHandler
private final RejectedExecutionHandler rejectedExecutionHandler
-
supportSuspension
private final boolean supportSuspension
-
accumulatedActiveTimeNanos
private volatile long accumulatedActiveTimeNanos
-
lastActivityTimeNanos
private volatile long lastActivityTimeNanos
-
consecutiveIdleCycles
private volatile int consecutiveIdleCycles
Tracks the number of consecutive monitor cycles this executor's utilization has been below the scale-down threshold.
-
consecutiveBusyCycles
private volatile int consecutiveBusyCycles
Tracks the number of consecutive monitor cycles this executor's utilization has been above the scale-up threshold.
-
lastExecutionTime
private long lastExecutionTime
-
state
private volatile int state
-
gracefulShutdownQuietPeriod
private volatile long gracefulShutdownQuietPeriod
-
gracefulShutdownTimeout
private volatile long gracefulShutdownTimeout
-
gracefulShutdownStartTime
private long gracefulShutdownStartTime
-
terminationFuture
private final Promise<?> terminationFuture
-
SCHEDULE_PURGE_INTERVAL
private static final long SCHEDULE_PURGE_INTERVAL
-
-
Constructor Detail
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.ThreadFactory threadFactory, boolean addTaskWakesUp)
Create a new instance- Parameters:
parent- theEventExecutorGroupwhich is the parent of this instance and belongs to itthreadFactory- theThreadFactorywhich will be used for the usedThreadaddTaskWakesUp-trueif and only if invocation ofaddTask(Runnable)will wake up the executor thread
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.ThreadFactory threadFactory, boolean addTaskWakesUp, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)
Create a new instance- Parameters:
parent- theEventExecutorGroupwhich is the parent of this instance and belongs to itthreadFactory- theThreadFactorywhich will be used for the usedThreadaddTaskWakesUp-trueif and only if invocation ofaddTask(Runnable)will wake up the executor threadmaxPendingTasks- the maximum number of pending tasks before new tasks will be rejected.rejectedHandler- theRejectedExecutionHandlerto use.
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.ThreadFactory threadFactory, boolean addTaskWakesUp, boolean supportSuspension, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)
Create a new instance- Parameters:
parent- theEventExecutorGroupwhich is the parent of this instance and belongs to itthreadFactory- theThreadFactorywhich will be used for the usedThreadaddTaskWakesUp-trueif and only if invocation ofaddTask(Runnable)will wake up the executor threadsupportSuspension-trueif suspension of thisSingleThreadEventExecutoris supported.maxPendingTasks- the maximum number of pending tasks before new tasks will be rejected.rejectedHandler- theRejectedExecutionHandlerto use.
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp)
Create a new instance- Parameters:
parent- theEventExecutorGroupwhich is the parent of this instance and belongs to itexecutor- theExecutorwhich will be used for executingaddTaskWakesUp-trueif and only if invocation ofaddTask(Runnable)will wake up the executor thread
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)
Create a new instance- Parameters:
parent- theEventExecutorGroupwhich is the parent of this instance and belongs to itexecutor- theExecutorwhich will be used for executingaddTaskWakesUp-trueif and only if invocation ofaddTask(Runnable)will wake up the executor threadmaxPendingTasks- the maximum number of pending tasks before new tasks will be rejected.rejectedHandler- theRejectedExecutionHandlerto use.
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, boolean supportSuspension, int maxPendingTasks, RejectedExecutionHandler rejectedHandler)
Create a new instance- Parameters:
parent- theEventExecutorGroupwhich is the parent of this instance and belongs to itexecutor- theExecutorwhich will be used for executingaddTaskWakesUp-trueif and only if invocation ofaddTask(Runnable)will wake up the executor threadsupportSuspension-trueif suspension of thisSingleThreadEventExecutoris supported.maxPendingTasks- the maximum number of pending tasks before new tasks will be rejected.rejectedHandler- theRejectedExecutionHandlerto use.
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, java.util.Queue<java.lang.Runnable> taskQueue, RejectedExecutionHandler rejectedHandler)
-
SingleThreadEventExecutor
protected SingleThreadEventExecutor(EventExecutorGroup parent, java.util.concurrent.Executor executor, boolean addTaskWakesUp, boolean supportSuspension, java.util.Queue<java.lang.Runnable> taskQueue, RejectedExecutionHandler rejectedHandler)
-
-
Method Detail
-
newTaskQueue
@Deprecated protected java.util.Queue<java.lang.Runnable> newTaskQueue()
Deprecated.Please use and overridenewTaskQueue(int).
-
newTaskQueue
protected java.util.Queue<java.lang.Runnable> newTaskQueue(int maxPendingTasks)
Create a newQueuewhich will holds the tasks to execute. This default implementation will return aLinkedBlockingQueuebut if your sub-class ofSingleThreadEventExecutorwill not do any blocking calls on the thisQueueit may make sense to@Overridethis and return some more performant implementation that does not support blocking operations at all.
-
interruptThread
protected void interruptThread()
Interrupt the current runningThread.
-
pollTask
protected java.lang.Runnable pollTask()
- See Also:
Queue.poll()
-
pollTaskFrom
protected static java.lang.Runnable pollTaskFrom(java.util.Queue<java.lang.Runnable> taskQueue)
-
takeTask
protected java.lang.Runnable takeTask()
Take the nextRunnablefrom the task queue and so will block if no task is currently present.Be aware that this method will throw an
UnsupportedOperationExceptionif the task queue, which was created vianewTaskQueue(), does not implementBlockingQueue.- Returns:
nullif the executor thread has been interrupted or waken up.
-
fetchFromScheduledTaskQueue
private boolean fetchFromScheduledTaskQueue()
-
executeExpiredScheduledTasks
private boolean executeExpiredScheduledTasks()
- Returns:
trueif at least one scheduled task was executed.
-
peekTask
protected java.lang.Runnable peekTask()
- See Also:
Queue.peek()
-
hasTasks
protected boolean hasTasks()
- See Also:
Collection.isEmpty()
-
pendingTasks
public int pendingTasks()
Return the number of tasks that are pending for processing.
-
addTask
protected void addTask(java.lang.Runnable task)
Add a task to the task queue, or throws aRejectedExecutionExceptionif this instance was shutdown before.
-
offerTask
final boolean offerTask(java.lang.Runnable task)
-
removeTask
protected boolean removeTask(java.lang.Runnable task)
- See Also:
Collection.remove(Object)
-
runAllTasks
protected boolean runAllTasks()
Poll all tasks from the task queue and run them viaRunnable.run()method.- Returns:
trueif and only if at least one task was run
-
runScheduledAndExecutorTasks
protected final boolean runScheduledAndExecutorTasks(int maxDrainAttempts)
Execute all expired scheduled tasks and all current tasks in the executor queue until both queues are empty, ormaxDrainAttemptshas been exceeded.- Parameters:
maxDrainAttempts- The maximum amount of times this method attempts to drain from queues. This is to prevent continuous task execution and scheduling from preventing the EventExecutor thread to make progress and return to the selector mechanism to process inbound I/O events.- Returns:
trueif at least one task was run.
-
runAllTasksFrom
protected final boolean runAllTasksFrom(java.util.Queue<java.lang.Runnable> taskQueue)
Runs all tasks from the passedtaskQueue.- Parameters:
taskQueue- To poll and execute all tasks.- Returns:
trueif at least one task was executed.
-
runExistingTasksFrom
private boolean runExistingTasksFrom(java.util.Queue<java.lang.Runnable> taskQueue)
What ever tasks are present intaskQueuewhen this method is invoked will beRunnable.run().- Parameters:
taskQueue- the task queue to drain.- Returns:
trueif at leastRunnable.run()was called.
-
runAllTasks
protected boolean runAllTasks(long timeoutNanos)
Poll all tasks from the task queue and run them viaRunnable.run()method. This method stops running the tasks in the task queue and returns if it ran longer thantimeoutNanos.
-
afterRunningAllTasks
protected void afterRunningAllTasks()
Invoked before returning fromrunAllTasks()andrunAllTasks(long).
-
delayNanos
protected long delayNanos(long currentTimeNanos)
Returns the amount of time left until the scheduled task with the closest dead line is executed.
-
deadlineNanos
protected long deadlineNanos()
Returns the absolute point in time (relative toAbstractScheduledEventExecutor.getCurrentTimeNanos()) at which the next closest scheduled task should run.
-
updateLastExecutionTime
protected void updateLastExecutionTime()
Updates the internal timestamp that tells when a submitted task was executed most recently.runAllTasks()andrunAllTasks(long)updates this timestamp automatically, and thus there's usually no need to call this method. However, if you take the tasks manually usingtakeTask()orpollTask(), you have to call this method at the end of task execution loop for accurate quiet period checks.
-
getNumOfRegisteredChannels
protected int getNumOfRegisteredChannels()
Returns the number of registered channels for auto-scaling related decisions. This is intended to be used byMultithreadEventExecutorGroupfor dynamic scaling.- Returns:
- The number of registered channels, or
-1if not applicable.
-
reportActiveIoTime
protected void reportActiveIoTime(long nanos)
Adds the given duration to the total active time for the current measurement window.Note: This method is not thread-safe and must only be called from the
event loop thread.- Parameters:
nanos- The active time in nanoseconds to add.
-
getAndResetAccumulatedActiveTimeNanos
protected long getAndResetAccumulatedActiveTimeNanos()
Returns the accumulated active time since the last call and resets the counter.
-
getLastActivityTimeNanos
protected long getLastActivityTimeNanos()
Returns the timestamp of the last known activity (tasks + I/O).
-
getAndIncrementIdleCycles
protected int getAndIncrementIdleCycles()
Atomically increments the counter for consecutive monitor cycles where utilization was below the scale-down threshold. This is used by the auto-scaling monitor to track sustained idleness.- Returns:
- The number of consecutive idle cycles before the increment.
-
resetIdleCycles
protected void resetIdleCycles()
Resets the counter for consecutive idle cycles to zero. This is typically called when the executor's utilization is no longer considered idle, breaking the streak.
-
getAndIncrementBusyCycles
protected int getAndIncrementBusyCycles()
Atomically increments the counter for consecutive monitor cycles where utilization was above the scale-up threshold. This is used by the auto-scaling monitor to track a sustained high load.- Returns:
- The number of consecutive busy cycles before the increment.
-
resetBusyCycles
protected void resetBusyCycles()
Resets the counter for consecutive busy cycles to zero. This is typically called when the executor's utilization is no longer considered busy, breaking the streak.
-
isSuspensionSupported
protected boolean isSuspensionSupported()
Returnstrueif thisSingleThreadEventExecutorsupports suspension.
-
run
protected abstract void run()
Runs the task-processing loop untilconfirmShutdown()returnstrue.Implementations must not let a
Throwablethrown by a task escape this method: any uncaughtThrowableterminates the executor (logged atWARNand surfaced viaterminationFuture()), at which point everyChannelregistered with this executor stops processing I/O and new task submissions are rejected. The supplied helpers -runAllTasks(),runAllTasks(long), andAbstractEventExecutor.safeExecute(Runnable)- catchThrowablefor you; custom loops built onpollTask()ortakeTask()are responsible for wrapping each task invocation accordingly.
-
cleanup
protected void cleanup()
Do nothing, sub-classes may override
-
wakeup
protected void wakeup(boolean inEventLoop)
-
inEventLoop
public boolean inEventLoop(java.lang.Thread thread)
Description copied from interface:EventExecutorReturntrueif the givenThreadis executed in the event loop,falseotherwise.- Specified by:
inEventLoopin interfaceEventExecutor
-
addShutdownHook
public void addShutdownHook(java.lang.Runnable task)
Add aRunnablewhich will be executed on shutdown of this instance
-
removeShutdownHook
public void removeShutdownHook(java.lang.Runnable task)
Remove a previous addedRunnableas a shutdown hook
-
runShutdownHooks
private boolean runShutdownHooks()
-
shutdown0
private void shutdown0(long quietPeriod, long timeout, int shutdownState)
-
shutdownGracefully
public Future<?> shutdownGracefully(long quietPeriod, long timeout, java.util.concurrent.TimeUnit unit)
Description copied from interface:EventExecutorGroupSignals this executor that the caller wants the executor to be shut down. Once this method is called,EventExecutorGroup.isShuttingDown()starts to returntrue, and the executor prepares to shut itself down. UnlikeEventExecutorGroup.shutdown(), graceful shutdown ensures that no tasks are submitted for 'the quiet period' (usually a couple seconds) before it shuts itself down. If a task is submitted during the quiet period, it is guaranteed to be accepted and the quiet period will start over.- Specified by:
shutdownGracefullyin interfaceEventExecutorGroup- Parameters:
quietPeriod- the quiet period as described in the documentationtimeout- the maximum amount of time to wait until the executor is EventExecutorGroup.shutdown() regardless if a task was submitted during the quiet periodunit- the unit ofquietPeriodandtimeout- Returns:
- the
EventExecutorGroup.terminationFuture()
-
terminationFuture
public Future<?> terminationFuture()
Description copied from interface:EventExecutorGroupReturns theFuturewhich is notified when allEventExecutors managed by thisEventExecutorGrouphave been terminated.- Specified by:
terminationFuturein interfaceEventExecutorGroup
-
shutdown
@Deprecated public void shutdown()
Deprecated.- Specified by:
shutdownin interfaceEventExecutorGroup- Specified by:
shutdownin interfacejava.util.concurrent.ExecutorService- Specified by:
shutdownin classAbstractEventExecutor
-
isShuttingDown
public boolean isShuttingDown()
Description copied from interface:EventExecutorGroupReturnstrueif and only if allEventExecutors managed by thisEventExecutorGroupare being shut down gracefully or was shut down.- Specified by:
isShuttingDownin interfaceEventExecutorGroup
-
isShutdown
public boolean isShutdown()
- Specified by:
isShutdownin interfacejava.util.concurrent.ExecutorService
-
isTerminated
public boolean isTerminated()
- Specified by:
isTerminatedin interfacejava.util.concurrent.ExecutorService
-
isSuspended
public boolean isSuspended()
Description copied from interface:EventExecutorReturnstrueif theEventExecutoris considered suspended.- Specified by:
isSuspendedin interfaceEventExecutor- Returns:
trueif suspended,falseotherwise.
-
trySuspend
public boolean trySuspend()
Description copied from interface:EventExecutorTry to suspend thisEventExecutorand returntrueif suspension was successful. Suspending anEventExecutorwill allow it to free up resources, like for example aThreadthat is backing theEventExecutor. Once anEventExecutorwas suspended it will be started again by submitting work to it via one of the following methods:Executor.execute(Runnable)EventExecutorGroup.schedule(Runnable, long, TimeUnit)EventExecutorGroup.schedule(Callable, long, TimeUnit)EventExecutorGroup.scheduleAtFixedRate(Runnable, long, long, TimeUnit)EventExecutorGroup.scheduleWithFixedDelay(Runnable, long, long, TimeUnit)
trueit might take some time for theEventExecutorto fully suspend itself.- Specified by:
trySuspendin interfaceEventExecutor- Returns:
trueif suspension was successful, otherwisefalse.
-
canSuspend
protected boolean canSuspend()
- Returns:
- if suspension is possible at the moment.
-
canSuspend
protected boolean canSuspend(int state)
Returnstrueif thisSingleThreadEventExecutorcan be suspended at the moment,falseotherwise. Subclasses might override this method to add extra checks.- Parameters:
state- the current internal state of theSingleThreadEventExecutor.- Returns:
- if suspension is possible at the moment.
-
confirmShutdown
protected boolean confirmShutdown()
Confirm that the shutdown if the instance should be done now!
-
awaitTermination
public boolean awaitTermination(long timeout, java.util.concurrent.TimeUnit unit) throws java.lang.InterruptedException- Specified by:
awaitTerminationin interfacejava.util.concurrent.ExecutorService- Throws:
java.lang.InterruptedException
-
execute
public void execute(java.lang.Runnable task)
- Specified by:
executein interfacejava.util.concurrent.Executor
-
lazyExecute
public void lazyExecute(java.lang.Runnable task)
Description copied from class:AbstractEventExecutorLikeExecutor.execute(Runnable)but does not guarantee the task will be run until either a non-lazy task is executed or the executor is shut down.The default implementation just delegates to
Executor.execute(Runnable).- Overrides:
lazyExecutein classAbstractEventExecutor
-
execute0
private void execute0(java.lang.Runnable task)
-
lazyExecute0
private void lazyExecute0(java.lang.Runnable task)
-
scheduleRemoveScheduled
void scheduleRemoveScheduled(ScheduledFutureTask<?> task)
- Overrides:
scheduleRemoveScheduledin classAbstractScheduledEventExecutor
-
execute
private void execute(java.lang.Runnable task, boolean immediate)
-
invokeAny
public <T> T invokeAny(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks) throws java.lang.InterruptedException, java.util.concurrent.ExecutionException- Specified by:
invokeAnyin interfacejava.util.concurrent.ExecutorService- Overrides:
invokeAnyin classjava.util.concurrent.AbstractExecutorService- Throws:
java.lang.InterruptedExceptionjava.util.concurrent.ExecutionException
-
invokeAny
public <T> T invokeAny(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks, long timeout, java.util.concurrent.TimeUnit unit) throws java.lang.InterruptedException, java.util.concurrent.ExecutionException, java.util.concurrent.TimeoutException- Specified by:
invokeAnyin interfacejava.util.concurrent.ExecutorService- Overrides:
invokeAnyin classjava.util.concurrent.AbstractExecutorService- Throws:
java.lang.InterruptedExceptionjava.util.concurrent.ExecutionExceptionjava.util.concurrent.TimeoutException
-
invokeAll
public <T> java.util.List<java.util.concurrent.Future<T>> invokeAll(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks) throws java.lang.InterruptedException- Specified by:
invokeAllin interfacejava.util.concurrent.ExecutorService- Overrides:
invokeAllin classjava.util.concurrent.AbstractExecutorService- Throws:
java.lang.InterruptedException
-
invokeAll
public <T> java.util.List<java.util.concurrent.Future<T>> invokeAll(java.util.Collection<? extends java.util.concurrent.Callable<T>> tasks, long timeout, java.util.concurrent.TimeUnit unit) throws java.lang.InterruptedException- Specified by:
invokeAllin interfacejava.util.concurrent.ExecutorService- Overrides:
invokeAllin classjava.util.concurrent.AbstractExecutorService- Throws:
java.lang.InterruptedException
-
throwIfInEventLoop
private void throwIfInEventLoop(java.lang.String method)
-
threadProperties
public final ThreadProperties threadProperties()
Returns theThreadPropertiesof theThreadthat powers theSingleThreadEventExecutor. If theSingleThreadEventExecutoris not started yet, this operation will start it and block until it is fully started.
-
wakesUpForTask
protected boolean wakesUpForTask(java.lang.Runnable task)
Can be overridden to control which tasks require waking theEventExecutorthread if it is waiting so that they can be run immediately.
-
reject
protected static void reject()
-
reject
protected final void reject(java.lang.Runnable task)
Offers the task to the associatedRejectedExecutionHandler.- Parameters:
task- to reject.
-
startThread
private void startThread()
-
ensureThreadStarted
private boolean ensureThreadStarted(int oldState)
-
doStartThread
private void doStartThread()
-
drainTasks
final int drainTasks()
-
-