Skip to content
Charles Shao
Go back

Netty 深挖(开篇)· Reactor 模型与 EventLoop 线程模型

views

广告竞价 Server 是典型的「海量长短连接 + 亚毫秒延迟 + 单机几万 QPS」场景:交易所(ADX)把请求灌进来,tmax 只有 80–150ms,晚一毫秒就是一次 no-bid。这种流量特征下,几乎没有人会去裸写 Java NIO——大厂的竞价引擎、网关、RPC 框架(Dubbo、gRPC-Java、RocketMQ)底层清一色是 Netty

这一篇先讲清最底层的两块——Reactor 模型EventLoop 线程模型。Netty 的性能与心智负担,都建立在「主从多 Reactor + 一个 EventLoop 绑定一个线程、串行无锁地跑完一条连接的一生」之上;后面的 ByteBuf、Pipeline、背压都落在这块地基上。

TL;DR

Table of contents

Open Table of contents

1. 为什么竞价 Server 用 Netty,而不是裸写 NIO

连接与 IO 模型 那篇里,我们把 BIO / NIO / 多路复用讲过一遍:BIO 是「一连接一线程」,几万连接就要几万线程,上下文切换和内存直接压垮机器;NIO 靠一个 Selector 用一个线程管成千上万条连接,是高并发服务器的唯一现实解。竞价 Server 面对的是海量并发连接,答案毫无悬念——必须走 NIO 多路复用。

问题是:Java 原生 NIO 能用,但极难用对。 真要拿 java.nio 裸写一个生产级服务器,你会依次踩到这些坑:

Netty 的价值,就是把上面这些一次性封装成稳定、久经生产考验的 API:Selector 空轮询它替你兜底、事件分发变成 pipeline 回调、ByteBuffer 换成好用得多的 ByteBuf(第 2 篇)、粘包拆包给你现成的解码器(第 3 篇)。你写的不再是「怎么操作 socket」,而是「收到一条消息我要干什么」。

对竞价这种 tmax 只有百毫秒级、单机几万 QPS 的场景,Netty 还额外重要在两点:极致的低延迟(零拷贝、对象池、无锁串行)和可控的线程模型(下面几节的主角)。竞价链路上任何抖动都会转化为 no-bid,容不得原生 NIO 那种「能跑但不稳」。

2. Reactor 模式的三级演进

Netty 的线程模型本质是 Reactor 模式的工业级实现。核心思想:用一个(或几个)专门的线程通过多路复用器监听事件,事件就绪后「反应式」地分发给对应的处理器。 它有三种形态,层层递进:

Reactor 三级演进。左:单 Reactor 单线程(accept+IO+业务全挤一线程,任一阻塞全停);中:单 Reactor 多线程(IO 仍单线程、业务丢线程池,IO 仍是瓶颈);右:主从多 Reactor(boss 只 accept、worker 组跑 IO,Netty 默认)。底部说明业务默认仍在 worker IO 线程,耗时逻辑仍要外移。

2.1 单 Reactor 单线程

一个线程包揽一切:accept 新连接、read/write 所有连接的 IO、以及业务处理,全在这一个线程的死循环里。

2.2 单 Reactor 多线程

把「业务处理」从 Reactor 线程里摘出来:Reactor 线程仍然单线程负责 accept + 所有 IO 读写,但读到完整请求后,把业务逻辑丢给一个 worker 线程池去跑,处理完再回到 Reactor 线程去 write

2.3 主从多 Reactor(Netty 默认)

再进一步:acceptread/write 也拆开,用两组 Reactor。

这就是 Netty 服务端 ServerBootstrap 的默认模型

主从多 Reactor 架构图:最左侧是三条客户端连接,通过 TCP connect 指向 bossGroup(Boss EventLoop 只监听 OP_ACCEPT);Boss 通过 accept → 轮询 register 把连接分给 workerGroup 内多个 Worker EventLoop,各自用线程+Selector 绑定若干 Channel。底部说明 boss 只接连接,此后连接一生不换 worker。

