乐于分享
好东西不私藏

Netty 写缓冲与背压源码:慢连接为什么会拖垮内存(源码系列第10篇)

Netty 写缓冲与背压源码:慢连接为什么会拖垮内存(源码系列第10篇)

Netty 不会因为对端慢就自动丢弃业务消息。写不出去的数据会留在 ChannelOutboundBuffer,pending bytes 超过高水位后 Channel.isWritable() 变为 false,并触发 channelWritabilityChanged。真正的背压必须由业务响应这个信号:暂停生产、停止读取或降级丢弃,否则慢连接会把直接内存拖爆。

很多 Netty 服务压测初期表现很好,上线后却在某些时刻直接内存飙升。日志看不到异常,CPU 也不一定高,最终出现 OutOfDirectMemoryError 或容器 OOM。根因常常不是 ByteBuf 泄漏,而是写侧积压:消息都还在等着发。

慢连接很危险,因为它不会立即失败。客户端还活着,只是读得慢;网络还能写,只是每次写一点;服务端业务还在生产,只是发送队列越来越长。Netty 提供了水位线和可写状态,但不会替业务决定“还要不要继续生产”。

源码基线:

  • • Netty netty-4.1.135.Final
  • • Commit:f05f765d81460799c53123a207f665bf3b465171
  • • 重点文件:ChannelOutboundBuffer.javaDefaultChannelConfig.javaWriteBufferWaterMark.java
  • • 官方仓库:https://github.com/netty/netty

这张流程图要抓住高低水位的滞回关系:pending bytes 超过 highWaterMark 后,Channel 变为不可写;只有下降到 lowWaterMark 以下才恢复可写。Netty 只提供信号,不会替业务决定丢弃、合并、暂停拉取还是关闭连接。真正的背压闭环必须把 isWritable=false 传回生产者。

1. 写缓冲积压不是泄漏,但同样会 OOM

ByteBuf 泄漏是“没有人再需要它,却没有 release”。写缓冲积压是“仍然有人需要它,但发送速度跟不上生产速度”。二者表现都可能是直接内存增长,但处理方式完全不同。

写积压路径:

业务持续 write  -> ChannelOutboundBuffer.addMessage  -> pendingOutboundBytes 增长  -> SocketChannel.write 返回 0 或写很少  -> 注册 OP_WRITE 等待可写  -> 业务继续 write  -> 队列继续增长

如果业务不看可写状态,Netty只能继续接收消息,直到内存耗尽或连接被关闭。

2. ChannelOutboundBuffer 统计待发送字节

ChannelOutboundBuffer.addMessage(msg, size, promise) 会创建 Entry,并增加待发送字节。简化逻辑:

Entry entry = Entry.newInstance(msg, size, total(msg), promise);if (tailEntry == null) {    flushedEntry = null;} else {    Entry tail = tailEntry;    tail.next = entry;}tailEntry = entry;if (unflushedEntry == null) {    unflushedEntry = entry;}incrementPendingOutboundBytes(entry.pendingSize, false);

pending size 不只是消息本体大小,还包括一定对象开销估算。它用于水位线判断,而不是精确内存会计。

源码位置:

  • • 文件:transport/src/main/java/io/netty/channel/ChannelOutboundBuffer.java
  • • 方法:addMessageincrementPendingOutboundBytes
  • • 链接:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/ChannelOutboundBuffer.java

当消息成功写完或失败移除时,decrementPendingOutboundBytes 会减少计数,并可能触发可写状态恢复。

3. 高低水位形成滞回,避免状态抖动

Netty 使用 WriteBufferWaterMark

new WriteBufferWaterMark(low, high)

默认值常见为低水位 32KB、高水位 64KB。语义:

  • • pending bytes > high:Channel 变为不可写。
  • • pending bytes < low:Channel 恢复可写。
  • • low 与 high 之间保持当前状态。

为什么不是超过 64KB 不可写,低于 64KB 立即可写?因为那会在边界附近频繁抖动。高低水位形成滞回区间,让业务有更稳定的暂停/恢复信号。

可写状态不是单一位。ChannelOutboundBuffer 内部维护 unwritable 标志位,除了写缓冲水位,用户也可设置自定义不可写位。最终 channel.isWritable() 综合这些状态。

4. 状态变化通过 channelWritabilityChanged 传播

pending bytes 超过高水位时,Netty 会触发:

