
来源: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>() {@Overridepublic Class<T> getInterface() {return invoker.getInterface();}@Overridepublic URL getUrl() {return invoker.getUrl();}@Overridepublic boolean isAvailable() {return invoker.isAvailable();}@Overridepublic 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);}}@Overridepublic void destroy() {invoker.destroy();}@Overridepublic 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 eventBoolean 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 invokeboolean 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 overrideserver.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进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。




