Class SingleThreadEventExecutor

All Implemented Interfaces:
EventExecutor, EventExecutorGroup, OrderedEventExecutor, ThreadAwareExecutor, Iterable<EventExecutor>, Executor, ExecutorService, ScheduledExecutorService
Direct Known Subclasses:
DefaultEventExecutor, SingleThreadEventLoop

public abstract class SingleThreadEventExecutor extends AbstractScheduledEventExecutor implements OrderedEventExecutor
Abstract base class for OrderedEventExecutor's that execute all its submitted tasks in a single thread.
  • Field Details

    • 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:
    • ST_SUSPENDING

      private static final int ST_SUSPENDING
      See Also:
    • ST_SUSPENDED

      private static final int ST_SUSPENDED
      See Also:
    • ST_STARTED

      private static final int ST_STARTED
      See Also:
    • ST_SHUTTING_DOWN

      private static final int ST_SHUTTING_DOWN
      See Also:
    • ST_SHUTDOWN

      private static final int ST_SHUTDOWN
      See Also:
    • ST_TERMINATED

      private static final int ST_TERMINATED
      See Also:
    • NOOP_TASK

      private static final Runnable NOOP_TASK
    • STATE_UPDATER

      private static final AtomicIntegerFieldUpdater<SingleThreadEventExecutor> STATE_UPDATER
    • PROPERTIES_UPDATER

      private static final AtomicReferenceFieldUpdater<SingleThreadEventExecutor, ThreadProperties> PROPERTIES_UPDATER
    • ACCUMULATED_ACTIVE_TIME_NANOS_UPDATER

      private static final AtomicLongFieldUpdater<SingleThreadEventExecutor> ACCUMULATED_ACTIVE_TIME_NANOS_UPDATER
    • CONSECUTIVE_IDLE_CYCLES_UPDATER

      private static final AtomicIntegerFieldUpdater<SingleThreadEventExecutor> CONSECUTIVE_IDLE_CYCLES_UPDATER
    • CONSECUTIVE_BUSY_CYCLES_UPDATER

      private static final AtomicIntegerFieldUpdater<SingleThreadEventExecutor> CONSECUTIVE_BUSY_CYCLES_UPDATER
    • taskQueue

      private final Queue<Runnable> taskQueue
    • thread

      private volatile Thread thread
    • threadProperties

      private volatile ThreadProperties threadProperties
    • executor

      private final Executor executor
    • interrupted

      private volatile boolean interrupted
    • processingLock

      private final Lock processingLock
    • threadLock

      private final CountDownLatch threadLock
    • shutdownHooks

      private final Set<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 Details

  • Method Details

    • newTaskQueue

      @Deprecated protected Queue<Runnable> newTaskQueue()
      Deprecated.
      Please use and override newTaskQueue(int).
    • newTaskQueue

      protected Queue<Runnable> newTaskQueue(int maxPendingTasks)
      Create a new Queue which will holds the tasks to execute. This default implementation will return a LinkedBlockingQueue but if your sub-class of SingleThreadEventExecutor will not do any blocking calls on the this Queue it may make sense to @Override this and return some more performant implementation that does not support blocking operations at all.
    • interruptThread

      protected void interruptThread()
      Interrupt the current running Thread.
    • pollTask

      protected Runnable pollTask()
      See Also:
    • pollTaskFrom

      protected static Runnable pollTaskFrom(Queue<Runnable> taskQueue)
    • takeTask

      protected Runnable takeTask()
      Take the next Runnable from the task queue and so will block if no task is currently present.

      Be aware that this method will throw an UnsupportedOperationException if the task queue, which was created via newTaskQueue(), does not implement BlockingQueue.

      Returns:
      null if the executor thread has been interrupted or waken up.
    • fetchFromScheduledTaskQueue

      private boolean fetchFromScheduledTaskQueue()
    • executeExpiredScheduledTasks

      private boolean executeExpiredScheduledTasks()
      Returns:
      true if at least one scheduled task was executed.
    • peekTask

      protected Runnable peekTask()
      See Also:
    • hasTasks

      protected boolean hasTasks()
      See Also:
    • pendingTasks

      public int pendingTasks()
      Return the number of tasks that are pending for processing.
    • addTask

      protected void addTask(Runnable task)
      Add a task to the task queue, or throws a RejectedExecutionException if this instance was shutdown before.
    • offerTask

      final boolean offerTask(Runnable task)
    • removeTask

      protected boolean removeTask(Runnable task)
      See Also:
    • runAllTasks

      protected boolean runAllTasks()
      Poll all tasks from the task queue and run them via Runnable.run() method.
      Returns:
      true if 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, or maxDrainAttempts has 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:
      true if at least one task was run.
    • runAllTasksFrom

      protected final boolean runAllTasksFrom(Queue<Runnable> taskQueue)
      Runs all tasks from the passed taskQueue.
      Parameters:
      taskQueue - To poll and execute all tasks.
      Returns:
      true if at least one task was executed.
    • runExistingTasksFrom

      private boolean runExistingTasksFrom(Queue<Runnable> taskQueue)
      What ever tasks are present in taskQueue when this method is invoked will be Runnable.run().
      Parameters:
      taskQueue - the task queue to drain.
      Returns:
      true if at least Runnable.run() was called.
    • runAllTasks

      protected boolean runAllTasks(long timeoutNanos)
      Poll all tasks from the task queue and run them via Runnable.run() method. This method stops running the tasks in the task queue and returns if it ran longer than timeoutNanos.
    • afterRunningAllTasks

      protected void afterRunningAllTasks()
      Invoked before returning from runAllTasks() and runAllTasks(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 to AbstractScheduledEventExecutor.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() and runAllTasks(long) updates this timestamp automatically, and thus there's usually no need to call this method. However, if you take the tasks manually using takeTask() or pollTask(), 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 by MultithreadEventExecutorGroup for dynamic scaling.
      Returns:
      The number of registered channels, or -1 if 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()
      Returns true if this SingleThreadEventExecutor supports suspension.
    • run

      protected abstract void run()
      Runs the task-processing loop until confirmShutdown() returns true.

      Implementations must not let a Throwable thrown by a task escape this method: any uncaught Throwable terminates the executor (logged at WARN and surfaced via terminationFuture()), at which point every Channel registered with this executor stops processing I/O and new task submissions are rejected. The supplied helpers - runAllTasks(), runAllTasks(long), and AbstractEventExecutor.safeExecute(Runnable) - catch Throwable for you; custom loops built on pollTask() or takeTask() 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(Thread thread)
      Description copied from interface: EventExecutor
      Return true if the given Thread is executed in the event loop, false otherwise.
      Specified by:
      inEventLoop in interface EventExecutor
    • addShutdownHook

      public void addShutdownHook(Runnable task)
      Add a Runnable which will be executed on shutdown of this instance
    • removeShutdownHook

      public void removeShutdownHook(Runnable task)
      Remove a previous added Runnable as 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, TimeUnit unit)
      Description copied from interface: EventExecutorGroup
      Signals this executor that the caller wants the executor to be shut down. Once this method is called, EventExecutorGroup.isShuttingDown() starts to return true, and the executor prepares to shut itself down. Unlike EventExecutorGroup.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:
      shutdownGracefully in interface EventExecutorGroup
      Parameters:
      quietPeriod - the quiet period as described in the documentation
      timeout - the maximum amount of time to wait until the executor is EventExecutorGroup.shutdown() regardless if a task was submitted during the quiet period
      unit - the unit of quietPeriod and timeout
      Returns:
      the EventExecutorGroup.terminationFuture()
    • terminationFuture

      public Future<?> terminationFuture()
      Description copied from interface: EventExecutorGroup
      Returns the Future which is notified when all EventExecutors managed by this EventExecutorGroup have been terminated.
      Specified by:
      terminationFuture in interface EventExecutorGroup
    • shutdown

      @Deprecated public void shutdown()
      Deprecated.
      Specified by:
      shutdown in interface EventExecutorGroup
      Specified by:
      shutdown in interface ExecutorService
      Specified by:
      shutdown in class AbstractEventExecutor
    • isShuttingDown

      public boolean isShuttingDown()
      Description copied from interface: EventExecutorGroup
      Returns true if and only if all EventExecutors managed by this EventExecutorGroup are being shut down gracefully or was shut down.
      Specified by:
      isShuttingDown in interface EventExecutorGroup
    • isShutdown

      public boolean isShutdown()
      Specified by:
      isShutdown in interface ExecutorService
    • isTerminated

      public boolean isTerminated()
      Specified by:
      isTerminated in interface ExecutorService
    • isSuspended

      public boolean isSuspended()
      Description copied from interface: EventExecutor
      Returns true if the EventExecutor is considered suspended.
      Specified by:
      isSuspended in interface EventExecutor
      Returns:
      true if suspended, false otherwise.
    • trySuspend

      public boolean trySuspend()
      Description copied from interface: EventExecutor
      Try to suspend this EventExecutor and return true if suspension was successful. Suspending an EventExecutor will allow it to free up resources, like for example a Thread that is backing the EventExecutor. Once an EventExecutor was suspended it will be started again by submitting work to it via one of the following methods: Even if this method returns true it might take some time for the EventExecutor to fully suspend itself.
      Specified by:
      trySuspend in interface EventExecutor
      Returns:
      true if suspension was successful, otherwise false.
    • canSuspend

      protected boolean canSuspend()
      Returns true if this SingleThreadEventExecutor can be suspended at the moment, false otherwise.
      Returns:
      if suspension is possible at the moment.
    • canSuspend

      protected boolean canSuspend(int state)
      Returns true if this SingleThreadEventExecutor can be suspended at the moment, false otherwise. Subclasses might override this method to add extra checks.
      Parameters:
      state - the current internal state of the SingleThreadEventExecutor.
      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, TimeUnit unit) throws InterruptedException
      Specified by:
      awaitTermination in interface ExecutorService
      Throws:
      InterruptedException
    • execute

      public void execute(Runnable task)
      Specified by:
      execute in interface Executor
    • lazyExecute

      public void lazyExecute(Runnable task)
      Description copied from class: AbstractEventExecutor
      Like Executor.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:
      lazyExecute in class AbstractEventExecutor
    • execute0

      private void execute0(Runnable task)
    • lazyExecute0

      private void lazyExecute0(Runnable task)
    • scheduleRemoveScheduled

      void scheduleRemoveScheduled(ScheduledFutureTask<?> task)
      Overrides:
      scheduleRemoveScheduled in class AbstractScheduledEventExecutor
    • execute

      private void execute(Runnable task, boolean immediate)
    • invokeAny

      public <T> T invokeAny(Collection<? extends Callable<T>> tasks) throws InterruptedException, ExecutionException
      Specified by:
      invokeAny in interface ExecutorService
      Overrides:
      invokeAny in class AbstractExecutorService
      Throws:
      InterruptedException
      ExecutionException
    • invokeAny

      public <T> T invokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException, ExecutionException, TimeoutException
      Specified by:
      invokeAny in interface ExecutorService
      Overrides:
      invokeAny in class AbstractExecutorService
      Throws:
      InterruptedException
      ExecutionException
      TimeoutException
    • invokeAll

      public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks) throws InterruptedException
      Specified by:
      invokeAll in interface ExecutorService
      Overrides:
      invokeAll in class AbstractExecutorService
      Throws:
      InterruptedException
    • invokeAll

      public <T> List<Future<T>> invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) throws InterruptedException
      Specified by:
      invokeAll in interface ExecutorService
      Overrides:
      invokeAll in class AbstractExecutorService
      Throws:
      InterruptedException
    • throwIfInEventLoop

      private void throwIfInEventLoop(String method)
    • threadProperties

      public final ThreadProperties threadProperties()
      Returns the ThreadProperties of the Thread that powers the SingleThreadEventExecutor. If the SingleThreadEventExecutor is not started yet, this operation will start it and block until it is fully started.
    • wakesUpForTask

      protected boolean wakesUpForTask(Runnable task)
      Can be overridden to control which tasks require waking the EventExecutor thread if it is waiting so that they can be run immediately.
    • reject

      protected static void reject()
    • reject

      protected final void reject(Runnable task)
      Offers the task to the associated RejectedExecutionHandler.
      Parameters:
      task - to reject.
    • startThread

      private void startThread()
    • ensureThreadStarted

      private boolean ensureThreadStarted(int oldState)
    • doStartThread

      private void doStartThread()
    • drainTasks

      final int drainTasks()