基于Netty实现TCP代理:如何触发处理已缓冲的初始请求数据
解决方案:Netty TCP代理缓冲数据与连接完成后触发转发
一、核心问题解决:缓冲入站数据 + 连接完成后批量转发
你的问题本质是出站连接建立前收到的入站数据未被缓存,且未在连接就绪后触发转发,可以通过以下步骤实现修复:
1. 在入站Handler中维护核心状态
- 保存出站连接的
ChannelFuture,用于判断连接状态 - 用线程安全队列缓存未处理的入站
ByteBuf - 保存已建立的出站
Channel,用于后续数据转发
2. 关键代码实现
public class ProxyInboundHandler extends ChannelInboundHandlerAdapter { private final String remoteHost; private final int remotePort; private ChannelFuture outboundFuture; private final Queue<ByteBuf> pendingData = new ConcurrentLinkedQueue<>(); private Channel outboundChannel; public ProxyInboundHandler(String remoteHost, int remotePort) { this.remoteHost = remoteHost; this.remotePort = remotePort; } @Override public void channelActive(ChannelHandlerContext ctx) { // 初始化出站连接Bootstrap Bootstrap bootstrap = new Bootstrap(); bootstrap.group(ctx.channel().eventLoop()) .channel(ctx.channel().getClass()) .handler(new ProxyOutboundHandler(ctx.channel())); // 发起连接并添加完成监听器 outboundFuture = bootstrap.connect(remoteHost, remotePort); outboundFuture.addListener((ChannelFutureListener) future -> { if (future.isSuccess()) { outboundChannel = future.channel(); // 连接就绪后,批量转发缓冲的数据 flushPendingData(); } else { // 连接失败,关闭入站通道并清理缓冲 ctx.channel().close(); clearPendingData(); } }); } @Override public void channelRead(ChannelHandlerContext ctx, Object msg) { if (msg instanceof ByteBuf buf) { if (outboundChannel != null && outboundChannel.isActive()) { // 出站连接就绪,直接转发 outboundChannel.writeAndFlush(buf) .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } else { // 缓冲数据:必须调用retain(),避免Netty自动释放数据 buf.retain(); pendingData.add(buf); } } else { ctx.fireChannelRead(msg); } } @Override public void channelInactive(ChannelHandlerContext ctx) { // 入站通道关闭时,清理资源 if (outboundChannel != null) { outboundChannel.close(); } clearPendingData(); } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); clearPendingData(); } // 批量转发缓冲数据 private void flushPendingData() { ByteBuf buf; while ((buf = pendingData.poll()) != null) { if (buf.isReadable()) { outboundChannel.writeAndFlush(buf) .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } else { buf.release(); } } } // 清理缓冲数据,避免内存泄漏 private void clearPendingData() { ByteBuf buf; while ((buf = pendingData.poll()) != null) { buf.release(); } } } // 出站Handler:负责将B端的数据转发回A端 public class ProxyOutboundHandler extends ChannelInboundHandlerAdapter { private final Channel inboundChannel; public ProxyOutboundHandler(Channel inboundChannel) { this.inboundChannel = inboundChannel; } @Override public void channelRead(ChannelHandlerContext ctx, Object msg) { if (msg instanceof ByteBuf buf) { inboundChannel.writeAndFlush(buf) .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } else { ctx.fireChannelRead(msg); } } @Override public void channelInactive(ChannelHandlerContext ctx) { inboundChannel.close(); } @Override public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } }
3. 关键细节说明
- 引用计数处理:调用
buf.retain()是因为Netty默认会在channelRead方法结束后自动释放ByteBuf的引用,缓存数据时必须手动增加引用计数,否则数据会被提前回收导致空指针或数据丢失。 - 线程安全队列:使用
ConcurrentLinkedQueue避免多线程环境下的队列操作异常(比如IO线程和Future监听器线程同时操作队列)。 - 资源清理:在通道关闭、连接失败时必须清空并释放缓冲的
ByteBuf,防止内存泄漏。
二、Netty Handler与Adapter的区别梳理
ChannelHandler:顶层接口,定义了所有入站、出站、生命周期相关的方法,实现该接口必须重写所有方法(哪怕是空实现),一般不直接使用。ChannelInboundHandlerAdapter:入站事件适配器,实现了ChannelInboundHandler接口,给所有入站方法提供了默认实现(比如channelRead默认调用ctx.fireChannelRead(msg)继续传递事件),只需重写你关心的入站方法。ChannelOutboundHandlerAdapter:出站事件适配器,实现了ChannelOutboundHandler接口,同理提供了出站方法的默认实现,用于处理write、connect、bind等出站操作。ChannelDuplexHandler:同时实现入站和出站适配器,适合需要同时处理双向事件的场景(比如某些编解码器)。
简单总结:优先使用对应的Adapter类,减少冗余代码,只重写需要的方法即可。
内容的提问来源于stack exchange,提问作者dty
相关产品推荐
相关产品推荐

