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

Netty实现RPC时channelRead未触发、writeAndFlush不生效问题求助

问题根因与修复方案

核心问题1:客户端连接刚建立就被关闭

你的ServiceProxy.startClient()方法逻辑错误,finally块在连接建立后直接执行group.shutdownGracefully(),直接销毁了IO线程组、断开了TCP连接,后续RPCClientHandler.call()里调用writeAndFlush()时连接已经失效,自然发不出消息。只有startClient()里连接刚建立、group还没关闭前调用的writeAndFlush("DDD")能短暂发出去。

核心问题2:服务端不会响应你发送的消息

RPCClientHandler.call()里写死发送的是字符串Ack,但服务端RPCServerHandler.channelRead()只有收到前缀为send?的消息才会返回响应,所以channelRead()永远收不到返回,wait()会一直阻塞,不会执行后续逻辑。

核心问题3:类语法错误

你贴的RPCClientHandler代码中,call()方法写在了类的闭合大括号外面,属于语法错误,实际运行会编译不通过。

核心问题4:动态代理逻辑失效

动态代理的invoke()方法里完全没有使用实际调用的方法名、参数、传输协议等信息,只是固定提交clientHandler,无法实现通用RPC调用。


修复步骤

1. 修复ServiceProxy的startClient逻辑

不要在startClient里直接关闭EventLoopGroup,要持有Channel和Group的引用,避免连接提前断开:

public class ServiceProxy {
    private static ExecutorService executor = Executors.newFixedThreadPool(Runtime.getRuntime().availableProcessors());
    private static RPCClientHandler clientHandler;
    // 新增持有Channel和EventLoopGroup,避免被提前销毁
    private static Channel channel;
    private static EventLoopGroup group;

    public Object getProxy(Class<?> serviceClass, String protocol){
        return Proxy.newProxyInstance(this.getClass().getClassLoader(), new Class<?>[]{serviceClass}, new InvocationHandler() {
            @Override
            public Object invoke(Object proxy, Method method, Object[] args) throws Throwable {
                if(clientHandler == null) {
                    startClient();
                    System.out.println("clientHandler exists:" + clientHandler);
                }
                // 把要发送的参数设置给clientHandler,不要写死
                clientHandler.setRequestParam(protocol + args[0]);
                return executor.submit(clientHandler).get();
            }
        });
    }

    public static void startClient() throws Exception{
        clientHandler = new RPCClientHandler();
        group = new NioEventLoopGroup();
        Bootstrap bootstrap = new Bootstrap();
        bootstrap.group(group)
                .channel(NioSocketChannel.class)
                .handler(new ChannelInitializer<SocketChannel>() {
                    @Override
                    protected void initChannel(SocketChannel ch) throws Exception {
                        ch.pipeline().addLast(new StringDecoder());
                        ch.pipeline().addLast(new StringEncoder());
                        ch.pipeline().addLast(clientHandler);
                    }
                });
        ChannelFuture cf = bootstrap.connect("127.0.0.1", 8890).sync();
        channel = cf.channel();
        // 移除原来的finally块关闭逻辑,不要在这里销毁group
    }
}

2. 修复RPCClientHandler逻辑

把call方法放回类内部,新增请求参数的成员变量,发送符合协议的请求内容:

public class RPCClientHandler extends ChannelInboundHandlerAdapter implements Callable {
    private ChannelHandlerContext context;
    private String response;
    // 新增请求参数成员变量
    private String requestParam;

    @Override
    public void channelActive(ChannelHandlerContext ctx) throws Exception {
        context = ctx;
    }

    @Override
    public synchronized void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
        System.out.println("invoked success");
        response = msg.toString();
        System.out.println("invoked success"+response);
        notify();
    }

    @Override
    public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws Exception {
        ctx.close();
    }

    public void setRequestParam(String requestParam) {
        this.requestParam = requestParam;
    }

    @Override
    public synchronized Object call() throws Exception {
        // 发送实际请求参数,不是写死的Ack
        ChannelFuture cf = context.writeAndFlush(requestParam);
        cf.addListener(future -> {
            if(!future.isSuccess()){
                System.out.println("发送失败:" + future.cause());
            }
        });
        System.out.println(cf);
        System.out.println("invoked1, ctx is: "+context);
        wait();
        System.out.println("invoked2");
        return response;
    }
}

3. 可选:服务端增加全量消息日志

方便调试,在RPCServerHandler的channelRead里增加非协议消息的日志输出:

@Override
public void channelRead(ChannelHandlerContext ctx, Object msg) throws Exception {
    System.out.println("msg in channelRead:" + msg +"the class of msg is:" + msg.getClass());
    if(msg.toString().startsWith("send?")){
        String res = new UserServiceImpl().say(msg.toString().substring(msg.toString().lastIndexOf("?")+1));
        ctx.writeAndFlush(res);
    } else {
        System.out.println("收到不符合协议的消息:" + msg);
    }
}

运行验证

修复后先启动服务端,再启动客户端,就能正常看到打印的RPC调用返回结果。

内容的提问来源于stack exchange,提问作者s666

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 16:36:00