Netty 的 ChannelPipeline
如果说 Reactor 负责“发现事件并分发”,那么 ChannelPipeline 负责“事件到了之后怎么一层层处理”。
为什么 Netty 要设计 Pipeline
网络请求处理通常不是一步完成的。
比如一个消息可能会经过:
- 字节流读取。
- 粘包拆包处理。
- 字符串解码。
- JSON 解析。
- 权限校验。
- 业务处理。
- 响应编码。
处理链模型
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 的“流水线和责任链”。
它让网络处理过程具备了三个优势:
- 分层清晰。
- 可插拔。
- 易于扩展和复用。
常见风险和排查
| 现象 | 常见原因 | 排查方式 |
|---|---|---|
| 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 就收不到这条消息。
