关于手写 RPC 的一个 Bug

前言

先简单说一下我目前的情况,我学到第六章快完了,前面几章基本跟着实现了,也跟着提示优化了一下,目前有的优化是:支持读取 .yml 格式的配置文件、配置文件支持中文、接口 Mock 支持更多返回类型的默认值、序列化器工厂和 SPI Loader 实现懒加载。 在手写 RPC 第六章 注册中心优化里,有一个问题。 就是写到 3、服务缓存更新 - 监听机制这里的测试,鱼皮哥提示的测试方式如下:

image.png

我给大家看看我的消费者和提供者的代码

java
复制代码
public static void main(String[] args) throws InterruptedException { // 获取代理 UserService userService = ServiceProxyFactory.getProxy(UserService.class); // 打印服务名称和版本,确认与提供者一致 System.out.println("查找服务:" + UserService.class.getName()); // 第一次调用服务发现 System.out.println("第一次调用服务发现:"); User user = new User(); user.setName("hsurosy"); User newUser1 = userService.getUser(user); if (newUser1 != null) { System.out.println("第一次服务发现结果: " + newUser1.getName()); } else { System.out.println("第一次服务发现结果: user == null"); } // 第二次调用服务发现 System.out.println("第二次调用服务发现:"); User newUser2 = userService.getUser(user); if (newUser2 != null) { System.out.println("第二次服务发现结果: " + newUser2.getName()); } else { System.out.println("第二次服务发现结果: user == null"); } // 增加延迟 System.out.println("等待10秒后进行第三次服务发现..."); Thread.sleep(10000); // 10秒的延迟时间,这期间可以关闭服务提供者以测试缓存失效 // 第三次调用服务发现 System.out.println("第三次调用服务发现:"); User newUser3 = userService.getUser(user); if (newUser3 != null) { System.out.println("第三次服务发现结果: " + newUser3.getName()); } else { System.out.println("第三次服务发现结果: user == null"); } }
java
复制代码
public static void main(String[] args) { // 初始化 RPC 框架 RpcApplication.init(); // 本地注册服务 String serviceName = UserService.class.getName(); LocalRegistry.register(serviceName, UserServiceImpl.class); // 注册服务到注册中心 RpcConfig rpcConfig = RpcApplication.getRpcConfig(); RegistryConfig registryConfig = rpcConfig.getRegistryConfig(); Registry registry = RegistryFactory.getInstance(registryConfig.getRegistry()); ServiceMetaInfo serviceMetaInfo = new ServiceMetaInfo(); serviceMetaInfo.setServiceName(serviceName); serviceMetaInfo.setServiceHost(rpcConfig.getServerHost()); serviceMetaInfo.setServicePort(rpcConfig.getServerPort()); try { registry.register(serviceMetaInfo); // 服务注册 System.out.println("服务注册成功:" + serviceMetaInfo.getServiceAddress()); } catch (Exception e) { throw new RuntimeException(e); } // 启动 Web 服务 HttpServer httpServer = new VertxHttpServer(); httpServer.doStart(RpcApplication.getRpcConfig().getServerPort()); }

我是结合评论区第一个大佬的提示 debug 的

image.png

也顺便贴出 EtcdRegistry 的代码,这里是已经解决问题后的代码

java
复制代码
package com.hsu.hsurpc.registry; import cn.hutool.core.collection.CollUtil; import cn.hutool.core.collection.ConcurrentHashSet; import cn.hutool.cron.CronUtil; import cn.hutool.cron.task.Task; import cn.hutool.json.JSONUtil; import com.hsu.hsurpc.config.RegistryConfig; import com.hsu.hsurpc.model.ServiceMetaInfo; import io.etcd.jetcd.*; import io.etcd.jetcd.options.GetOption; import io.etcd.jetcd.options.PutOption; import io.etcd.jetcd.watch.WatchEvent; import java.nio.charset.StandardCharsets; import java.time.Duration; import java.util.HashSet; import java.util.List; import java.util.Set; import java.util.stream.Collectors; /** * Etcd 注册中心实现类 * @Author Hsu琛君珩 * @Date 2024-09-20 16:55 * @Description * @Version: v1.0.0 */ public class EtcdRegistry implements Registry { /** * Etcd 客户端实例 */ private Client client; /** * KV 客户端,用于与 Etcd 进行键值操作 */ private KV kvClient; /** * Etcd 根节点路径,防止不同项目的数据冲突 */ private static final String ETCD_ROOT_PATH = "/rpc/"; /** * 本机注册的节点 key 集合(用于维护续期) */ private final Set<String> localRegisterNodeKeySet = new HashSet<>(); /** * 注册中心服务缓存(支持多个服务键) */ private final RegistryServiceMultiCache registryServiceMultiCache = new RegistryServiceMultiCache(); /** * 正在监听的 key 集合 */ private final Set<String> watchingKeySet = new ConcurrentHashSet<>(); /** * 初始化 Etcd 客户端,连接到注册中心 * * @param registryConfig 注册中心的配置信息,包含连接地址、超时时间等 */ @Override public void init(RegistryConfig registryConfig) { client = Client.builder() .endpoints(registryConfig.getAddress()) // 使用配置的地址 .connectTimeout(Duration.ofMillis(registryConfig.getTimeout())) // 设置超时时间 .build(); kvClient = client.getKVClient(); // 获取 KV 操作客户端 // 启动心跳机制,定期检查并续签服务节点的租约 heartBeat(); } /** * 注册服务到 Etcd 注册中心 * 服务注册通过键值对的方式存储,并为服务设置租约(过期时间), * 同时将注册的服务节点信息添加到本地缓存以便后续管理。 * * @param serviceMetaInfo 服务的元数据信息(包括服务名称、版本、地址、端口等) * @throws Exception 如果注册失败,抛出异常 */ @Override public void register(ServiceMetaInfo serviceMetaInfo) throws Exception { // 创建 Lease 和 KV 客户端 Lease leaseClient = client.getLeaseClient(); // 创建 30 秒租约,用于设置服务节点的过期时间 long leaseId = leaseClient.grant(30).get().getID(); // 设置要存储的键值对,键为服务节点的唯一键名 String registerKey = ETCD_ROOT_PATH + serviceMetaInfo.getServiceNodeKey(); // /rpc/com.hsu.example.common.service.UserService:1.0/localhost:8080 ByteSequence key = ByteSequence.from(registerKey, StandardCharsets.UTF_8); ByteSequence value = ByteSequence.from(JSONUtil.toJsonStr(serviceMetaInfo), StandardCharsets.UTF_8); // 使用租约设置键值对并存储到 Etcd,并设置过期时间 PutOption putOption = PutOption.builder().withLeaseId(leaseId).build(); kvClient.put(key, value, putOption).get(); // 添加节点信息到本地缓存 localRegisterNodeKeySet.add(registerKey); } /** * 注销服务节点 * 从 Etcd 注册中心中删除已注册的服务节点信息,并从本地缓存中移除该节点的键名。 * * @param serviceMetaInfo 要注销的服务的元数据信息 */ @Override public void unRegister(ServiceMetaInfo serviceMetaInfo) { // 构建服务节点的唯一键名,格式为 "/rpc/{serviceName}:{version}/{host}:{port}" String registerKey = ETCD_ROOT_PATH + serviceMetaInfo.getServiceNodeKey(); // 从 Etcd 中删除与该服务节点对应的键值对 kvClient.delete(ByteSequence.from(registerKey, StandardCharsets.UTF_8)); // 从本地缓存移除 localRegisterNodeKeySet.remove(registerKey); } /** * 服务发现,根据服务键名前缀,查找已注册的所有服务节点信息 * * @param serviceKey 服务的唯一键名(通常由服务名称和版本号组成) * @return 返回该服务下的所有节点列表 */ @Override public List<ServiceMetaInfo> serviceDiscovery(String serviceKey) { // 优先从缓存获取服务 List<ServiceMetaInfo> cachedServiceMetaInfoList = registryServiceMultiCache.readCache(serviceKey); if (cachedServiceMetaInfoList != null) { System.out.println("发现服务 [" + serviceKey + "] 使用缓存,服务节点信息:" + cachedServiceMetaInfoList); return cachedServiceMetaInfoList; } // 从注册中心查询服务节点 String searchPrefix = ETCD_ROOT_PATH + serviceKey; try { GetOption getOption = GetOption.builder().isPrefix(true).build(); List<KeyValue> keyValues = kvClient.get( ByteSequence.from(searchPrefix, StandardCharsets.UTF_8), getOption) .get() .getKvs(); List<ServiceMetaInfo> serviceMetaInfoList = keyValues.stream() .map(keyValue -> { String key = keyValue.getKey().toString(StandardCharsets.UTF_8); watch(key); // 监听 key 的变化,这里监听的是完整的 key:/rpc/com.hsu.example.common.service.UserService:1.0:localhost:8080 String value = keyValue.getValue().toString(StandardCharsets.UTF_8); return JSONUtil.toBean(value, ServiceMetaInfo.class); }) .collect(Collectors.toList()); System.out.println("发现服务 [" + serviceKey + "] 不使用缓存,服务节点信息:" + serviceMetaInfoList); registryServiceMultiCache.writeCache(serviceKey, serviceMetaInfoList); return serviceMetaInfoList; } catch (Exception e) { throw new RuntimeException("获取服务列表失败", e); } } /** * 销毁注册中心,释放 Etcd 客户端资源 * 在项目关闭或服务停止时调用此方法,断开与 Etcd 注册中心的连接, * 下线所有已注册的服务节点,并释放相关的客户端资源。 */ @Override public void destroy() { System.out.println("当前节点下线"); // 下线节点:遍历本节点所有的 key 并删除 for (String key : localRegisterNodeKeySet) { try { kvClient.delete(ByteSequence.from(key, StandardCharsets.UTF_8)).get(); } catch (Exception e) { throw new RuntimeException(key + "节点下线失败", e); } } // 释放 KV 客户端和 Etcd 客户端的资源 if (kvClient != null) { kvClient.close(); } if (client != null) { client.close(); } } /** * 定期发送心跳信号,保持服务节点的存活状态 * 通过定时任务,每隔 10 秒遍历本地注册的所有服务节点,检查它们在 Etcd 中的存活状态, * 如果节点未过期,则重新注册该节点(续签租约),以延长其在 Etcd 中的存活时间。 * 如果节点已过期,则不进行续签处理。 */ @Override public void heartBeat() { // 使用 CronUtil 设置一个每 10 秒执行一次的定时任务 CronUtil.schedule("*/10 * * * * *", new Task() { @Override public void execute() { // 遍历本地缓存中所有已注册的服务节点键 for (String key : localRegisterNodeKeySet) { try { // 从 Etcd 中获取对应的服务节点信息 List<KeyValue> keyValues = kvClient.get(ByteSequence.from(key, StandardCharsets.UTF_8)) .get() .getKvs(); // 检查该节点是否已过期 if (CollUtil.isEmpty(keyValues)) { // 如果节点已过期,跳过续签 continue; } // 如果节点未过期,进行重新注册以续签租约 KeyValue keyValue = keyValues.get(0); String value = keyValue.getValue().toString(StandardCharsets.UTF_8); ServiceMetaInfo serviceMetaInfo = JSONUtil.toBean(value, ServiceMetaInfo.class); register(serviceMetaInfo); } catch (Exception e) { throw new RuntimeException(key + "续签失败", e); } } } }); // 支持秒级精度的定时任务 CronUtil.setMatchSecond(true); CronUtil.start(); } /** * 监听指定服务节点的变化 * 通过 Etcd 的 Watch 机制,监控指定服务节点的状态变化(如删除事件), * 当服务节点发生变化时,执行相应的处理逻辑。 * * @param serviceNodeKey 需要监听的服务节点的键名 */ @Override public void watch(String serviceNodeKey) { Watch watchClient = client.getWatchClient(); // 之前未被监听,开启监听 boolean newWatch = watchingKeySet.add(serviceNodeKey); if (newWatch) { watchClient.watch(ByteSequence.from(serviceNodeKey, StandardCharsets.UTF_8), response -> { for (WatchEvent event : response.getEvents()) { switch (event.getEventType()) { // key 被删除时触发 case DELETE: String serviceKey = extractServiceKey(serviceNodeKey); // 清理注册服务缓存 System.out.println("服务节点 [" + serviceKey + "] 被删除,清理缓存"); registryServiceMultiCache.clearCache(serviceKey); break; case PUT: default: break; } } }); } } // 提取出serviceKey的方法 private String extractServiceKey(String serviceNodeKey) { // 假设 serviceNodeKey 的格式是 /rpc/serviceName:version:host:port // 我们需要提取出 serviceName:version int startIndex = ETCD_ROOT_PATH.length(); // 去掉 /rpc/ 前缀 int endIndex = serviceNodeKey.lastIndexOf(':'); // 找到最后一个冒号 if (endIndex > startIndex) { endIndex = serviceNodeKey.lastIndexOf(':', endIndex - 1); // 找到倒数第二个冒号 return serviceNodeKey.substring(startIndex, endIndex); // 提取出 serviceKey } else { return serviceNodeKey.substring(startIndex); // 如果没有找到冒号,返回去掉前缀的部分 } } }

问题

按照测试的要求里,我们的三次调用中,第一次查询缓存没查到就去注册中心查了,第二次是查询缓存查到直接返回,第三次开始前(我们是通过 sleep 10秒期间)停掉提供者,可以通过 Etcd 可视化工具看见服务确实下线了,按道理应该是查询缓存查不到了的,但是如果你按照鱼皮哥原来的代码写的话,还是会进入查询缓存。这应该就不对了。

我看了一下鱼皮哥 gitee 上的代码,在 watch 方法里好像说了传入的应该是 serviceKey,但是没有改。 这里举例说一下 serviceKey = com.hsu.example.common.service.UserService:1.0serviceNodeKey = /rpc/com.hsu.example.common.service.UserService:1.0:localhost,形如这样的。

这里感觉也不太对,应该是 serviceNodeKey = com.hsu.example.common.service.UserService:1.0:localhost,上面说的是 watch 方法里的 serviceNodeKey

上面为什么不对的原因就是,写入缓存的时候存入的 key 是 serviceKey,但是删除缓存的时候的 key 是 serviceNodeKey,自然就不对了。我开始是通过日志的方式看是否删除成功,每次都显示成功,但我 debug 看缓存还在就红温了,折磨我好几个小时。

解决方法

我开始是想把传入 watch 里的 key 直接传入 serviceKey 就行了,但是这样不行,因为 watch 方法里监听的是完整的 serviceNodeKey,因此只能改 serviceNodeKey 改为 serviceKey,按照下面的方法改就能解决上面的问题了

java
复制代码
/** * 监听指定服务节点的变化 * 通过 Etcd 的 Watch 机制,监控指定服务节点的状态变化(如删除事件), * 当服务节点发生变化时,执行相应的处理逻辑。 * * @param serviceNodeKey 需要监听的服务节点的键名 */ @Override public void watch(String serviceNodeKey) { Watch watchClient = client.getWatchClient(); // 之前未被监听,开启监听 boolean newWatch = watchingKeySet.add(serviceNodeKey); if (newWatch) { watchClient.watch(ByteSequence.from(serviceNodeKey, StandardCharsets.UTF_8), response -> { for (WatchEvent event : response.getEvents()) { switch (event.getEventType()) { // key 被删除时触发 case DELETE: String serviceKey = extractServiceKey(serviceNodeKey); // 清理注册服务缓存 System.out.println("服务节点 [" + serviceKey + "] 被删除,清理缓存"); registryServiceMultiCache.clearCache(serviceKey); break; case PUT: default: break; } } }); } } // 提取出serviceKey的方法 private String extractServiceKey(String serviceNodeKey) { // 假设 serviceNodeKey 的格式是 /rpc/serviceName:version:host:port // 我们需要提取出 serviceName:version int startIndex = ETCD_ROOT_PATH.length(); // 去掉 /rpc/ 前缀 int endIndex = serviceNodeKey.lastIndexOf(':'); // 找到最后一个冒号 if (endIndex > startIndex) { endIndex = serviceNodeKey.lastIndexOf(':', endIndex - 1); // 找到倒数第二个冒号 return serviceNodeKey.substring(startIndex, endIndex); // 提取出 serviceKey } else { return serviceNodeKey.substring(startIndex); // 如果没有找到冒号,返回去掉前缀的部分 } }

这样之后测试就能成功,只不过我觉得也不是很优雅,不过我也没有更好的办法了,就这样吧,这里仅个人而言,我觉得挺烦人的,不过本身实力一般,自然解决起来也就比较麻烦了,大家就看一乐吧,如果对你有帮助那就最好啦。如果你有不太懂的地方,那可能是你没做过这项目,可以先跟着做做,看看你会不会有遇到这个问题。

0个评论
点击登录,快来和大家讨论吧~
表情
图片
暂无评论
下载 APP