入门第一课,我们搭建了一个简单的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®istry=zookeeper×tamp=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缓存起来。
第三步:暴露远程服务
- 协议在接收请求时,应记录请求来源方地址信息:RpcContext.getContext().setRemoteAddress();
- export()必须是幂等的,也就是暴露同一个URL的Invoker两次,和暴露一次没有区别。
- 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×tamp=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×tamp=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×tamp=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×tamp=1548921845226
(2)dubbo协议的服务器端实现类型默认为Netty.
(3)Exchangers负责数据交换和网络通信的组件
第五步:注册中心注册服务
注册需处理契约:
- 当URL设置了check=false时,注册失败后不报错,在后台定时重试,否则抛出异常。
- 当URL设置了dynamic=false参数,则需持久存储,否则,当注册者出现断电等情况异常退出时,需自动删除。
- 当URL设置了category=routers时,表示分类存储,缺省类别为providers,可按分类部分通知数据。
- 当注册中心重启,网络抖动,不能丢失数据,包括断线自动删除数据。
- 允许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);
}
}
第六步:订阅符合条件的已注册数据,当有注册数据变更时自动推送
订阅需处理契约:
- 当URL设置了check=false时,订阅失败后不报错,在后台定时重试。
- 当URL设置了category=routers,只通知指定分类的数据,多个分类用逗号分隔,并允许星号通配,表示订阅所有分类数据。
- 允许以interface,group,version,classifier作为条件查询,如:interface=com.alibaba.foo.BarService&version=1.0.0
- 并且查询条件允许星号通配,订阅所有接口的所有分组的所有版本,或:interface=&group=&version=&classifier=
- 当注册中心重启,网络抖动,需自动恢复订阅请求。
- 允许URI相同但参数不同的URL并存,不能覆盖。
- 必须阻塞订阅过程,等第一次通知完后再返回。
@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件事:
- 将配置文件的属性对转换成URL(dubbo中间传输都是根据URL,URL 采用标准格式:protocol://username:password@host:port/path?key=value&key=value)
- 服务的实现转为Invoker,暴露服务(主要是Invoker)
- 注册中心注册服务
- 订阅与通知服务
- 返回新的exporter实例