
前言
上篇文章《Dubbo之服务暴露》分析 Dubbo 服务是如何暴露的,本文接着分析 Dubbo 服务的消费流程。主要从以下几个方面进行分析:注册中心的暴露;通过注册中心进行服务消费通知;直连服务进行消费。 服务消费端启动时,将自身的信息注册到注册中心的目录,同时还订阅服务提供方的目录,当服务提供方的 URL 发生更改时,实时获取新的数据。
服务消费端流程
下面是一个服务消费的流程图:
上图中可以看到,服务消费的流程与服务暴露的流程有点类似逆向的。同样,Dubbo 服务也是分为两个大步骤:第一步就是将远程服务通过Protocol转换成Invoker(概念在上篇文章中有解释)。第二步通过动态代理将Invoker转换成消费服务需要的接口。
org.apache.dubbo.config.ReferenceConfig类是ReferenceBean的父类,与生产端服务的ServiceBean一样,存放着解析出来的 XML 和注解信息。类关系如下:

服务初始化中转换的入口
当我们消费端调用本地接口就能实现远程服务的调用,这是怎么实现的呢?根据上面的流程图,来分析消费原理。 在消费端进行初始化时ReferenceConfig#init,会执行ReferenceConfig#createProxy来完成这一系列操作。以下为ReferenceConfig#createProxy主要的代码部分:
1private T createProxy(Map<String, String> map) { 2 // 判断是否为 Jvm 本地引用 3 if (shouldJvmRefer(map)) { 4 // 通过 injvm 协议,获取本地服务 5 URL url = new URL(LOCAL_PROTOCOL, LOCALHOST_VALUE, 0, interfaceClass.getName()).addParameters(map); 6 invoker = REF_PROTOCOL.refer(interfaceClass, url); 7 } else { 8 urls.clear(); 9 // 判断是否有自定义的直连地址,或注册中心地址 10 if (url != null && url.length() > 0) { 11 String[] us = SEMICOLON_SPLIT_PATTERN.split(url); 12 if (us != null && us.length > 0) { 13 for (String u : us) { 14 URL url = URL.valueOf(u); 15 if (StringUtils.isEmpty(url.getPath())) { 16 url = url.setPath(interfaceName); 17 } 18 if (UrlUtils.isRegistry(url)) { 19 // 如果是注册中心Protocol类型,则向地址中添加 refer 服务消费元数据 20 urls.add(url.addParameterAndEncoded(REFER_KEY, StringUtils.toQueryString(map))); 21 } else { 22 // 直连服务提供端 23 urls.add(ClusterUtils.mergeUrl(url, map)); 24 } 25 } 26 } 27 } else { 28 // 组装注册中心的配置 29 if (!LOCAL_PROTOCOL.equalsIgnoreCase(getProtocol())) { 30 // 检查配置中心 31 checkRegistry(); 32 List<URL> us = ConfigValidationUtils.loadRegistries(this, false); 33 if (CollectionUtils.isNotEmpty(us)) { 34 for (URL u : us) { 35 URL monitorUrl = ConfigValidationUtils.loadMonitor(this, u); 36 if (monitorUrl != null) { 37 // 监控上报信息 38 map.put(MONITOR_KEY, URL.encode(monitorUrl.toFullString())); 39 } 40 // 注册中心地址添加 refer 服务消费元数据 41 urls.add(u.addParameterAndEncoded(REFER_KEY, StringUtils.toQueryString(map))); 42 } 43 } 44 } 45 } 46 47 // 只有一条注册中心数据,即单注册中心 48 if (urls.size() == 1) { 49 // 将远程服务转化成 Invoker 50 invoker = REF_PROTOCOL.refer(interfaceClass, urls.get(0)); 51 } else { 52 // 因为多注册中心就会存在多个 Invoker,这里用保存在 List 中 53 List<Invoker<?>> invokers = new ArrayList<Invoker<?>>(); 54 URL registryURL = null; 55 for (URL url : urls) { 56 // 将每个注册中心转换成 Invoker 数据 57 invokers.add(REF_PROTOCOL.refer(interfaceClass, url)); 58 if (UrlUtils.isRegistry(url)) { 59 // 会覆盖前遍历的注册中心,使用最后一条注册中心数据 60 registryURL = url; 61 } 62 } 63 if (registryURL != null) { 64 // 默认使用 zone-aware 策略来处理多个订阅 65 URL u = registryURL.addParameterIfAbsent(CLUSTER_KEY, ZoneAwareCluster.NAME); 66 // 将转换后的多个 Invoker 合并成一个 67 invoker = CLUSTER.join(new StaticDirectory(u, invokers)); 68 } else { 69 invoker = CLUSTER.join(new StaticDirectory(invokers)); 70 } 71 } 72 } 73 // 利用动态代理,将 Invoker 转换成本地接口代理 74 return (T) PROXY_FACTORY.getProxy(invoker); 75} 76
上面转换的过程中,主要可概括为:先分为本地引用和远程引用两类。本地就是以 inJvm 协议的获取本地服务,这不做过多说明;远程引用分为直连服务和通过注册中心。注册中心分为单注册中心和多注册中心的情况,单注册中心好解决,直接使用即可,多注册中心时,将转换后的 Invoker 合并成一个 Invoker。最后通过动态代理将 Invoker 转换成本地接口代理。
获取 Invoker 实例
由于本地服务时直接从缓存中获取,这里就注册中心的消费进行分析,上面代码片段中使用的是REF_PROTOCOL.refer进行转换,该方法代码:
1public <T> Invoker<T> refer(Class<T> type, URL url) throws RpcException { 2 // 获取服务的注册中心url,里面会设置注册中心的协议和移除 registry 的参数 3 url = getRegistryUrl(url); 4 // 获取注册中心实例 5 Registry registry = registryFactory.getRegistry(url); 6 if (RegistryService.class.equals(type)) { 7 return proxyFactory.getInvoker((T) registry, type, url); 8 } 9 10 // 获取服务消费元数据 11 Map<String, String> qs = StringUtils.parseQueryString(url.getParameterAndDecoded(REFER_KEY)); 12 // 从服务消费元数据中获取分组信息 13 String group = qs.get(GROUP_KEY); 14 if (group != null && group.length() > 0) { 15 if ((COMMA_SPLIT_PATTERN.split(group)).length > 1 || "*".equals(group)) { 16 // 执行 Invoker 转换工作 17 return doRefer(getMergeableCluster(), registry, type, url); 18 } 19 } 20 // 执行 Invoker 转换工作 21 return doRefer(cluster, registry, type, url); 22} 23
上面主要是获取服务消费的注册中心实例和进行服务分组,最后调用doRefer方法进行转换工作,以下为doRefer的代码:
1private <T> Invoker<T> doRefer(Cluster cluster, Registry registry, Class<T> type, URL url) { 2 // 创建 RegistryDirectory 对象 3 RegistryDirectory<T> directory = new RegistryDirectory<T>(type, url); 4 // 设置注册中心 5 directory.setRegistry(registry); 6 // 设置协议 7 directory.setProtocol(protocol); 8 // directory.getUrl().getParameters() 是服务消费元数据 9 Map<String, String> parameters = new HashMap<String, String>(directory.getUrl().getParameters()); 10 URL subscribeUrl = new URL(CONSUMER_PROTOCOL, parameters.remove(REGISTER_IP_KEY), 0, type.getName(), parameters); 11 if (!ANY_VALUE.equals(url.getServiceInterface()) && url.getParameter(REGISTER_KEY, true)) { 12 directory.setRegisteredConsumerUrl(getRegisteredConsumerUrl(subscribeUrl, url)); 13 // 消费消息注册到注册中心 14 registry.register(directory.getRegisteredConsumerUrl()); 15 } 16 17 directory.buildRouterChain(subscribeUrl); 18 // 服务消费者订阅:服务提供端,动态配置,路由的通知 19 directory.subscribe(subscribeUrl.addParameter(CATEGORY_KEY, 20 PROVIDERS_CATEGORY + "," + CONFIGURATORS_CATEGORY + "," + ROUTERS_CATEGORY)); 21 22 // 多个Invoker合并为一个 23 Invoker invoker = cluster.join(directory); 24 return invoker; 25} 26
上面实现主要是完成创建 RegistryDirectory 对象,将消费服务元数据注册到注册中心,通过 RegistryDirectory 对象里的信息,实现服务提供端,动态配置及路由的订阅相关功能。
RegistryDirectory 这个类实现了 NotifyListener 这个通知监听接口,当订阅的服务,配置或路由发生变化时,会接收到通知,进行相应改变:
1public synchronized void notify(List<URL> urls) { 2 // 将服务提供方配置,路由配置,服务提供方的服务分别以不同的 key 保存在 Map 中 3 Map<String, List<URL>> categoryUrls = urls.stream() 4 .filter(Objects::nonNull) 5 .filter(this::isValidCategory) 6 .filter(this::isNotCompatibleFor26x) 7 .collect(Collectors.groupingBy(url -> { 8 if (UrlUtils.isConfigurator(url)) { 9 return CONFIGURATORS_CATEGORY; 10 } else if (UrlUtils.isRoute(url)) { 11 return ROUTERS_CATEGORY; 12 } else if (UrlUtils.isProvider(url)) { 13 return PROVIDERS_CATEGORY; 14 } 15 return ""; 16 })); 17 18 // 更新服务提供方配置 19 List<URL> configuratorURLs = categoryUrls.getOrDefault(CONFIGURATORS_CATEGORY, Collections.emptyList()); 20 this.configurators = Configurator.toConfigurators(configuratorURLs).orElse(this.configurators); 21 22 // 更新路由配置 23 List<URL> routerURLs = categoryUrls.getOrDefault(ROUTERS_CATEGORY, Collections.emptyList()); 24 toRouters(routerURLs).ifPresent(this::addRouters); 25 26 // 加载服务提供方的服务信息 27 List<URL> providerURLs = categoryUrls.getOrDefault(PROVIDERS_CATEGORY, Collections.emptyList()); 28 /** 29 * 3.x added for extend URL address 30 */ 31 ExtensionLoader<AddressListener> addressListenerExtensionLoader = ExtensionLoader.getExtensionLoader(AddressListener.class); 32 List<AddressListener> supportedListeners = addressListenerExtensionLoader.getActivateExtension(getUrl(), (String[]) null); 33 if (supportedListeners != null && !supportedListeners.isEmpty()) { 34 for (AddressListener addressListener : supportedListeners) { 35 providerURLs = addressListener.notify(providerURLs, getUrl(),this); 36 } 37 } 38 // 重新加载 Invoker 实例 39 refreshOverrideAndInvoker(providerURLs); 40} 41
RegistryDirectory#notify里面最后会刷新 Invoker 进行重新加载,下面是核心代码的实现:
1private void refreshOverrideAndInvoker(List<URL> urls) { 2 // mock zookeeper://xxx?mock=return null 3 overrideDirectoryUrl(); 4 // 刷新 invoker 5 refreshInvoker(urls); 6} 7 8private void refreshInvoker(List<URL> invokerUrls) { 9 Assert.notNull(invokerUrls, "invokerUrls should not be null"); 10 11 if (invokerUrls.size() == 1 12 && invokerUrls.get(0) != null 13 && EMPTY_PROTOCOL.equals(invokerUrls.get(0).getProtocol())) { 14 15 ...... 16 17 } else { 18 // 刷新之前的 Invoker 19 Map<String, Invoker<T>> oldUrlInvokerMap = this.urlInvokerMap; // local reference 20 // 加载新的 Invoker Map 21 Map<String, Invoker<T>> newUrlInvokerMap = toInvokers(invokerUrls);// Translate url list to Invoker map 22 // 获取新的 Invokers 23 List<Invoker<T>> newInvokers = Collections.unmodifiableList(new ArrayList<>(newUrlInvokerMap.values())); 24 // 缓存新的 Invokers 25 routerChain.setInvokers(newInvokers); 26 this.invokers = multiGroup ? toMergeInvokerList(newInvokers) : newInvokers; 27 this.urlInvokerMap = newUrlInvokerMap; 28 29 try { 30 // 通过新旧 Invokers 对比,销毁无用的 Invokers 31 destroyUnusedInvokers(oldUrlInvokerMap, newUrlInvokerMap); // Close the unused Invoker 32 } catch (Exception e) { 33 logger.warn("destroyUnusedInvokers error. ", e); 34 } 35 } 36} 37
获取刷新前后的 Invokers,将新的 Invokers 重新缓存起来,通过对比,销毁无用的 Invoker。
上面将 URL 转换 Invoker 是在RegistryDirectory#toInvokers中进行。
1private Map<String, Invoker<T>> toInvokers(List<URL> urls) { 2 Map<String, Invoker<T>> newUrlInvokerMap = new HashMap<>(); 3 4 Set<String> keys = new HashSet<>(); 5 String queryProtocols = this.queryMap.get(PROTOCOL_KEY); 6 for (URL providerUrl : urls) { 7 8 // 过滤消费端不匹配的协议,及非法协议 9 ...... 10 11 // 合并服务提供端配置数据 12 URL url = mergeUrl(providerUrl); 13 // 过滤重复的服务提供端配置数据 14 String key = url.toFullString(); 15 if (keys.contains(key)) { 16 continue; 17 } 18 keys.add(key); 19 20 // 缓存键是不与使用者端参数合并的url,无论使用者如何合并参数,如果服务器url更改,则再次引用 21 Map<String, Invoker<T>> localUrlInvokerMap = this.urlInvokerMap; // local reference 22 Invoker<T> invoker = localUrlInvokerMap == null ? null : localUrlInvokerMap.get(key); 23 24 // 缓存无对应 invoker,再次调用 protocol#refer 是否有数据 25 if (invoker == null) { 26 try { 27 boolean enabled = true; 28 if (url.hasParameter(DISABLED_KEY)) { 29 enabled = !url.getParameter(DISABLED_KEY, false); 30 } else { 31 enabled = url.getParameter(ENABLED_KEY, true); 32 } 33 if (enabled) { 34 invoker = new InvokerDelegate<>(protocol.refer(serviceType, url), url, providerUrl); 35 } 36 } catch (Throwable t) { 37 logger.error("Failed to refer invoker for interface:" + serviceType + ",url:(" + url + ")" + t.getMessage(), t); 38 } 39 // 将新的 Invoker 缓存起来 40 if (invoker != null) { // Put new invoker in cache 41 newUrlInvokerMap.put(key, invoker); 42 } 43 } else { 44 // 缓存里有数据,则进行重新覆盖 45 newUrlInvokerMap.put(key, invoker); 46 } 47 } 48 keys.clear(); 49 return newUrlInvokerMap; 50} 51
总结
通过《Dubbo之服务暴露》和本文两篇文章对 Dubbo 服务暴露和服务消费原理的了解。我们可以看到,不管是暴露还是消费,Dubbo 都是以 Invoker 为数据交换主体进行,通过对 Invoker 发起调用,实现一个远程或本地的实现。
个人博客: https://ytao.top
关注公众号 【ytao】,更多原创好文