pipeline.fireChannelWritabilityChanged();

业务 Handler 可以处理:

@Overridepublic void channelWritabilityChanged(ChannelHandlerContext ctx) {    if (ctx.channel().isWritable()) {        producer.resume(ctx.channel());        ctx.read();    } else {        producer.pause(ctx.channel());    }    ctx.fireChannelWritabilityChanged();}

注意一定要把信号传给真正的生产者。只打印日志没有背压效果。对于请求响应服务,生产者可能就是入站 read:不可写时停止 auto read 或不再调用 read;对于推送服务,生产者可能是消息队列订阅、定时广播或上游 RPC,需要暂停拉取或丢弃低优先级消息。

5. bytesBeforeUnwritable/bytesBeforeWritable 可做细粒度控制

Channel 提供:

channel.bytesBeforeUnwritable();channel.bytesBeforeWritable();

前者表示还能写多少字节就会变不可写,后者表示还需要发送多少字节才恢复可写。

批量推送时可以在写前判断:

while (channel.isWritable() && hasMore()) {    Message msg = next();    if (channel.bytesBeforeUnwritable() < estimateSize(msg)) {        break;    }    channel.write(msg);}channel.flush();

这样比写到不可写再暂停更温和。尤其在单条消息很大时,如果完全不估算,可能一次写入就越过水位很多。

6. OP_WRITE 只说明内核缓冲曾经写不下

写不完时,Netty 注册 OP_WRITE 等待 Socket 可写。很多人误以为 OP_WRITE 是背压信号。它只是内核层面的“之前写不下,现在可能能写了”。用户态积压和业务生产速度仍需要水位线控制。

可能出现:

  • • OP_WRITE 频繁触发,但每次只能写很少。
  • • pending bytes 仍持续增长。
  • • Channel 长期不可写。
  • • EventLoop 花大量时间尝试 flush 慢连接。

此时正确策略不是无限等,而是业务超时、限流、丢弃或关闭慢连接。不同业务选择不同:聊天消息可能保留关键消息丢弃弱提示;行情推送可能只保留最新快照;文件传输则应自然等待但限制并发。

7. 入站背压要把写侧状态传回 read

请求响应模型里,如果响应写不出去,还继续读取新请求,就会形成双重积压:入站业务任务和出站响应队列同时增长。

一种简单策略:

@Overridepublic void channelReadComplete(ChannelHandlerContext ctx) {    if (ctx.channel().isWritable()) {        ctx.read();    }}@Overridepublic void channelWritabilityChanged(ChannelHandlerContext ctx) {    if (ctx.channel().isWritable()) {        ctx.read();    }}

前提是关闭 AUTO_READ。这样当写侧不可写时,不再主动读取更多请求;当发送缓冲下降到低水位,再恢复读取。

这会把压力传回 TCP 接收窗口和客户端发送端,是端到端背压链的一部分。

8. 推送系统需要按用户/连接维度做队列策略

长连接推送服务中,业务生产者通常不由客户端请求驱动,而由后端事件驱动。写侧不可写时,必须定义队列策略:

策略
适合消息
代价
暂停上游拉取
可反压 MQ 或流系统
增加端到端延迟
丢弃低优先级
在线状态、弱提醒
需要业务可接受
合并保留最新
行情、位置、进度
丢中间状态
断开慢连接
强实时服务
用户体验受影响
落盘补偿
重要通知
系统复杂度高

不要把 Netty 的 ChannelOutboundBuffer 当业务消息队列。它缺少优先级、过期、合并、持久化和可观测性,只是网络发送队列。

9. 水位线不是越大越好

调大 highWaterMark 可以减少不可写状态,但也允许更多内存积压。假设 10 万连接,每条连接允许 1MB 写缓冲,理论上就是 100GB 用户态积压。即使大多数连接不满,峰值风险仍很高。

水位线应结合:

  • • 单连接最大可接受积压。
  • • 消息平均大小和最大大小。
  • • 连接数量。
  • • 直接内存上限。
  • • 客户端网络质量。
  • • 业务是否允许丢弃或合并。

通常应先建立业务背压策略,再微调水位线。只调大水位线是一种延迟爆炸。

