NOTE

Netty Source: NioEventLoop run Method

Code analysis of the NioEventLoop run loop: select, selector rebuild, selected-key optimization, processSelectedKeys, and asynchronous task execution.

JavaCreated Updated 1 min readhistorical

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

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 an openSelector operation:
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;
}

5. Reference

5.1. JDK Empty-Polling Bug

Discussion

Sign in with GitHub to comment. Discussions are stored as GitHub Issues.View on GitHub