本节摘要:ChannelPipeline 是每条连接专属的双向工序表,入站事件从 head 流向 tail,出站事件反向。本节讲透事件传播机制、Handler 基类的选择、动态热插拔与用户自定义事件,并用 EmbeddedChannel 演示不启网络的流水线测试法。
每个 Channel 内部都有一条 Pipeline,两端是 Netty 内置的 head 与 tail 上下文,中间是你 addLast/addFirst 塞进去的 Handler 序列。它最容易被忽视也最关键的特性是双向:
注意图中出站事件的路径:write 从业务代码发出后,只会穿过位于它之前加入的出站工位。把编码器 addLast 在业务 Handler 之后,出站消息根本不会经过它——"编码器不生效"十有八九是位置问题。
事件是否继续流动,完全取决于当前工位有没有"放行":
public class AuditHandler extends ChannelInboundHandlerAdapter { @Override public void channelRead(ChannelHandlerContext ctx, Object msg) { System.out.println("工位一审计到数据"); ctx.fireChannelRead(msg); // 放行:下一个入站工位才能收到 // 若此处不调用,事件在此终止,后面的业务 Handler 全部失聪 } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.fireExceptionCaught(cause); // 异常同样靠 fire 向下游传递 } }
忘记 fire 是新手第一大坑:现象通常是"日志显示读到了数据,但业务没反应"。反过来,SimpleChannelInboundHandler 的 channelRead0 会自动释放消息并继续传播(细节涉及引用计数,第四章展开)。
基类怎么选,一张表说清:
| 基类 | 方向 | 自动释放 msg | 适合 |
|---|---|---|---|
| ChannelInboundHandlerAdapter | 入站 | 否,需手动 release | 中间审计、透传工位 |
| SimpleChannelInboundHandler | 入站 | 是 | 终点业务工位 |
| ChannelOutboundHandlerAdapter | 出站 | 否 | 出站过滤、改写 |
| MessageToByteEncoder | 出站 | 是 | 编码器基类(第五章) |
| ByteToMessageDecoder | 入站 | 是 | 解码器基类(第五章) |
另一个高频注解是 @ChannelSharable。默认每个 Handler 实例只服务一条连接(有状态);无状态的审计、日志类 Handler 标注该注解后可以全局单例,省内存又免去重复创建。有状态的 Handler 加了 @Sharable 是并发事故——第六章从线程角度还会再审判一次这条纪律。
Pipeline 是运行时可变的。最经典的应用是"首次握手协议成功后,撤掉握手 Handler、换上正式工序":
// 握手完成后:先移除自己,再装上正式业务工序 public class HandshakeHandler extends SimpleChannelInboundHandler<String> { @Override protected void channelRead0(ChannelHandlerContext ctx, String msg) { if ("HELLO".equals(msg.trim())) { ChannelPipeline p = ctx.pipeline(); p.remove(this); // 拆掉握手工位 p.addLast(new LengthFieldDecoder(...), // 之后才挂正式解码 new BusinessHandler()); ctx.writeAndFlush("WELCOME\n"); } else { ctx.close(); // 握手失败直接停线 } } }
可用的动作还有 addBefore、addAfter、replace、get(按名取)、remove(按名删)。因为修改操作会被投递到该 Channel 的 EventLoop 执行,跨线程调用也是安全的。
除了读写在流水线上流动,你也可以让自定义信号走同一条路。典型场景:空闲检测(第七章心跳的前置知识)——Netty 自带的 IdleStateHandler 在连接空闲时发射 IdleStateEvent,业务工位接住它做断连或发心跳:
ch.pipeline() .addLast(new IdleStateHandler(60, 0, 0)) // 60 秒没读到数据就发空闲事件 .addLast(new IdleEventHandler()); public class IdleEventHandler extends ChannelInboundHandlerAdapter { @Override public void userEventTriggered(ChannelHandlerContext ctx, Object evt) { if (evt instanceof IdleStateEvent) { System.out.println("连接空闲,主动断开"); ctx.close(); } else { ctx.fireUserEventTriggered(evt); // 不认识的事件继续传递 } } }
自定义事件同样能携带数据:定义一个普通类(比如 HandshakeSuccessEvent),在某个工位 ctx.pipeline().fireUserEventTriggered(event),下游所有工位的 userEventTriggered 都有机会接住。它比在 Handler 之间互相持有引用干净得多,是流水线内部的标准通信方式。
Pipeline 的行为可以完全离线验证,这是写单测的正道:
EmbeddedChannel ch = new EmbeddedChannel( new StringDecoder(), new StringEncoder(), new EchoHandler()); // 模拟字节进站:写字节到入站方向 ch.writeInbound(Unpooled.copiedBuffer("hello\n", StandardCharsets.UTF_8)); // 读取流到终点的入站消息 assert ch.readInbound().equals("hello"); // 模拟出站:业务 write 的内容会穿过编码器 ch.writeOutbound("world\n"); ByteBuf out = ch.readOutbound(); System.out.println(out.toString(StandardCharsets.UTF_8)); // world ch.finishAndReleaseAll();
解码器、业务 Handler、自定义事件都能这样在毫秒级跑完,不用起端口、不用等握手。第五章写完解码器后,这就是你的主力验证工具。
⚠️ 常见坑:在 Handler 里做耗时操作(查库、调下游接口)会卡住整条流水线——因为整条 Pipeline 跑在唯一的 EventLoop 线程上。重活必须挪到业务线程池,第六章给出完整方案。
💡 关键直觉:Pipeline 的本质是"责任链 + 双向过滤器"。入站像流水线前段的车床逐层加工原料,出站像包装段逐层打包成品;位置决定经历,fire 决定生死。