乐于分享
好东西不私藏

Netty 新连接接入源码:Boss 如何把 Channel 交给 Worker(源码系列第3篇)

Netty 新连接接入源码:Boss 如何把 Channel 交给 Worker(源码系列第3篇)

Netty 的主从 Reactor 不是两组线程直接传递 Socket,而是 Boss 把 accept() 得到的连接包装成子 Channel,再通过父 Pipeline 的 ServerBootstrapAcceptor 完成配置和 Worker 注册。父子 Channel、父子 Pipeline 与两组 EventLoop 的边界,共同保证了“接入轻量、连接固定归属、业务互不串线”。

上一章走完了服务端 bind()NioServerSocketChannel 已注册到 Boss EventLoop 的 Selector,并完成端口监听。接下来客户端发起 TCP 连接,内核完成握手,Selector 返回 OP_ACCEPT。从这一刻到业务 Handler 收到 channelActive,Netty 要完成一场精确的连接接力。

这条路径最容易出现三个误解:

  1. 1. Boss EventLoop 会一直负责该连接的读写。
  2. 2. childHandler 在服务启动时已经加入父 Pipeline。
  3. 3. JDK SocketChannel 被直接交给 Worker Selector。

实际情况是:Boss 只负责监听 Channel 的 Accept 事件;每次 Accept 都会创建一个新的 NioSocketChannel;这个子 Channel 先作为一条入站消息经过父 Pipeline,再由 ServerBootstrapAcceptor 配置并注册到 Worker EventLoop。注册之后,后续读写才由 Worker 负责。

源码基线:

  • • Netty netty-4.1.135.Final
  • • Commit:f05f765d81460799c53123a207f665bf3b465171
  • • 重点文件:AbstractNioMessageChannel.javaNioServerSocketChannel.javaServerBootstrap.java
  • • 官方仓库:https://github.com/netty/netty

1. 父 Channel 只代表监听端口

服务端绑定 8080 后,存在一个 NioServerSocketChannel。它包装 JDK ServerSocketChannel,关注 SelectionKey.OP_ACCEPT。它不是某个客户端连接,也不会承载客户端请求字节。

当三个客户端连入时,运行时对象大致是:

NioServerSocketChannel :8080  ├── NioSocketChannel client-A  ├── NioSocketChannel client-B  └── NioSocketChannel client-C

父 Channel 的 parent() 为 null,子 Channel 的 parent() 指向 ServerChannel。父子关系主要表达连接来源和生命周期上下文,并不意味着所有 I/O 都在同一线程执行。

父 Pipeline 通常只有启动日志、接入控制和系统安装的 ServerBootstrapAcceptor。子 Pipeline 才会包含 Decoder、Encoder、鉴权和业务 Handler。把这个边界想清楚,handler() 与 childHandler() 的区别就非常自然。

2. OP_ACCEPT 最终进入 NioMessageUnsafe.read()

Boss EventLoop 调用 processSelectedKey() 时,发现 SelectionKey 包含 OP_ACCEPT。对消息型 NIO Channel,最终会调用 AbstractNioMessageChannel.NioMessageUnsafe.read()

源码位置:

  • • 文件:transport/src/main/java/io/netty/channel/nio/AbstractNioMessageChannel.java
  • • 方法:NioMessageUnsafe.read()
  • • 链接:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/nio/AbstractNioMessageChannel.java

主逻辑可简化为:

public void read() {    ChannelConfig config = config();    ChannelPipeline pipeline = pipeline();    RecvByteBufAllocator.Handle allocHandle = unsafe().recvBufAllocHandle();    allocHandle.reset(config);    boolean closed = false;    Throwable exception = null;    try {        do {            int localRead = doReadMessages(readBuf);            if (localRead == 0) {                break;            }            if (localRead < 0) {                closed = true;                break;            }            allocHandle.incMessagesRead(localRead);        } while (allocHandle.continueReading());    } catch (Throwable t) {        exception = t;    }    for (Object message : readBuf) {        pipeline.fireChannelRead(message);    }    pipeline.fireChannelReadComplete();    if (exception != null) {        pipeline.fireExceptionCaught(exception);    }}

