Skip to content

Netty 的 ChannelPipeline

如果说 Reactor 负责“发现事件并分发”,那么 ChannelPipeline 负责“事件到了之后怎么一层层处理”。

为什么 Netty 要设计 Pipeline

网络请求处理通常不是一步完成的。

比如一个消息可能会经过:

  1. 字节流读取。
  2. 粘包拆包处理。
  3. 字符串解码。
  4. JSON 解析。
  5. 权限校验。
  6. 业务处理。
  7. 响应编码。

处理链模型

mermaid
flowchart TD
    A["原始字节数据"] --> B["入站 Handler 1<br/>拆包"]
    B --> C["入站 Handler 2<br/>解码"]
    C --> D["业务 Handler<br/>校验和处理"]
    D --> E["出站 Handler 1<br/>编码"]
    E --> F["出站 Handler 2<br/>压缩 / 加密"]
    F --> G["写回客户端"]

入站和出站怎么理解

入站

数据从网络进到应用。

出站

数据从应用写回网络。

核心原理

Pipeline 是一条双向责任链。入站事件从 Head 往 Tail 方向传播,出站事件从当前节点往 Head 方向传播。理解这个方向非常重要,因为很多“Handler 不生效”的问题都不是代码没写,而是放错了位置或没有继续传播事件。

mermaid
flowchart TD
    A["网络读事件"] --> B["HeadContext"]
    B --> C["Inbound Handler<br/>Decoder / Auth / Biz"]
    C --> D["TailContext"]
    E["业务写响应"] --> F["Outbound Handler<br/>Encoder / Flush"]
    F --> G["HeadContext 写入 Socket"]

关键规则:

规则为什么如果不会怎样
入站 Handler 需要调用 ctx.fireChannelRead 才会继续往后传Netty 不知道你处理完后是否还要交给后续 Handler后面的业务 Handler 收不到消息
出站写通常用 ctx.writeAndFlush从当前 Handler 的前一个出站节点开始找编码器位置不对时,响应可能没有被编码
解码器放在业务 Handler 前面业务代码应该处理对象,不应该处理原始字节业务层充满字节解析逻辑,难维护
异常要在 Pipeline 中兜底网络异常、解码异常、业务异常都可能发生连接泄漏、日志缺失、客户端无响应

一次消息处理的流转

mermaid
sequenceDiagram
    participant Client as 客户端
    participant Decoder as 解码器
    participant Auth as 鉴权Handler
    participant Biz as 业务Handler
    participant Encoder as 编码器

    Client->>Decoder: 发送字节流
    Decoder->>Auth: 转成消息对象
    Auth->>Biz: 鉴权通过后继续传递
    Biz->>Encoder: 生成响应对象
    Encoder-->>Client: 编码后写回

本章小结

Pipeline 可以理解成 Netty 的“流水线和责任链”。

它让网络处理过程具备了三个优势:

  1. 分层清晰。
  2. 可插拔。
  3. 易于扩展和复用。

常见风险和排查

现象常见原因排查方式
Handler 没有收到消息前一个入站 Handler 没有 fireChannelRead在每个 Handler 打日志确认传播位置
编码器没有执行出站 Handler 顺序放错,或写出位置不对确认 Encoder 是否在写出路径上
消息被重复释放SimpleChannelInboundHandler 自动释放后又手动释放明确是否需要 retain,不要重复 release
鉴权通过后业务无响应鉴权 Handler 消费消息后没有继续传播调用 ctx.fireChannelRead 传递处理后的对象
异常后连接不关闭没有实现 exceptionCaught记录日志并按协议决定是否 ctx.close()

代码 Demo:自定义入站 Handler

java
public class AuthHandler extends ChannelInboundHandlerAdapter {
    @Override
    public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
        String text = (String) msg;
        if (!text.startsWith("token:")) {
            ctx.writeAndFlush("unauthorized");
            ctx.close();
            return;
        }
        ctx.fireChannelRead(text.substring("token:".length()));
    }
}

注册到 Pipeline:

java
ch.pipeline().addLast(new StringDecoder());
ch.pipeline().addLast(new StringEncoder());
ch.pipeline().addLast(new AuthHandler());
ch.pipeline().addLast(new BizHandler());

ctx.fireChannelRead(...) 表示把处理后的消息继续交给下一个入站 Handler。如果不调用,后面的 Handler 就收不到这条消息。