这张时序图展示了慢客户端故障的完整扩散路径:生产速度高于发送速度,内核发送缓冲写满,Netty 注册 OP_WRITE,用户态 ChannelOutboundBuffer 持续增长。正确分支是在不可写时暂停生产或停止读取;错误分支是继续写,把网络慢问题放大成直接内存问题。

10. ChannelOutboundBuffer 还影响 writability 事件顺序

当 pending bytes 跨越水位线,Netty 不一定在当前调用栈同步执行用户 Handler。它可能调用 fireChannelWritabilityChanged(invokeLater),把事件延后到 EventLoop 后续任务中,避免重入和复杂状态交错。

这意味着业务不能假设刚 write 完马上同步收到不可写回调。更稳妥的方式是在写循环中主动检查 isWritable(),同时在回调中恢复生产。

11. 慢连接排查指标

建议至少观测:

  • • 每个 Channel 的 pending outbound bytes。
  • • 不可写 Channel 数量。
  • • channelWritabilityChanged 频率。
  • • OP_WRITE 注册数量或持续时间。
  • • write Future 完成延迟。
  • • EventLoop flush 耗时。
  • • 业务生产速率与实际发送速率。
  • • 连接维度消息丢弃/合并数量。

如果不可写连接数上升且 direct memory 增长,优先判断慢客户端或网络拥塞。如果不可写不多但内存增长,可能是少数超级连接或 ByteBuf 泄漏。如果 write Future 大量失败,可能是连接关闭或编码异常。

12. 一个可落地的背压模板

连接初始化:

ch.config().setAutoRead(false);ch.config().setWriteBufferWaterMark(    new WriteBufferWaterMark(32 * 1024, 128 * 1024));

Handler:

@Overridepublic void channelActive(ChannelHandlerContext ctx) {    ctx.read();}@Overridepublic void channelRead(ChannelHandlerContext ctx, Object msg) {    handleAsync(msg).whenComplete((resp, err) -> {        ctx.executor().execute(() -> {            if (err != null) {                ctx.close();                return;            }            ctx.writeAndFlush(resp);            if (ctx.channel().isWritable()) {                ctx.read();            }        });    });}@Overridepublic void channelWritabilityChanged(ChannelHandlerContext ctx) {    if (ctx.channel().isWritable()) {        ctx.read();    }    ctx.fireChannelWritabilityChanged();}

这不是唯一答案,但体现了闭环:处理完成、写侧可用、才读下一批。实际项目还要加超时、异常释放、业务线程池容量和协议层错误返回。

13. 和应用层限流的关系

Netty 背压解决的是单连接网络发送能力;应用限流解决的是系统整体处理能力。二者不能互相替代。

如果系统总 QPS 超过数据库能力,即使每条 Channel 都可写,也需要全局限流。如果只有某个客户端网络慢,不能因为全局限流让所有用户受影响,应在该 Channel 维度暂停或关闭。

最佳实践通常是多层背压:

全局入口限流  -> 业务线程池有界队列  -> Channel 写水位线  -> AUTO_READ 控制  -> 协议级降级/丢弃/关闭

每一层都要有明确指标和策略,否则压力只会转移到下一层。

14. 总结:Netty 给信号,业务做决策

写侧背压主链:

write  -> ChannelOutboundBuffer.addMessage  -> pending bytes 增长  -> 超过 highWaterMark  -> isWritable=false  -> channelWritabilityChanged  -> 业务暂停生产/读取  -> doWrite 逐步发送  -> 低于 lowWaterMark  -> isWritable=true  -> 业务恢复

Netty 能告诉你“这条连接现在承载不了更多写入”,但它不知道你的消息能不能丢、能不能合并、能不能延迟、要不要断开。背压的最后一公里必须由业务协议完成。

至此前 10 篇完成了一条完整链路:启动、接入、事件循环、读、写、Pipeline、内存、解码和背压。下一阶段可以继续深入 Future/Promise、定时任务、Selector 重建、IdleStateHandler、Native Transport 和 HTTP/WebSocket 等专项。

参考资料

  1. 1. ChannelOutboundBuffer.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/ChannelOutboundBuffer.java
  2. 2. WriteBufferWaterMark.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/WriteBufferWaterMark.java
  3. 3. DefaultChannelConfig.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/DefaultChannelConfig.java
  4. 4. AbstractNioByteChannel.java:https://github.com/netty/netty/blob/netty-4.1.135.Final/transport/src/main/java/io/netty/channel/nio/AbstractNioByteChannel.java