这里的“read”不是读取业务字节,而是读取一批消息。对于 ServerSocketChannel,一个消息就是一个新连接。readBuf 因此保存的是新创建的子 Channel。

Netty 会在一次 Accept 就绪中尽可能处理多个连接,但又由 RecvByteBufAllocator.Handle 控制循环,避免单个 ServerChannel 长时间霸占 Boss EventLoop。

3. doReadMessages() 完成 JDK accept 和 Channel 包装

NioServerSocketChannel.doReadMessages() 是真正调用 JDK accept() 的地方:

protected int doReadMessages(List<Object> buf) throws Exception {    SocketChannel ch = SocketUtils.accept(javaChannel());    try {        if (ch != null) {            buf.add(new NioSocketChannel(this, ch));            return 1;        }    } catch (Throwable t) {        if (ch != null) {            try {                ch.close();            } catch (Throwable t2) {                logger.warn("Failed to close a socket.", t2);            }        }        throw t;    }    return 0;}

源码位置:

  • • 文件:transport/src/main/java/io/netty/channel/socket/nio/NioServerSocketChannel.java
  • • 方法:doReadMessages(List<Object>)
  • • 链接:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/socket/nio/NioServerSocketChannel.java

这段代码做了两件关键的事。

第一,调用 JDK ServerSocketChannel.accept() 得到 SocketChannel。由于监听 Channel 是非阻塞模式,即使 SelectionKey 表示可接入,循环中后续 accept 也可能返回 null。

第二,用 new NioSocketChannel(this, ch) 把 JDK Socket 包装成 Netty 子 Channel。构造函数会关联父 Channel,设置非阻塞模式,创建 Pipeline、Config、Unsafe 和写缓冲等运行时结构。

此时子 Channel 还没有注册 Worker EventLoop。它只是一个已经持有底层 Socket、但尚未建立线程归属的 Netty 对象。

4. 新 Channel 为什么要先经过父 Pipeline

NioMessageUnsafe.read() 不直接调用 Worker Group,而是对 readBuf 中每个子 Channel 执行:

pipeline.fireChannelRead(readBuf.get(i));

这里的 Pipeline 是父 NioServerSocketChannel 的 Pipeline。新连接被当成一条入站消息向后传播。最终处理它的是 ServerBootstrapAcceptor

这个设计非常优雅:Accept 仍然遵守 Netty 的统一事件模型。用户可以在 Acceptor 之前加入父 Channel Handler,实现连接接入日志、IP 黑白名单、限流或统计,而不需要修改底层 Transport。

当然,父 Pipeline 中的 Handler 必须足够轻。它运行在 Boss EventLoop,一次阻塞会直接影响后续连接接入。复杂鉴权、数据库查询和协议握手应放在子 Channel 的 Worker Pipeline 中,而不是 Boss 接入链。

5. ServerBootstrapAcceptor 应用 child 配置

ServerBootstrapAcceptor.channelRead() 收到的 msg 就是子 Channel:

public void channelRead(ChannelHandlerContext ctx, Object msg) {    final Channel child = (Channel) msg;    child.pipeline().addLast(childHandler);    setChannelOptions(child, childOptions, logger);    setAttributes(child, childAttrs);    try {        childGroup.register(child).addListener(new ChannelFutureListener() {            @Override            public void operationComplete(ChannelFuture future) {                if (!future.isSuccess()) {                    forceClose(child, future.cause());                }            }        });    } catch (Throwable t) {        forceClose(child, t);    }}

