You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

基于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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.29 02:28:18