主从多 Reactor:boss 只接连接、worker 组各跑一个 Selector 复用大量连接。

它的好处是各司其职、水平扩展:boss 专心接连接(accept 很轻、1~2 个线程够了),worker 组按 CPU 核数并行处理海量连接的 IO。竞价 Server 每秒要接大量短连接又要处理海量并发 IO,正需要这种「接入」与「IO」分离的架构。

注意一个常见误解:主从多 Reactor 里,业务处理默认仍然在 worker 的 IO 线程上跑(在 pipeline 的 handler 里)。如果业务耗时,还是要像 §2.2 那样再外移到业务线程池——这正是本篇 war story 的核心。

3. Netty 核心组件与它们的关系

在讲线程模型细节前,先把 Netty 的几个核心抽象和它们的关系理清。它们环环相扣,理解了关系,代码就读得懂了。

它们的从属关系:

核心组件从属关系。自左向右:EventLoopGroup(boss/worker)→ EventLoop(1 线程+Selector)→ Channel(一生绑一个 EL)→ ChannelPipeline / ChannelHandler+Context。箭头标注含多个、绑定 1:N、含 1 条。底部强调一个 Channel 只绑一个 EventLoop。

关键数量关系:一个 EventLoopGroup 有多个 EventLoop;一个 EventLoop 可以绑定多个 Channel(1:N);但一个 Channel 只绑定一个 EventLoop(N:1)。 最后这条是整套无锁模型的关键,下一节展开。

4. EventLoop 线程模型:一线程、串行、无锁

这是 Netty 最精髓、也最容易用错的一节。三个词概括:一个 EventLoop 绑定一个线程、串行处理、无锁化。

EventLoop 线程模型图:顶部一条横向流程展示 NioEventLoop.run() 的死循环三步——① select()(等待 IO 就绪事件)→ ② processSelectedKeys()(处理就绪 Channel 的读写)→ ③ runAllTasks()(taskQueue + 到期定时任务),并有一条虚线返回箭头标注「死循环:三步跑完再来一轮,ioRatio 控制 IO 与任务的耗时配比」。下方左侧绿色框「绑定此 EventLoop 的 Channel · 串行、无锁」里有 Channel A、Channel B、Channel C 三个组件,并说明「同一时刻只处理一个 Channel → 无需加锁保护 Channel 状态」;下方右侧橙色框「EventLoop 自己的两个任务队列」里有 taskQueue(MPSC 无锁队列)和 scheduledTaskQueue(定时任务·小顶堆)两个组件,说明由 execute()/schedule() 提交、③ runAllTasks 执行。

一个 EventLoop 就是一个死循环线程:轮流处理就绪 Channel 的 IO 与自己队列里的任务;它绑定的多个 Channel 在这一个线程里串行流转,因此天然无锁。

4.1 一个 EventLoop 绑定一个线程

NioEventLoop 继承自 SingleThreadEventExecutor——顾名思义,每个 EventLoop 内部只有一个线程,这个线程被懒启动(第一个任务提交时才 start),此后终生只跑这一个 run() 死循环。它的骨架(简化)大致是:

// io.netty.channel.nio.NioEventLoop#run(高度简化)
protected void run() {
    for (;;) {
        // ① 多路复用:等待就绪的 IO 事件(带空轮询 bug 的兜底)
        select(...);

        // ② 处理就绪 Channel 的 IO:accept / read / write
        processSelectedKeys();

        // ③ 跑普通任务 + 到期的定时任务;ioRatio 控制这一步与 ② 的耗时配比
        runAllTasks(...);
    }
}

三步循环往复,就是一个 EventLoop 线程的全部人生。ioRatio(默认 50)控制 ② 和 ③ 的耗时比例:ioRatio=50 表示花在处理任务上的时间约等于花在 IO 上的时间;调高更偏向 IO,调低更偏向任务。

