Подробное объяснение модели многопоточности Netty Reactor

Java

1. Введение

1. Что такое реактор?

Шаблон      Reactor — это шаблон обработки событий для обработки запросов на обслуживание, которые передаются на сервер одновременно через один или несколько входов. Обработчик службы мультиплексирует входящие запросы и синхронно отправляет их соответствующему обработчику. Ключевые моменты:

(1) Движимый событиями

(2) Обработка нескольких входов

(3) Используйте мультиплексирование для передачи событий соответствующему обработчику для обработки.

2. Основные компоненты реактора

(1) Reactor

     отвечает за реагирование на события и привязку распределения событий к обработчику события. Соответствует NioEventLoop.run(), processSelectedKeys() для netty.

(2) Handler

     Обработчик события, привязанный к определенному типу события, отвечает за выполнение задачи соответствующего события для обработки события. Соответствует IdleStateHandler netty и т. д.

(3) Acceptor

     Акцептор относится к одному из обработчиков, потому что он более специальный и независимый, это класс приема событий реактора, который отвечает за инициализацию селектора и получение буферной очереди. ServerBootstrapAcceptor, соответствующий netty.

2. Процесс

     Каждый поток Reactor в пуле потоков Reactor будет иметь собственный селектор, поток и логику цикла отправленных событий. Может быть только один mainReactor, но обычно есть несколько subReactors. Поток mainReacto в основном отвечает за получение запроса на подключение от клиента, а затем за передачу полученного SocketChannel в subReactor, который завершает связь с клиентом. Анализ исходного кода

1. Создайте пул потоков mainReactor и пул потоков subReactor.

	bossGroup = new NioEventLoopGroup();
	workGroup = new NioEventLoopGroup(4);
protected MultithreadEventExecutorGroup(int nThreads, ThreadFactory threadFactory, Object... args) {
     children = new SingleThreadEventExecutor[nThreads];
     ...
     for (int i = 0; i < nThreads; i ++) {
            ...
            children[i] = newChild(threadFactory, args);
            ...
     }
}
@Override
protected EventExecutor newChild(
        ThreadFactory threadFactory, Object... args) throws Exception {
    return new NioEventLoop(this, threadFactory, (SelectorProvider) args[0]);
}

     Здесь создаются пулы потоков mainReactor и subReactor, а также поток eventLoop.

NioEventLoop(NioEventLoopGroup parent, ThreadFactory threadFactory, SelectorProvider selectorProvider) {
    super(parent, threadFactory, false);
    if (selectorProvider == null) {
        throw new NullPointerException("selectorProvider");
    }
    provider = selectorProvider;
    selector = openSelector();
}

     Каждый поток eventLoop будет иметь свой собственный селектор, здесь поток eventLoop еще не запущен, после его запуска будет выполняться selector.select в run().

2. mainReactor привязывает селектор события OP_ACCEPT и запускает цикл потока для выполнения selector.select();

ChannelFuture regFuture = group().register(channel); Здесь group() — это bossGroup

@Override
public ChannelFuture register(Channel channel) {
    return next().register(channel);
}

     выполнение следующего()

@Override
public EventLoop next() {
    return (EventLoop) super.next();
}
private final class PowerOfTwoEventExecutorChooser implements EventExecutorChooser {
    @Override
    public EventExecutor next() {
        return children[childIndex.getAndIncrement() & children.length - 1];
    }
}

     Возьмите первый цикл событий из пула потоков