源码位置:

  • • 文件:transport/src/main/java/io/netty/bootstrap/ServerBootstrap.java
  • • 类:ServerBootstrapAcceptor
  • • 方法:channelRead(ChannelHandlerContext, Object)
  • • 链接:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/bootstrap/ServerBootstrap.java

先把 childHandler 加入子 Pipeline,再应用 childOption 和 Attribute,最后调用 childGroup.register(child)。通常 childHandler 是 ChannelInitializer,它会在子 Channel 注册过程中执行 initChannel(),把真正的业务 Handler 加入 Pipeline,随后把自己移除。

这里的失败处理很重要:如果注册 Worker 失败,必须强制关闭已 accept 的 Socket。否则内核连接已经建立,Netty 却没有线程负责它,会形成资源泄漏和客户端悬挂。

6. Worker 注册决定连接的长期线程归属

childGroup.register(child) 会选择一个 Worker EventLoop。选择结果写入 child.eventLoop,并在该 EventLoop 中执行 register0(),把 JDK SocketChannel 注册到 Worker Selector。

从此通常满足:

child.eventLoop() == selectedWorkerEventLoop

连接不会在每次请求时重新选择线程。固定归属带来三个后果:

  1. 1. 同一 Channel 的事件顺序稳定。
  2. 2. Handler 可以保存连接级状态而少用锁。
  3. 3. EventLoop 阻塞会影响它管理的全部 Channel。

Worker 的连接分配一般按轮询进行,不会根据每条连接的实际流量动态迁移。因此长连接业务里,即使连接数分布均匀,流量也可能不均匀:某个 Worker 恰好绑定了几个超级活跃连接,就可能比其他线程繁忙。排查时不能只看总连接数,还要观察各 EventLoop 的任务队列、CPU 和处理延迟。

7. ChannelInitializer 在注册期间完成业务 Pipeline

子 Channel 的 childHandler 常见写法是:

new ChannelInitializer<SocketChannel>() {    @Override    protected void initChannel(SocketChannel ch) {        ChannelPipeline p = ch.pipeline();        p.addLast(new LengthFieldBasedFrameDecoder(...));        p.addLast(new MessageDecoder());        p.addLast(new MessageEncoder());        p.addLast(new BusinessHandler());    }}

ChannelInitializer 本身是一个临时 Handler。它在 handlerAdded 或 channelRegistered 时检查是否需要初始化,调用 initChannel() 后将自身从 Pipeline 移除。这样每个子 Channel 都得到独立的 Handler 链。

要注意 Handler 的共享规则。标记 @Sharable 的无状态 Handler 可以复用同一实例;带可变成员状态且未设计并发安全的 Handler,通常应该为每条 Channel 创建新实例。否则多个连接可能共享解析状态、用户身份或计数器,产生非常隐蔽的数据串线。

Decoder 往往不是 Sharable,因为它可能持有累积 Buffer。业务 Handler 是否可共享,要看状态是否放在 Channel Attribute、请求对象或线程安全结构中。

8. 注册完成后触发 registered、active 与首次 read

子 Channel 已经由 TCP accept 创建,因此其底层 Socket 通常处于连接状态。Worker 注册完成后,register0() 触发:

handlerAdded  -> channelRegistered  -> channelActive  -> beginRead / 设置 OP_READ

channelActive 表示连接已经可用,业务可以发送欢迎消息、发起 TLS 握手或登记在线会话。随后,如果 AUTO_READ 开启,Netty 会调用 beginRead(),最终把 OP_READ 加入 SelectionKey 的 interestOps。

这也解释了为什么关闭 AUTO_READ 可以实现入站背压:连接注册和激活仍然发生,但没有显式 channel.read() 时,Netty不会持续关注或触发下一轮读取。业务可以在处理完当前消息后再请求下一次 read。

9. Accept 异常为什么会临时关闭 AutoRead

ServerBootstrapAcceptor 对连接接入异常有一段防护逻辑:如果父 Channel 出现某些异常,会临时把 autoRead 设为 false,并在短暂延迟后恢复。