4.2 一个 Channel 一生绑定同一个 EventLoop

这是最关键的设计:一条连接(Channel)在 register 到某个 EventLoop 之后,它整个生命周期内的所有 IO 事件、pipeline 传播、以及提交给它的任务,都在这个 EventLoop 的唯一线程里执行——一生不换线程。

带来的直接后果是串行:同一个 Channel 上的 channelReadwritechannelInactive 等回调,永远由同一个线程按顺序执行,不可能并发。于是:

4.3 inEventLoop():跨线程写的安全保证

那如果别的线程(比如业务线程池处理完,要把结果写回连接)调用了 channel.writeAndFlush(resp) 呢?Netty 用一个巧妙的判断保证「同一个 Channel 永远只被它绑定的那个 EventLoop 线程操作」:

// AbstractChannelHandlerContext#write(示意)
final EventExecutor executor = ctx.executor();
if (executor.inEventLoop()) {
    // 当前就是该 Channel 绑定的 EventLoop 线程 → 直接执行写
    write(msg, ...);
} else {
    // 是别的线程 → 封装成任务,塞进该 EventLoop 的 taskQueue,排队执行
    executor.execute(() -> write(msg, ...));
}

inEventLoop() 就是判断「当前线程 == 本 EventLoop 的线程」。任意线程都能安全地调 channel.write():如果恰好在 IO 线程内就直接写,否则自动转成任务排队到 IO 线程去写。这样对同一个 Channel 的所有操作最终都收敛到同一个线程串行执行——开发者无感知,就得到了线程安全。这也解释了为什么 taskQueue 用的是 MPSC(多生产者单消费者) 无锁队列:多个业务线程可能同时往里塞任务,但只有 EventLoop 这一个线程消费。

5. 服务端启动骨架代码

把前面的组件拼起来,一个 Netty 服务端的标准启动代码长这样。竞价 Server 的接入层骨架基本就是它的变体:

// boss:只接连接,1 个线程足够;worker:处理 IO,默认 2 × CPU
EventLoopGroup bossGroup = new NioEventLoopGroup(1);
EventLoopGroup workerGroup = new NioEventLoopGroup();   // 不传参 = 2 × CPU 核数

try {
    ServerBootstrap b = new ServerBootstrap();
    b.group(bossGroup, workerGroup)                     // 配置主从两组 Reactor
     .channel(NioServerSocketChannel.class)             // 用 NIO 的服务端 Channel
     .option(ChannelOption.SO_BACKLOG, 1024)            // 作用于 boss 的监听 socket
     .childOption(ChannelOption.TCP_NODELAY, true)      // 作用于每条被 accept 的连接
     .childHandler(new ChannelInitializer<SocketChannel>() {
         @Override
         protected void initChannel(SocketChannel ch) {
             // 每来一条新连接,就为它初始化一条 pipeline
             ChannelPipeline p = ch.pipeline();
             p.addLast(new HttpServerCodec());          // 编解码(第 3 篇)
             p.addLast(new HttpObjectAggregator(65536));
             p.addLast(new BidRequestHandler());        // 业务:处理竞价请求
         }
     });

    // 绑定端口,sync() 阻塞等绑定完成
    ChannelFuture f = b.bind(8080).sync();
    f.channel().closeFuture().sync();                   // 阻塞等服务端关闭
} finally {
    workerGroup.shutdownGracefully();
    bossGroup.shutdownGracefully();
}

读这段代码要抓住三个容易混淆的点:

6. EventLoop 不只跑 IO:普通任务与定时任务

回看 §4.1 的 run() 循环,第 ③ 步 runAllTasks() 说明 EventLoop 不止处理 IO,还处理提交给它的任务。这有两个入口:

定时任务在网络编程里无处不在,而且由 EventLoop 自己调度,天然线程安全、无需外部 ScheduledExecutorService

