Class SingleThreadIoEventLoop
- All Implemented Interfaces:
EventLoop, EventLoopGroup, IoEventLoop, IoEventLoopGroup, EventExecutor, EventExecutorGroup, OrderedEventExecutor, ThreadAwareExecutor, Iterable<EventExecutor>, Executor, ExecutorService, ScheduledExecutorService
- Direct Known Subclasses:
EpollEventLoop, NioEventLoop
IoEventLoop implementation that execute all its submitted tasks in a single thread using the provided
IoHandler.-
Nested Class Summary
Nested classes/interfaces inherited from class SingleThreadEventLoop
SingleThreadEventLoop.ChannelsReadOnlyIterator<T>Modifier and TypeClassDescriptionprotected static final classNested classes/interfaces inherited from class SingleThreadEventExecutor
SingleThreadEventExecutor.NonWakeupRunnableModifier and TypeClassDescriptionprotected static interfaceDeprecated.Nested classes/interfaces inherited from class AbstractEventExecutor
AbstractEventExecutor.LazyRunnableModifier and TypeClassDescriptionstatic interfaceDeprecated.overrideSingleThreadEventExecutor.wakesUpForTask(Runnable)to re-create this behaviour -
Field Summary
Fields inherited from class SingleThreadEventLoop
DEFAULT_MAX_PENDING_TASKS -
Constructor Summary
ConstructorsModifierConstructorDescriptionSingleThreadIoEventLoop(IoEventLoopGroup parent, Executor executor, IoHandlerFactory ioHandlerFactory) Creates a new instanceSingleThreadIoEventLoop(IoEventLoopGroup parent, Executor executor, IoHandlerFactory ioHandlerFactory, int maxPendingTasks, RejectedExecutionHandler rejectedExecutionHandler, long maxTaskProcessingQuantumMs) Creates a new instanceprotectedSingleThreadIoEventLoop(IoEventLoopGroup parent, Executor executor, IoHandlerFactory ioHandlerFactory, Queue<Runnable> taskQueue, Queue<Runnable> tailTaskQueue, RejectedExecutionHandler rejectedExecutionHandler) Creates a new instanceSingleThreadIoEventLoop(IoEventLoopGroup parent, ThreadFactory threadFactory, IoHandlerFactory ioHandlerFactory) Creates a new instanceSingleThreadIoEventLoop(IoEventLoopGroup parent, ThreadFactory threadFactory, IoHandlerFactory ioHandlerFactory, int maxPendingTasks, RejectedExecutionHandler rejectedExecutionHandler, long maxTaskProcessingQuantumMs) Creates a new instance -
Method Summary
Modifier and TypeMethodDescriptionprotected booleancanSuspend(int state) protected final voidcleanup()Do nothing, sub-classes may overrideprotected intReturns the number of registered channels for auto-scaling related decisions.protected final IoHandlerbooleanisCompatible(Class<? extends IoHandle> handleType) Returnstrueif the given type is compatible with thisIoEventLoopGroupand so can be registered to the containedIoEventLoops,falseotherwise.booleannewTaskQueue(int maxPendingTasks) Create a newQueuewhich will holds the tasks to execute.newTaskQueue0(int maxPendingTasks) next()Returns one of theEventExecutors managed by thisEventExecutorGroup.final Future<IoRegistration> protected voidrun()Runs the task-processing loop untilSingleThreadEventExecutor.confirmShutdown()returnstrue.protected intrunIo()Called when IO will be processed for all theIoHandles on thisSingleThreadIoEventLoop.protected final voidwakeup(boolean inEventLoop) Methods inherited from class SingleThreadEventLoop
afterRunningAllTasks, executeAfterEventLoopIteration, hasTasks, parent, pendingTasks, register, register, register, registeredChannels, registeredChannelsIteratorModifier and TypeMethodDescriptionprotected voidInvoked before returning fromSingleThreadEventExecutor.runAllTasks()andSingleThreadEventExecutor.runAllTasks(long).final voidAdds a task to be run once at the end of next (or current)eventloopiteration.protected booleanhasTasks()parent()Return theEventExecutorGroupwhich is the parent of thisEventExecutor,intReturn the number of tasks that are pending for processing.register(ChannelPromise promise) register(Channel channel, ChannelPromise promise) Deprecated.intMethods inherited from class SingleThreadEventExecutor
addShutdownHook, addTask, awaitTermination, canSuspend, confirmShutdown, deadlineNanos, delayNanos, execute, getAndIncrementBusyCycles, getAndIncrementIdleCycles, getAndResetAccumulatedActiveTimeNanos, getLastActivityTimeNanos, inEventLoop, interruptThread, invokeAll, invokeAll, invokeAny, invokeAny, isShutdown, isShuttingDown, isSuspended, isSuspensionSupported, isTerminated, lazyExecute, newTaskQueue, peekTask, pollTask, pollTaskFrom, reject, reject, removeShutdownHook, removeTask, reportActiveIoTime, resetBusyCycles, resetIdleCycles, runAllTasks, runAllTasks, runAllTasksFrom, runScheduledAndExecutorTasks, shutdown, shutdownGracefully, takeTask, terminationFuture, threadProperties, trySuspend, updateLastExecutionTime, wakesUpForTaskModifier and TypeMethodDescriptionvoidaddShutdownHook(Runnable task) Add aRunnablewhich will be executed on shutdown of this instanceprotected voidAdd a task to the task queue, or throws aRejectedExecutionExceptionif this instance was shutdown before.booleanawaitTermination(long timeout, TimeUnit unit) protected booleanprotected booleanConfirm that the shutdown if the instance should be done now!protected longReturns 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.voidprotected intAtomically increments the counter for consecutive monitor cycles where utilization was above the scale-up threshold.protected intAtomically increments the counter for consecutive monitor cycles where utilization was below the scale-down threshold.protected longReturns the accumulated active time since the last call and resets the counter.protected longReturns the timestamp of the last known activity (tasks + I/O).booleaninEventLoop(Thread thread) protected voidInterrupt the current runningThread.invokeAll(Collection<? extends Callable<T>> tasks) invokeAll(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) <T> TinvokeAny(Collection<? extends Callable<T>> tasks) <T> TinvokeAny(Collection<? extends Callable<T>> tasks, long timeout, TimeUnit unit) booleanbooleanReturnstrueif and only if allEventExecutors managed by thisEventExecutorGroupare being shut down gracefully or was shut down.booleanReturnstrueif theEventExecutoris considered suspended.protected booleanReturnstrueif thisSingleThreadEventExecutorsupports suspension.booleanvoidlazyExecute(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.Deprecated.Please use and overrideSingleThreadEventExecutor.newTaskQueue(int).protected RunnablepeekTask()protected RunnablepollTask()protected static RunnablepollTaskFrom(Queue<Runnable> taskQueue) protected static voidreject()protected final voidOffers the task to the associatedRejectedExecutionHandler.voidremoveShutdownHook(Runnable task) Remove a previous addedRunnableas a shutdown hookprotected booleanremoveTask(Runnable task) protected voidreportActiveIoTime(long nanos) Adds the given duration to the total active time for the current measurement window.protected voidResets the counter for consecutive busy cycles to zero.protected voidResets the counter for consecutive idle cycles to zero.protected booleanPoll 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 final booleanrunAllTasksFrom(Queue<Runnable> taskQueue) Runs all tasks from the passedtaskQueue.protected final booleanrunScheduledAndExecutorTasks(int maxDrainAttempts) Execute all expired scheduled tasks and all current tasks in the executor queue until both queues are empty, ormaxDrainAttemptshas been exceeded.voidshutdown()Deprecated.Future<?> shutdownGracefully(long quietPeriod, long timeout, TimeUnit unit) Signals this executor that the caller wants the executor to be shut down.protected RunnabletakeTask()Take the nextRunnablefrom the task queue and so will block if no task is currently present.Future<?> Returns theFuturewhich is notified when allEventExecutors managed by thisEventExecutorGrouphave been terminated.final ThreadPropertiesbooleanTry to suspend thisEventExecutorand returntrueif suspension was successful.protected voidUpdates the internal timestamp that tells when a submitted task was executed most recently.protected booleanwakesUpForTask(Runnable task) Can be overridden to control which tasks require waking theEventExecutorthread if it is waiting so that they can be run immediately.Methods inherited from class AbstractScheduledEventExecutor
afterScheduledTaskSubmitted, beforeScheduledTaskSubmitted, cancelScheduledTasks, deadlineToDelayNanos, delayNanos, fetchFromScheduledTaskQueue, getCurrentTimeNanos, hasScheduledTasks, initialNanoTime, nanoTime, nextScheduledTaskDeadlineNanos, nextScheduledTaskNano, pollScheduledTask, pollScheduledTask, schedule, schedule, scheduleAtFixedRate, scheduleWithFixedDelay, ticker, validateScheduledModifier and TypeMethodDescriptionprotected booleanafterScheduledTaskSubmitted(long deadlineNanos) protected booleanbeforeScheduledTaskSubmitted(long deadlineNanos) Called from arbitrary non-EventExecutorthreads prior to scheduled task submission.protected voidCancel all scheduled tasks.protected static longdeadlineToDelayNanos(long deadlineNanos) Deprecated.UseAbstractScheduledEventExecutor.ticker()insteadprotected longdelayNanos(long currentTimeNanos, long scheduledPurgeInterval) Returns the amount of time left until the scheduled task with the closest dead line is executed.protected booleanfetchFromScheduledTaskQueue(Queue<Runnable> taskQueue) Fetch scheduled tasks from the internal queue and add these to the givenQueue.protected longDeprecated.Please use (or override)AbstractScheduledEventExecutor.ticker()instead.protected final booleanReturnstrueif a scheduled task is ready for processing.protected static longDeprecated.UseAbstractScheduledEventExecutor.ticker()insteadprotected static longnanoTime()Deprecated.Use the non-staticAbstractScheduledEventExecutor.ticker()instead.protected final longReturn the deadline (in nanoseconds) when the next scheduled task is ready to be run or-1if no task is scheduled.protected final longReturn the nanoseconds until the next scheduled task is ready to be run or-1if no task is scheduled.protected final Runnableprotected final RunnablepollScheduledTask(long nanoTime) Return theRunnablewhich is ready to be executed with the givennanoTime.<V> ScheduledFuture<V> scheduleAtFixedRate(Runnable command, long initialDelay, long period, TimeUnit unit) scheduleWithFixedDelay(Runnable command, long initialDelay, long delay, TimeUnit unit) ticker()The ticker for this executor.protected voidvalidateScheduled(long amount, TimeUnit unit) Deprecated.will be removed in the future.Methods inherited from class AbstractEventExecutor
iterator, newTaskFor, newTaskFor, runTask, safeExecute, shutdownGracefully, shutdownNow, submit, submit, submitModifier and TypeMethodDescriptioniterator()protected final <T> RunnableFuture<T> newTaskFor(Runnable runnable, T value) protected final <T> RunnableFuture<T> newTaskFor(Callable<T> callable) protected static voidprotected static voidsafeExecute(Runnable task) Future<?> Shortcut method forEventExecutorGroup.shutdownGracefully(long, long, TimeUnit)with sensible default values.Future<?> <T> Future<T> <T> Future<T> Methods inherited from class Object
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, waitMethods inherited from interface EventExecutor
inEventLoop, isExecutorThread, newFailedFuture, newProgressivePromise, newPromise, newSucceededFutureModifier and TypeMethodDescriptiondefault booleanCallsEventExecutor.inEventLoop(Thread)withThread.currentThread()as argumentdefault booleanisExecutorThread(Thread thread) default <V> Future<V> newFailedFuture(Throwable cause) Create a newFuturewhich is marked as failed already.default <V> ProgressivePromise<V> Create a newProgressivePromise.default <V> Promise<V> Return a newPromise.default <V> Future<V> newSucceededFuture(V result) Create a newFuturewhich is marked as succeeded already.Methods inherited from interface IoEventLoopGroup
register, registerModifier and TypeMethodDescriptiondefault ChannelFutureDeprecated.default ChannelFutureregister(ChannelPromise promise) Deprecated.Methods inherited from interface Iterable
forEach, spliterator
-
Constructor Details
-
SingleThreadIoEventLoop
public SingleThreadIoEventLoop(IoEventLoopGroup parent, ThreadFactory threadFactory, IoHandlerFactory ioHandlerFactory) Creates a new instance- Parameters:
parent- the parent that holds thisIoEventLoop.threadFactory- theThreadFactorythat is used to create the underlyingThread.ioHandlerFactory- theIoHandlerFactorythat should be used to obtainIoHandlerto handle IO.
-
SingleThreadIoEventLoop
public SingleThreadIoEventLoop(IoEventLoopGroup parent, Executor executor, IoHandlerFactory ioHandlerFactory) Creates a new instance- Parameters:
parent- the parent that holds thisIoEventLoop.executor- theExecutorthat is used for dispatching the work.ioHandlerFactory- theIoHandlerFactorythat should be used to obtainIoHandlerto handle IO.
-
SingleThreadIoEventLoop
public SingleThreadIoEventLoop(IoEventLoopGroup parent, ThreadFactory threadFactory, IoHandlerFactory ioHandlerFactory, int maxPendingTasks, RejectedExecutionHandler rejectedExecutionHandler, long maxTaskProcessingQuantumMs) Creates a new instance- Parameters:
parent- the parent that holds thisIoEventLoop.threadFactory- theThreadFactorythat is used to create the underlyingThread.ioHandlerFactory- theIoHandlerFactorythat should be used to obtainIoHandlerto handle IO.maxPendingTasks- the maximum pending tasks that are allowed beforeRejectedExecutionHandler.rejected(Runnable, SingleThreadEventExecutor)is called to handle it.rejectedExecutionHandler- theRejectedExecutionHandlerthat handles when more tasks are added then allowed permaxPendingTasks.maxTaskProcessingQuantumMs- the maximum number of milliseconds that will be spent to run tasks before trying to run IO again.
-
SingleThreadIoEventLoop
public SingleThreadIoEventLoop(IoEventLoopGroup parent, Executor executor, IoHandlerFactory ioHandlerFactory, int maxPendingTasks, RejectedExecutionHandler rejectedExecutionHandler, long maxTaskProcessingQuantumMs) Creates a new instance- Parameters:
parent- the parent that holds thisIoEventLoop.ioHandlerFactory- theIoHandlerFactorythat should be used to obtainIoHandlerto handle IO.maxPendingTasks- the maximum pending tasks that are allowed beforeRejectedExecutionHandler.rejected(Runnable, SingleThreadEventExecutor)is called to handle it.rejectedExecutionHandler- theRejectedExecutionHandlerthat handles when more tasks are added then allowed permaxPendingTasks.maxTaskProcessingQuantumMs- the maximum number of milliseconds that will be spent to run tasks before trying to run IO again.
-
SingleThreadIoEventLoop
protected SingleThreadIoEventLoop(IoEventLoopGroup parent, Executor executor, IoHandlerFactory ioHandlerFactory, Queue<Runnable> taskQueue, Queue<Runnable> tailTaskQueue, RejectedExecutionHandler rejectedExecutionHandler) Creates a new instance- Parameters:
parent- the parent that holds thisIoEventLoop.executor- theExecutorthat is used for dispatching the work.ioHandlerFactory- theIoHandlerFactorythat should be used to obtainIoHandlerto handle IO.taskQueue- theQueueused for storing pending tasks.tailTaskQueue- theQueueused for storing tail pending tasks.rejectedExecutionHandler- theRejectedExecutionHandlerthat handles when more tasks are added then allowed.
-
-
Method Details
-
run
protected void run()Description copied from class:SingleThreadEventExecutorRuns the task-processing loop untilSingleThreadEventExecutor.confirmShutdown()returnstrue.Implementations must not let a
Throwablethrown by a task escape this method: any uncaughtThrowableterminates the executor (logged atWARNand surfaced viaSingleThreadEventExecutor.terminationFuture()), at which point everyChannelregistered with this executor stops processing I/O and new task submissions are rejected. The supplied helpers -SingleThreadEventExecutor.runAllTasks(),SingleThreadEventExecutor.runAllTasks(long), andAbstractEventExecutor.safeExecute(Runnable)- catchThrowablefor you; custom loops built onSingleThreadEventExecutor.pollTask()orSingleThreadEventExecutor.takeTask()are responsible for wrapping each task invocation accordingly.- Specified by:
runin classSingleThreadEventExecutor
-
ioHandler
-
canSuspend
protected boolean canSuspend(int state) Description copied from class:SingleThreadEventExecutorReturnstrueif thisSingleThreadEventExecutorcan be suspended at the moment,falseotherwise. Subclasses might override this method to add extra checks.- Overrides:
canSuspendin classSingleThreadEventExecutor- Parameters:
state- the current internal state of theSingleThreadEventExecutor.- Returns:
- if suspension is possible at the moment.
-
runIo
protected int runIo()Called when IO will be processed for all theIoHandles on thisSingleThreadIoEventLoop. This method returns the number ofIoHandles for which IO was processed. This method must be called from theEventLoopthread. -
next
Description copied from interface:EventExecutorGroupReturns one of theEventExecutors managed by thisEventExecutorGroup.- Specified by:
nextin interfaceEventExecutorGroup- Specified by:
nextin interfaceEventLoopGroup- Specified by:
nextin interfaceIoEventLoop- Specified by:
nextin interfaceIoEventLoopGroup- Overrides:
nextin classSingleThreadEventLoop
-
register
Description copied from interface:IoEventLoop- Specified by:
registerin interfaceIoEventLoop- Specified by:
registerin interfaceIoEventLoopGroup- Parameters:
handle- theIoHandleto register.- Returns:
- the
Futurethat is notified once the operations completes.
-
getNumOfRegisteredChannels
protected int getNumOfRegisteredChannels()Description copied from class:SingleThreadEventExecutorReturns the number of registered channels for auto-scaling related decisions. This is intended to be used byMultithreadEventExecutorGroupfor dynamic scaling.- Overrides:
getNumOfRegisteredChannelsin classSingleThreadEventExecutor- Returns:
- The number of registered channels, or
-1if not applicable.
-
wakeup
protected final void wakeup(boolean inEventLoop) - Overrides:
wakeupin classSingleThreadEventExecutor
-
cleanup
protected final void cleanup()Description copied from class:SingleThreadEventExecutorDo nothing, sub-classes may override- Overrides:
cleanupin classSingleThreadEventExecutor
-
isCompatible
Description copied from interface:IoEventLoopGroupReturnstrueif the given type is compatible with thisIoEventLoopGroupand so can be registered to the containedIoEventLoops,falseotherwise.- Specified by:
isCompatiblein interfaceIoEventLoop- Specified by:
isCompatiblein interfaceIoEventLoopGroup- Parameters:
handleType- the type of theIoHandle.- Returns:
- if compatible of not.
-
isIoType
Description copied from interface:IoEventLoopGroup- Specified by:
isIoTypein interfaceIoEventLoop- Specified by:
isIoTypein interfaceIoEventLoopGroup- Parameters:
handlerType- the type of theIoHandler.- Returns:
- if used or not.
-
newTaskQueue
Description copied from class:SingleThreadEventExecutorCreate 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.- Overrides:
newTaskQueuein classSingleThreadEventExecutor
-
newTaskQueue0
-
SingleThreadEventExecutor.wakesUpForTask(Runnable)to re-create this behaviour