dubbo入门第二课-服务暴露

入门第一课,我们搭建了一个简单的dubbo应用,这节课主要来介绍下提供者是如何服务暴露的。

第一步:获取所有的注册中心URL

其实在dubbo服务整个阶段,不管是服务提供者还是服务消费者,传递的都是URL。

 private void doExportUrls() {
        List<URL> registryURLs = loadRegistries(true);                  -------------------------------------(1)
        for (ProtocolConfig protocolConfig : protocols) {                  -------------------------------------(2)
            doExportUrlsFor1Protocol(protocolConfig, registryURLs);
        }
    }

(1)获取注册中心URL,主要是针对配置文件中:
<dubbo:registry protocol="zookeeper" address="172.59.72.126:2181,172.59.72.184:2181,172.59.72.135:2181" check="true" id="com.alibaba.dubbo.config.RegistryConfig" />
经过loadRegistries(true)之后返回的:
[registry://172.59.72.126:2181/com.alibaba.dubbo.registry.RegistryService?application=dubbo-provider&backup=172.59.72.184:2181,172.59.72.135:2181&check=true&dubbo=2.8.4&owner=dubbo&pid=12956&registry=zookeeper&timestamp=1548730068596]
用统一数据模型URL表示:

  • registry表示这是一个注册中心的URL
  • 注册中心地址 172.59.72.126:2181
  • 调用注册中心的服务 com.alibaba.dubbo.registry.RegistryService
  • 其他key-value项。
    (2)多协议暴露服务。for循环处理多协议暴露。

第二步:获取服务暴露的URL

接下来主要看 doExportUrlsFor1Protocol(protocolConfig, registryURLs);

private void doExportUrlsFor1Protocol(ProtocolConfig protocolConfig, List<URL> registryURLs) {
        String name = protocolConfig.getName();
        if (name == null || name.length() == 0) {
            name = "dubbo";                     -------------------------------------------------------------------(1)
        }

        String host = protocolConfig.getHost();
        if (provider != null && (host == null || host.length() == 0)) {
            host = provider.getHost();
        }
        boolean anyhost = false;
        if (NetUtils.isInvalidLocalHost(host)) {
            anyhost = true;
            try {
                host = InetAddress.getLocalHost().getHostAddress();
            } catch (UnknownHostException e) {
                logger.warn(e.getMessage(), e);
            }
            if (NetUtils.isInvalidLocalHost(host)) {
                if (registryURLs != null && registryURLs.size() > 0) {
                    for (URL registryURL : registryURLs) {
                        try {
                            Socket socket = new Socket();
                            try {
                                SocketAddress addr = new InetSocketAddress(registryURL.getHost(), registryURL.getPort());
                                socket.connect(addr, 1000);
                                host = socket.getLocalAddress().getHostAddress();
                                break;
                            } finally {
                                try {
                                    socket.close();
                                } catch (Throwable e) {}
                            }
                        } catch (Exception e) {
                            logger.warn(e.getMessage(), e);
                        }
                    }
                }
                if (NetUtils.isInvalidLocalHost(host)) {
                    host = NetUtils.getLocalHost();
                }
            }
        }

        Integer port = protocolConfig.getPort();
        if (provider != null && (port == null || port == 0)) {
            port = provider.getPort();
        }
        final int defaultPort = ExtensionLoader.getExtensionLoader(Protocol.class).getExtension(name).getDefaultPort();
        if (port == null || port == 0) {
            port = defaultPort;
        }
        if (port == null || port <= 0) {
            port = getRandomPort(name);
            if (port == null || port < 0) {
                port = NetUtils.getAvailablePort(defaultPort);
                putRandomPort(name, port);
            }
            logger.warn("Use random available port(" + port + ") for protocol " + name);
        }

        Map<String, String> map = new HashMap<String, String>();
        if (anyhost) {
            map.put(Constants.ANYHOST_KEY, "true");
        }
        map.put(Constants.SIDE_KEY, Constants.PROVIDER_SIDE);
        map.put(Constants.DUBBO_VERSION_KEY, Version.getVersion());
        map.put(Constants.TIMESTAMP_KEY, String.valueOf(System.currentTimeMillis()));
        if (ConfigUtils.getPid() > 0) {
            map.put(Constants.PID_KEY, String.valueOf(ConfigUtils.getPid()));
        }
        appendParameters(map, application);
        appendParameters(map, module);
        appendParameters(map, provider, Constants.DEFAULT_KEY);
        appendParameters(map, protocolConfig);
        appendParameters(map, this);
        if (methods != null && methods.size() > 0) {
            for (MethodConfig method : methods) {
                appendParameters(map, method, method.getName());
                String retryKey = method.getName() + ".retry";
                if (map.containsKey(retryKey)) {
                    String retryValue = map.remove(retryKey);
                    if ("false".equals(retryValue)) {
                        map.put(method.getName() + ".retries", "0");
                    }
                }
                List<ArgumentConfig> arguments = method.getArguments();
                if (arguments != null && arguments.size() > 0) {
                    for (ArgumentConfig argument : arguments) {
                        //类型自动转换.
                        if(argument.getType() != null && argument.getType().length() >0){
                            Method[] methods = interfaceClass.getMethods();
                            //遍历所有方法
                            if(methods != null && methods.length > 0){
                                for (int i = 0; i < methods.length; i++) {
                                    String methodName = methods[i].getName();
                                    //匹配方法名称,获取方法签名.
                                    if(methodName.equals(method.getName())){
                                        Class<?>[] argtypes = methods[i].getParameterTypes();
                                        //一个方法中单个callback
                                        if (argument.getIndex() != -1 ){
                                            if (argtypes[argument.getIndex()].getName().equals(argument.getType())){
                                                appendParameters(map, argument, method.getName() + "." + argument.getIndex());
                                            }else {
                                                throw new IllegalArgumentException("argument config error : the index attribute and type attirbute not match :index :"+argument.getIndex() + ", type:" + argument.getType());
                                            }
                                        } else {
                                            //一个方法中多个callback
                                            for (int j = 0 ;j<argtypes.length ;j++) {
                                                Class<?> argclazz = argtypes[j];
                                                if (argclazz.getName().equals(argument.getType())){
                                                    appendParameters(map, argument, method.getName() + "." + j);
                                                    if (argument.getIndex() != -1 && argument.getIndex() != j){
                                                        throw new IllegalArgumentException("argument config error : the index attribute and type attirbute not match :index :"+argument.getIndex() + ", type:" + argument.getType());
                                                    }
                                                }
                                            }
                                        }
                                    }
                                }
                            }
                        }else if(argument.getIndex() != -1){
                            appendParameters(map, argument, method.getName() + "." + argument.getIndex());
                        }else {
                            throw new IllegalArgumentException("argument config must set index or type attribute.eg: <dubbo:argument index='0' .../> or <dubbo:argument type=xxx .../>");
                        }

                    }
                }
            } // end of methods for
        }

        if (ProtocolUtils.isGeneric(generic)) {
            map.put("generic", generic);
            map.put("methods", Constants.ANY_VALUE);
        } else {
            String revision = Version.getVersion(interfaceClass, version);
            if (revision != null && revision.length() > 0) {
                map.put("revision", revision);
            }

            String[] methods = Wrapper.getWrapper(interfaceClass).getMethodNames();
            if(methods.length == 0) {
                logger.warn("NO method found in service interface " + interfaceClass.getName());
                map.put("methods", Constants.ANY_VALUE);
            }
            else {
                map.put("methods", StringUtils.join(new HashSet<String>(Arrays.asList(methods)), ","));
            }
        }
        if (! ConfigUtils.isEmpty(token)) {
            if (ConfigUtils.isDefault(token)) {
                map.put("token", UUID.randomUUID().toString());
            } else {
                map.put("token", token);
            }
        }
        if ("injvm".equals(protocolConfig.getName())) {
            protocolConfig.setRegister(false);
            map.put("notify", "false");------------------------------------------------------(2)
        }
        // 导出服务
        String contextPath = protocolConfig.getContextpath();
        if ((contextPath == null || contextPath.length() == 0) && provider != null) {
            contextPath = provider.getContextpath();
        }
        URL url = new URL(name, host, port, (contextPath == null || contextPath.length() == 0 ? "" : contextPath + "/") + path, map);           ---------------------------------------------------------(3)

        if (ExtensionLoader.getExtensionLoader(ConfiguratorFactory.class)
                .hasExtension(url.getProtocol())) {
            url = ExtensionLoader.getExtensionLoader(ConfiguratorFactory.class)
                    .getExtension(url.getProtocol()).getConfigurator(url).configure(url);
        }

        String scope = url.getParameter(Constants.SCOPE_KEY);
        //配置为none不暴露
        if (! Constants.SCOPE_NONE.toString().equalsIgnoreCase(scope)) {

            //配置不是remote的情况下做本地暴露 (配置为remote,则表示只暴露远程服务)
            if (!Constants.SCOPE_REMOTE.toString().equalsIgnoreCase(scope)) {  ------------------(4)
                exportLocal(url);                                                  
            }
            //如果配置不是local则暴露为远程服务.(配置为local,则表示只暴露远程服务)
            if (! Constants.SCOPE_LOCAL.toString().equalsIgnoreCase(scope) ){  ---------------------------(5)                                              -                                                                                                                                
                if (logger.isInfoEnabled()) {
                    logger.info("Export dubbo service " + interfaceClass.getName() + " to url " + url);
                }
                if (registryURLs != null && registryURLs.size() > 0
                        && url.getParameter("register", true)) {
                    for (URL registryURL : registryURLs) {
                        url = url.addParameterIfAbsent("dynamic", registryURL.getParameter("dynamic"));
                        URL monitorUrl = loadMonitor(registryURL);
                        if (monitorUrl != null) {
                            url = url.addParameterAndEncoded(Constants.MONITOR_KEY, monitorUrl.toFullString());
                        }
                        if (logger.isInfoEnabled()) {
                            logger.info("Register dubbo service " + interfaceClass.getName() + " url " + url + " to registry " + registryURL);
                        }
                        Invoker<?> invoker = proxyFactory.getInvoker(ref, (Class) interfaceClass, registryURL.addParameterAndEncoded(Constants.EXPORT_KEY, url.toFullString()));

                        Exporter<?> exporter = protocol.export(invoker);
                        exporters.add(exporter);
                    }
                } else {
                    Invoker<?> invoker = proxyFactory.getInvoker(ref, (Class) interfaceClass, url);

                    Exporter<?> exporter = protocol.export(invoker);
                    exporters.add(exporter);
                }
            }
        }
        this.urls.add(url);
    }

(1)默认协议为dubbo.
(2)将所有的配置文件中对应的属性值全部转为map中作为key-value,如果是使用injvm协议暴露,则不需要开启网络,也不开启远程调用。这里不需要去注册中心注册,也不需要通知服务变更。服务提供者的配置如下:

<dubbo:protocol  name="injvm"  port="20880" threads="200" accesslog="true" /> 
<!--    <dubbo:protocol name="hessian" port="20881" threads="200" accesslog="true" /> -->
    <dubbo:protocol name="custom" port="20881" threads="200" accesslog="true" />
    <!-- 服务提供者暴露服务,当一个接口有多种实现时,可以用group区分 -->
    <dubbo:service interface="com.wang.api.service.LogService" ref="logService" timeout="3000" loadbalance="random" protocol="injvm" />

(3)map转为URL,URL中包括:protocol,username,password,host,port,path,Map<String, String> para

{owner=dubbo, side=provider, methods=add,delete, dubbo=2.8.4, threads=200, loadbalance=random, pid=9752, interface=com.wang.api.service.LogService, generic=false, timeout=3000, accesslog=true, revision=0.0.1-SNAPSHOT, delay=-1, application=dubbo-provider, default.delay=-1, anyhost=true, 
timestamp=1548732142496}meters,

(4)本地服务暴露,可以看出如果我们没有配置,则先本地服务暴露,再远程服务暴露。
(5)远程服务暴露。
主要看一下远程服务暴露

//如果配置不是local则暴露为远程服务.(配置为local,则表示只暴露远程服务)
            if (! Constants.SCOPE_LOCAL.toString().equalsIgnoreCase(scope) ){
                if (logger.isInfoEnabled()) {
                    logger.info("Export dubbo service " + interfaceClass.getName() + " to url " + url);
                }
                if (registryURLs != null && registryURLs.size() > 0
                        && url.getParameter("register", true)) {
                    for (URL registryURL : registryURLs) {
                        url = url.addParameterIfAbsent("dynamic", registryURL.getParameter("dynamic"));
                        URL monitorUrl = loadMonitor(registryURL);
                        if (monitorUrl != null) {
                            url = url.addParameterAndEncoded(Constants.MONITOR_KEY, monitorUrl.toFullString());
                        }
                        if (logger.isInfoEnabled()) {
                            logger.info("Register dubbo service " + interfaceClass.getName() + " url " + url + " to registry " + registryURL);
                        }
                        Invoker<?> invoker = proxyFactory.getInvoker(ref, (Class) interfaceClass, registryURL.addParameterAndEncoded(Constants.EXPORT_KEY, url.toFullString()));----------------(1)

                        Exporter<?> exporter = protocol.export(invoker);               -----------------------------(2)
                        exporters.add(exporter);                     -----------------------------------------(3)
                    }
                } else {
                    Invoker<?> invoker = proxyFactory.getInvoker(ref, (Class) interfaceClass, url);

                    Exporter<?> exporter = protocol.export(invoker);
                    exporters.add(exporter);
                }
            }

(1)获取Invoker过程。根据服务具体实现,实现接口,以及registryUrl通过ProxyFactory将LogServiceImpl封装成一个本地执行的Invoker,这就是所谓的RPC实现了像调用本地服务一样调用远程服务,其实invoker是对具体实现的一种代理。
(2)将Invoker转换成Exporter,然后进行服务暴露。
(3)将exporter缓存起来。

第三步:暴露远程服务

  1. 协议在接收请求时,应记录请求来源方地址信息:RpcContext.getContext().setRemoteAddress();
  2. export()必须是幂等的,也就是暴露同一个URL的Invoker两次,和暴露一次没有区别。
  3. export()传入的Invoker由框架实现并传入,协议不需要关心。

    我们来看看RegistryProtocol.class
 public <T> Exporter<T> export(final Invoker<T> originInvoker) throws RpcException {
        //export invoker
        final ExporterChangeableWrapper<T> exporter = doLocalExport(originInvoker);-------------(1)
        //registry provider
        final Registry registry = getRegistry(originInvoker);         -------------------------------------------(2)
        final URL registedProviderUrl = getRegistedProviderUrl(originInvoker);  --------------------------------(3)
        registry.register(registedProviderUrl);------------------------------(4)
        // 订阅override数据
        // FIXME 提供者订阅时,会影响同一JVM即暴露服务,又引用同一服务的的场景,因为subscribed以服务名为缓存的key,导致订阅信息覆盖。
        final URL overrideSubscribeUrl = getSubscribedOverrideUrl(registedProviderUrl);
        final OverrideListener overrideSubscribeListener = new OverrideListener(overrideSubscribeUrl);
        overrideListeners.put(overrideSubscribeUrl, overrideSubscribeListener);
        registry.subscribe(overrideSubscribeUrl, overrideSubscribeListener);-----------------------------(5)
        //保证每次export都返回一个新的exporter实例
        return new Exporter<T>() {--------------------------------------------(6)
            public Invoker<T> getInvoker() {
                return exporter.getInvoker();
            }
            public void unexport() {
                try {
                    exporter.unexport();
                } catch (Throwable t) {
                    logger.warn(t.getMessage(), t);
                }
                try {
                    registry.unregister(registedProviderUrl);
                } catch (Throwable t) {
                    logger.warn(t.getMessage(), t);
                }
                try {
                    overrideListeners.remove(overrideSubscribeUrl);
                    registry.unsubscribe(overrideSubscribeUrl, overrideSubscribeListener);
                } catch (Throwable t) {
                    logger.warn(t.getMessage(), t);
                }
            }
        };
    }

(1)具体的协议去暴露服务,我们使用的是dubbo,所以是用DubboProtocol.class暴露服务。
(2)根据invoker中的url获取Registry实例
(3)返回注册到注册中心的URL,对URL参数进行一次过滤。
(4)注册数据。比如:提供者地址,消费者地址,路由规则,覆盖规则,等数据。若有消费者订阅此服务,则推送消息让消费者引用此服务。这里注册的是提供者地址,如下所示:

dubbo://192.168.0.104:20880/com.wang.api.service.LogService?accesslog=true&anyhost=true&application=dubbo-provider&default.delay=-1&delay=-1&dubbo=2.8.4&generic=false&interface=com.wang.api.service.LogService&loadbalance=random&methods=add,delete&owner=dubbo&pid=13496&revision=0.0.1-SNAPSHOT&side=provider&threads=200&timeout=3000&timestamp=1548734246686

构建服务发布的URL的统一模型:dubbo表示服务发布的URL,192.168.0.104:20880服务发布的地址,com.wang.api.service.LogService发布的具体服务,然后就是其他key-value对。
(5)提供者向注册中心订阅所有注册服务的覆盖配置,当注册中心有此服务的覆盖配置注册进来时,推送消息给提供者,重新暴露服务,这由管理页面完成。
(6)保证每次export都返回一个新的exporter实例。
现在我们来分析如何去注册中心注册呢?

第四步:交给具体的协议进行服务暴露

 @SuppressWarnings("unchecked")
    private <T> ExporterChangeableWrapper<T>  doLocalExport(final Invoker<T> originInvoker){
        String key = getCacheKey(originInvoker);      ---------------------------------(1)
        ExporterChangeableWrapper<T> exporter = (ExporterChangeableWrapper<T>) bounds.get(key);
        if (exporter == null) {
            synchronized (bounds) {
                exporter = (ExporterChangeableWrapper<T>) bounds.get(key); ---------------------(2)
                if (exporter == null) {
                    final Invoker<?> invokerDelegete = new InvokerDelegete<T>(originInvoker, getProviderUrl(originInvoker));
                    exporter = new ExporterChangeableWrapper<T>((Exporter<T>)protocol.export(invokerDelegete), originInvoker);
                    bounds.put(key, exporter);
                }
            }
        }
        return (ExporterChangeableWrapper<T>) exporter;
    }

(1)从原始的invoker中获取key:dubbo://192.168.0.100:20880/com.wang.api.service.LogService?accesslog=true&anyhost=true&application=dubbo-provider&default.delay=-1&delay=-1&dubbo=2.8.4&generic=false&interface=com.wang.api.service.LogService&loadbalance=random&methods=add,delete&owner=dubbo&pid=13944&revision=0.0.1-SNAPSHOT&side=provider&threads=200&timeout=3000&timestamp=1548921263732
(2)从缓存中获取exporter,如果缓存中没有,说明服务么有暴露过,则暴露服务。用于解决rmi重复暴露端口冲突的问题,已经暴露过的服务不再重新暴露。
正式暴露服务:

 public <T> Exporter<T> export(Invoker<T> invoker) throws RpcException {
        URL url = invoker.getUrl();    ------------------------------(1)
        
        // export service.
        String key = serviceKey(url);     -------------------------(2)
        DubboExporter<T> exporter = new DubboExporter<T>(invoker, key, exporterMap);
        exporterMap.put(key, exporter);           ------------------------------------(3)
        
        //export an stub service for dispaching event
        Boolean isStubSupportEvent = url.getParameter(Constants.STUB_EVENT_KEY,Constants.DEFAULT_STUB_EVENT);
        Boolean isCallbackservice = url.getParameter(Constants.IS_CALLBACK_SERVICE, false);
        if (isStubSupportEvent && !isCallbackservice){
            String stubServiceMethods = url.getParameter(Constants.STUB_EVENT_METHODS_KEY);
            if (stubServiceMethods == null || stubServiceMethods.length() == 0 ){
                if (logger.isWarnEnabled()){
                    logger.warn(new IllegalStateException("consumer [" +url.getParameter(Constants.INTERFACE_KEY) +
                            "], has set stubproxy support event ,but no stub methods founded."));
                }
            } else {
                stubServiceMethodsMap.put(url.getServiceKey(), stubServiceMethods);
            }
        }

        openServer(url);                ------------------------------------(4)

        // modified by lishen
        optimizeSerialization(url);

        return exporter;
    }

(1)获取URL:dubbo://192.168.0.100:20880/com.wang.api.service.LogService?accesslog=true&anyhost=true&application=dubbo-provider&default.delay=-1&delay=-1&dubbo=2.8.4&generic=false&interface=com.wang.api.service.LogService&loadbalance=random&methods=add,delete&owner=dubbo&pid=13944&revision=0.0.1-SNAPSHOT&side=provider&threads=200&timeout=3000&timestamp=1548921263732
(2)获取key:com.wang.api.service.LogService:20880
(3)将invoker转化为DubboExporter,并将(key,exporter)缓存在map中。
(4)开启服务: 根据URL绑定IP与端口,建立NIO框架的Server。
如何开启服务?

private void openServer(URL url) {
        // find server.
        String key = url.getAddress();
        //client 也可以暴露一个只有server可以调用的服务。
        boolean isServer = url.getParameter(Constants.IS_SERVER_KEY,true);
        if (isServer) {
            ExchangeServer server = serverMap.get(key);
            if (server == null) {
                serverMap.put(key, createServer(url));
            } else {
                //server支持reset,配合override功能使用
                server.reset(url);
            }
        }
    }

(1)获取key:ip+端口号(192.168.0.100:20880)
(2)同一JVM中,同协议的服务,共享同一个Server,第一个暴露服务的时候创建server,以后相同协议的服务都使用同一个server。
(3)创建服务。

    private ExchangeServer createServer(URL url) {
        //默认开启server关闭时发送readonly事件
        url = url.addParameterIfAbsent(Constants.CHANNEL_READONLYEVENT_SENT_KEY, Boolean.TRUE.toString());
        //默认开启heartbeat
        url = url.addParameterIfAbsent(Constants.HEARTBEAT_KEY, String.valueOf(Constants.DEFAULT_HEARTBEAT));
        String str = url.getParameter(Constants.SERVER_KEY, Constants.DEFAULT_REMOTING_SERVER);

        if (str != null && str.length() > 0 && ! ExtensionLoader.getExtensionLoader(Transporter.class).hasExtension(str))
            throw new RpcException("Unsupported server type: " + str + ", url: " + url);

        url = url.addParameter(Constants.CODEC_KEY, Version.isCompatibleVersion() ? COMPATIBLE_CODEC_NAME : DubboCodec.NAME);
        ExchangeServer server;
        try {
            server = Exchangers.bind(url, requestHandler);
        } catch (RemotingException e) {
            throw new RpcException("Fail to start server(url: " + url + ") " + e.getMessage(), e);
        }
        str = url.getParameter(Constants.CLIENT_KEY);
        if (str != null && str.length() > 0) {
            Set<String> supportedTypes = ExtensionLoader.getExtensionLoader(Transporter.class).getSupportedExtensions();
            if (!supportedTypes.contains(str)) {
                throw new RpcException("Unsupported client type: " + str);
            }
        }
        return server;
    }

(1)获取url:
dubbo://192.168.0.100:20880/com.wang.api.service.LogService?accesslog=true&anyhost=true&application=dubbo-provider&channel.readonly.sent=true&default.delay=-1&delay=-1&dubbo=2.8.4&generic=false&interface=com.wang.api.service.LogService&loadbalance=random&methods=add,delete&owner=dubbo&pid=15748&revision=0.0.1-SNAPSHOT&side=provider&threads=200&timeout=3000&timestamp=1548921845226
(2)dubbo协议的服务器端实现类型默认为Netty.
(3)Exchangers负责数据交换和网络通信的组件

第五步:注册中心注册服务

注册需处理契约:

  1. 当URL设置了check=false时,注册失败后不报错,在后台定时重试,否则抛出异常。
  2. 当URL设置了dynamic=false参数,则需持久存储,否则,当注册者出现断电等情况异常退出时,需自动删除。
  3. 当URL设置了category=routers时,表示分类存储,缺省类别为providers,可按分类部分通知数据。
  4. 当注册中心重启,网络抖动,不能丢失数据,包括断线自动删除数据。
  5. 允许URI相同但参数不同的URL并存,不能覆盖。
@Override
    public void register(URL url) {
        super.register(url);
        failedRegistered.remove(url);
        failedUnregistered.remove(url);
        try {
            // 向服务器端发送注册请求
            doRegister(url);
        } catch (Exception e) {
            Throwable t = e;

            // 如果开启了启动时检测,则直接抛出异常
            boolean check = getUrl().getParameter(Constants.CHECK_KEY, true)
                    && url.getParameter(Constants.CHECK_KEY, true)
                    && ! Constants.CONSUMER_PROTOCOL.equals(url.getProtocol());
            boolean skipFailback = t instanceof SkipFailbackWrapperException;
            if (check || skipFailback) {
                if(skipFailback) {
                    t = t.getCause();
                }
                throw new IllegalStateException("Failed to register " + url + " to registry " + getUrl().getAddress() + ", cause: " + t.getMessage(), t);
            } else {
                logger.error("Failed to register " + url + ", waiting for retry, cause: " + t.getMessage(), t);
            }

            // 将失败的注册请求记录到失败列表,定时重试
            failedRegistered.add(url);
        }
    }

第六步:订阅符合条件的已注册数据,当有注册数据变更时自动推送

订阅需处理契约:

  1. 当URL设置了check=false时,订阅失败后不报错,在后台定时重试。
  2. 当URL设置了category=routers,只通知指定分类的数据,多个分类用逗号分隔,并允许星号通配,表示订阅所有分类数据。
  3. 允许以interface,group,version,classifier作为条件查询,如:interface=com.alibaba.foo.BarService&version=1.0.0
  4. 并且查询条件允许星号通配,订阅所有接口的所有分组的所有版本,或:interface=&group=&version=&classifier=
  5. 当注册中心重启,网络抖动,需自动恢复订阅请求。
  6. 允许URI相同但参数不同的URL并存,不能覆盖。
  7. 必须阻塞订阅过程,等第一次通知完后再返回。
 @Override
    public void subscribe(URL url, NotifyListener listener) {
        super.subscribe(url, listener);
        removeFailedSubscribed(url, listener);
        try {
            // 向服务器端发送订阅请求
            doSubscribe(url, listener);
        } catch (Exception e) {
            Throwable t = e;

            List<URL> urls = getCacheUrls(url);
            if (urls != null && urls.size() > 0) {
                notify(url, listener, urls);
                logger.error("Failed to subscribe " + url + ", Using cached list: " + urls + " from cache file: " + getUrl().getParameter(Constants.FILE_KEY, System.getProperty("user.home") + "/dubbo-registry-" + url.getHost() + ".cache") + ", cause: " + t.getMessage(), t);
            } else {
                // 如果开启了启动时检测,则直接抛出异常
                boolean check = getUrl().getParameter(Constants.CHECK_KEY, true)
                        && url.getParameter(Constants.CHECK_KEY, true);
                boolean skipFailback = t instanceof SkipFailbackWrapperException;
                if (check || skipFailback) {
                    if(skipFailback) {
                        t = t.getCause();
                    }
                    throw new IllegalStateException("Failed to subscribe " + url + ", cause: " + t.getMessage(), t);
                } else {
                    logger.error("Failed to subscribe " + url + ", waiting for retry, cause: " + t.getMessage(), t);
                }
            }

            // 将失败的订阅请求记录到失败列表,定时重试
            addFailedSubscribed(url, listener);
        }
    }

订阅其实主要就是用来当数据变化通知重新暴露,已让服务消费者更新注册数据。

总结

大概可以看到服务暴露的流程:
获取注册中心URL——获取服务提供者暴露的URL——获取invoker——暴露服务准备工作——本地暴露服务——注册中心注册服务——注册中心订阅服务。
其实dubbo服务提供者暴露服务主要做了4件事:

最后编辑于
©著作权归作者所有,转载或内容合作请联系作者
【社区内容提示】社区部分内容疑似由AI辅助生成,浏览时请结合常识与多方信息审慎甄别。
平台声明:文章内容(如有图片或视频亦包括在内)由作者上传并发布,文章内容仅代表作者本人观点,简书系信息发布平台,仅提供信息存储服务。

相关阅读更多精彩内容

友情链接更多精彩内容