上一篇【第29篇】Netty OP_READ事件处理源码解析(下)——数据如何在Pipeline中流转 下一篇【第31篇】031 Netty零拷贝深度解析(明日更新,敬请期待)
开篇故事
某IM系统,频繁出现"消息发送失败"。排查发现:写到Netty的write()调用成功后没有flush(),数据一直积压在缓冲区没有发送!
write() = 把数据放入发送队列(不发送) flush() = 把队列中的数据真正写入Socket writeAndFlush() = write() + flush()
一、ChannelOutboundBuffer双链表结构
write()操作:
unflushedEntry →
[Msg1] → [Msg2] → [Msg3] → null
↑ tailEntry
flush()操作后:
flushedEntry →
[Msg1] → [Msg2] → [Msg3] → null
↑ tailEntry
doWrite():从flushedEntry开始写入,写完后移除
二、write()源码
// Pipeline传播:tail.write(msg) → HeadContext.write(msg) → unsafe.write(msg)
public void write(Object msg, ChannelPromise promise) {
ChannelOutboundBuffer outboundBuffer = this.outboundBuffer;
int size = msg instanceof ByteBuf ? ((ByteBuf) msg).readableBytes() : 0;
outboundBuffer.addMessage(msg, size, promise); // 加入unflushed队列
}
// addMessage:添加消息到unflushed队列
public void addMessage(Object msg, int size, ChannelPromise promise) {
Entry entry = Entry.newInstance(msg, size, total(msg), promise);
if (tailEntry == null) {
flushedEntry = null;
tailEntry = entry;
} else {
Entry tail = tailEntry;
tail.next = entry;
tailEntry = entry;
}
// 检查水位线
incrementPendingOutboundBytes(entry.pendingSize, false);
}
三、flush()源码
public void flush() {
ChannelOutboundBuffer outboundBuffer = this.outboundBuffer;
outboundBuffer.addFlush(); // 标记unflushed→flushed
flush0(); // 调用doWrite()
}
// addFlush:将unflushedEntry指向的队列标记为flushed
public void addFlush() {
Entry entry = unflushedEntry;
if (entry != null) {
do {
entry = entry.next;
} while (entry != null);
flushedEntry = unflushedEntry;
unflushedEntry = null;
}
}
四、水位线机制
// 默认水位线:高水位64KB,低水位32KB
int highWaterMark = 64 * 1024;
int lowWaterMark = 32 * 1024;
public void incrementPendingOutboundBytes(long size, boolean invokeLater) {
long newValue = TOTAL_PENDING_SIZE_UPDATER.addAndGet(this, size);
if (newValue > highWaterMark) {
setUnwritable(invokeLater); // 设置不可写
}
}
// 写出数据后减少计数
public void removeBytes(long writtenBytes) {
long newValue = TOTAL_PENDING_SIZE_UPDATER.addAndGet(this, –writtenBytes);
if (newValue < lowWaterMark && !isWritable()) {
setWritable(); // 恢复可写
}
}
水位线状态转换:
可写状态 不可写状态
| ^
| 待发送>高水位(64KB) |
+————————–> |
| |
| 待发送<低水位(32KB) |
+ <————————-+
五、实战:水位线监听
import io.netty.channel.ChannelInboundHandlerAdapter;
import io.netty.channel.ChannelHandlerContext;
public class WaterMarkHandler extends ChannelInboundHandlerAdapter {
@Override
public void channelWritabilityChanged(ChannelHandlerContext ctx) {
boolean writable = ctx.channel().isWritable();
System.out.println("可写状态变化: " + writable);
if (writable) {
// 恢复发送积压数据
resumeSending();
} else {
// 暂停发送,等恢复
pauseSending();
}
ctx.fireChannelWritabilityChanged();
}
}
六、总结
| write() | 数据进入unflushed队列,不发送 |
| flush() | 标记为flushed,调用doWrite()写入Socket |
| 高水位 | 待发送数据超过高水位→channelWritabilityChanged(false) |
| 低水位 | 待发送数据低于低水位→channelWritabilityChanged(true) |
上一篇【第29篇】Netty OP_READ事件处理源码解析(下)——数据如何在Pipeline中流转 下一篇【第31篇】031 Netty零拷贝深度解析(明日更新,敬请期待)

