源码版本是 2.7.8
什么是注册中心? 服务治理框架中大致分为服务通信和服务管理两部分,服务管理可以分为服务注册、服务发现以及服务被热加工介入,服务提供者Provider会往注册中心注册服务,而消费者Consumer会从注册中心订阅相关服务,并不会订阅全部的服务。
dubbo-registry 模块 在dubbo中,注册中心相关的代码在dubbo-registry模块下,子模块dubbo-registry-api中定义了注册中心相关的基础代码,而在dubbo-registry-xxx模块中则定义了具体的注册中心类型实现代码,例如dubbo-registry-zookeeper模块则存放了zookeeper注册中心的实现代码。
类关系图:
dubbo-registry-api 相关实现 通过Registry的实现管理,分析下面各个接口类:
RegistryService 注册中心模块的服务接口,提供注册、取消注册、订阅、取消订阅、查询符合条件的已注册数据。
下面的注释,是官方的解释:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 public interface RegistryService { void register (URL url) ; void unregister (URL url) ; void subscribe (URL url, NotifyListener listener) ; void unsubscribe (URL url, NotifyListener listener) ; List <URL> lookup (URL url) ; }
Node Node 接口不是registry中,而是在common模块中的一个接口。Node 接口中主要声明了一些节点的操作方法,获取节点Url、是否可用、销毁节点。
Registry Registry 接口主要是继承了 Node 接口和 RegistryService 接口,将这两个接口的内容都统一在一起。同时 Registry 接口也提供了两个方法。
1 2 3 4 5 6 7 8 9 10 public interface Registry extends Node , RegistryService { default void reExportRegister (URL url) { register(url); } default void reExportUnregister (URL url) { unregister(url); } }
AbstractRegistry AbstractRegistry 实现 Registry 接口,实现了接口中定义的注册、订阅等方法。
抽象类的属性 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 private static final char URL_SEPARATOR = ' ' ;private static final String URL_SPLIT = "\\s+" ;private static final int MAX_RETRY_TIMES_SAVE_PROPERTIES = 3 ;protected final Logger logger = LoggerFactory.getLogger(getClass());private final Properties properties = new Properties ();private final ExecutorService registryCacheExecutor = Executors.newFixedThreadPool(1 , new NamedThreadFactory ("DubboSaveRegistryCache" , true ));private boolean syncSaveFile;private final AtomicLong lastCacheChanged = new AtomicLong ();private final AtomicInteger savePropertiesRetryTimes = new AtomicInteger ();private final Set <URL> registered = new ConcurrentHashSet <>();private final ConcurrentMap <URL, Set<NotifyListener>> subscribed = new ConcurrentHashMap <>();private final ConcurrentMap <URL, Map<String, List<URL>>> notified = new ConcurrentHashMap <>();private URL registryUrl;private File file;
构造方法 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 public AbstractRegistry (URL url) { setUrl(url); if (url.getParameter(REGISTRY__LOCAL_FILE_CACHE_ENABLED, true )) { syncSaveFile = url.getParameter(REGISTRY_FILESAVE_SYNC_KEY, false ); String defaultFilename = System.getProperty("user.home" ) + "/.dubbo/dubbo-registry-" + url.getParameter(APPLICATION_KEY) + "-" + url.getAddress().replaceAll(":" , "-" ) + ".cache" ; String filename = url.getParameter(FILE_KEY, defaultFilename); File file = null ; if (ConfigUtils.isNotEmpty(filename)) { file = new File (filename); if (!file.exists() && file.getParentFile() != null && !file.getParentFile().exists()) { if (!file.getParentFile().mkdirs()) { throw new IllegalArgumentException ("Invalid registry cache file " + file + ", cause: Failed to create directory " + file.getParentFile() + "!" ); } } } this .file = file; loadProperties(); notify(url.getBackupUrls()); } } private void loadProperties () { if (file != null && file.exists()) { InputStream in = null ; try { in = new FileInputStream (file); properties.load(in); if (logger.isInfoEnabled()) { logger.info("Load registry cache file " + file + ", data: " + properties); } } catch (Throwable e) { logger.warn("Failed to load registry cache file " + file, e); } finally { if (in != null ) { try { in.close(); } catch (IOException e) { logger.warn(e.getMessage(), e); } } } } }
lookup 获取消费者URL订阅的服务URL列表。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 @Override public List <URL> lookup (URL url) { List <URL> result = new ArrayList <>(); Map <String, List<URL>> notifiedUrls = getNotified().get(url); if (CollectionUtils.isNotEmptyMap(notifiedUrls)) { for (List<URL> urls : notifiedUrls.values()) { for (URL u : urls) { if (!EMPTY_PROTOCOL.equals(u.getProtocol())) { result.add(u); } } } } else { final AtomicReference <List<URL>> reference = new AtomicReference <>(); NotifyListener listener = reference::set; subscribe(url, listener); List <URL> urls = reference.get(); if (CollectionUtils.isNotEmpty(urls)) { for (URL u : urls) { if (!EMPTY_PROTOCOL.equals(u.getProtocol())) { result.add(u); } } } } return result; }
register and unregister URL 注册和取消注册,代码的主要逻辑就是从registered的内存缓存中添加或者删除URL。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 @Override public void register (URL url) { if (url == null ) { throw new IllegalArgumentException ("register url == null" ); } if (logger.isInfoEnabled()) { logger.info("Register: " + url); } registered.add(url); } @Override public void unregister (URL url) { if (url == null ) { throw new IllegalArgumentException ("unregister url == null" ); } if (logger.isInfoEnabled()) { logger.info("Unregister: " + url); } registered.remove(url); }
notify 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 protected void notify (List<URL> urls) { if (CollectionUtils.isEmpty(urls)) { return ; } for (Map.Entry<URL, Set<NotifyListener>> entry : getSubscribed().entrySet()) { URL url = entry.getKey(); if (!UrlUtils.isMatch(url, urls.get(0 ))) { continue ; } Set <NotifyListener> listeners = entry.getValue(); if (listeners != null ) { for (NotifyListener listener : listeners) { try { notify(url, listener, filterEmpty(url, urls)); } catch (Throwable t) { logger.error("Failed to notify registry event, urls: " + urls + ", cause: " + t.getMessage(), t); } } } } } protected void notify (URL url, NotifyListener listener, List<URL> urls) { if (url == null ) { throw new IllegalArgumentException ("notify url == null" ); } if (listener == null ) { throw new IllegalArgumentException ("notify listener == null" ); } if ((CollectionUtils.isEmpty(urls)) && !ANY_VALUE.equals(url.getServiceInterface())) { logger.warn("Ignore empty notify urls for subscribe url " + url); return ; } if (logger.isInfoEnabled()) { logger.info("Notify urls for subscribe url " + url + ", urls: " + urls); } Map <String, List<URL>> result = new HashMap <>(); for (URL u : urls) { if (UrlUtils.isMatch(url, u)) { String category = u.getParameter(CATEGORY_KEY, DEFAULT_CATEGORY); List <URL> categoryList = result.computeIfAbsent(category, k -> new ArrayList <>()); categoryList.add(u); } } if (result.size() == 0 ) { return ; } Map <String, List<URL>> categoryNotified = notified.computeIfAbsent(url, u -> new ConcurrentHashMap <>()); for (Map.Entry<String, List<URL>> entry : result.entrySet()) { String category = entry.getKey(); List <URL> categoryList = entry.getValue(); categoryNotified.put(category, categoryList); listener.notify(categoryList); saveProperties(url); } }
FailbackRegistry FailbackRegistry 是继承AbstractRegistry,增加了失败重试的机制作为抽象能力,后面不同的注册中心具体实现继承了这个类,就可以直接使用这个能力。
抽象类的属性 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 private final ConcurrentMap <URL, FailedRegisteredTask> failedRegistered = new ConcurrentHashMap <URL, FailedRegisteredTask>();private final ConcurrentMap <URL, FailedUnregisteredTask> failedUnregistered = new ConcurrentHashMap <URL, FailedUnregisteredTask>();private final ConcurrentMap <Holder, FailedSubscribedTask> failedSubscribed = new ConcurrentHashMap <Holder, FailedSubscribedTask>();private final ConcurrentMap <Holder, FailedUnsubscribedTask> failedUnsubscribed = new ConcurrentHashMap <Holder, FailedUnsubscribedTask>();private final int retryPeriod;private final HashedWheelTimer retryTimer;
构造方法 1 2 3 4 5 6 7 public FailbackRegistry (URL url) { super (url); this .retryPeriod = url.getParameter(REGISTRY_RETRY_PERIOD_KEY, DEFAULT_REGISTRY_RETRY_PERIOD); retryTimer = new HashedWheelTimer (new NamedThreadFactory ("DubboRegistryRetryTimer" , true ), retryPeriod, TimeUnit.MILLISECONDS, 128 ); }
register 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 @Override public void register (URL url) { if (!acceptable(url)) { logger.info("URL " + url + " will not be registered to Registry. Registry " + url + " does not accept service of this protocol type." ); return ; } super .register(url); removeFailedRegistered(url); removeFailedUnregistered(url); try { doRegister(url); } catch (Exception e) { Throwable t = e; boolean check = getUrl().getParameter(Constants.CHECK_KEY, true ) && url.getParameter(Constants.CHECK_KEY, true ) && !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); } addFailedRegistered(url); } }
另外的几个方法,unregister、subscribe、unsubscribe都类似,里面会调用一个doxxx方法,底层就是不同的注册中心具体实现方法。
notify 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 protected void notify (URL url, NotifyListener listener, List<URL> urls) { if (url == null ) { throw new IllegalArgumentException ("notify url == null" ); } if (listener == null ) { throw new IllegalArgumentException ("notify listener == null" ); } try { doNotify(url, listener, urls); } catch (Exception t) { logger.error("Failed to notify addresses for subscribe " + url + ", cause: " + t.getMessage(), t); } } protected void doNotify (URL url, NotifyListener listener, List<URL> urls) { super .notify(url, listener, urls); }
RegistryFactory RegistryFactory 接口定义只有一个getRegistry方法。URL 为dubbo封装的统一资源定位符,其中定义了协议protocol、用户名username、密码password、host主机、path路径等等属性。
RegistryFactory 是一个工厂方法,根据具体的注册协议,比如zookeeper,获取具体的注册中心实现。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 @SPI("dubbo") public interface RegistryFactory { @Adaptive({"protocol"}) Registry getRegistry (URL url) ; }
NotifyListener 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 public interface NotifyListener { void notify (List<URL> urls) ; default void addServiceListener (ServiceInstancesChangedListener instanceListener) { } }
dubbo-registry-zookeeper 对于dubbo的注册中心,这里只介绍下zookeeper,这个也是dubbo默认的注册中心。
ZookeeperRegistry 属性 1 2 3 4 5 6 7 8 9 10 private final static String DEFAULT_ROOT = "dubbo" ;private final String root;private final Set <String> anyServices = new ConcurrentHashSet <>();private final ConcurrentMap <URL, ConcurrentMap<NotifyListener, ChildListener>> zkListeners = new ConcurrentHashMap <>();private final ZookeeperClient zkClient;
构造方法 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 public ZookeeperRegistry (URL url, ZookeeperTransporter zookeeperTransporter) { super (url); if (url.isAnyHost()) { throw new IllegalStateException ("registry address == null" ); } String group = url.getParameter(GROUP_KEY, DEFAULT_ROOT); if (!group.startsWith(PATH_SEPARATOR)) { group = PATH_SEPARATOR + group; } this .root = group; zkClient = zookeeperTransporter.connect(url); zkClient.addStateListener((state) -> { if (state == StateListener.RECONNECTED) { logger.warn("Trying to fetch the latest urls, in case there're provider changes during connection loss.\n" + " Since ephemeral ZNode will not get deleted for a connection lose, " + "there's no need to re-register url of this instance." ); ZookeeperRegistry.this .fetchLatestAddresses(); } else if (state == StateListener.NEW_SESSION_CREATED) { logger.warn("Trying to re-register urls and re-subscribe listeners of this instance to registry..." ); try { ZookeeperRegistry.this .recover(); } catch (Exception e) { logger.error(e.getMessage(), e); } } else if (state == StateListener.SESSION_LOST) { logger.warn("Url of this instance will be deleted from registry soon. " + "Dubbo client will try to re-register once a new session is created." ); } else if (state == StateListener.SUSPENDED) { } else if (state == StateListener.CONNECTED) { } }); }
服务注册发布和下线取消注册 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 public void doRegister (URL url) { try { zkClient.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); } } public void doUnregister (URL url) { try { zkClient.delete(toUrlPath(url)); } catch (Throwable e) { throw new RpcException ("Failed to unregister " + url + " to zookeeper " + getUrl() + ", cause: " + e.getMessage(), e); } }
Reference