Skip to content

Netty 商业场景训练营

这页把 Netty 放到真实商业系统里学习。目标不是背 BootstrapChannelEventLoop 这些类名,而是能解释一个长连接怎么接入、一条 TCP 字节流怎么变成业务消息、业务处理为什么不能阻塞 IO 线程、写缓冲为什么会堆积、ByteBuf 为什么会泄漏、线上连接假死和消息错乱怎么排查。

训练目标

学完这页,你要能做到:

  1. 解释 Netty 服务端从绑定端口、接收连接、注册 Channel 到处理消息的全过程。
  2. 解释 Reactor、BossGroup、WorkerGroup、EventLoop、Channel、Pipeline、Handler 的关系。
  3. 解释 TCP 粘包半包为什么一定要用协议边界解决。
  4. 写出长度字段协议的编解码 Demo。
  5. 解释为什么不能在 EventLoop 里查数据库、调 HTTP 或做大计算。
  6. 解释心跳、空闲检测、重连、连接管理和用户会话绑定。
  7. 解释 ByteBuf 池化、引用计数和泄漏排查。
  8. 排查连接数上不去、消息错乱、延迟升高、写缓冲堆积、直接内存 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["等待读写事件"]

关键点:

  1. BossGroup 只负责接入连接,不负责处理所有业务。
  2. WorkerGroup 负责已建立连接的读写事件。
  3. 一个 Channel 通常固定绑定一个 EventLoop。
  4. Pipeline 在 Channel 初始化时建立,后续消息按 Handler 顺序流动。
  5. 如果 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请求和响应关联
lengthbody 长度,解决粘包半包
bodyJSON、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;
        }
    }
}

这个解码器体现了几个关键点:

  1. 不够头部长度就返回,等待更多字节。
  2. markReaderIndexresetReaderIndex 处理半包。
  3. 校验 magic 和 length,防止乱流量和超大包攻击。
  4. 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();
        }
    }
}

为什么要会话绑定?

  1. 平台要知道设备当前连接在哪个 Channel。
  2. 下行指令需要找到设备连接。
  3. 设备断开时要清理状态,避免给旧 Channel 发消息。
  4. 多实例部署时还要结合 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["继续被持有"]

常见泄漏:

  1. retain 后忘记 release
  2. 异常分支没有释放。
  3. ByteBuf 放入异步队列,所有权不清楚。
  4. 自定义 Decoder 中切片后没有管理引用计数。
  5. 出站失败后没有清理缓存。

排查建议:

bash
-Dio.netty.leakDetection.level=advanced
-Dio.netty.maxDirectMemory=512m

不要长期在高流量生产环境开最高级别泄漏检测,它有性能成本。可以在压测、灰度或复现环境打开。

商业场景一:设备采集网关

设计要点:

  1. 每个设备登录后绑定 deviceId -> Channel
  2. 数据上报必须有 requestId 和时间戳。
  3. 上报数据进入 MQ,避免直接在 EventLoop 写库。
  4. 业务线程池满时要返回限流或断开低优先级连接。
  5. 心跳超时关闭连接,设备端自动重连。
  6. 协议版本要兼容旧设备。
  7. 原始报文要按采集批次归档,方便排查。
mermaid
flowchart TD
    A["设备上报"] --> B["Netty 解码"]
    B --> C["设备鉴权"]
    C --> D["投递业务线程池"]
    D --> E["写 MQ"]
    E --> F["消费者清洗入库"]
    F --> G["ACK 或补偿"]

商业场景二:IM/通知长连接

设计要点:

  1. 用户登录后绑定 userId -> Channel
  2. 多端登录要区分 deviceId。
  3. 消息要有 messageId,客户端 ACK 后才算送达。
  4. 离线消息进入数据库或 MQ。
  5. 写缓冲过高时低优先级消息降级。
  6. 多实例部署要有用户路由表。

生产排查流程

连接数上不去

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 等待更多字节。

关联知识点