主要是介绍`Netty`在新连接接入后的相关处理。
新连接建立可以分为以下三个步骤:
1.检测到有新的连接
2.将新的连接注册到Worker线程组
3.注册新连接的读事件
在`Reactor`线程模型详解中得知当服务端读取到`IO`事件(新连接接入事件)后,会调用`processSelectedKey`方法对事件进行处理,此处以新连接接入事件为例它最后会调用底层的`unsafe`进行`read`操作。
public void read() {assert eventLoop().inEventLoop();final ChannelPipeline pipeline = pipeline();final RecvByteBufAllocator.Handle allocHandle = unsafe().recvBufAllocHandle();do {int localRead = doReadMessages(readBuf);if (localRead == 0) {break;}if (localRead < 0) {closed = true;break;}} while (allocHandle.continueReading());int size = readBuf.size();for (int i = 0; i < size; i ++) {pipeline.fireChannelRead(readBuf.get(i));}readBuf.clear();pipeline.fireChannelReadComplete();}
这里有两个主要的方法:
1.doReadMessages
2.pipeline.fireChannelRead
protected int doReadMessages(List<Object> buf) throws Exception {SocketChannel ch = javaChannel().accept();try {if (ch != null) {buf.add(new NioSocketChannel(this, ch));return 1;}} catch (Throwable t) {// ...}return 0;}
该方法主要作用是通过`JDK`底层的`API`获取到`SocketChannel`,然后包装成`Netty`的`NioSocketChannel`。`NioSocketChannel`与服务端启动时创建的`NioServerSocketChannel`最主要的区别在于它们关注的事件不同,
`NioSocketChannel`的构造方法如下:
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);}
这里我们看到一个`SelectionKey.OP_READ`,说明这个`Channel`关心读事件而服务端的`Channel`关心`ACCEPT`事件。接下来调用父类`AbstractNioChannel`构造,后续过程与服务端启动流程一致此处不再赘述.
接着来看`pipeline.fireChannelRead(readBuf.get(i))`方法,关于`Pipeline`我们将在下一篇博客中详细介绍。
我们知道服务端在启动的过程中会往`Pipeline`中添加一个`ServerBootstrapAcceptor`(连接处理器),所以到这里服务端`Channel`对应的`Pipeline`的数据结构为:`Hea⇋ServerBootstrapAcceptor⇋Tail`,因此在调用`pipeline.fireChannelRead`时会依次触发这三个节点上的`channelRead`方法.
`ServerBootstrapAcceptor`的`channelRead`方法,代码如下:
public void channelRead(ChannelHandlerContext ctx, Object msg) {final Channel child = (Channel) msg;child.pipeline().addLast(childHandler);setChannelOptions(child, childOptions, logger);for (Entry<AttributeKey<?>, Object> e: childAttrs) {child.attr((AttributeKey<Object>) e.getKey()).set(e.getValue());}try {childGroup.register(child).addListener(new ChannelFutureListener() {@Overridepublic void operationComplete(ChannelFuture future) throws Exception {if (!future.isSuccess()) {forceClose(child, future.cause());}}});} catch (Throwable t) {forceClose(child, t);}}
首先获取我们之前实例化的`NioSocketChannel`,然后将我们设置的`chlidHandler`添加到`NioSocketChannel`对应的`Pipeline`中(这里的`chlidHandler`对应用户通过`.childHandler()`设置的`Handler`),代码执行到这里`NioSocketChannel`中`Pipeline`对应的数据结构为: `head⇋ChannelInitializer⇋tail`,接着设置对应的`attr`和`option`,然后进入到`childGroup.register(child)`(这里的`childGroup`就是`WorkerGroup`),接下来我们进入`NioEventLoopGroup`的`register`方法:
public ChannelFuture register(Channel channel) {return next().register(channel);}public ChannelFuture register(Channel channel) {return register(new DefaultChannelPromise(channel, this));}public ChannelFuture register(final ChannelPromise promise) {ObjectUtil.checkNotNull(promise, "promise");promise.channel().unsafe().register(this, promise);return promise;}
这段代码和服务端启动的时候像`BossGroup`注册`NioServerSocketChannel`是类似的,通过`next()`方法获取到`NioEventLoop`然后将`Channel`注册到该`NioEventLoop`上(即将该`Channel`与`NioEventLoop`的`Selector`进行绑定)。注册的逻辑最终是交给`Unsafe`对象完成的,我们继续跟进`Unsafe`的`register`方法代码如下:
public final void register(EventLoop eventLoop, final ChannelPromise promise) {//...AbstractChannel.this.eventLoop = eventLoop;if (eventLoop.inEventLoop()) {register0(promise);} else {try {eventLoop.execute(new Runnable() {@Overridepublic void run() {register0(promise);}});} catch (Throwable t) {//...}}}
由于是在`Boss`线程中执行的`IO`操作所以不会是跟`Worker`线程是同一个线程,故`eventLoop.inEventLoop()`返回`false`,最后会通过`eventLoop.execute`的方式去执行注册任务。在`Reactor`线程模型中我们讲到在调用`execute`的时候,如果是首次添加任务那这个`NioEventLoop`线程会被启动,所以从此`Worker`线程开始执行,接下来看下具体的注册逻辑:
private void register0(ChannelPromise promise) {try {boolean firstRegistration = neverRegistered;doRegister();neverRegistered = false;registered = true;pipeline.invokeHandlerAddedIfNeeded();safeSetSuccess(promise);pipeline.fireChannelRegistered();if (isActive()) {if (firstRegistration) {pipeline.fireChannelActive();} else if (config().isAutoRead()) {beginRead();}}} catch (Throwable t) {//...}}
和服务端启动过程一样,先是调用`doRegister()`执行真正的注册过程
protected void doRegister() throws Exception {boolean selected = false;for (;;) {try {selectionKey = javaChannel().register(eventLoop().selector, 0, this);return;} catch (CancelledKeyException e) {//...}}}
该`Channel`绑定到`NioEventLoop`对应的`Selector`上去,后续该`Channel`的事件轮询、事件处理、异步`Task`执行都由此线程负责,绑定完`Reactor`线程之后调用`pipeline.invokeHandlerAddedIfNeeded()`代码如下:
final void invokeHandlerAddedIfNeeded() {assert channel.eventLoop().inEventLoop();if (firstRegistration) {firstRegistration = false;callHandlerAddedForAllHandlers();}}
往下跟`callHandlerAddedForAllHandlers`方法:
private void callHandlerAddedForAllHandlers() {final PendingHandlerCallback pendingHandlerCallbackHead;synchronized (this) {assert !registered;registered = true;pendingHandlerCallbackHead = this.pendingHandlerCallbackHead;this.pendingHandlerCallbackHead = null;}PendingHandlerCallback task = pendingHandlerCallbackHead;while (task != null) {task.execute();task = task.next;}}
这里有个对象叫`pendingHandlerCallbackHead`,它是在`callHandlerCallbackLater`方法中被初始化的
private void callHandlerCallbackLater(AbstractChannelHandlerContext ctx, boolean added) {assert !registered;PendingHandlerCallback task = added ? new PendingHandlerAddedTask(ctx) : new PendingHandlerRemovedTask(ctx);PendingHandlerCallback pending = pendingHandlerCallbackHead;if (pending == null) {pendingHandlerCallbackHead = task;} else {// Find the tail of the linked-list.while (pending.next != null) {pending = pending.next;}pending.next = task;}}
在`Channel`注册到之前添加或删除`Handler`时没有`EventExecutor`可执行`HandlerAdd`或`HandlerRemove`事件,所以`Netty`为此事件生成一个相应任务等注册完成后在调用执行任务。添加或删除任务可能会有很多个,所以`DefaultChannelPipeline`使用一个链表存储,链表头部为先前的字段`pendingHandlerCallbackHead`接下来我们继续分析`task.execute`方法, 它主要是完成`NioSocketChannel`对应的`Pipeline`的初始化.
void execute() {// ...callHandlerAdded0(ctx);// ...}
通过以上对`pendingHandlerCallbackHead`的分析,这里会调用`ChannelInitializer`的`handlerAdded`方法
public void handlerAdded(ChannelHandlerContext ctx) throws Exception {if (ctx.channel().isRegistered()) {initChannel(ctx);}}private boolean initChannel(ChannelHandlerContext ctx) throws Exception {if (initMap.putIfAbsent(ctx, Boolean.TRUE) == null) {try {initChannel((C) ctx.channel());} catch (Throwable cause) {exceptionCaught(ctx, cause);} finally {remove(ctx);}return true;}return false;}
`ChannelInitializer`的`initChannel`主要完成两个功能以下两个功能
调用`initChannel((C) ctx.channel())`进入用户自定义的代码完成`Pipeline`的初始化.
.childHandler(new ChannelInitializer<SocketChannel>() {@Overridepublic void initChannel(SocketChannel ch) throws Exception {}})
- 在`finally`中调用`remove`方法将`ChannelInitializer`删除
private void remove(ChannelHandlerContext ctx) {try {ChannelPipeline pipeline = ctx.pipeline();if (pipeline.context(this) != null) {pipeline.remove(this);}} finally {initMap.remove(ctx);
执行该方法前`NioSocketChannel`对应的`Pipeline`的数据结构为:`head⇋ChannelInitializer⇋tail`,执行该方法后`ChannelInitializer`被删除,`NioSocketChannel`对应的`Pipeline`的数据结构为:`head⇋自定义的HandlerContext⇋tail`。到目前为止我们完成了新连接的注册、`pipeline`的绑定,但是新连接注册的时候的感兴趣事件还是0还无法进行读写操作,新连接对读事件的绑定是在`pipeline.fireChannelActive`方法中完成的,它最后会调用到`AbstractNioChannel`的`doBeginRead`
protected void doBeginRead() throws Exception {final SelectionKey selectionKey = this.selectionKey;if (!selectionKey.isValid()) {return;}readPending = true;final int interestOps = selectionKey.interestOps();if ((interestOps & readInterestOp) == 0) {selectionKey.interestOps(interestOps | readInterestOp);}}
前面`register0()`方法的时候向`selector`注册的事件代码是0,而`readInterestOp`对应的事件代码是`SelectionKey.OP_READ`,所以本段代码的用处是将`SelectionKey.OP_READ`事件注册到`Selector`中去,`fireChannelActive`的执行逻辑在服务端启动过程中有详细描述,至此已完成客户端新连接接入的操作.
下一客将介绍`Pipeline`相关的源码解析




