NOTE
Netty Source: NioEventLoop run Method
Code analysis of the NioEventLoop run loop: select, selector rebuild, selected-key optimization, processSelectedKeys, and asynchronous task execution.
This is a historical learning note and may contain outdated or incomplete understanding.
1. Code to Analyze
protected void run() {
for (;;) {
try {
switch (selectStrategy.calculateStrategy(selectNowSupplier, hasTasks())) {
case SelectStrategy.CONTINUE:
continue;
case SelectStrategy.SELECT:
// Poll the I/O events registered with this NioEventLoop's selector
select(wakenUp.getAndSet(false));
if (wakenUp.get()) {
selector.wakeup();
}
// fall through
default:
}
cancelledKeys = 0;
needsToSelectAgain = false;
final int ioRatio = this.ioRatio;// 50 by default
// ioRatio balances the time spent on the two kinds of work
if (ioRatio == 100) {
try {
// Process I/O events
processSelectedKeys();
} finally {
// Process tasks submitted by external threads into taskQueue
// Ensure we always run tasks.
runAllTasks();
}
} else {
//...
// See below
}
} catch (Throwable t) {
handleLoopException(t);
}
// Always handle shutdown even if the loop processing threw an exception.
try {
if (isShuttingDown()) {
closeAll();
if (confirmShutdown()) {
return;
}
}
} catch (Throwable t) {
handleLoopException(t);
}
}
}
The code above mainly contains three pieces of logic:
- An infinite loop
- Detect whether there are I/O events:
select - Process I/O events:
processSelectedKeys - Process the asynchronous task queue:
runAllTasks
- Detect whether there are I/O events:
The time spent processing I/O events and the time spent processing asynchronous tasks from external threads are balanced by ioRatio. The default is 50, meaning half of the time is spent executing processSelectedKeys and the other half executing runAllTasks.
That is exactly what the else branch above does:
else {
// Start time for processing I/O
final long ioStartTime = System.nanoTime();
try {
processSelectedKeys();
} finally {
// Ensure we always run tasks.
// Time after I/O processing - start time = time spent on I/O
final long ioTime = System.nanoTime() - ioStartTime;
// When ioRatio is 50, the parameter passed is ioTime * (100 - 50) / 50 == ioTime
// That is, runAllTasks is also given ioTime
runAllTasks(ioTime * (100 - ioRatio) / ioRatio);
}
}
2. Detect Whether There Are I/O Events: select
private void select(boolean oldWakenUp) throws IOException {
Selector selector = this.selector;
try {
// The key to solving empty polling
int selectCnt = 0;// Number of empty-polling loops that have executed
long currentTimeNanos = System.nanoTime();// Start time
long selectDeadLineNanos = currentTimeNanos + delayNanos(currentTimeNanos);// Time by which normal execution should finish
for (;;) {
// Calculate timeout duration
long timeoutMillis = (selectDeadLineNanos - currentTimeNanos + 500000L) / 1000000L;
// If timed out
if (timeoutMillis <= 0) {
// And select has not been executed even once
if (selectCnt == 0) {
// Execute non-blocking select
selector.selectNow();
selectCnt = 1;
}
break;
}
// If not timed out
// There is a task in the task queue -- that is, an external thread put a task into the task queue
if (hasTasks() && wakenUp.compareAndSet(false, true)) {
// Execute non-blocking select
selector.selectNow();
selectCnt = 1;
break;
}
// There is no task in the queue. Execute blocking select for timeout.
// [If empty polling occurs, this will not block for the full timeout duration]
int selectedKeys = selector.select(timeoutMillis);
selectCnt ++;
// If an event was selected || current select needs waking up ||
// select was awakened by an external thread || there is a task in the queue ||
// there is a scheduled task
if (selectedKeys != 0 || oldWakenUp || wakenUp.get() || hasTasks() || hasScheduledTasks()) {
// End this select operation
break;
}
if (Thread.interrupted()) {
// Thread was interrupted so reset selected keys and break so we not run into a busy loop.
// As this is most likely a bug in the handler of the user or it's client library we will
// also log it.
//
// See https://github.com/netty/netty/issues/2426
if (logger.isDebugEnabled()) {
logger.debug("Selector.select() returned prematurely because " +
"Thread.currentThread().interrupt() was called. Use " +
"NioEventLoop.shutdownGracefully() to shutdown the NioEventLoop.");
}
selectCnt = 1;
break;
}
// Reaching here each time means one blocking select operation was performed
// Current time - start polling time > timeout
// This means select(timeout) above really did block for timeout, so empty polling did not occur
// Reset selectCnt to 1
long time = System.nanoTime();
if (time - TimeUnit.MILLISECONDS.toNanos(timeoutMillis) >= currentTimeNanos) {
// timeoutMillis elapsed without anything selected.
selectCnt = 1;
// Reaching here means empty polling occurred. If the number of polls is > 512, rebuild the selector
} else if (SELECTOR_AUTO_REBUILD_THRESHOLD > 0 &&
selectCnt >= SELECTOR_AUTO_REBUILD_THRESHOLD) {
logger.warn(
"Selector.select() returned prematurely {} times in a row; rebuilding Selector {}.",
selectCnt, selector);
// Register the SelectionKeys from the old selector onto a new selector
rebuildSelector();
selector = this.selector;
// Select again to populate selectedKeys.
selector.selectNow();
selectCnt = 1;
break;
}
currentTimeNanos = time;
}
if (selectCnt > MIN_PREMATURE_SELECTOR_RETURNS) {
if (logger.isDebugEnabled()) {
logger.debug("Selector.select() returned prematurely {} times in a row for Selector {}.",
selectCnt - 1, selector);
}
}
} catch (CancelledKeyException e) {
if (logger.isDebugEnabled()) {
logger.debug(CancelledKeyException.class.getSimpleName() + " raised by a Selector {} - JDK bug?",
selector, e);
}
// Harmless exception - log anyway
}
}
2.1. Rebuild the Selector
private void rebuildSelector0() {
final Selector oldSelector = selector;
final SelectorTuple newSelectorTuple;
if (oldSelector == null) {
return;
}
try {
// Create a new selector
newSelectorTuple = openSelector();
} catch (Exception e) {
logger.warn("Failed to create a new Selector.", e);
return;
}
// Register all channels to the new Selector.
int nChannels = 0;
// For every key on the old selector
for (SelectionKey key: oldSelector.keys()) {
// Get its attachment
Object a = key.attachment();
try {
if (!key.isValid() || key.channel().keyFor(newSelectorTuple.unwrappedSelector) != null) {
continue;
}
// Get the events it is interested in
int interestOps = key.interestOps();
// Cancel the old key
key.cancel();
// Re-register a new key using the interested events and attachment
SelectionKey newKey = key.channel().register(newSelectorTuple.unwrappedSelector, interestOps, a);
if (a instanceof AbstractNioChannel) {
// Update SelectionKey
// Associate it with the channel
((AbstractNioChannel) a).selectionKey = newKey;
}
nChannels ++;
} catch (Exception e) {
logger.warn("Failed to re-register a Channel to the new Selector.", e);
if (a instanceof AbstractNioChannel) {
AbstractNioChannel ch = (AbstractNioChannel) a;
ch.unsafe().close(ch.unsafe().voidPromise());
} else {
@SuppressWarnings("unchecked")
NioTask<SelectableChannel> task = (NioTask<SelectableChannel>) a;
invokeChannelUnregistered(task, key, e);
}
}
}
selector = newSelectorTuple.selector;
unwrappedSelector = newSelectorTuple.unwrappedSelector;
try {
// time to close the old selector as everything else is registered to the new one
oldSelector.close();
} catch (Throwable t) {
if (logger.isWarnEnabled()) {
logger.warn("Failed to close the old Selector.", t);
}
}
logger.info("Migrated " + nChannels + " channel(s) to the new Selector.");
}
3. Process I/O Events: processSelectedKeys
3.1. Selected-Key Set Optimization
It essentially replaces the implementation of HashSet.add with an array so that add can achieve O(1) time complexity.
- Go back to the constructor that creates
NioEventLoop. There is anopenSelectoroperation:
private SelectorTuple openSelector() {
final Selector unwrappedSelector;
try {
// Call the JDK to create a selector
unwrappedSelector = provider.openSelector();
} catch (IOException e) {
throw new ChannelException("failed to open a new selector", e);
}
// If optimization is disabled, directly return the JDK selector
if (DISABLE_KEYSET_OPTIMIZATION) {
return new SelectorTuple(unwrappedSelector);
}
// Optimized set data structure -- implemented with an array
final SelectedSelectionKeySet selectedKeySet = new SelectedSelectionKeySet();
Object maybeSelectorImplClass = AccessController.doPrivileged(new PrivilegedAction<Object>() {
@Override
public Object run() {
try {
// Obtain the sun.nio.ch.SelectorImpl Class object through reflection
return Class.forName(
"sun.nio.ch.SelectorImpl",
false,
PlatformDependent.getSystemClassLoader());
} catch (Throwable cause) {
return cause;
}
}
});
// After obtaining the sun.nio.ch.SelectorImpl Class object, verify that it was actually obtained
if (!(maybeSelectorImplClass instanceof Class) ||
// ensure the current selector implementation is what we can instrument.
// And whether selector is an implementation of sun.nio.ch.SelectorImpl
!((Class<?>) maybeSelectorImplClass).isAssignableFrom(unwrappedSelector.getClass())) {
if (maybeSelectorImplClass instanceof Throwable) {
Throwable t = (Throwable) maybeSelectorImplClass;
logger.trace("failed to instrument a special java.util.Set into: {}", unwrappedSelector, t);
}
// Otherwise return the native selector
return new SelectorTuple(unwrappedSelector);
}
final Class<?> selectorImplClass = (Class<?>) maybeSelectorImplClass;
Object maybeException = AccessController.doPrivileged(new PrivilegedAction<Object>() {
@Override
public Object run() {
try {
// Get the two most important fields: selectedKeys and publicSelectedKeys.
// By default they are HashSet
Field selectedKeysField = selectorImplClass.getDeclaredField("selectedKeys");
Field publicSelectedKeysField = selectorImplClass.getDeclaredField("publicSelectedKeys");
Throwable cause = ReflectionUtil.trySetAccessible(selectedKeysField, true);
if (cause != null) {
return cause;
}
cause = ReflectionUtil.trySetAccessible(publicSelectedKeysField, true);
if (cause != null) {
return cause;
}
// Standard reflection flow: set them to our array-based implementation
selectedKeysField.set(unwrappedSelector, selectedKeySet);
publicSelectedKeysField.set(unwrappedSelector, selectedKeySet);
return null;
} catch (NoSuchFieldException e) {
return e;
} catch (IllegalAccessException e) {
return e;
}
}
});
if (maybeException instanceof Exception) {
selectedKeys = null;
Exception e = (Exception) maybeException;
logger.trace("failed to instrument a special java.util.Set into: {}", unwrappedSelector, e);
return new SelectorTuple(unwrappedSelector);
}
selectedKeys = selectedKeySet;
logger.trace("instrumented a special java.util.Set into: {}", unwrappedSelector);
return new SelectorTuple(unwrappedSelector,
new SelectedSelectionKeySetSelector(unwrappedSelector, selectedKeySet));
}
SelectedSelectionKeySet
final class SelectedSelectionKeySet extends AbstractSet<SelectionKey> {
// Implemented with an array and size
SelectionKey[] keys;
int size;
SelectedSelectionKeySet() {
keys = new SelectionKey[1024];// Default length 1024
}
@Override
public boolean add(SelectionKey o) {
if (o == null) {
return false;
}
// Direct assignment (O(1))
keys[size++] = o;
if (size == keys.length) {
increaseCapacity();// Expand to twice the size. SelectionKey[] newKeys = new SelectionKey[keys.length << 1];
}
return true;
}
//........
// Other operations are not implemented
}
3.2. processSelectedKeysOptimized
Return to processSelectedKeysOptimized in NioEventLoop.run.
private void processSelectedKeysOptimized() {
// Traverse our array implementation to obtain all keys
for (int i = 0; i < selectedKeys.size; ++i) {
final SelectionKey k = selectedKeys.keys[i];
// null out entry in the array to allow to have it GC'ed once the Channel close
// See https:github.com/netty/netty/issues/2363
selectedKeys.keys[i] = null;
// Get the attachment corresponding to the key
final Object a = k.attachment();
// If it is AbstractNioChannel, process it
if (a instanceof AbstractNioChannel) {
processSelectedKey(k, (AbstractNioChannel) a);
} else {
@SuppressWarnings("unchecked")
NioTask<SelectableChannel> task = (NioTask<SelectableChannel>) a;
processSelectedKey(k, task);
}
if (needsToSelectAgain) {
// null out entries in the array to allow to have it GC'ed once the Channel close
// See https:github.com/netty/netty/issues/2363
selectedKeys.reset(i + 1);
selectAgain();
i = -1;
}
}
}
private void processSelectedKey(SelectionKey k, AbstractNioChannel ch) {
final AbstractNioChannel.NioUnsafe unsafe = ch.unsafe();
// Handling an invalid key
if (!k.isValid()) {
final EventLoop eventLoop;
try {
eventLoop = ch.eventLoop();
} catch (Throwable ignored) {
// If the channel implementation throws an exception because there is no event loop, we ignore this
// because we are only trying to determine if ch is registered to this event loop and thus has authority
// to close ch.
return;
}
// Only close ch if ch is still registered to this EventLoop. ch could have deregistered from the event loop
// and thus the SelectionKey could be cancelled as part of the deregistration process, but the channel is
// still healthy and should not be closed.
// See https://github.com/netty/netty/issues/5125
if (eventLoop != this || eventLoop == null) {
return;
}
// close the channel if the key is not valid anymore
unsafe.close(unsafe.voidPromise());
return;
}
try {
// Obtain all events
int readyOps = k.readyOps();
// We first need to call finishConnect() before try to trigger a read(...) or write(...) as otherwise
// the NIO JDK channel implementation may throw a NotYetConnectedException.
// Handle events such as OP_CONNECT
if ((readyOps & SelectionKey.OP_CONNECT) != 0) {
// remove OP_CONNECT as otherwise Selector.select(..) will always return without blocking
// See https://github.com/netty/netty/issues/924
int ops = k.interestOps();
ops &= ~SelectionKey.OP_CONNECT;
k.interestOps(ops);
unsafe.finishConnect();
}
// Process OP_WRITE first as we may be able to write some queued buffers and so free memory.
if ((readyOps & SelectionKey.OP_WRITE) != 0) {
// Call forceFlush which will also take care of clear the OP_WRITE once there is nothing left to write
ch.unsafe().forceFlush();
}
// Also check for readOps of 0 to workaround possible JDK bug which may otherwise lead
// to a spin loop
// If this is bossGroup, the polled event is OP_READ; if this is workerGroup, the polled event is OP_ACCEPT
if ((readyOps & (SelectionKey.OP_READ | SelectionKey.OP_ACCEPT)) != 0 || readyOps == 0) {
unsafe.read();
}
} catch (CancelledKeyException ignored) {
unsafe.close(unsafe.voidPromise());
}
}
4. Process the Asynchronous Task Queue: runAllTasks
4.1. Task Categories and Addition
There are two kinds of tasks: scheduled tasks and ordinary tasks.
4.1.1. Ordinary Tasks
protected SingleThreadEventExecutor(EventExecutorGroup parent, Executor executor,
boolean addTaskWakesUp, int maxPendingTasks,
RejectedExecutionHandler rejectedHandler) {
super(parent);
this.addTaskWakesUp = addTaskWakesUp;
this.maxPendingTasks = Math.max(16, maxPendingTasks);
this.executor = ObjectUtil.checkNotNull(executor, "executor");
taskQueue = newTaskQueue(this.maxPendingTasks);// Here
rejectedExecutionHandler = ObjectUtil.checkNotNull(rejectedHandler, "rejectedHandler");
}
4.1.1.1. When Are Ordinary Tasks Added?
When an external thread calls NioEventLoop.execute:
public void execute(Runnable task) {
if (task == null) {
throw new NullPointerException("task");
}
boolean inEventLoop = inEventLoop();
addTask(task);// Add directly into the task queue
// (which means it is thread-safe: PlatformDependent.<Runnable>newMpscQueue): taskQueue.offer(task);
if (!inEventLoop) {// A thread not in NioEventLoop
startThread();// Start a new thread to process it
if (isShutdown() && removeTask(task)) {
reject();
}
}
if (!addTaskWakesUp && wakesUpForTask(task)) {
wakeup(inEventLoop);
}
}
4.1.2. Scheduled Tasks
AbstractScheduledEventExecutor#schedule
public <V> ScheduledFuture<V> schedule(Callable<V> callable, long delay, TimeUnit unit) {
ObjectUtil.checkNotNull(callable, "callable");
ObjectUtil.checkNotNull(unit, "unit");
if (delay < 0) {
delay = 0;
}
validateScheduled(delay, unit);
// Wrap it as ScheduledFutureTask
return schedule(new ScheduledFutureTask<V>(
this, callable, ScheduledFutureTask.deadlineNanos(unit.toNanos(delay))));
}
<V> ScheduledFuture<V> schedule(final ScheduledFutureTask<V> task) {
// If this is the thread in NioEventLoop, add it directly
if (inEventLoop()) {
scheduledTaskQueue().add(task);
} else {
// Otherwise add it through an indirect thread -- why?
// Because this queue, DefaultPriorityQueue, is not thread-safe
execute(new Runnable() {
@Override
public void run() {
scheduledTaskQueue().add(task);
}
});
}
return task;
}
4.2. Aggregate Tasks
Return to the first operation of runAllTasks, fetchFromScheduledTaskQueue.
private boolean fetchFromScheduledTaskQueue() {
long nanoTime = AbstractScheduledEventExecutor.nanoTime();
// Take tasks from the scheduled-task queue:
// tasks are ordered from earliest to latest deadline by ScheduledFutureTask.compareTo
Runnable scheduledTask = pollScheduledTask(nanoTime);
while (scheduledTask != null) {
// Put it into the ordinary task queue
if (!taskQueue.offer(scheduledTask)) {
// No space left in the task queue add it back to the scheduledTaskQueue so we pick it up again.
// If it fails, add it back to the scheduled queue
scheduledTaskQueue().add((ScheduledFutureTask<?>) scheduledTask);
return false;
}
// Continue
scheduledTask = pollScheduledTask(nanoTime);
}
return true;
}
4.3. Execute Tasks
protected boolean runAllTasks(long timeoutNanos) {
// See the analysis above
fetchFromScheduledTaskQueue();
// Take one task from the ordinary task queue
Runnable task = pollTask();
if (task == null) {
afterRunningAllTasks();
return false;
}
final long deadline = ScheduledFutureTask.nanoTime() + timeoutNanos;
long runTasks = 0;
long lastExecutionTime;
// Continuously execute tasks
for (;;) {
// task.run
safeExecute(task);
runTasks ++;
// Check timeout every 64 tasks because nanoTime() is relatively expensive.
// XXX: Hard-coded value - will make it configurable if it is really a problem.
// Once 64 tasks have accumulated, calculate the current time.
// If it exceeds the deadline, stop executing
if ((runTasks & 0x3F) == 0) {
lastExecutionTime = ScheduledFutureTask.nanoTime();
if (lastExecutionTime >= deadline) {
break;
}
}
// If it has not exceeded the deadline, take another task
task = pollTask();
if (task == null) {
lastExecutionTime = ScheduledFutureTask.nanoTime();
break;
}
}
afterRunningAllTasks();
this.lastExecutionTime = lastExecutionTime;
return true;
}
Discussion
Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub