From b575b08bde1bfae80caed01252820ac53cb336e1 Mon Sep 17 00:00:00 2001 From: "william.liangf" Date: Fri, 11 Nov 2011 05:16:20 +0000 Subject: [PATCH] =?UTF-8?q?DUBBO-3=09=E5=AE=9E=E7=8E=B0=E4=B8=8Ezookeeper?= =?UTF-8?q?=E6=B3=A8=E5=86=8C=E4=B8=AD=E5=BF=83=E7=9A=84=E6=A1=A5=E6=8E=A5?= =?UTF-8?q?=EF=BC=8C=E5=A2=9E=E5=8A=A0=E9=87=8D=E8=BF=9E=E6=81=A2=E5=A4=8D?= =?UTF-8?q?=E6=95=B0=E6=8D=AE=E9=80=BB=E8=BE=91=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit git-svn-id: http://code.alibabatech.com/svn/dubbo/trunk@248 1a56cb94-b969-4eaa-88fa-be21384802f2 --- .../registry/zookeeper/ZookeeperRegistry.java | 62 ++++++++++++++++++- .../registry/support/FailbackRegistry.java | 14 ++--- 2 files changed, 66 insertions(+), 10 deletions(-) diff --git a/dubbo-registry-zookeeper/src/main/java/com/alibaba/dubbo/registry/zookeeper/ZookeeperRegistry.java b/dubbo-registry-zookeeper/src/main/java/com/alibaba/dubbo/registry/zookeeper/ZookeeperRegistry.java index 37d4a9a78e..e022aee858 100644 --- a/dubbo-registry-zookeeper/src/main/java/com/alibaba/dubbo/registry/zookeeper/ZookeeperRegistry.java +++ b/dubbo-registry-zookeeper/src/main/java/com/alibaba/dubbo/registry/zookeeper/ZookeeperRegistry.java @@ -16,6 +16,7 @@ package com.alibaba.dubbo.registry.zookeeper; import java.util.ArrayList; +import java.util.HashSet; import java.util.List; import java.util.Map; import java.util.Set; @@ -34,6 +35,7 @@ import com.alibaba.dubbo.common.Constants; import com.alibaba.dubbo.common.URL; import com.alibaba.dubbo.common.logger.Logger; import com.alibaba.dubbo.common.logger.LoggerFactory; +import com.alibaba.dubbo.common.utils.ConcurrentHashSet; import com.alibaba.dubbo.common.utils.UrlUtils; import com.alibaba.dubbo.registry.NotifyListener; import com.alibaba.dubbo.registry.support.FailbackRegistry; @@ -56,6 +58,8 @@ public class ZookeeperRegistry extends FailbackRegistry { private final ReentrantLock zookeeperLock = new ReentrantLock(); + private final Set failedWatched = new ConcurrentHashSet(); + private volatile ZooKeeper zookeeper; public ZookeeperRegistry(URL url) { @@ -75,6 +79,54 @@ public class ZookeeperRegistry extends FailbackRegistry { @Override protected void doRetry() { initZookeeper(); + if (failedWatched.size() > 0) { + Set failed = new HashSet(failedWatched); + if (failed.size() > 0) { + if (logger.isInfoEnabled()) { + logger.info("Retry watch " + failed); + } + for (String service : failed) { + try { + zookeeper.getChildren(service, true); + failedWatched.remove(service); + } catch (Throwable t) { + logger.warn("Failed to retry register " + failed + ", waiting for again, cause: " + t.getMessage(), t); + } + } + } + } + } + + private List watch(String service) { + try { + ZooKeeper zk = ZookeeperRegistry.this.zookeeper; + if (zk != null) { + List result = zookeeper.getChildren(service, true); + failedWatched.remove(service); + return result; + } + } catch (Throwable e) { + logger.warn(e.getMessage(), e); + } + failedWatched.add(service); + return new ArrayList(0); + } + + private void recover() { + for (String url : new HashSet(getRegistered())) { + try { + register(URL.valueOf(url)); + } catch (Throwable e) { + logger.warn(e.getMessage(), e); + } + } + for (String url : new HashSet(getSubscribed().keySet())) { + try { + watch(toServicePath(URL.valueOf(url))); + } catch (Throwable e) { + logger.warn(e.getMessage(), e); + } + } } private void initZookeeper() { @@ -85,6 +137,7 @@ public class ZookeeperRegistry extends FailbackRegistry { zk = this.zookeeper; if (zk == null || zk.getState() == null || ! zk.getState().isAlive()) { this.zookeeper = createZookeeper(); + recover(); } if (zk != null) { zk.close(); @@ -110,6 +163,9 @@ public class ZookeeperRegistry extends FailbackRegistry { try { if (event.getState() == KeeperState.Expired) { initZookeeper(); + } else if (event.getState() == KeeperState.SyncConnected + && event.getType() == EventType.None) { + recover(); } if (event.getType() != EventType.NodeChildrenChanged) { return; @@ -118,8 +174,10 @@ public class ZookeeperRegistry extends FailbackRegistry { if (path == null || path.length() == 0) { return; } - ZooKeeper zk = ZookeeperRegistry.this.zookeeper; - List providers = zk.getChildren(path, true); + List providers = watch(path); + if (providers == null || providers.size() == 0) { + return; + } String service = path; int i = service.lastIndexOf(SEPARATOR); if (i >= 0) { diff --git a/dubbo-registry/src/main/java/com/alibaba/dubbo/registry/support/FailbackRegistry.java b/dubbo-registry/src/main/java/com/alibaba/dubbo/registry/support/FailbackRegistry.java index c012f52242..828cdcd82b 100644 --- a/dubbo-registry/src/main/java/com/alibaba/dubbo/registry/support/FailbackRegistry.java +++ b/dubbo-registry/src/main/java/com/alibaba/dubbo/registry/support/FailbackRegistry.java @@ -84,15 +84,13 @@ public abstract class FailbackRegistry extends AbstractRegistry { if (logger.isInfoEnabled()) { logger.info("Retry register " + failed); } - if (! failed.isEmpty()) { - try { - for (String url : failed) { - doRegister(URL.valueOf(url)); - failedRegistered.remove(url); - } - } catch (Throwable t) { // 忽略所有异常,等待下次重试 - logger.warn("Failed to retry register " + failed + ", waiting for again, cause: " + t.getMessage(), t); + try { + for (String url : failed) { + doRegister(URL.valueOf(url)); + failedRegistered.remove(url); } + } catch (Throwable t) { // 忽略所有异常,等待下次重试 + logger.warn("Failed to retry register " + failed + ", waiting for again, cause: " + t.getMessage(), t); } } }