使用 Netty 从零实现一个 RPC(类似 Dubbo)是一个非常经典的工程实践案例。为了遵循由简单到复杂的工程原则,我们将按以下五个步骤逐步进行深化。
Stage 1:跑通 Netty 最基础的通信,用反射完成一个最基本的调用。
Stage 2:引入 RequestId 和 CompletableFuture,完成单连接多路复用。
Stage 3:引入心跳和重连,保证了长连接的网络健壮性。
Stage 4:重构掉废弃 API,设计自定义二进制私有协议头和高性能 Kryo 序列化。
Stage 5:跳出单机思维,引入服务注册发现中心与 软负载均衡轮询算法。
基础案例 首先我们从最小可行性产品开始。在这个阶段,我们不考虑复杂的路由、服务注册中心(ZooKeeper/Nacos)、多线程模型优化和复杂的序列化,只聚焦于核心:如何通过 Netty 实现客户端发送一个请求对象,服务端接收、执行并返回结果。一个最基础的 RPC 包含以下四个要素:
API 接口:服务提供方和服务消费方共同约定的契约。
网络传输(Netty):负责数据的发送与接收。
协议与序列化:定义数据传输格式。为了简化直接使用 Netty 自带的 ObjectEncoder/Decoder。
反射调用:服务端收到请求后,通过反射调用本地真正的方法。
接口和契约类 首先,定义一个服务接口,以及传输的请求/响应对象,传输对象必须实现 Serializable 接口。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 public interface HelloService { String sayHello (String msg) ; } @Data public class RpcRequest implements Serializable { private String className; private String methodName; private Class<?>[] parameterTypes; private Object[] parameters; } @Data public class RpcResponse implements Serializable { private Object result; private Throwable error; }
服务端实现 服务端需要实现接口,并启动一个 Netty 服务监听端口。收到 RpcRequest 后,反射调用实现类,再把 RpcResponse 写回。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 public class HelloServiceImpl implements HelloService { @Override public String sayHello (String msg) { return "【Server 回应】: 你好,我已经收到了你的消息 -> " + msg; } } public class RpcServerHandler extends SimpleChannelInboundHandler <RpcRequest> { @Override protected void channelRead0 (ChannelHandlerContext ctx, RpcRequest request) throws Exception { RpcResponse response = new RpcResponse (); try { Object serviceBean = new HelloServiceImpl (); Method method = serviceBean.getClass().getMethod(request.getMethodName(), request.getParameterTypes()); Object result = method.invoke(serviceBean, request.getParameters()); response.setResult(result); } catch (Throwable t) { response.setError(t); } ctx.writeAndFlush(response); } } public class RpcServer { private int port; public RpcServer (int port) { this .port = port; } public void start () throws Exception { EventLoopGroup bossGroup = new NioEventLoopGroup (1 ); EventLoopGroup workerGroup = new NioEventLoopGroup (); try { ServerBootstrap b = new ServerBootstrap (); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializer <SocketChannel>() { @Override protected void initChannel (SocketChannel ch) throws Exception { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new ObjectDecoder (Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null ))); pipeline.addLast(new ObjectEncoder ()); pipeline.addLast(new RpcServerHandler ()); } }); ChannelFuture f = b.bind(port).sync(); System.out.println("RPC Server 成功启动,监听端口: " + port); f.channel().closeFuture().sync(); } finally { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); } } public static void main (String[] args) throws Exception { new RpcServer (8080 ).start(); } }
客户端实现 客户端需要使用动态代理,让开发者像调用本地方法一样调用远程服务。同时利用 Netty 发送数据并同步阻塞等待结果。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 public class RpcClientHandler extends SimpleChannelInboundHandler <RpcResponse> { private final SynchronousQueue<RpcResponse> queue = new SynchronousQueue <>(); @Override protected void channelRead0 (ChannelHandlerContext ctx, RpcResponse msg) throws Exception { queue.put(msg); } public RpcResponse getResponse () throws InterruptedException { return queue.take(); } @Override public void exceptionCaught (ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } } public class RpcClientProxy { @SuppressWarnings("unchecked") public static <T> T createProxy (Class<T> serviceClass, String host, int port) { return (T) newProxyInstance( serviceClass.getClassLoader(), new Class <?>[]{serviceClass}, (proxy, method, args) -> { RpcRequest request = new RpcRequest (); request.setClassName(serviceClass.getName()); request.setMethodName(method.getName()); request.setParameterTypes(method.getParameterTypes()); request.setParameters(args); EventLoopGroup group = new NioEventLoopGroup (); RpcClientHandler clientHandler = new RpcClientHandler (); try { Bootstrap b = new Bootstrap (); b.group(group) .channel(NioSocketChannel.class) .handler(new ChannelInitializer <SocketChannel>() { @Override protected void initChannel (SocketChannel ch) throws Exception { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new ObjectEncoder ()); pipeline.addLast(new ObjectDecoder (Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null ))); pipeline.addLast(clientHandler); } }); ChannelFuture f = b.connect(host, port).sync(); f.channel().writeAndFlush(request).sync(); RpcResponse response = clientHandler.getResponse(); if (response.getError() != null ) { throw response.getError(); } return response.getResult(); } finally { group.shutdownGracefully(); } } ); } public static void main (String[] args) { HelloService helloService = RpcClientProxy.createProxy(HelloService.class, "127.0.0.1" , 8080 ); String result = helloService.sayHello("Hello Netty RPC!" ); System.out.println("客户端收到远程调用结果: " + result); } }
测试验证 依次执行以下步骤:
长连接复用与多服务路由 需要解决的问题 在 RPC 框架中,网络通信模型的演进(从短连接到长连接复用)才是最核心、最硬核的架构部分。这里我们需要实现长连接复用与多服务路由。它要解决的核心痛点如下:
痛点一:短连接开销巨大。 上面的基础案例,客户端每调用一次方法,就要 new NioEventLoopGroup(),建连、发数据、断连、销毁线程池,开销巨大。我们需要客户端与服务端建立一条(或几条)常驻的长连接,所有 RPC 请求都复用这个通道(Channel)发送。
痛点二:高并发下的请求错乱(多路复用) 。当多个线程复用同一个 Channel 发送请求时,Netty 是异步返回的。比如线程 A 发了请求 1,线程 B 发了请求 2;服务器可能先返回响应 2,再返回响应 1。客户端如何让线程 A 拿到响应 1,线程 B 拿到响应 2 就是个棘手的问题。这里我们需要引入请求 ID (RequestId) 和 全局期约(CompletableFuture) 映射表。
痛点三:服务端无法支持多个服务 。在基础案例中,Server 端写死了 new HelloServiceImpl()。这里我们在服务端引入一个本地服务注册表(Service Registry),用一个 Map 维护服务名与实现类的映射。
代码实现演进 为了实现多路复用,我们需要改造 RpcRequest 和 RpcResponse,加入一个 requestId。客户端发送请求时,把一个 CompletableFuture 存入全局 Map,然后业务线程阻塞在 future.get() 上。当 Netty 收到响应时,根据响应里的 requestId 找到对应的 Future 并赋值,从而唤醒业务线程。
契约升级 加入 RequestId:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 @Data public class RpcRequest implements Serializable { private String requestId; private String className; private String methodName; private Class<?>[] parameterTypes; private Object[] parameters; } @Data public class RpcResponse implements Serializable { private String requestId; private Object result; private Throwable error; }
服务端升级 服务端不再硬编码,而是提供一个注册方法,支持多服务路由。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 public class RpcServerHandler extends SimpleChannelInboundHandler <RpcRequest> { private final Map<String, Object> serviceMap; public RpcServerHandler (Map<String, Object> serviceMap) { this .serviceMap = serviceMap; } @Override protected void channelRead0 (ChannelHandlerContext ctx, RpcRequest request) throws Exception { RpcResponse response = new RpcResponse (); response.setRequestId(request.getRequestId()); try { Object serviceBean = serviceMap.get(request.getClassName()); if (serviceBean == null ) { throw new RuntimeException ("未找到服务实现类: " + request.getClassName()); } Method method = serviceBean.getClass().getMethod(request.getMethodName(), request.getParameterTypes()); Object result = method.invoke(serviceBean, request.getParameters()); response.setResult(result); } catch (Throwable t) { response.setError(t); } ctx.writeAndFlush(response); } }
服务端启动类
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 public class RpcServer { private int port; public RpcServer (int port) { this .port = port; } private final Map<String, Object> serviceMap = new ConcurrentHashMap <>(); public void registerService (Class<?> interfaceClass, Object serviceImpl) { serviceMap.put(interfaceClass.getName(), serviceImpl); } public void start () throws Exception { EventLoopGroup bossGroup = new NioEventLoopGroup (1 ); EventLoopGroup workerGroup = new NioEventLoopGroup (); try { ServerBootstrap b = new ServerBootstrap (); b.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializer <SocketChannel>() { @Override protected void initChannel (SocketChannel ch) throws Exception { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new ObjectDecoder (Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null ))); pipeline.addLast(new ObjectEncoder ()); pipeline.addLast(new RpcServerHandler (serviceMap)); } }); ChannelFuture f = b.bind(port).sync(); System.out.println("RPC Server 成功启动,监听端口: " + port); f.channel().closeFuture().sync(); } finally { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); } } }
客户端升级 客户端支持长连接与多路复用。客户端需要保持长连接,并使用 CompletableFuture 解决异步响应的匹配问题。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 public class RpcClientHandler extends SimpleChannelInboundHandler <RpcResponse> { private final Map<String, CompletableFuture<RpcResponse>> pendingResponseMap = new ConcurrentHashMap <>(); public CompletableFuture<RpcResponse> registerFuture (String requestId) { CompletableFuture<RpcResponse> future = new CompletableFuture <>(); pendingResponseMap.put(requestId, future); return future; } @Override protected void channelRead0 (ChannelHandlerContext ctx, RpcResponse response) throws Exception { CompletableFuture<RpcResponse> future = pendingResponseMap.remove(response.getRequestId()); if (future != null ) { future.complete(response); } } @Override public void exceptionCaught (ChannelHandlerContext ctx, Throwable cause) { cause.printStackTrace(); ctx.close(); } }
客户端连接管理器与代理(单例长连接)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 public class RpcClientContainer { private final String host; private final int port; private Channel channel; private RpcClientHandler clientHandler; private final EventLoopGroup group = new NioEventLoopGroup (); public RpcClientContainer (String host, int port) { this .host = host; this .port = port; } public void init () throws Exception { clientHandler = new RpcClientHandler (); Bootstrap b = new Bootstrap (); b.group(group) .channel(NioSocketChannel.class) .handler(new ChannelInitializer <SocketChannel>() { @Override protected void initChannel (SocketChannel ch) { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new ObjectEncoder ()); pipeline.addLast(new ObjectDecoder (Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null ))); pipeline.addLast(clientHandler); } }); ChannelFuture f = b.connect(host, port).sync(); this .channel = f.channel(); System.out.println("RPC 客户端长连接建立成功!" ); } public void close () { group.shutdownGracefully(); } @SuppressWarnings("unchecked") public <T> T createProxy (Class<T> serviceClass) { return (T) Proxy.newProxyInstance( serviceClass.getClassLoader(), new Class <?>[]{serviceClass}, (proxy, method, args) -> { RpcRequest request = new RpcRequest (); request.setRequestId(UUID.randomUUID().toString()); request.setClassName(serviceClass.getName()); request.setMethodName(method.getName()); request.setParameterTypes(method.getParameterTypes()); request.setParameters(args); CompletableFuture<RpcResponse> future = clientHandler.registerFuture(request.getRequestId()); if (channel == null || !channel.isActive()) { throw new RuntimeException ("RPC 调用失败:当前网络连接已断开,正在尝试重连中..." ); } channel.writeAndFlush(request); RpcResponse response = future.get(); if (response.getError() != null ) { throw response.getError(); } return response.getResult(); } ); } }
测试验证 增加一个新接口 OrderService,验证 “多服务路由”。
1 2 3 4 5 6 7 8 9 10 public interface OrderService { String getOrderSymbol (String userId) ; } public class OrderServiceImpl implements OrderService { @Override public String getOrderSymbol (String userId) { return "【OrderServer】用户 " + userId + " 的最新订单号为:TX20260705001" ; } }
服务端和客户端的启动类:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 public class ServerMain { public static void main (String[] args) throws Exception { RpcServer server = new RpcServer (8080 ); server.registerService(HelloService.class, new HelloServiceImpl ()); server.registerService(OrderService.class, new OrderServiceImpl ()); server.start(); } } public class ClientMain { public static void main (String[] args) throws Exception { RpcClientContainer container = new RpcClientContainer ("127.0.0.1" , 8080 ); container.init(); HelloService helloService = container.createProxy(HelloService.class); OrderService orderService = container.createProxy(OrderService.class); for (int i = 0 ; i < 3 ; i++) { final int index = i; new Thread (() -> { String res1 = helloService.sayHello("线程-" + index); System.out.println(res1); String res2 = orderService.getOrderSymbol("User-" + index); System.out.println(res2); }).start(); } Thread.sleep(2000 ); container.close(); } }
客户端打印:
1 2 3 4 5 6 7 RPC 客户端长连接建立成功! 【Server 回应】: 你好,我已经收到了你的消息 -> 线程-0 【Server 回应】: 你好,我已经收到了你的消息 -> 线程-1 【Server 回应】: 你好,我已经收到了你的消息 -> 线程-2 【OrderServer】用户 User-0 的最新订单号为:TX20260705001 【OrderServer】用户 User-1 的最新订单号为:TX20260705001 【OrderServer】用户 User-2 的最新订单号为:TX20260705001
在这个阶段中,我们完成了 RPC 框架最核心的蜕变:用一条长连接 + Map,支撑起了并发多线程、多接口的远程调用。性能比基础案例提升了几个数量级。
不过,细心的你可能发现了另一个隐患:在 RpcClientContainer 中,我们所有的请求都共用一条 channel。如果高并发极端场景下,大量的读写数据导致这条 Channel 堵塞了(或者长连接因网络闪断意外关闭了),整个客户端就会陷入瘫痪。为了解决这个问题,在接下来的案例升级中,业界(包括 Dubbo)通常会引入连接池(Channel Pool) 或者心跳检测与断线重连机制来解决这个问题。
实现连接池与心跳重连机制 需要解决的问题 在 Stage 2 中,我们实现了一条长连接的多路复用。但如果在生产环境中,这条常驻的长连接因为网络闪断、运营商抖动或者长时间没有数据传输被防火墙切断,客户端的后续请求就会全部失败。为了让这个 RPC 框架达到 “工业级” 的健壮性,Stage 3 我们需要攻克两大核心机制:
心跳检测与自动重连(Keep-Alive & Reconnect) :在空闲时发送微小的探测包,发现连接断开立刻自动重新连接。
连接池(Channel Pool) :高并发下如果单条 Channel 成为瓶颈或意外受损,连接池可以提供多条通道分担压力,并提供自动剔除坏连接的能力。
Netty 提供了开箱即用的 IdleStateHandler(空闲状态处理器)。
客户端 :如果在一段时间内(比如 5 秒)没有向服务端发任何数据,就自动发送一个“心跳契约包”给服务端。
服务端 :如果长时间(比如 15 秒)没有收到客户端的任何请求或心跳,则认为客户端已掉线,主动关闭该 Channel 释放资源。
重连 :客户端如果检测到 Channel 断开了,触发监听器自动发起异步重连。
代码实现演进 契约升级 为了支持心跳,我们先定义一个轻量级的 “心跳乒乓包”,不走复杂的业务反射。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 @Data public class RpcRequest implements Serializable { private String requestId; private String className; private String methodName; private Class<?>[] parameterTypes; private Object[] parameters; private boolean heartbeat; public static RpcRequest createHeartbeat () { RpcRequest r = new RpcRequest (); r.setHeartbeat(true ); return r; } }
服务端升级 服务端开启空闲超时,踢掉死连接。使用 IdleStateHandler 监听读空闲。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 public class RpcServerHandler extends SimpleChannelInboundHandler <RpcRequest> { private final Map<String, Object> serviceMap; public RpcServerHandler (Map<String, Object> serviceMap) { this .serviceMap = serviceMap; } @Override protected void channelRead0 (ChannelHandlerContext ctx, RpcRequest request) throws Exception { if (request.isHeartbeat()) { System.out.println("【Server】收到客户端心跳维持包..." ); return ; } RpcResponse response = new RpcResponse (); response.setRequestId(request.getRequestId()); try { Object serviceBean = serviceMap.get(request.getClassName()); if (serviceBean == null ) { throw new RuntimeException ("未找到服务实现类: " + request.getClassName()); } Method method = serviceBean.getClass().getMethod(request.getMethodName(), request.getParameterTypes()); Object result = method.invoke(serviceBean, request.getParameters()); response.setResult(result); } catch (Throwable t) { response.setError(t); } ctx.writeAndFlush(response); } @Override public void userEventTriggered (ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { System.out.println("【Server警告】长时间未收到客户端消息,判定为无效连接,断开。" ); ctx.close(); } else { super .userEventTriggered(ctx, evt); } } }
服务端 Pipeline 配置修改:
1 2 3 4 5 6 7 8 9 10 11 @Override protected void initChannel (SocketChannel ch) throws Exception { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new IdleStateHandler (15 , 0 , 0 , TimeUnit.SECONDS)); pipeline.addLast(new ObjectDecoder (Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null ))); pipeline.addLast(new ObjectEncoder ()); pipeline.addLast(new RpcServerHandler (serviceMap)); }
客户端升级 客户端需要在空闲时发心跳,并且在连接断开时,利用 Netty 的 ChannelFutureListener 发起自动重连。
初次启动守住大门:当你在 ClientMain 中调用 container.init() 时,主线程会卡在 future.sync() 这一行,直到网络通道彻底打通、this.channel 被成功赋值,主线程才会被放行去创建代理和启动多线程测试。
重连不影响线程池:当后面网络意外断开时,触发 channelInactive,此时我们调用 doConnect(false) 走异步分支。在重连成功的几秒钟空窗期内,如果业务线程来调用,会被我们的 “健壮性补丁” 拦截,抛出可控的业务异常,而不会报底层无厘头的 NullPointerException。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 public class RpcClientContainer { private final String host; private final int port; private Channel channel; private RpcClientHandler clientHandler; private Bootstrap bootstrap; private final EventLoopGroup group = new NioEventLoopGroup (); public RpcClientContainer (String host, int port) { this .host = host; this .port = port; } public void init () throws Exception { clientHandler = new RpcClientHandler (); bootstrap = new Bootstrap (); bootstrap.group(group) .channel(NioSocketChannel.class) .option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 5000 ) .handler(new ChannelInitializer <SocketChannel>() { @Override protected void initChannel (SocketChannel ch) { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new IdleStateHandler (0 , 5 , 0 , TimeUnit.SECONDS)); pipeline.addLast(new ObjectEncoder ()); pipeline.addLast(new ObjectDecoder (Integer.MAX_VALUE, ClassResolvers.cacheDisabled(null ))); pipeline.addLast(new RpcClientHandlerInner ()); } }); doConnect(true ); } private void doConnect (boolean isInit) { if (channel != null && channel.isActive()) { return ; } System.out.println("尝试连接 RPC 服务器 [" + host + ":" + port + "]..." ); try { ChannelFuture future = bootstrap.connect(host, port); if (isInit) { future.sync(); this .channel = future.channel(); System.out.println("RPC 客户端初次连接成功!" ); } else { future.addListener((ChannelFutureListener) f -> { if (f.isSuccess()) { this .channel = f.channel(); System.out.println("RPC 客户端自动重连成功!" ); } else { System.out.println("RPC 客户端重连失败,3秒后准备再次重连..." ); f.channel().eventLoop().schedule(() -> doConnect(false ), 3 , TimeUnit.SECONDS); } }); } } catch (Exception e) { if (isInit) { System.err.println("RPC 客户端初次连接服务器失败!原因: " + e.getMessage()); throw new RuntimeException ("RPC 客户端初始化失败,无法连接服务器" , e); } } } public void close () { group.shutdownGracefully(); } private class RpcClientHandlerInner extends SimpleChannelInboundHandler <RpcResponse> { @Override protected void channelRead0 (ChannelHandlerContext ctx, RpcResponse msg) throws Exception { clientHandler.channelRead0(ctx, msg); } @Override public void userEventTriggered (ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { ctx.writeAndFlush(RpcRequest.createHeartbeat()); } else { super .userEventTriggered(ctx, evt); } } @Override public void channelInactive (ChannelHandlerContext ctx) throws Exception { System.out.println("【客户端警告】检测到连接断开,触发自动重连机制!" ); super .channelInactive(ctx); ctx.channel().eventLoop().schedule(() -> doConnect(false ), 3 , TimeUnit.SECONDS); } } @SuppressWarnings("unchecked") public <T> T createProxy (Class<T> serviceClass) { return (T) Proxy.newProxyInstance( serviceClass.getClassLoader(), new Class <?>[]{serviceClass}, (proxy, method, args) -> { RpcRequest request = new RpcRequest (); request.setRequestId(UUID.randomUUID().toString()); request.setClassName(serviceClass.getName()); request.setMethodName(method.getName()); request.setParameterTypes(method.getParameterTypes()); request.setParameters(args); CompletableFuture<RpcResponse> future = clientHandler.registerFuture(request.getRequestId()); if (channel == null || !channel.isActive()) { throw new RuntimeException ("RPC 调用失败:当前网络连接已断开,正在尝试重连中..." ); } channel.writeAndFlush(request); RpcResponse response = future.get(); if (response.getError() != null ) { throw response.getError(); } return response.getResult(); } ); } }
测试健壮性 微调 ClientMain:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 public class ClientMain { public static void main (String[] args) throws Exception { RpcClientContainer container = new RpcClientContainer ("127.0.0.1" , 8080 ); container.init(); HelloService helloService = container.createProxy(HelloService.class); OrderService orderService = container.createProxy(OrderService.class); for (int i = 0 ; i < 3 ; i++) { final int index = i; new Thread (() -> { String res1 = helloService.sayHello("线程-" + index); System.out.println(res1); String res2 = orderService.getOrderSymbol("User-" + index); System.out.println(res2); }).start(); } } }
依次进行执行验证:
正常启动 ServerMain 和 ClientMain。
客户端开始发送请求。此时你可以看到服务端定期打印:【Server】收到客户端心跳维持包…
断网模拟:直接在控制台强行把 ServerMain 关闭(Kill 掉)。
客户端控制台会立刻报警:【客户端警告】检测到连接断开,触发自动重连机制! 并每隔 3 秒持续打印 尝试连接 RPC 服务器…
恢复模拟:重新把 ServerMain 启动起来。
客户端会在几秒内捕捉到服务的复活,打印 RPC 客户端连接成功!,整条链路完好如初,不需要重启客户端!
自定义协议与高性能序列化升级 自定义私有协议 在这一阶段,我们将彻底告别已经被废弃且漏洞百出的 ObjectDecoder/Encoder。我们将像 Dubbo 一样,设计一套专属的自定义二进制私有协议。我们设计的协议头共占 8 个字节,结构如下:
1 2 3 4 5 6 +---------------------------------------------------------------+ | 魔数 (4B) | 序列化类型 (1B) | 消息类型 (1B) | 数据长度 (2B) | +---------------------------------------------------------------+ | 真正的数据 (Payload) | +---------------------------------------------------------------+
魔数 (Magic Number) - 4字节:固定为 0xCAFEBABE(或者你喜欢的任意 4 字节)。用来筛选和过滤非本框架的非法请求,防止恶意攻击。
序列化类型 (Serialization Type) - 1字节:标记使用的是哪种序列化方式。例如:0x01 代表 Kryo,0x02 代表 JSON。这为以后的多序列化扩展打下基础。
消息类型 (MessageType) - 1字节:标记这是一个请求、响应还是心跳包。例如:0x01 请求,0x02 响应,0x03 心跳。这样心跳包就不需要复杂的业务反序列化了。
数据长度 (Length) - 2字节:标记后面真正的数据体(Payload)占用了多少字节。最大支持 $2^{16}-1 = 65535$ 字节(如果需要传大对象,可改为 4 字节)。
引入高性能序列化组件 Kryo Kryo 是一个快速、高效的 Java 对象图形序列化框架,性能和压缩率远超 JDK 原生序列化。
1 2 3 4 5 <dependency > <groupId > com.esotericsoftware</groupId > <artifactId > kryo</artifactId > <version > 5.5.0</version > </dependency >
Kryo 工具类封装 由于 Kryo 不是线程安全的,必须使用 ThreadLocal 来保证每个线程拥有独立的 Kryo 实例。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 import com.esotericsoftware.kryo.Kryo;import com.esotericsoftware.kryo.io.Input;import com.esotericsoftware.kryo.io.Output;import com.owlias.janus.comet.client.zdemo.common.RpcRequest;import com.owlias.janus.comet.client.zdemo.common.RpcResponse;import java.io.ByteArrayInputStream;import java.io.ByteArrayOutputStream;public class KryoSerializer { private static final ThreadLocal<Kryo> kryoThreadLocal = ThreadLocal.withInitial(() -> { Kryo kryo = new Kryo (); kryo.register(RpcRequest.class); kryo.register(RpcResponse.class); kryo.register(Class[].class); kryo.register(Class.class); kryo.register(Object[].class); kryo.setReferences(true ); kryo.setRegistrationRequired(false ); return kryo; }); public static byte [] serialize(Object obj) { Kryo kryo = kryoThreadLocal.get(); ByteArrayOutputStream out = new ByteArrayOutputStream (); Output output = new Output (out); kryo.writeClassAndObject(output, obj); output.close(); return out.toByteArray(); } public static Object deserialize (byte [] bytes) { Kryo kryo = kryoThreadLocal.get(); ByteArrayInputStream in = new ByteArrayInputStream (bytes); Input input = new Input (in); Object obj = kryo.readClassAndObject(input); input.close(); return obj; } }
编写自定义编解码器 编码器:将对象打包成二进制流(Outbound)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 import com.owlias.janus.comet.client.zdemo.common.RpcRequest;import com.owlias.janus.comet.client.zdemo.serialize.KryoSerializer;import io.netty.buffer.ByteBuf;import io.netty.channel.ChannelHandlerContext;import io.netty.handler.codec.MessageToByteEncoder;public class MyProtocolEncoder extends MessageToByteEncoder <Object> { private static final int MAGIC_NUMBER = 0xCAFEBABE ; @Override protected void encode (ChannelHandlerContext ctx, Object msg, ByteBuf out) throws Exception { out.writeInt(MAGIC_NUMBER); out.writeByte((byte ) 0x01 ); if (msg instanceof RpcRequest) { RpcRequest req = (RpcRequest) msg; if (req.isHeartbeat()) { out.writeByte((byte ) 0x03 ); } else { out.writeByte((byte ) 0x01 ); } } else { out.writeByte((byte ) 0x02 ); } byte [] bodyBytes = KryoSerializer.serialize(msg); out.writeShort(bodyBytes.length); out.writeBytes(bodyBytes); } }
解码器:解决粘包半包,剥离协议头(Inbound)
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 import com.owlias.janus.comet.client.zdemo.common.RpcRequest;import com.owlias.janus.comet.client.zdemo.serialize.KryoSerializer;import io.netty.buffer.ByteBuf;import io.netty.channel.ChannelHandlerContext;import io.netty.handler.codec.LengthFieldBasedFrameDecoder;public class MyProtocolDecoder extends LengthFieldBasedFrameDecoder { private static final int MAGIC_NUMBER = 0xCAFEBABE ; private static final int HEADER_LENGTH = 8 ; public MyProtocolDecoder () { super (65535 , 6 , 2 , 0 , 0 ); } @Override protected Object decode (ChannelHandlerContext ctx, ByteBuf in) throws Exception { ByteBuf frame = (ByteBuf) super .decode(ctx, in); if (frame == null ) return null ; try { int magic = frame.readInt(); if (magic != MAGIC_NUMBER) { throw new IllegalArgumentException ("非法协议魔数: " + Integer.toHexString(magic)); } byte serializeType = frame.readByte(); byte messageType = frame.readByte(); int length = frame.readShort(); if (messageType == (byte ) 0x03 ) { return RpcRequest.createHeartbeat(); } byte [] bodyBytes = new byte [length]; frame.readBytes(bodyBytes); return KryoSerializer.deserialize(bodyBytes); } finally { frame.release(); } } }
替换原编解码器 编解码器写好后,我们在服务端和客户端的 initChannel 中,把之前废弃的 ObjectDecoder/Encoder 移出队列,换上我们帅气的专属私有协议。
服务端 RpcServer 修改:
1 2 3 4 pipeline.addLast(new IdleStateHandler (15 , 0 , 0 , TimeUnit.SECONDS)); pipeline.addLast(new MyProtocolDecoder ()); pipeline.addLast(new MyProtocolEncoder ()); pipeline.addLast(new RpcServerHandler (serviceMap));
客户端 RpcClientContainer 修改:
1 2 3 4 pipeline.addLast(new IdleStateHandler (0 , 5 , 0 , TimeUnit.SECONDS)); pipeline.addLast(new MyProtocolEncoder ()); pipeline.addLast(new MyProtocolDecoder ()); pipeline.addLast(new RpcClientHandlerInner ());
测试验证 按照阶段3测试,我们依然能够正常稳定运行。下面是客户端部分日志:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 尝试连接 RPC 服务器 [127.0.0.1:8080]... RPC 客户端初次连接成功! 【Server 回应】: 你好,我已经收到了你的消息 -> 线程-0 【Server 回应】: 你好,我已经收到了你的消息 -> 线程-1 【Server 回应】: 你好,我已经收到了你的消息 -> 线程-2 【OrderServer】用户 User-0 的最新订单号为:TX20260705001 【OrderServer】用户 User-1 的最新订单号为:TX20260705001 【OrderServer】用户 User-2 的最新订单号为:TX20260705001 【客户端警告】检测到连接断开,触发自动重连机制! 尝试连接 RPC 服务器 [127.0.0.1:8080]... RPC 客户端重连失败,3秒后准备再次重连... 尝试连接 RPC 服务器 [127.0.0.1:8080]... RPC 客户端重连失败,3秒后准备再次重连... 尝试连接 RPC 服务器 [127.0.0.1:8080]... RPC 客户端自动重连成功!
服务发现与软负载均衡 需要解决的问题 在前面的阶段中,我们的客户端(Consumer)在 RpcClientContainer 里把服务端的 IP 和端口(127.0.0.1:8080)给硬编码死了。但在实际的生产分布式系统中,服务端往往是一个集群。当某个服务端节点宕机或者新加了机器时,客户端必须能够动态感知,并将请求均匀地分摊到各个健康的节点上。这就是服务注册中心与负载均衡的核心价值。
为了遵循由浅入深的原则,我们不急于立刻引入复杂的 ZooKeeper 或 Nacos 外部集群,而是先在本地使用一个标准的接口设计,抽象出一个 虚拟的/轻量级的分布式注册中心模型。
服务提供者 (Provider):启动时,将自己提供的接口名和自己的 IP:Port 注册到注册中心。
服务消费者 (Consumer):调用方法前,先去注册中心“拉取”该接口对应的全部可用地址列表(如 [192.168.1.10:8080, 192.168.1.11:8080])。
负载均衡 (Load Balance):客户端从地址列表中,通过一定的算法(如轮询 RoundRobin 或随机 Random)挑选出一个地址,然后发起网络调用。
代码实现演进 注册与负载接口 首先,我们需要解耦,定义好服务注册、发现以及负载均衡的通用标准。
1 2 3 4 5 6 7 8 9 10 11 12 13 public interface ServiceRegistry { void register (String serviceName, String serverAddress) ; List<String> discover (String serviceName) ; } public interface LoadBalance { String select (List<String> addressList) ; }
注册中心实现 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 public class LocalMockRegistry implements ServiceRegistry { private static final Map<String, List<String>> registryCenter = new ConcurrentHashMap <>(); @Override public void register (String serviceName, String serverAddress) { registryCenter.computeIfAbsent(serviceName, k -> new CopyOnWriteArrayList <>()).add(serverAddress); System.out.println("【注册中心】服务 [" + serviceName + "] 成功注册新节点 -> " + serverAddress); } @Override public List<String> discover (String serviceName) { return registryCenter.getOrDefault(serviceName, new ArrayList <>()); } }
负载均衡器实现 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 public class RoundRobinLoadBalance implements LoadBalance { private final AtomicInteger index = new AtomicInteger (0 ); @Override public String select (List<String> addressList) { if (addressList == null || addressList.isEmpty()) { return null ; } int pos = Math.abs(index.getAndIncrement() % addressList.size()); return addressList.get(pos); } }
客户端升级 现在,客户端 RpcClientContainer 不再只盯着一个 host:port 了,它需要在运行时,根据你要调用的接口,动态去注册中心找地址!为了不频繁建连,我们内部通过一个 Map<String, Channel> 维护针对不同服务器地址的长连接。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 public class RpcClientRouter { private final ServiceRegistry registry; private final LoadBalance loadBalance = new RoundRobinLoadBalance (); private final Map<String, Channel> channelMap = new ConcurrentHashMap <>(); private final EventLoopGroup group = new NioEventLoopGroup (); private final RpcClientHandler clientHandler = new RpcClientHandler (); public RpcClientRouter (ServiceRegistry registry) { this .registry = registry; } private Channel getChannel (String address) throws Exception { if (channelMap.containsKey(address)) { Channel ch = channelMap.get(address); if (ch.isActive()) return ch; } String[] parts = address.split(":" ); String host = parts[0 ]; int port = Integer.parseInt(parts[1 ]); Bootstrap b = new Bootstrap (); b.group(group).channel(NioSocketChannel.class) .handler(new ChannelInitializer <>() { @Override protected void initChannel (Channel ch) { ch.pipeline().addLast(new MyProtocolEncoder ()); ch.pipeline().addLast(new MyProtocolDecoder ()); ch.pipeline().addLast(clientHandler); } }); ChannelFuture f = b.connect(host, port).sync(); Channel channel = f.channel(); channelMap.put(address, channel); return channel; } @SuppressWarnings("unchecked") public <T> T createProxy (Class<T> serviceClass) { return (T) Proxy.newProxyInstance( serviceClass.getClassLoader(), new Class <?>[]{serviceClass}, (proxy, method, args) -> { String serviceName = serviceClass.getName(); List<String> addresses = registry.discover(serviceName); if (addresses.isEmpty()) { throw new RuntimeException ("RPC 错误:没有找到可用的服务节点 [" + serviceName + "]" ); } String selectedAddress = loadBalance.select(addresses); System.out.println("负载均衡 selectedAddress:" + selectedAddress); Channel channel = getChannel(selectedAddress); RpcRequest request = new RpcRequest (); request.setRequestId(UUID.randomUUID().toString()); request.setClassName(serviceName); request.setMethodName(method.getName()); request.setParameterTypes(method.getParameterTypes()); request.setParameters(args); var future = clientHandler.registerFuture(request.getRequestId()); channel.writeAndFlush(request); return future.get().getResult(); } ); } }
由于这里的 clientHandler 是一个共享的对象 ,它允许同一个实例被同时配置到多个 Channel 中,这种情况下你必须声明它是可被共享的,否则当你两次 rpcClientRouter.getChannel 就会报错:
1 com.xxx.client.zdemo.client.RpcClientHandler is not a @Sharable handler, so can't be added or removed multiple times.
声明 RpcClientHandler 是共享的(它实际上是无状态的):
1 2 3 4 @ChannelHandler .Sharable public class RpcClientHandler extends SimpleChannelInboundHandler <RpcResponse> { }
测试验证 为了方便测试,我们做如下新增或修改。
第一,RpcServer 增加 异步启动方法:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 public void startAsync () { Thread serverThread = new Thread (() -> { try { this .start(); } catch (Exception e) { System.err.println("【Server 错误】端口 " + port + " 异步启动失败: " + e.getMessage()); } }, "RpcServer-Binder-Thread-" + port); serverThread.start(); }
第二,新增服务端和客户端测试类:
ClusterServerMain
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 public class ClusterServerMain { public static void main (String[] args) throws Exception { ServiceRegistry registry = new LocalMockRegistry (); RpcServer server1 = new RpcServer (8080 ); server1.registerService(HelloService.class, new HelloServiceImpl ()); server1.startAsync(); registry.register(HelloService.class.getName(), "127.0.0.1:8080" ); RpcServer server2 = new RpcServer (8081 ); server2.registerService(HelloService.class, new HelloServiceImpl ()); server2.startAsync(); registry.register(HelloService.class.getName(), "127.0.0.1:8081" ); } }
ClusterClientMain
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 public class ClusterClientMain { public static void main (String[] args) throws Exception { ServiceRegistry registry = new LocalMockRegistry (); registry.register(HelloService.class.getName(), "127.0.0.1:8080" ); registry.register(HelloService.class.getName(), "127.0.0.1:8081" ); RpcClientRouter router = new RpcClientRouter (registry); HelloService helloService = router.createProxy(HelloService.class); for (int i = 1 ; i <= 4 ; i++) { String result = helloService.sayHello("集群测试 " + i); System.out.println("客户端收到结果: " + result); } } }
分别启动服务端和客户端:
服务端日志:
1 2 3 4 【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8080 【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8081 RPC Server 成功启动,监听端口: 8080 RPC Server 成功启动,监听端口: 8081
客户端日志:
1 2 3 4 5 6 7 8 9 10 【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8080 【注册中心】服务 [com.xxx.client.zdemo.common.HelloService] 成功注册新节点 -> 127.0.0.1:8081 负载均衡 selectedAddress:127.0.0.1:8080 客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 1 负载均衡 selectedAddress:127.0.0.1:8081 客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 2 负载均衡 selectedAddress:127.0.0.1:8080 客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 3 负载均衡 selectedAddress:127.0.0.1:8081 客户端收到结果: 【Server 回应】: 你好,我已经收到了你的消息 -> 集群测试 4
标题:
Java NIO - 使用 netty 模拟 dubbo RPC 远程服务调用