1. 服务引用简介
在dubbo中,有两种方式引用远程服务。第一种是通过服务直连的方式引用服务;第二种是基于注册中心引用。服务直连的方式,一般用于开发环境的调试阶段,不建议在线上使用。服务引用的过程,总体可以分为3个步骤: 一,前置工作,加载Consumer端的配置元数据信息;二,将远程服务信息封装为Invoker对象,如果有多个服务注册中心,则会获取到一组Invoker实例,并通过集群管理类Cluster将多个Invoker合并为实例;三,创建代理类,用于调用服务。
2. 前置工作,加载Consumer配置元数据
服务的引用对应着ServiceBean这个类,ServiceBean继承了ServiceConfig,实现了InitializingBean、FactoryBean、ApplicationContextAware等接口。通过重写afterPropertiesSet方法,设置了配置元数据信息;通过重写setApplicationContext方法,设置spring上下文;通过重写getObject方法, 触发服务引用。
其实触发服务引用有两个地方。
第一个地方在afterPropertiesSet方法最后,会判断是否进行init操作,也就是判断<dubbo:reference/>
或者<dubbo:consumer/>
标签配置的init是否为true,默认为false。
<dubbo:reference id="demoService" check="false" interface="org.apache.dubbo.demo.DemoService" retries="0" init="true"/>// 饿汉式,如果配置了init的参数为true,则直接进行服务引用if (shouldInit()) {getObject();}
第二个地方在getObject方法,这由spring容器触发,在ReferenceBean对应的服务被注入到其它类中时,会触发引用。
这两种方式的区别在于前者为饿汉式,后者为懒汉式。dubbo默认使用懒汉式引用服务,如果要使用饿汉式,可以配置<dubbo:reference/>
或者<dubbo:consumer/>
标签的init值为true。
getObject方法触发服务引用会调用get方法, get方法在父类ReferenceConfig中实现。
public synchronized T get() {checkAndUpdateSubConfigs(); // ①if (destroyed) {throw new IllegalStateException("The invoker of ReferenceConfig(" + url + ") has already destroyed!");}if (ref == null) { // ②init();}return ref;}
①: 首先对配置信息进行处理和检查,来保证配置信息的准确性。
②: 判断ref是否为null,如果为null,表示服务还未进行引用,调用init方法进行初始化;反之表示已经完成了服务引用。ref是ReferenceConfig的成员变量,代表着接口引用的代理类。
2.1 init方法实现
init方法主要是将配置信息解析出来,使用map容器存放所有配置参数和值。在后续的服务引用过程中,会将map中的参数信息组装为URL,URL最终会贯穿整个服务引用过程。

图中红框部分就是配置信息解析的逻辑。黄色框中的是创建代理类的逻辑。在createProxy方法中会完成Invoker的转换过程,以及创建代理类。
3. 获取远程服务信息,转换为Invoker对象
createProxy方法不只是用于创建代理对象,还有创建Invoker对象。创建Invoker对象也分三种情况。首先如果配置了<dubbo reference/>
标签配置了injvm,表示进行本地服务引用,会通过InjvmProtocol的refer方法生成InjvmInvoker实例;如果获取到的提供者地址只有一个,则直接使用对应的Protocol实例生成Invoker实例即可;如果有多个注册中心地址,会将创建的Invoker对象放入集合中,然后通过Cluster合并为一个Invoker对象。
private T createProxy(Map<String, String> map) {// 是否打开本地引用if (shouldJvmRefer(map)) {URL url = new URL(LOCAL_PROTOCOL, LOCALHOST_VALUE, 0, interfaceClass.getName()).addParameters(map);// 调用InjvmProtocol的refer方法生成InjvmInvoker实例invoker = REF_PROTOCOL.refer(interfaceClass, url);if (logger.isInfoEnabled()) {logger.info("Using injvm service " + interfaceClass.getName());}} else {urls.clear(); // reference retry init will add url to urls, lead to OOM// 用户是否指定服务提供方地址:可以使服务提供方的ip地址,也就是服务直连的方式if (url != null && url.length() > 0) { // ①String[] us = SEMICOLON_SPLIT_PATTERN.split(url);if (us != null && us.length > 0) {for (String u : us) {URL url = URL.valueOf(u);if (StringUtils.isEmpty(url.getPath())) {url = url.setPath(interfaceName);}if (REGISTRY_PROTOCOL.equals(url.getProtocol())) {urls.add(url.addParameterAndEncoded(REFER_KEY, StringUtils.toQueryString(map)));} else { // ②urls.add(ClusterUtils.mergeUrl(url, map));}}}} else { // ③// if protocols not injvm checkRegistryif (!LOCAL_PROTOCOL.equalsIgnoreCase(getProtocol())){// 检查注册中心的配置是否存在,然后将注册信息转换为Registry对象checkRegistry(); // ④List<URL> us = loadRegistries(false); // ⑤if (CollectionUtils.isNotEmpty(us)) {for (URL u : us) { // ⑥URL monitorUrl = loadMonitor(u);if (monitorUrl != null) {map.put(MONITOR_KEY, URL.encode(monitorUrl.toFullString()));}urls.add(u.addParameterAndEncoded(REFER_KEY, StringUtils.toQueryString(map)));}}if (urls.isEmpty()) {throw new IllegalStateException("No such any registry to reference " + interfaceName + " on the consumer " + NetUtils.getLocalHost() + " use dubbo version " + Version.getVersion() + ", please config <dubbo:registry address=\"...\" /> to your spring config.");}}}// 只有一个服务提供者地址的时候if (urls.size() == 1) { // ⑦invoker = REF_PROTOCOL.refer(interfaceClass, urls.get(0));} else {// ⑧List<Invoker<?>> invokers = new ArrayList<Invoker<?>>();URL registryURL = null;for (URL url : urls) { // ⑨ 有多个提供者地址invokers.add(REF_PROTOCOL.refer(interfaceClass, url));if (REGISTRY_PROTOCOL.equals(url.getProtocol())) { // ⑩registryURL = url; // use last registry url}}if (registryURL != null) {// ⑪URL u = registryURL.addParameter(CLUSTER_KEY, RegistryAwareCluster.NAME);// ⑫invoker = CLUSTER.join(new StaticDirectory(u, invokers));} else { // ⑬invoker = CLUSTER.join(new StaticDirectory(invokers));}}}// ⑭if (shouldCheck() && !invoker.isAvailable()) {throw new IllegalStateException("Failed to check the status of the service " + interfaceName + ". No provider available for the service " + (group == null ? "" : group + "/") + interfaceName + (version == null ? "" : ":" + version) + " from the url " + invoker.getUrl() + " to the consumer " + NetUtils.getLocalHost() + " use dubbo version " + Version.getVersion());}if (logger.isInfoEnabled()) {logger.info("Refer dubbo service " + interfaceClass.getName() + " from url " + invoker.getUrl());}// ⑮return (T) PROXY_FACTORY.getProxy(invoker);}
①: 如果<dubbo: reference/>
标签配置了url的值,也就是直连提供者的方式,则对配置的值进行解析,将解析出来的地址放入到urls集合中。
②: 这里会将map中的参数配置合并到URL中。
③: 如果没有使用直连的方式,则从注册中心加载URL信息。
④: 检查有没有配置registry,如果本地没有配置,会构建出一个默认的配置。然后会将registry的配置信息转换为Registry实例。
⑤: 将Registry实例转换为URL对象,并会添加其它和服务引用相关的参数信息。
⑥: 循环每一个注册信息的url,在url中添加refer参数,refer的参数值是map中保存的consumer端配置元数据信息。
⑦: 如果只用一个url对象,则调用Protocol的refer方法会获取一个Invoker对象。
⑧: 如果有多个url对象,则定义一个invokers列表,存放多个Invoker实例。
⑨: 循环每一个url对象,调用Protocol的refer方法会获取一个Invoker对象。并放入invokers集合中。
⑩: 如果url的protocol协议是registry,则赋值给registryUrl用于后续的Cluster实例封装。
⑪: 如果registryUrl不为null,将添加cluster参数到url中,值为registryaware。
⑫: 将多个Invoker实例合并,并使用RegistryAwareClusterInvoker包装一层。
⑬: 如果url的protocol不是registry,则就是直接调用的方式,直接合并多个Invoker即可。
⑭: 检查<dubbo:reference/>
标签是否配置了check的值,如果配置check=false,则不会检查提供者的可用性,正常启动consumer端。否则会检查提供者的可用性,如果不可用,则会抛出异常。check的缺省配置为true。一般在测试环境对某些服务不关心时,会设置为false。生成环境不建议。
⑮: 调用ProxyFactory的getProxy方法,创建调用代理对象。
这个方法的细节步骤有点多,但是转换Invoker的逻辑,主要关注Protocol的refer的方法和CLUSTER.join方法的实现。创建代理对象在下一节分析。
3.1 Protocol的refer方法实现
REF_PROTOCOL是ReferenceConfig的成员变量,定义如下:
private static final Protocol REF_PROTOCOL = ExtensionLoader.getExtensionLoader(Protocol.class).getAdaptiveExtension();
即REF_PROTOCOL是Protocol的自适应类Protocol$Adaptive实例,refer方法会根据URL中的protocol的值选择具体的扩展实现类,此时URL中的protocol有两种情况,dubbo和registry。一般的如果是直连的方式,就直接会配置<dubbo: reference url="dubbo://host:ip">
,也可以指定选择注册中心地址。或者是配置<dubbo: registry address="{protocol}://host:ip">
。所以Protocol的refer方法有两种实现方式,基于DubboProtocol和RegistryProtocol的实现方式。
3.1.1 基于DubboProtocol的refer方法实现
如果是以dubbo直连的方式,配置提供者的地址,那么创建Invoker的对象的方法在DubboProtocol中实现。DubboProtocol被ProtocolFilterWrapper和ProtocolListenerWrapper包装了,所以会先调用这两个类的refer方法。
ProtocolFilterWrapper的refer方法
public <T> Invoker<T> refer(Class<T> type, URL url) throws RpcException {if (REGISTRY_PROTOCOL.equals(url.getProtocol())) {return protocol.refer(type, url);}return buildInvokerChain(protocol.refer(type, url), REFERENCE_FILTER_KEY, CommonConstants.CONSUMER);}
由于当前protocol是dubbo,所有在完成Invoker对象创建之后,会调用buildInvokerChain方法构建调用链。这个调用链的构建和export方法中的一致,只是此时的构建方式consumer端,也就是group为consumer。这就意味着在buildInvokerChain方法在获取Filter的自动激活扩展类时,只会获取被标注为Consumer端的扩展类。构建的调用链如下图所示:

ProtocolListenerWrapper的refer方法
该方法和export方法一样,都是注册一个监听器,监听器可用用户自己定义,通过spi机制进行加载。在完成服务的引用之后,会调用该监听器的相关方法,做一些操作。
分析完两个包装类的refer方法之后,再来看看DubboProtocol的refer方法实现。因为DubboProtocol继承了AbstractProtocol,AbstractProtocol的refer做了实现,内部逻辑比较简单,首先会调用ProtocolBindingRefer方法,该方法是一个抽象方法,具体要看子类DubboProtocol的实现,然后会将子类返回的结果封装成AsyncToSyncInvoker对象返回。
DubboProtocol的protocolBindingRefer方法
public <T> Invoker<T> protocolBindingRefer(Class<T> serviceType, URL url) throws RpcException {optimizeSerialization(url);// create rpc invoker.DubboInvoker<T> invoker = new DubboInvoker<T>(serviceType, url, getClients(url), invokers);invokers.add(invoker);return invoker;}
该方法我们需要关注的调用getClients方法,这个方法用于获取客户端实例,实例类型为ExchangeClient。ExchangeClient实际上并不具备通信能力,需要基于更底层的客户端实例进行通信。比如NettyClient、MinaClient等,默认情况下,dubbo使用NettyClient进行通信。下面看下getClients方法的实现。
DubboProtocol的getClients方法
private ExchangeClient[] getClients(URL url) {// 是否使用共享连接boolean useShareConnect = false;// 获取连接数,默认为0。表示未配置int connections = url.getParameter(CONNECTIONS_KEY, 0);List<ReferenceCountExchangeClient> shareClients = null;// 如果没有配置connections,则共享连接,否则一个连接供一个service使用if (connections == 0) {useShareConnect = true;String shareConnectionsStr = url.getParameter(SHARE_CONNECTIONS_KEY, (String) null);// 如果没有配置shareconnections的值,默认为connections的值默认为1connections = Integer.parseInt(StringUtils.isBlank(shareConnectionsStr) ? ConfigUtils.getProperty(SHARE_CONNECTIONS_KEY, DEFAULT_SHARE_CONNECTIONS) : shareConnectionsStr);// 获取共享客户端shareClients = getSharedClient(url, connections);}ExchangeClient[] clients = new ExchangeClient[connections];for (int i = 0; i < clients.length; i++) {if (useShareConnect) {clients[i] = shareClients.get(i);} else {// 初始化新的客户端clients[i] = initClient(url);}}return clients;}
该方法会根据connections的值来决定是获取共享客户端还是创建新的客户端实例,默认情况下,使用共享客户端实例。getShareClient方法中也会调用initClient方法。
DubboProtocol的getShareClient方法
private List<ReferenceCountExchangeClient> getSharedClient(URL url, int connectNum) {String key = url.getAddress();// 从换从中获取List<ReferenceCountExchangeClient> clients = referenceClientMap.get(key); // ①// 检查客户端是否可用,只要有一个不可用,就需要将不可用的客户端替换成可以用的客户端if (checkClientCanUse(clients)) { // ②// 因为创建的共享客户端是带有引用计数功能的ReferenceCountExchangeClient,所以这里会对所有的客户端引用计数+1batchClientRefIncr(clients);return clients;}locks.putIfAbsent(key, new Object());synchronized (locks.get(key)) { // ③clients = referenceClientMap.get(key);// dubbo checkif (checkClientCanUse(clients)) {batchClientRefIncr(clients);return clients;}// 连接数必须大于等于1connectNum = Math.max(connectNum, 1);// 如果clients集合为空,则进行初始化,并放入到缓存if (CollectionUtils.isEmpty(clients)) { // ④clients = buildReferenceCountExchangeClientList(url, connectNum);referenceClientMap.put(key, clients);} else {for (int i = 0; i < clients.size(); i++) { // ⑤ReferenceCountExchangeClient referenceCountExchangeClient = clients.get(i);// 如果客户端集合中某个客户端不可用,创建一个新的客户端来替代// If there is a client in the list that is no longer available, create a new one to replace him.if (referenceCountExchangeClient == null || referenceCountExchangeClient.isClosed()) {clients.set(i, buildReferenceCountExchangeClient(url));continue;}referenceCountExchangeClient.incrementAndGetCount();}}locks.remove(key);return clients;}}
①: 根据地址从referenceClientMap缓存中获取客户端列表。
②: 判断客户端是否可用,如果集合中存在一个不可用,则视为不可用,后续会将不可用的客户端替换掉。如果都可用,则对所有的客户端引用计数+1,并返回客户端列表。
③: 通过加锁的方式,确保一个地址只会创建一次客户端。
④: 如果clients结合为空,则调用buildReferenceCountExchangeClientList方法创建connectNum格式的客户端实例。
⑤: 循环clients集合中的每个客户端,如果当前客户端不可用,则调用buildReferenceCountExchangeClient方法创建新的客户端来替代。最后引用计数+1。
buildReferenceCountExchangeClient方法中会调用initClient方法,返回一个ExchangeClient实例,最后通过ReferenceCountExchangeClient对结果进行封装返回。
DubboProtocol的initClient方法
private ExchangeClient initClient(URL url) {// 获取配置的客户端类型配置,默认为nettyString str = url.getParameter(CLIENT_KEY, url.getParameter(SERVER_KEY, DEFAULT_REMOTING_CLIENT));url = url.addParameter(CODEC_KEY, DubboCodec.NAME);// 默认开启心跳机制url = url.addParameterIfAbsent(HEARTBEAT_KEY, String.valueOf(DEFAULT_HEARTBEAT));// BIO is not allowed since it has severe performance issue.if (str != null && str.length() > 0 && !ExtensionLoader.getExtensionLoader(Transporter.class).hasExtension(str)) {throw new RpcException("Unsupported client type: " + str + "," +" supported client type is " + StringUtils.join(ExtensionLoader.getExtensionLoader(Transporter.class).getSupportedExtensions(), " "));}ExchangeClient client;try {// 是否进行懒加载,也就是当有请求发生时,再去创建客户端if (url.getParameter(LAZY_CONNECT_KEY, false)) {client = new LazyConnectExchangeClient(url, requestHandler);} else {// 创建客户端实例client = Exchangers.connect(url, requestHandler);}} catch (RemotingException e) {throw new RpcException("Fail to create remoting client for service(" + url + "): " + e.getMessage(), e);}return client;}
该方法首先会根据配置的client值来确定客户端的类型,默认是netty。还有判断是否设置了lazy=true,默认为false。如果lazy为true,则会返回LazyConnectExchangeClient客户端,区别在于当有request请求的时候,才会调用Exchangers.connect方法创建客户端实例。
Exhangers.connect方法和bind方法一样,首先获取Exchanger实例,这里还是HeaderExchanger。然后在调用HeaderExchanger的connect方法回去ExchangeClient客户端,同样的ExchangeClient的connect方法和bind一样,会构建一个ChannelHandler链。在ExchangeClient的connect方法内部,会调用Transporters的connect方法来创建一个NettyClient,最终会调用NettyTransporter的connect方法,返回一个NettyClient实例。NettyClient的构造函数,同样调用父类的构造函数完成实例化,在此之前,会调用wrapChannelHandler方法,对传过来的ChannelHandler进行一次封装,封装完之后的ChannelHandler链如下图所示:

完成ChannelHandler的包装之后,就会调用父类AbstractClient的构造函数,在父类构造函数中,会调用doOpen方法,开启客户端。doOpen方法是抽象方法,有子类NettyClient实现。
protected void doOpen() throws Throwable {// 实例化处理客户端业务的Handlerfinal NettyClientHandler nettyClientHandler = new NettyClientHandler(getUrl(), this);// 实例化bootstrapbootstrap = new Bootstrap();bootstrap.group(nioEventLoopGroup).option(ChannelOption.SO_KEEPALIVE, true).option(ChannelOption.TCP_NODELAY, true).option(ChannelOption.ALLOCATOR, PooledByteBufAllocator.DEFAULT).channel(NioSocketChannel.class);// 设置超时时间,默认3sif (getConnectTimeout() < 3000) {bootstrap.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, 3000);} else {bootstrap.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, getConnectTimeout());}bootstrap.handler(new ChannelInitializer() {@Overrideprotected void initChannel(Channel ch) throws Exception {int heartbeatInterval = UrlUtils.getHeartbeat(getUrl());NettyCodecAdapter adapter = new NettyCodecAdapter(getCodec(), getUrl(), NettyClient.this);ch.pipeline()//.addLast("logging",new LoggingHandler(LogLevel.INFO))//for debug.addLast("decoder", adapter.getDecoder()) // 添加解码器Handler到管道.addLast("encoder", adapter.getEncoder()) // 添加编码码器Handler到管道.addLast("client-idle-handler", new IdleStateHandler(heartbeatInterval, 0, 0, MILLISECONDS)).addLast("handler", nettyClientHandler); // 添加处理业务的handler到管道String socksProxyHost = ConfigUtils.getProperty(SOCKS_PROXY_HOST);if(socksProxyHost != null) {int socksProxyPort = Integer.parseInt(ConfigUtils.getProperty(SOCKS_PROXY_PORT, DEFAULT_SOCKS_PROXY_PORT));Socks5ProxyHandler socks5ProxyHandler = new Socks5ProxyHandler(new InetSocketAddress(socksProxyHost, socksProxyPort));ch.pipeline().addFirst(socks5ProxyHandler);}}});}
子类完成doOpen的逻辑之后,会调用connect方法,完成和服务端连接。connect方法中做了一个加锁的操作,保证只会进行一次连接。真正进行连接的操作是doConnect方法,doConnect方法是抽象方法,由子类实现,这里看NettyClient方法的doConnect方法实现。在doConnect方法的内部逻辑中,通过ChannelFuture来实现异步的连接,通过future.awaitUninterruptibly方法来控制等待时间,超时则会关闭连接通道。
3.1.2 基于RegistryProtocol的refer方法实现
这里我们开始分析RegistryProtocol的refer方法实现。在进入RegistryProctocol的refer方式内部之前,也会经过ProtocolFilterWrapper和ProtocolListenerWrapper这两个包装类的refer方法。
public <T> Invoker<T> refer(Class<T> type, URL url) throws RpcException {url = URLBuilder.from(url).setProtocol(url.getParameter(REGISTRY_KEY, DEFAULT_REGISTRY)).removeParameter(REGISTRY_KEY).build(); // ①Registry registry = registryFactory.getRegistry(url); // ②if (RegistryService.class.equals(type)) { // ③return proxyFactory.getInvoker((T) registry, type, url);}// group="a,b" or group="*"Map<String, String> qs = StringUtils.parseQueryString(url.getParameterAndDecoded(REFER_KEY));String group = qs.get(GROUP_KEY);if (group != null && group.length() > 0) {if ((COMMA_SPLIT_PATTERN.split(group)).length > 1 || "*".equals(group)) {return doRefer(getMergeableCluster(), registry, type, url);}}// ④return doRefer(cluster, registry, type, url);}
①: 从url中获取registry的值,作为url的protocol值,并去掉registry参数。这是url的protocol就是配置的<dubbo:registry/>
标签配置的协议名称,一般是zookeeper。
②: registryFactory是一个扩展点,getRegistry方法会根据url中的protocol值选择具体的扩展实现类。因为第①中protocol的值已经更新为zookeeper,所以获取的Registry实例为ZookeeperRegistry。
③: 如果被引用的服务是RegistryService,则调用ProxyFactory的getInvoker方法获取Invoker实例。
④: 调用doRefer方法完成Invoker实例的创建。
RegistryProtocol的doRefer方法
private <T> Invoker<T> doRefer(Cluster cluster, Registry registry, Class<T> type, URL url) {// 实例化RegistryDirectory并设置注册中心实例和协议RegistryDirectory<T> directory = new RegistryDirectory<T>(type, url); // ①directory.setRegistry(registry);directory.setProtocol(protocol);// 从directory中获取所有和refer相关的参数Map<String, String> parameters = new HashMap<String, String>(directory.getUrl().getParameters()); // ②// 构建服务消费者的URL对象URL subscribeUrl = new URL(CONSUMER_PROTOCOL, parameters.remove(REGISTER_IP_KEY), 0, type.getName(), parameters); // ③// 将服务消费者信息注册到注册中心if (!ANY_VALUE.equals(url.getServiceInterface()) && url.getParameter(REGISTER_KEY, true)) { // ④directory.setRegisteredConsumerUrl(getRegisteredConsumerUrl(subscribeUrl, url));registry.register(directory.getRegisteredConsumerUrl());}directory.buildRouterChain(subscribeUrl);// 订阅providers、configurators、routers等节点数据directory.subscribe(subscribeUrl.addParameter(CATEGORY_KEY, PROVIDERS_CATEGORY + "," + CONFIGURATORS_CATEGORY + "," + ROUTERS_CATEGORY)); // ⑤// 一个注册中心会有多个提供者,这里需要将多个提供者合并成一个。Invoker invoker = cluster.join(directory); // ⑥ProviderConsumerRegTable.registerConsumer(invoker, url, subscribeUrl, directory);return invoker;}
①: 实例化RegistryDirectory,并设置注册中心和协议。
②: 从directorys实例中获取refer相关的参数。
③: 根据相关参数构建消费者的URL对象。
④: 根据配置,判断是否需要将消费者信息注册到注册中心。如果需要,调用Registry的register方法完成信息注册。
⑤: RegistryDirectory实例订阅providers、configurators、routes等节点的信息,如果这几个节点及子节点信息发生变化,RegistryDirectory实例会接收通知,并进行更新。
⑥: 因为provider一般不是单机部署的,所以一个注册中心会有多个提供者信息,所有这里会将多个provider信息合并为一个Invoker对象。
3.1.2.1 Registry的register方法
因为配置的registry协议是zookeeper,所以register的实现在ZookeeperRegistry中,但是ZookeeperRegistry继承自FailbackRegistry,并且FailbackRegistry对register方法进行了实现。FailbackRegistry的register方法没有做实质性的操作,而是通过抽象方法doRegister方法扔给了子类去实现。那么就看下ZookeeperRegistry的doRegister方法的实现。
public void doRegister(URL url) {try {// 通过调用zkClient.create方法注册服务节点信息到zookeeperzkClient.create(toUrlPath(url), url.getParameter(DYNAMIC_KEY, true));} catch (Throwable e) {throw new RpcException("Failed to register " + url + " to zookeeper " + getUrl() + ", cause: " + e.getMessage(), e);}}
该方法调用zkClient的create方法在Zookeeper服务端创建节点,这个和服务暴露过程的创建服务节点是一样的逻辑。
3.1.2.2 RegistryDirectory的subscribe方法
消费者端完成注册之后,会订阅注册中心的providers、configurators、routes等节点的信息。并通过注册监听器,对节点信息进行监听。
public void subscribe(URL url) {setConsumerUrl(url);CONSUMER_CONFIGURATION_LISTENER.addNotifyListener(this); // 添加consumer端配置的监听器serviceConfigurationListener = new ReferenceConfigurationListener(this, url); // 添加服务引用相关配置的监听器registry.subscribe(url, this); // 执行订阅操作}
该方法会添加两个监听器: consumer端配置的监听器—ConsumerConfigurationListener;服务引用相关配置的监听器—ReferenceConfigurationListener。这两个监听器都继承自AbstractConfiguratorListener,重写了notifyOverrides方法,当notifyOverrides方法被触发时,会调用RegistryDirectory的refreshInvoker方法完成配置更新。最后调用Registry的subscribe方法订阅相关节点信息。
Registry的subscribe方法
同样的,因为ZookeeperRegistry继承自FailbackRegistry,subscribe方法在FailbackRegistry中进行了实现,但是执行实际的订阅操作扔给了抽象方法doSubscribe,具体逻辑看子类的实现。
public void doSubscribe(final URL url, final NotifyListener listener) {try {// 如果配置的是*,那就会为root节点以及子节点添加监听器if (ANY_VALUE.equals(url.getServiceInterface())) {String root = toRootPath();ConcurrentMap<NotifyListener, ChildListener> listeners = zkListeners.get(url);if (listeners == null) {zkListeners.putIfAbsent(url, new ConcurrentHashMap<>());listeners = zkListeners.get(url);}ChildListener zkListener = listeners.get(listener);if (zkListener == null) {listeners.putIfAbsent(listener, (parentPath, currentChilds) -> {for (String child : currentChilds) {child = URL.decode(child);if (!anyServices.contains(child)) {anyServices.add(child);subscribe(url.setPath(child).addParameters(INTERFACE_KEY, child, Constants.CHECK_KEY, String.valueOf(false)), listener);}}});zkListener = listeners.get(listener);}zkClient.create(root, false);List<String> services = zkClient.addChildListener(root, zkListener);if (CollectionUtils.isNotEmpty(services)) {for (String service : services) {service = URL.decode(service);anyServices.add(service);subscribe(url.setPath(service).addParameters(INTERFACE_KEY, service,Constants.CHECK_KEY, String.valueOf(false)), listener);}}} else {// 如果配置的不是*,那么就会获取providers、configurators、routes这几个节点的url,然后为这几个节点及子节点添加监听器List<URL> urls = new ArrayList<>();for (String path : toCategoriesPath(url)) {ConcurrentMap<NotifyListener, ChildListener> listeners = zkListeners.get(url);if (listeners == null) {zkListeners.putIfAbsent(url, new ConcurrentHashMap<>());listeners = zkListeners.get(url);}ChildListener zkListener = listeners.get(listener);if (zkListener == null) {listeners.putIfAbsent(listener, (parentPath, currentChilds) -> ZookeeperRegistry.this.notify(url, listener, toUrlsWithEmpty(url, parentPath, currentChilds)));zkListener = listeners.get(listener);}zkClient.create(path, false);List<String> children = zkClient.addChildListener(path, zkListener);if (children != null) {urls.addAll(toUrlsWithEmpty(url, path, children));}}notify(url, listener, urls);}} catch (Throwable e) {throw new RpcException("Failed to subscribe " + url + " to zookeeper " + getUrl() + ", cause: " + e.getMessage(), e);}}
该方法的实现,分为两个分支,通过服务引用配置的service地址,来判断需要给哪些节点创建监听器。如果配置的是*
,那么会将给root节点以及其所有的子节点都创建监听器。如果配置的不是*
,则只会给providers、configurators、routes这几个节点以及子节点添加监听器。
3.1.2.3 Cluster的join方法
Cluster的join方法,会将多个Provider合并为一个Invoker实例。Cluster是一个spi扩展点,默认扩展实现类是FailoverCluster。但是Cluster还有个扩展包装类——MockClusterWrapper,它会对扩展实现实现类进行包装,也就说在调用扩展实现类的join方法之前,会先调用包装类的join方法。
// 包装类,对Cluster进行包装public class MockClusterWrapper implements Cluster {private Cluster cluster;public MockClusterWrapper(Cluster cluster) {this.cluster = cluster;}@Overridepublic <T> Invoker<T> join(Directory<T> directory) throws RpcException {// 在方法内部,先调用被包装的cluster实例的join方法获取Invoker实例,然后再通过MockClusterInvoker进行封装return new MockClusterInvoker<T>(directory, this.cluster.join(directory));}}// 默认扩展实现类 FailoverClusterpublic class FailoverCluster implements Cluster {public final static String NAME = "failover";@Overridepublic <T> Invoker<T> join(Directory<T> directory) throws RpcException {// 返回一个FailoverClusterInvoker对象return new FailoverClusterInvoker<T>(directory);}}
在FailoverClusterInvoker的构造函数中,会调用父类的构造函数,并设置其directory属性为传入的Directory实例。在进行集群调用实现的时候,会从directory中获取一个提供者进行调用。最后返回的Invoker实例为MockClusterInvoker,在进行远程接口调用过程中,MockClusterInvoker会对实际的Invoker实例功能进行增加,增加了mock调用的功能。会根据服务引用的配置,来判断是否执行mock调用,以及mock的策略是哪些。
3.2 CLUSTER.join 合并多个Invoker实例
针对每个地址,都会通过Protocol的refer方法将引用的远程服务转换为Invoker对象,如果是默认的集群实现,那么Invoker实例为FailoverClusterInvoker。多个注册地址会生成多个FailoverClusterInvoker实例,放入到invokers集合中,最后会调用CLUSTER.join方法将多个FailoverClusterInvoker实例进行一次包装。不同的服务引用方式,封装方式也不用。
如果是基于dubbo直连的方式,则会使用StaticDirectory对象将invokers集合进行封装,最后返回的还是FailoverClusterInvoker实例,也就是将多个FailoverClusterInvoker最后归并为一个FailoverClusterInvoker实例。
如果是基于注册中心的引用,那么会在服务引用的URL中添加cluster参数,默认值为registryaware。也就是说使用StaticDirectory对象将invokers集合进行封装,join方法会根据cluster的值选择对应的实现, 这里也就是RegistryAwareCluster,最终返回RegistryAwareClusterInvoker实例。
4. 创建代理类
Invoker创建完毕之后,接下来要做的事情就是要创建服务代理对象。有了代理对象,就可以进行远程调用。PROXY_FACTORY是ProxyFactory的自适应扩展来ProxyFactory$Adaptive,getProxy方法会根据URL中的proxy值来确定具体的扩展类。ProxyFactory的子类都继承了AbstractProxyFactory,AbstractProxyFactory实现了getProxy方法,在AbstractProxyFactory的getProxy方法实现中,主要是解析interfaces参数,然后把目标服务以及EchoService.class放入到一个Class数组中,如果目标服务类支持泛化调用,也会将GenericService.class放入到数组中。
public <T> T getProxy(Invoker<T> invoker, boolean generic) throws RpcException {Class<?>[] interfaces = null;// 获取url中的interfaces参数值String config = invoker.getUrl().getParameter(INTERFACES);if (config != null && config.length() > 0) {// 解析配置,将服务接口类和EchoService.class方法如到interfaces数组String[] types = COMMA_SPLIT_PATTERN.split(config);if (types != null && types.length > 0) {interfaces = new Class<?>[types.length + 2];interfaces[0] = invoker.getInterface();interfaces[1] = EchoService.class;for (int i = 0; i < types.length; i++) {interfaces[i + 2] = ReflectUtils.forName(types[i]);}}}if (interfaces == null) {interfaces = new Class<?>[]{invoker.getInterface(), EchoService.class};}// 泛化调用相关配置if (!GenericService.class.isAssignableFrom(invoker.getInterface()) && generic) {int len = interfaces.length;Class<?>[] temp = interfaces;interfaces = new Class<?>[len + 1];System.arraycopy(temp, 0, interfaces, 0, len);interfaces[len] = com.alibaba.dubbo.rpc.service.GenericService.class;}// 获取代理类,由子类实现return getProxy(invoker, interfaces);}
方法最后会调用getProxy方法生成代理类,这个方法是抽象方法,由子类实现。ProxyFactory$Adaptive类会根据proxy的值选择具体的实现,默认是JavassistProxyFactory。
JavassistProxyFactory的getProxy方法
该方法首先会通过然后通过Proxy.getProxy方法回去interfaces的子类,然后调用newInstance方法生成代理类的实例之前,会创建InvokerInvocationHandler对象,InvokerInvocationHandler实现了jdk的InvocationHandler,是用于拦截目标接口类的调用。这里使用的Proxy了是dubbo自己定义的,而非jdk中的Proxy类,从而避免使用反射,减少性能消耗。
public <T> T getProxy(Invoker<T> invoker, Class<?>[] interfaces) {return (T) Proxy.getProxy(interfaces).newInstance(new InvokerInvocationHandler(invoker));}
JdkProxyFactory的getProxy方法
该方法是使用java.lang.reflect包下的Proxy类,来生成代理类。同样的也会创建InvokerInvocationHandler类,用于拦截目标接口类的调用。
public <T> T getProxy(Invoker<T> invoker, Class<?>[] interfaces) {return (T) Proxy.newProxyInstance(Thread.currentThread().getContextClassLoader(), interfaces, new InvokerInvocationHandler(invoker));}