目的不是“修复异常”,而是避免资源耗尽时陷入高速失败循环。典型场景是文件描述符耗尽:Selector 不断报告可 Accept,accept() 又不断失败。如果继续无节制接入,会占满 CPU 和日志,服务无法恢复。

临时暂停 Accept 给系统留下回收资源、运维介入和已有连接完成的时间。这体现了 Netty 很多稳定性设计的共同思路:发生系统性错误时,先降低事件产生速度,而不是在热循环中反复尝试。

生产环境仍应配置:

  • • 足够的文件描述符上限。
  • • 合理的 SO_BACKLOG
  • • 连接建立速率限制。
  • • 连接数和 Accept 失败监控。
  • • 对 Too many open files 等异常单独告警。

10. 主从 Reactor 的收益与边界

Boss/Worker 分离的收益是让连接接入与数据处理互不直接阻塞。Boss 的工作非常短:accept、包装 Channel、沿 Pipeline 传播、触发注册。Worker 承担连接生命周期内的大部分读写。

但它并不是万能隔离。

如果 childHandler 初始化很重,它仍然在注册路径中拖慢 Worker。如果父 Pipeline 做远程鉴权,它会拖慢 Boss。如果 Worker 数量设置过少,所有连接仍会竞争少数 EventLoop。如果业务任务直接运行在 EventLoop,慢操作仍会造成队头阻塞。

对于普通服务,Boss 线程数通常不需要很多,因为 accept 工作量小;Worker 数量则要结合 CPU、协议处理成本、连接活跃度和是否把业务卸载到独立线程池评估。简单照搬“CPU 核数乘二”并不可靠。

11. 排查连接接入问题的源码化方法

当客户端“TCP 能连但业务无响应”,可以按以下边界定位:

  1. 1. 端口是否已经 bind 成功,ServerChannel 是否 Active。
  2. 2. Boss Selector 是否收到 OP_ACCEPT
  3. 3. doReadMessages() 是否创建 NioSocketChannel。
  4. 4. 父 Pipeline 是否把 Channel 传播到 Acceptor。
  5. 5. childHandler 初始化是否抛异常。
  6. 6. childGroup.register 是否成功。
  7. 7. 子 Channel 是否触发 channelRegistered/channelActive
  8. 8. SelectionKey 是否关注 OP_READ

如果连接数达到阈值后新连接超时,检查 backlog、文件描述符、Boss EventLoop 是否阻塞。如果连接已建立但 Handler 没有 active 日志,重点检查 Worker 注册和 ChannelInitializer。如果 active 正常但读不到数据,再进入下一篇的 EventLoop 与读路径。

12. 总结:新连接是一条经过 Pipeline 的消息

Netty 连接接入链最值得记住的设计,是“新 Channel 先作为消息经过父 Pipeline”。它把底层 Accept 与上层扩展统一到同一套事件机制中。

完整主链是:

Boss Selector OP_ACCEPT  -> NioMessageUnsafe.read  -> NioServerSocketChannel.doReadMessages  -> new NioSocketChannel  -> parentPipeline.fireChannelRead(child)  -> ServerBootstrapAcceptor.channelRead  -> child pipeline/options/attrs  -> childGroup.register  -> Worker Selector  -> channelActive + OP_READ

下一篇会进入 Worker 的心脏 NioEventLoop.run(),解释它如何在一个线程里同时处理 Selector I/O、普通任务、定时任务和唤醒信号,以及 ioRatio 为什么会直接影响延迟与吞吐。

参考资料

  1. 1. AbstractNioMessageChannel.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/nio/AbstractNioMessageChannel.java
  2. 2. NioServerSocketChannel.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/socket/nio/NioServerSocketChannel.java
  3. 3. ServerBootstrap.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/bootstrap/ServerBootstrap.java
  4. 4. ChannelInitializer.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/ChannelInitializer.java