@Override
public ChannelFuture register(final Channel channel, final ChannelPromise promise) {
     ...
    channel.unsafe().register(this, promise);
    return promise;
}
@Override
public final void register(EventLoop eventLoop, final ChannelPromise promise) {
    ...
    AbstractChannel.this.eventLoop = eventLoop;

    if (eventLoop.inEventLoop()) {
        register0(promise);
    } else {
    try {
        eventLoop.execute(new OneTimeTask() {
            @Override
            public void run() {
                register0(promise);
            }
        });
    } catch (Throwable t) {
}

     Здесь eventLoop mainReactor привязан к NioServerSocketChannel сервера. Поскольку основной поток запускается в начале, выполняется eventLoop.execute, где mainReactor запускает только один поток.

@Override
public void execute(Runnable task) {
    boolean inEventLoop = inEventLoop();
    if (inEventLoop) {
        addTask(task);
    } else {
        startThread();
        addTask(task);
        ...
    }
        ...
}

     Выполните startThread() в execute, чтобы официально запустить цикл потока mainReactor, и добавьте Task register0 (обещание) в taskQueue, чтобы позволить циклу mainReactor выполняться.

private void register0(ChannelPromise promise) {
    doRegister();
    neverRegistered = false;
    registered = true;
    safeSetSuccess(promise);
    pipeline.fireChannelRegistered();
    if (firstRegistration && isActive()) {
        pipeline.fireChannelActive();
    }
}
@Override
protected void doRegister() throws Exception {
    boolean selected = false;
    for (;;) {
        ...
        selectionKey = javaChannel().register(eventLoop().selector, 0, this);
	...	
    }
}

     Здесь селектор eventLoop в mainReactor регистрирует бит прослушивания операции, равный 0, и привязывает NioServerSocketChannel сервера к потоку mainSubReactor. В doBind()-->doBind0()-->channel.bind()-->…-->next.invokeBind()-->HeadContext.Bind()-->unsafe.bind()-->конвейер .fireChannelActive()-->channel.read()-->…-->doBeginRead() изменен на бит контроля операции OP_ACCEPT(16).

@Override
protected void doBeginRead() throws Exception {
…
    final int interestOps = selectionKey.interestOps();
    if ((interestOps & readInterestOp) == 0) {
        selectionKey.interestOps(interestOps | readInterestOp);
    }
}
  自此,mainReactor的eventLoop从run开始循环执行selector.select。
  注:readInterestOp的值来自于创建NioServerSocketChannel的构造函数
public NioServerSocketChannel(ServerSocketChannel channel) {
    super(null, channel, SelectionKey.OP_ACCEPT);
    config = new NioServerSocketChannelConfig(this, javaChannel().socket());
}

3. subReactor регистрирует событие OP_READ

     После получения клиентского соединения Канал клиента будет зарегистрирован в потоке subReactor в ServerBootstrapAcceptor, и канал будет привязан к селектору потока subReactor, а событие OP_READ клиентского канала будет отслеживаться.

if ((readyOps & (SelectionKey.OP_READ | SelectionKey.OP_ACCEPT)) != 0 || readyOps == 0) {
    unsafe.read();
}

     При прослушивании клиентского соединения выполните read() сервера AbstractNioUnsafe;

@Override
public void read() {
    ...
    int localRead = doReadMessages(readBuf);
    ...
    for (int i = 0; i < size; i ++) {
        pipeline.fireChannelRead(readBuf.get(i));
    }
    ...
    pipeline.fireChannelReadComplete();
    ...
}

(1) doReadMessage

@Override
protected int doReadMessages(List<Object> buf) throws Exception {
    SocketChannel ch = javaChannel().accept();
    ...
    buf.add(new NioSocketChannel(this, ch));
    ...
}
public NioSocketChannel(Channel parent, SocketChannel socket) {
    super(parent, socket);
    config = new NioSocketChannelConfig(this, socket.socket());
}
protected AbstractNioByteChannel(Channel parent, SelectableChannel ch) {
    super(parent, ch, SelectionKey.OP_READ);
}

     Установите значение бита прослушивания клиентского канала в OP_READ(1)

(2) pipeline.fireChannelRead()

private static class ServerBootstrapAcceptor extends ChannelInboundHandlerAdapter {
public void channelRead(ChannelHandlerContext ctx, Object msg) {
   final Channel child = (Channel) msg;
   child.pipeline().addLast(childHandler);

   for (Entry<ChannelOption<?>, Object> e: childOptions) {
       try {
           if (!child.config().setOption((ChannelOption<Object>) e.getKey(), e.getValue())) {
               logger.warn("Unknown channel option: " + e);
           }
       } catch (Throwable t) {
           logger.warn("Failed to set a channel option: " + child, t);
       }
   }

   for (Entry<AttributeKey<?>, Object> e: childAttrs) {
       child.attr((AttributeKey<Object>) e.getKey()).set(e.getValue());
   }

   try {
        childGroup.register(child).addListener(new ChannelFutureListener() {
        @Override
        public void operationComplete(ChannelFuture future) throws Exception {
            if (!future.isSuccess()) {
                forceClose(child, future.cause());
            }
         }
       });
     } catch (Throwable t) {
         forceClose(child, t);
     }
}
}

     ServerBootstrapAcceptor не только привязывает subReactor к клиентскому каналу, но и инициализирует некоторые параметры для клиентского канала.

@Override
public ChannelFuture register(Channel channel) {
    return next().register(channel);
}

     То же, что и регистр выше, за исключением того, что пул потоков mainReactor изменяется на пул потоков subReactor. Здесь селектор потока берется из пула потоков subReactor и привязывается к каналу клиента, а событие клиента 0 отслеживается.

(3) pipeline.fireChannelReadComplete()

@Override
public ChannelPipeline fireChannelReadComplete() {
    head.fireChannelReadComplete();
    if (channel.config().isAutoRead()) {
        read();
    }
    return this;
}

     read() --> tail.read() --> next.invokeRead() --> HeadContext.read() -->… --> doBeginRead()

@Override
protected void doBeginRead() throws Exception {
    ...
    final int interestOps = selectionKey.interestOps();
    if ((interestOps & readInterestOp) == 0) {
        selectionKey.interestOps(interestOps | readInterestOp);
    }
}

     Измените бит прослушивания на OP_READ(1) ​​здесь

4. SubReactor обрабатывает событие чтения

	if ((readyOps & (SelectionKey.OP_READ | SelectionKey.OP_ACCEPT)) != 0 || readyOps == 0) {
            unsafe.read();
            ...
        }

     Перейдите к методу read() NioByteUnsafe

@Override
public final void read() {
    ...
    final ChannelPipeline pipeline = pipeline();
    final ByteBufAllocator allocator = config.getAllocator();
    ...
    byteBuf = allocHandle.allocate(allocator);
    ...
    pipeline.fireChannelRead(byteBuf);
    ...
}
@Override
public ChannelPipeline fireChannelRead(Object msg) {
    head.fireChannelRead(msg);
    return this;
}
private void invokeChannelRead(Object msg) {
    try {
        ((ChannelInboundHandler) handler()).channelRead(this, msg);
    } catch (Throwable t) {
        notifyHandlerException(t);
    }
}
public class InBoundHandlerB extends ChannelInboundHandlerAdapter {
    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
        System.out.println("InBoundHandlerB: " + msg);
        super.channelRead(ctx, msg);
    }
}

     Здесь будет обрабатываться клиентское сообщение

Суммировать

     Сколько селекторов будет создано, так как есть пулы потоков Reactor, цикл событий mainReactor будет привязан к каналу сервера, и обратите внимание только на событие ACCEPT канала сервера, цикл событий subReactor будет привязан к каналу клиента, и обращайте внимание только на клиентское событие READ канала.

     mainReactor и subReactor зацикливают свои соответствующие селекторы, mainReactor зацикливает селектор события ACCEPT, subReactor зацикливает селектор события READ, после того как mainReactor получает клиентское соединение, он выполняет метод channelRead ServerBootstrapAcceptor для привязки клиентского соединения к subReactor.