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

