Netty 商业场景训练营
这页把 Netty 放到真实商业系统里学习。目标不是背 Bootstrap、Channel、EventLoop 这些类名,而是能解释一个长连接怎么接入、一条 TCP 字节流怎么变成业务消息、业务处理为什么不能阻塞 IO 线程、写缓冲为什么会堆积、ByteBuf 为什么会泄漏、线上连接假死和消息错乱怎么排查。
训练目标
学完这页,你要能做到:
- 解释 Netty 服务端从绑定端口、接收连接、注册 Channel 到处理消息的全过程。
- 解释 Reactor、BossGroup、WorkerGroup、EventLoop、Channel、Pipeline、Handler 的关系。
- 解释 TCP 粘包半包为什么一定要用协议边界解决。
- 写出长度字段协议的编解码 Demo。
- 解释为什么不能在 EventLoop 里查数据库、调 HTTP 或做大计算。
- 解释心跳、空闲检测、重连、连接管理和用户会话绑定。
- 解释 ByteBuf 池化、引用计数和泄漏排查。
- 排查连接数上不去、消息错乱、延迟升高、写缓冲堆积、直接内存 OOM。
商业场景总览
以“医疗设备长连接采集网关”为例,设备通过 TCP 长连接把采集数据推给平台:
mermaid
flowchart TD
A["医疗设备"] --> B["Netty TCP 网关"]
B --> C["协议拆包和解码"]
C --> D["设备鉴权和会话绑定"]
D --> E["业务校验"]
E --> F["投递业务线程池"]
F --> G["写入 MQ 或数据库"]
G --> H["返回 ACK"]
B --> I["心跳检测和连接管理"]这里涉及的不是普通 HTTP CRUD,而是高并发连接、协议边界、异步处理、背压和连接生命周期。
一次连接从进入到可读的全过程
mermaid
flowchart TD
A["ServerBootstrap bind 端口"] --> B["BossGroup 监听 ServerSocketChannel"]
B --> C["客户端发起 TCP 连接"]
C --> D["Boss EventLoop accept"]
D --> E["创建 SocketChannel"]
E --> F["注册到某个 Worker EventLoop"]
F --> G["初始化 ChannelPipeline"]
G --> H["触发 channelActive"]
H --> I["等待读写事件"]关键点:
- BossGroup 只负责接入连接,不负责处理所有业务。
- WorkerGroup 负责已建立连接的读写事件。
- 一个 Channel 通常固定绑定一个 EventLoop。
- Pipeline 在 Channel 初始化时建立,后续消息按 Handler 顺序流动。
- 如果 BossGroup 被阻塞,新连接接入会变慢;如果 Worker EventLoop 被阻塞,它负责的多个连接都会变慢。
一条消息从字节到业务对象的全过程
mermaid
flowchart TD
A["网络收到字节"] --> B["ByteBuf"]
B --> C["拆包器判断完整帧"]
C --> D{"是否完整"}
D -- "否" --> E["缓存等待更多字节"]
D -- "是" --> F["Decoder 解码业务对象"]
F --> G["鉴权和协议校验"]
G --> H["业务 Handler"]
H --> I["Encoder 编码响应"]
I --> J["写回 Socket"]TCP 是字节流协议,不会保留应用层消息边界。应用层必须自己定义“一个消息到哪里结束”。
协议设计:长度字段协议
商业系统最常用的是长度字段协议:
text
magic(2字节) | version(1字节) | type(1字节) | requestId(8字节) | length(4字节) | body(N字节)| 字段 | 作用 |
|---|---|
| magic | 快速识别是否是本协议 |
| version | 协议版本,方便升级 |
| type | 登录、心跳、数据上报、ACK 等消息类型 |
| requestId | 请求和响应关联 |
| length | body 长度,解决粘包半包 |
| body | JSON、Protobuf 或自定义二进制内容 |
为什么要有 length?因为一次 channelRead 可能读到半条消息,也可能读到多条消息。没有长度字段,解码器不知道什么时候可以把字节交给业务层。
Demo 一:长度字段解码器
java
public class DeviceFrameDecoder extends ByteToMessageDecoder {
private static final short MAGIC = (short) 0xCAFE;
private static final int HEADER_LENGTH = 16;
@Override
protected void decode(ChannelHandlerContext ctx, ByteBuf in, List<Object> out) {
if (in.readableBytes() < HEADER_LENGTH) {
return;
}
in.markReaderIndex();
short magic = in.readShort();
if (magic != MAGIC) {
ctx.close();
return;
}
byte version = in.readByte();
byte type = in.readByte();
long requestId = in.readLong();
int length = in.readInt();
if (length < 0 || length > 1024 * 1024) {
ctx.close();
return;
}
if (in.readableBytes() < length) {
in.resetReaderIndex();
return;
}
ByteBuf body = in.readRetainedSlice(length);
try {
DeviceMessage message = new DeviceMessage(version, type, requestId, body);
out.add(message);
} catch (Exception ex) {
body.release();
throw ex;
}
}
}这个解码器体现了几个关键点:
- 不够头部长度就返回,等待更多字节。
- 用
markReaderIndex和resetReaderIndex处理半包。 - 校验 magic 和 length,防止乱流量和超大包攻击。
readRetainedSlice会增加引用计数,后续必须明确释放。
Demo 二:Pipeline 顺序
java
public class DeviceChannelInitializer extends ChannelInitializer<SocketChannel> {
private final EventExecutorGroup businessGroup;
public DeviceChannelInitializer(EventExecutorGroup businessGroup) {
this.businessGroup = businessGroup;
}
@Override
protected void initChannel(SocketChannel ch) {
ch.pipeline()
.addLast("idle", new IdleStateHandler(60, 30, 0))
.addLast("frameDecoder", new DeviceFrameDecoder())
.addLast("messageDecoder", new DeviceMessageDecoder())
.addLast("encoder", new DeviceMessageEncoder())
.addLast("auth", new DeviceAuthHandler())
.addLast(businessGroup, "business", new DeviceBusinessHandler())
.addLast("exception", new DeviceExceptionHandler());
}
}顺序为什么重要:
| 顺序 | 作用 | 放错后果 |
|---|---|---|
| IdleStateHandler | 先检测空闲连接 | 心跳不生效或连接假死 |
| FrameDecoder | 先解决消息边界 | 后续 Decoder 拿到乱字节 |
| MessageDecoder | 字节转业务对象 | 业务 Handler 无法理解消息 |
| Encoder | 响应编码 | 写回对象无法变成字节 |
| AuthHandler | 鉴权 | 未认证设备可能进入业务 |
| BusinessHandler | 业务处理 | 不能阻塞 IO 线程 |
| ExceptionHandler | 统一异常 | 异常分支资源泄漏或连接不关闭 |
EventLoop 为什么不能阻塞
一个 EventLoop 可能负责多个 Channel:
mermaid
flowchart TD
A["EventLoop-1"] --> B["Channel A"]
A --> C["Channel B"]
A --> D["Channel C"]
B --> E["业务 Handler 慢 SQL"]
E --> F["EventLoop 被阻塞"]
F --> G["B/C/D 都不能及时读写"]错误示例:
java
protected void channelRead0(ChannelHandlerContext ctx, DeviceData data) {
deviceDataMapper.insert(data); // 慢数据库写入
String result = httpClient.call(); // 同步 HTTP 调用
ctx.writeAndFlush(result);
}正确做法:慢业务投递到有界业务线程池。
java
DefaultEventExecutorGroup businessGroup =
new DefaultEventExecutorGroup(16);
pipeline.addLast(businessGroup, "business", new DeviceBusinessHandler());业务线程池也不能无限大。要设置队列、拒绝策略、耗时监控和降级,否则只是把阻塞从 EventLoop 转移到另一个黑洞里。
Demo 三:设备登录、会话绑定和心跳
java
public class DeviceAuthHandler extends SimpleChannelInboundHandler<DeviceLoginMessage> {
@Override
protected void channelRead0(ChannelHandlerContext ctx, DeviceLoginMessage msg) {
if (!checkSignature(msg.deviceId(), msg.timestamp(), msg.signature())) {
ctx.writeAndFlush(DeviceResponse.fail("AUTH_FAILED"));
ctx.close();
return;
}
ctx.channel().attr(DeviceAttrs.DEVICE_ID).set(msg.deviceId());
DeviceSessionRegistry.bind(msg.deviceId(), ctx.channel());
ctx.writeAndFlush(DeviceResponse.ok("LOGIN_OK"));
}
@Override
public void channelInactive(ChannelHandlerContext ctx) {
String deviceId = ctx.channel().attr(DeviceAttrs.DEVICE_ID).get();
if (deviceId != null) {
DeviceSessionRegistry.unbind(deviceId, ctx.channel());
}
}
}心跳处理:
java
public class HeartbeatHandler extends ChannelInboundHandlerAdapter {
@Override
public void userEventTriggered(ChannelHandlerContext ctx, Object evt) {
if (evt instanceof IdleStateEvent event
&& event.state() == IdleState.READER_IDLE) {
ctx.close();
}
}
}为什么要会话绑定?
- 平台要知道设备当前连接在哪个 Channel。
- 下行指令需要找到设备连接。
- 设备断开时要清理状态,避免给旧 Channel 发消息。
- 多实例部署时还要结合 Redis、注册中心或路由表记录设备在哪个网关实例。
写缓冲堆积和背压
如果服务端写得比客户端收得快,Netty 的 outbound buffer 会堆积。
mermaid
flowchart TD
A["服务端持续写消息"] --> B["客户端接收慢"]
B --> C["Socket 发送缓冲区满"]
C --> D["ChannelOutboundBuffer 堆积"]
D --> E["直接内存上涨"]
E --> F["延迟升高或 OOM"]处理方式:
java
if (!ctx.channel().isWritable()) {
metrics.incrementDropCount();
return;
}
ctx.writeAndFlush(response);配置水位:
java
bootstrap.childOption(
ChannelOption.WRITE_BUFFER_WATER_MARK,
new WriteBufferWaterMark(32 * 1024, 64 * 1024)
);业务策略:
| 消息类型 | 处理 |
|---|---|
| 心跳 ACK | 可以丢弃旧的 |
| 普通状态推送 | 可以限流或合并 |
| 关键控制指令 | 入可靠队列,失败告警 |
| 大文件传输 | 分片、限速、断点续传 |
背压的本质是:当下游消费能力不足时,上游不能无限生产,否则内存一定会被打满。
ByteBuf 泄漏怎么理解
Netty 为了性能大量使用池化和堆外内存。ByteBuf 常用引用计数管理生命周期。
mermaid
flowchart TD
A["创建或读取 ByteBuf"] --> B["refCnt = 1"]
B --> C["retain 后 refCnt + 1"]
C --> D["使用完成 release"]
D --> E{"refCnt 是否为 0"}
E -- "是" --> F["回收到内存池"]
E -- "否" --> G["继续被持有"]常见泄漏:
retain后忘记release。- 异常分支没有释放。
- 把
ByteBuf放入异步队列,所有权不清楚。 - 自定义 Decoder 中切片后没有管理引用计数。
- 出站失败后没有清理缓存。
排查建议:
bash
-Dio.netty.leakDetection.level=advanced
-Dio.netty.maxDirectMemory=512m不要长期在高流量生产环境开最高级别泄漏检测,它有性能成本。可以在压测、灰度或复现环境打开。
商业场景一:设备采集网关
设计要点:
- 每个设备登录后绑定
deviceId -> Channel。 - 数据上报必须有
requestId和时间戳。 - 上报数据进入 MQ,避免直接在 EventLoop 写库。
- 业务线程池满时要返回限流或断开低优先级连接。
- 心跳超时关闭连接,设备端自动重连。
- 协议版本要兼容旧设备。
- 原始报文要按采集批次归档,方便排查。
mermaid
flowchart TD
A["设备上报"] --> B["Netty 解码"]
B --> C["设备鉴权"]
C --> D["投递业务线程池"]
D --> E["写 MQ"]
E --> F["消费者清洗入库"]
F --> G["ACK 或补偿"]商业场景二:IM/通知长连接
设计要点:
- 用户登录后绑定
userId -> Channel。 - 多端登录要区分 deviceId。
- 消息要有 messageId,客户端 ACK 后才算送达。
- 离线消息进入数据库或 MQ。
- 写缓冲过高时低优先级消息降级。
- 多实例部署要有用户路由表。
生产排查流程
连接数上不去
mermaid
flowchart TD
A["连接数上不去"] --> B["检查 ulimit 文件句柄"]
B --> C["检查端口和 backlog"]
C --> D["检查 Boss/Worker 线程"]
D --> E["检查握手和鉴权耗时"]
E --> F["检查内存和 direct memory"]消息错乱
mermaid
flowchart TD
A["消息错乱"] --> B["抓原始报文"]
B --> C["检查协议 magic 和 length"]
C --> D["检查粘包半包解码"]
D --> E["检查 Handler 是否共享可变状态"]
E --> F["检查异步处理是否改变顺序"]延迟升高
mermaid
flowchart TD
A["延迟升高"] --> B["看 EventLoop 线程栈"]
B --> C["是否有慢 SQL/HTTP/日志阻塞"]
C --> D["看业务线程池队列"]
D --> E["看写缓冲和客户端接收速度"]
E --> F["看 GC 和 direct memory"]直接内存 OOM
mermaid
flowchart TD
A["Direct memory OOM"] --> B["看 ByteBuf 泄漏日志"]
B --> C["看写缓冲 pending bytes"]
C --> D["看是否大包或无限缓存"]
D --> E["检查 retain/release 所有权"]
E --> F["压测复现并打开 leak detection"]面试标准回答
Netty 一次连接怎么走
text
服务端通过 ServerBootstrap 绑定端口,BossGroup 负责监听和接收新连接。客户端连接到来后,Boss EventLoop accept 创建 SocketChannel,并把它注册到某个 Worker EventLoop。随后初始化 ChannelPipeline,后续这个连接的读写事件通常都由绑定的 Worker EventLoop 处理,消息再按 Pipeline 中的 Handler 顺序完成拆包、解码、鉴权、业务处理、编码和写回。为什么不能在 EventLoop 里做慢业务
text
一个 EventLoop 通常负责多个 Channel 的 IO 事件。如果在 EventLoop 里执行慢 SQL、同步 HTTP、大计算或大量同步日志,这个 EventLoop 上的其他连接也会被拖慢,表现为多个客户端同时延迟升高甚至心跳超时。正确做法是把慢业务投递到有界业务线程池,并监控队列、耗时和拒绝策略。粘包半包怎么解决
text
TCP 是字节流协议,不保留应用层消息边界,所以一次读取可能拿到半条消息,也可能拿到多条消息。解决方式是应用层定义协议边界,常见有固定长度、分隔符和长度字段。商业系统最常用长度字段协议,解码器先读取头部和 body 长度,不够完整消息就 resetReaderIndex 等待更多字节。