欢迎光临
我们一直在努力

【Netty源码解读和权威指南】第30篇:Netty写数据源码解析——write/flush背后的双队列设计

上一篇【第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零拷贝深度解析(明日更新,敬请期待)


赞(0)
未经允许不得转载:171主机测评 » 【Netty源码解读和权威指南】第30篇:Netty写数据源码解析——write/flush背后的双队列设计
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址