暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Dubbo源码分析之提供端(2)

徘徊笔记 2019-05-24
174

来源:https://github.com/apache/incubator-dubbo



获取注册URL,还原原来的协议,因为设置的临时协议只是为了导向到该类。

    private URL getRegistryUrl(Invoker<?> originInvoker) {
    URL registryUrl = originInvoker.getUrl();
    if (REGISTRY_PROTOCOL.equals(registryUrl.getProtocol())) {
    String protocol = registryUrl.getParameter(REGISTRY_KEY, DEFAULT_DIRECTORY);
    registryUrl = registryUrl.setProtocol(protocol).removeParameter(REGISTRY_KEY);
    }
    return registryUrl;
    }


    获取提供者URL,之前该URL被当做参数设置到注册URL上。

      private URL getProviderUrl(final Invoker<?> originInvoker) {
      String export = originInvoker.getUrl().getParameterAndDecoded(EXPORT_KEY);
      if (export == null || export.length() == 0) {
      throw new IllegalArgumentException("The registry export url is null! registry: " + originInvoker.getUrl());
      }
      return URL.valueOf(export);
      }


      暴露本地代理invoker

        private <T> ExporterChangeableWrapper<T> doLocalExport(final Invoker<T> originInvoker, URL providerUrl) {
        String key = getCacheKey(originInvoker);


        return (ExporterChangeableWrapper<T>) bounds.computeIfAbsent(key, s -> {
        Invoker<?> invokerDelegate = new InvokerDelegate<>(originInvoker, providerUrl);
        return new ExporterChangeableWrapper<>((Exporter<T>) protocol.export(invokerDelegate), originInvoker);
        });
        }


        private String getCacheKey(final Invoker<?> originInvoker) {
        URL providerUrl = getProviderUrl(originInvoker);
        String key = providerUrl.removeParameters("dynamic", "enabled").toFullString();
        return key;
        }


        提供者URL默认为dubbo协议,这次也会走那三个包装类,QosProtocolWrapper

        不做处理,ProtocolFilterWrapper对代理进行过滤器包装,他会获取提供端组设置的所有需要激活的过滤器类,ProtocolListenerWrapper获取URL中设置的和自动激活的监听器类。

          private static <T> Invoker<T> buildInvokerChain(final Invoker<T> invoker, String key, String group) {
          Invoker<T> last = invoker;
          List<Filter> filters = ExtensionLoader.getExtensionLoader(Filter.class).getActivateExtension(invoker.getUrl(), key, group);
          if (!filters.isEmpty()) {
          for (int i = filters.size() - 1; i >= 0; i--) {
          final Filter filter = filters.get(i);
          final Invoker<T> next = last;
          last = new Invoker<T>() {


          @Override
          public Class<T> getInterface() {
          return invoker.getInterface();
          }


          @Override
          public URL getUrl() {
          return invoker.getUrl();
          }


          @Override
          public boolean isAvailable() {
          return invoker.isAvailable();
          }


          @Override
          public Result invoke(Invocation invocation) throws RpcException {
          Result result = filter.invoke(next, invocation);
          if (result instanceof AsyncRpcResult) {
          AsyncRpcResult asyncResult = (AsyncRpcResult) result;
          asyncResult.thenApplyWithContext(r -> filter.onResponse(r, invoker, invocation));
          return asyncResult;
          } else {
          return filter.onResponse(result, invoker, invocation);
          }
          }


          @Override
          public void destroy() {
          invoker.destroy();
          }


          @Override
          public String toString() {
          return invoker.toString();
          }
          };
          }
          }
          return last;
          }


          通过dubbo协议暴露服务,获取具体的提供端URL,组装dubbo暴露者类

            public <T> Exporter<T> export(Invoker<T> invoker) throws RpcException {
            URL url = invoker.getUrl();


            // export service.
            String key = serviceKey(url);
            DubboExporter<T> exporter = new DubboExporter<T>(invoker, key, exporterMap);
            exporterMap.put(key, exporter);


            //export an stub service for dispatching 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);
            optimizeSerialization(url);


            return exporter;
            }


            组装服务key,由服务组名,服务类名,版本号以及端口组成

              protected static String serviceKey(URL url) {
              int port = url.getParameter(Constants.BIND_PORT_KEY, url.getPort());
              return serviceKey(port, url.getPath(), url.getParameter(Constants.VERSION_KEY),
              url.getParameter(Constants.GROUP_KEY));
              }


              protected static String serviceKey(int port, String serviceName, String serviceVersion, String serviceGroup) {
              return ProtocolUtils.serviceKey(port, serviceName, serviceVersion, serviceGroup);
              }


              public static String serviceKey(int port, String serviceName, String serviceVersion, String serviceGroup) {
              StringBuilder buf = new StringBuilder();
              if (StringUtils.isNotEmpty(serviceGroup)) {
              buf.append(serviceGroup);
              buf.append("/");
              }
              buf.append(serviceName);
              if (serviceVersion != null && serviceVersion.length() > 0 && !"0.0.0".equals(serviceVersion)) {
              buf.append(":");
              buf.append(serviceVersion);
              }
              buf.append(":");
              buf.append(port);
              return buf.toString();
              }


              接收消费端的请求,所以需要启动服务端

                private void openServer(URL url) {
                // find server.
                String key = url.getAddress();
                //client can export a service which's only for server to invoke
                boolean isServer = url.getParameter(Constants.IS_SERVER_KEY, true);
                if (isServer) {
                ExchangeServer server = serverMap.get(key);
                if (server == null) {
                synchronized (this) {
                server = serverMap.get(key);
                if (server == null) {
                serverMap.put(key, createServer(url));
                }
                }
                } else {
                // server supports reset, use together with override
                server.reset(url);
                }
                }
                }


                二重检查后创建交换服务端,设置心跳时间以及编码协议

                  private ExchangeServer createServer(URL url) {
                  url = URLBuilder.from(url)
                  // send readonly event when server closes, it's enabled by default
                  .addParameterIfAbsent(Constants.CHANNEL_READONLYEVENT_SENT_KEY, Boolean.TRUE.toString())
                  // enable heartbeat by default
                  .addParameterIfAbsent(Constants.HEARTBEAT_KEY, String.valueOf(Constants.DEFAULT_HEARTBEAT))
                  .addParameter(Constants.CODEC_KEY, DubboCodec.NAME)
                  .build();
                  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);
                  }


                  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;
                  }


                  获取对应的交换器,默认为header,即HeaderExchanger

                    public static ExchangeServer bind(URL url, ExchangeHandler handler) throws RemotingException {
                    if (url == null) {
                    throw new IllegalArgumentException("url == null");
                    }
                    if (handler == null) {
                    throw new IllegalArgumentException("handler == null");
                    }
                    url = url.addParameterIfAbsent(Constants.CODEC_KEY, "exchange");
                    return getExchanger(url).bind(url, handler);
                    }


                    public static Exchanger getExchanger(URL url) {
                    String type = url.getParameter(Constants.EXCHANGER_KEY, Constants.DEFAULT_EXCHANGER);
                    return getExchanger(type);
                    }


                    public static Exchanger getExchanger(String type) {
                    return ExtensionLoader.getExtensionLoader(Exchanger.class).getExtension(type);
                    }


                    文章转载自徘徊笔记,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

                    评论