// 在某个 handler 里,给这条连接设一个「空闲超时」定时任务
ctx.executor().schedule(() -> {
    if (System.nanoTime() - lastActiveNanos > TIMEOUT_NANOS) {
        ctx.close();     // 超时未收到数据,主动关连接
    }
}, 5, TimeUnit.SECONDS);

竞价链路里,这类定时任务用得极多:连接空闲检测IdleStateHandler 底层就是它)、心跳保活响应超时兜底tmax 到点没出价就返回 no-bid)、断线重连的退避调度。因为它们和 IO 事件共用同一个 EventLoop 线程,所以定时任务里同样不能阻塞——否则又回到 §8 的坑。ioRatio 就是用来平衡「IO 处理」和「任务处理」谁多占点时间的旋钮。

7. 一个请求在 Netty 里的完整流转

把前面所有组件串成一条时间线,看一条连接从建立到收发的完整生命周期:

请求流转时序图,四个参与方从左到右为:客户端 Client、Boss EventLoop(Acceptor)、Worker EventLoop(1 线程·Selector)、ChannelPipeline(Handler 链)。步骤依次为:1. 客户端与 Boss 完成 TCP 三次握手,服务端 OP_ACCEPT 就绪;2. Boss 自身 accept() 得到一条新的 SocketChannel;3. Boss 按轮询把 Channel 注册到某个 worker(此后一生绑定它);随后一个说明框指出「这一步之后 boss 彻底退出,该连接的全部 IO 只由这一个 worker EventLoop 处理」;4. 客户端发送请求字节到 Worker,OP_READ 就绪、读进 ByteBuf;5. Worker 触发 fireChannelRead,字节从 head 沿 inbound 向后传播到 Pipeline;6. Pipeline 自身用解码器把 ByteBuf 解成业务消息、业务 Handler 处理请求;7. Pipeline 通过 ctx.writeAndFlush(响应) 沿 outbound 从 tail 回到 head,交回 Worker;8. Worker 把响应编码为 ByteBuf 写回 socket、flush 出栈到客户端。底部红色说明框强调:整条链路都在同一个 worker EventLoop 线程里串行完成——绝不能在 Handler 里阻塞它。

一条连接的一生:boss 只在 accept + register 时出现,之后所有读、pipeline 传播、写都在同一个 worker EventLoop 线程里串行完成。

拆开来讲:

  1. 握手 + accept:客户端 TCP 三次握手完成,boss 的 Selector OP_ACCEPT 就绪,boss accept() 拿到一条新的 NioSocketChannel
  2. register(关键分界点):boss 从 workerGroup 里按轮询挑一个 EventLoop,把这条 Channel register 上去。这一步之后 boss 就彻底退场——连接的一切都交给这个 worker EventLoop,一生不换。
  3. read:客户端发来数据,worker 的 Selector OP_READ 就绪,worker 线程把字节从 socket 读进一个 ByteBuf
  4. inbound 传播:worker 触发 pipeline.fireChannelRead(byteBuf),数据从 pipeline 的 head 开始,沿入站方向依次经过解码器、业务 handler;解码器把字节流拆成完整消息(第 3 篇),业务 handler 处理请求。
  5. outbound 回写:业务 handler 调 ctx.writeAndFlush(response),数据沿出站方向tail 回流到 head,经过编码器把消息编成 ByteBuf
  6. write + flush:到达 head 后,worker 把 ByteBuf 写回 socket、flush 出栈,客户端收到响应。

全程(第 3 步到第 6 步)都在同一个 worker EventLoop 的唯一线程里串行完成。 记住这句话,下一节的事故就是它的反面教材。

延伸阅读


views
Share this post on:

Previous Post
Netty 深挖 · ByteBuf、引用计数与零拷贝内存模型
Next Post
并发协作与异步编排:闭锁、信号量与 CompletableFuture