暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

死磕Netty源码之新连接接入源码解析

纯洁的明依 2019-11-24
161

主要是介绍`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() {
@Override
public 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() {
@Override
public 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>() {
@Override
public 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`相关的源码解析


文章转载自纯洁的明依,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

评论