From 09321afc6f888e2d1eeb8846451ee6f95dd972b0 Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sat, 4 Aug 2018 00:44:39 +0800 Subject: [PATCH 1/8] =?UTF-8?q?=E5=A2=9E=E5=8A=A0zookeeper=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E6=BA=90?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pom.xml | 5 + sentinel-demo/pom.xml | 1 + .../pom.xml | 49 +++++++++ .../zookeeper/ZookeeperConfigSender.java | 37 +++++++ .../zookeeper/ZookeeperDataSourceDemo.java | 50 +++++++++ sentinel-extension/pom.xml | 1 + .../sentinel-datasource-zookeeper/pom.xml | 31 ++++++ .../zookeeper/ZookeeperDataSource.java | 100 ++++++++++++++++++ 8 files changed, 274 insertions(+) create mode 100644 sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml create mode 100644 sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java create mode 100644 sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java create mode 100644 sentinel-extension/sentinel-datasource-zookeeper/pom.xml create mode 100644 sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java diff --git a/pom.xml b/pom.xml index 6c2f6d1672..cff702d472 100755 --- a/pom.xml +++ b/pom.xml @@ -91,6 +91,11 @@ sentinel-datasource-nacos ${project.version} + + com.alibaba.csp + sentinel-datasource-zookeeper + ${project.version} + com.alibaba.csp sentinel-adapter diff --git a/sentinel-demo/pom.xml b/sentinel-demo/pom.xml index 8accce646c..d168be48cb 100755 --- a/sentinel-demo/pom.xml +++ b/sentinel-demo/pom.xml @@ -18,6 +18,7 @@ sentinel-demo-rocketmq sentinel-demo-dubbo sentinel-demo-nacos-datasource + sentinel-demo-zookeeper-datasource diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml b/sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml new file mode 100644 index 0000000000..b320ecba59 --- /dev/null +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml @@ -0,0 +1,49 @@ + + + + sentinel-demo + com.alibaba.csp + 0.1.1-SNAPSHOT + + 4.0.0 + + sentinel-demo-zookeeper-datasource + + + + com.alibaba.csp + sentinel-core + + + com.alibaba.csp + sentinel-datasource-extension + + + com.alibaba.csp + sentinel-datasource-zookeeper + + + + com.alibaba + fastjson + + + + + + + org.apache.maven.plugins + maven-compiler-plugin + ${maven.compiler.version} + + 1.8 + 1.8 + ${java.encoding} + + + + + + \ No newline at end of file diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java new file mode 100644 index 0000000000..e6853f79b8 --- /dev/null +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java @@ -0,0 +1,37 @@ +package com.alibaba.csp.sentinel.demo.datasource.zookeeper; + +import org.I0Itec.zkclient.ZkClient; + +/** + * Zookeeper config sender for demo + * + * @author guonanjun + */ +public class ZookeeperConfigSender { + + public static void main(String[] args) throws Exception { + + + + final String remoteAddress = "127.0.0.1:2181"; + final String groupId = "Sentinel-Demo"; + final String dataId = "SYSTEM-CODE-DEMO-FLOW"; + final String rule = "[\n" + + " {\n" + + " \"resource\": \"TestResource\",\n" + + " \"controlBehavior\": 0,\n" + + " \"count\": 10.0,\n" + + " \"grade\": 1,\n" + + " \"limitApp\": \"default\",\n" + + " \"strategy\": 0\n" + + " }\n" + + "]"; + + ZkClient zkClient = new ZkClient(remoteAddress, 5000); + String path = "/" + groupId + "/" + dataId; + if (!zkClient.exists(path)) { + zkClient.createPersistent(path, true); + } + zkClient.writeData(path, rule); + } +} diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java new file mode 100644 index 0000000000..be5e39f527 --- /dev/null +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java @@ -0,0 +1,50 @@ +package com.alibaba.csp.sentinel.demo.datasource.zookeeper; + +import java.util.List; + +import com.alibaba.csp.sentinel.datasource.DataSource; +import com.alibaba.csp.sentinel.datasource.zookeeper.ZookeeperDataSource; +import com.alibaba.csp.sentinel.slots.block.degrade.DegradeRule; +import com.alibaba.csp.sentinel.slots.block.degrade.DegradeRuleManager; +import com.alibaba.csp.sentinel.slots.block.flow.FlowRule; +import com.alibaba.csp.sentinel.slots.block.flow.FlowRuleManager; +import com.alibaba.csp.sentinel.slots.system.SystemRule; +import com.alibaba.csp.sentinel.slots.system.SystemRuleManager; +import com.alibaba.fastjson.JSON; +import com.alibaba.fastjson.TypeReference; + +/** + * Zookeeper DataSource Demo + * + * @author guonanjun + */ +public class ZookeeperDataSourceDemo { + + private static final String KEY = "TestResource"; + + public static void main(String[] args) { + loadRules(); + } + + private static void loadRules() { + final String remoteAddress = "127.0.0.1:2181"; + final String groupId = "Sentinel-Demo"; + final String flowDataId = "SYSTEM-CODE-DEMO-FLOW"; + final String degradeDataId = "SYSTEM-CODE-DEMO-DEGRADE"; + final String systemDataId = "SYSTEM-CODE-DEMO-SYSTEM"; + + + DataSource> flowRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, flowDataId, + source -> JSON.parseObject(source, new TypeReference>() {})); + FlowRuleManager.register2Property(flowRuleDataSource.getProperty()); + + DataSource> degradeRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, degradeDataId, + source -> JSON.parseObject(source, new TypeReference>() {})); + DegradeRuleManager.register2Property(degradeRuleDataSource.getProperty()); + + DataSource> systemRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, systemDataId, + source -> JSON.parseObject(source, new TypeReference>() {})); + SystemRuleManager.register2Property(systemRuleDataSource.getProperty()); + + } +} diff --git a/sentinel-extension/pom.xml b/sentinel-extension/pom.xml index ba9b729fa5..e72add499a 100755 --- a/sentinel-extension/pom.xml +++ b/sentinel-extension/pom.xml @@ -14,6 +14,7 @@ sentinel-datasource-extension sentinel-datasource-nacos + sentinel-datasource-zookeeper diff --git a/sentinel-extension/sentinel-datasource-zookeeper/pom.xml b/sentinel-extension/sentinel-datasource-zookeeper/pom.xml new file mode 100644 index 0000000000..00fe649612 --- /dev/null +++ b/sentinel-extension/sentinel-datasource-zookeeper/pom.xml @@ -0,0 +1,31 @@ + + + + sentinel-extension + com.alibaba.csp + 0.1.1-SNAPSHOT + + 4.0.0 + + sentinel-datasource-zookeeper + jar + + + 0.10 + + + + + com.alibaba.csp + sentinel-datasource-extension + + + + com.101tec + zkclient + ${zkclient.version} + + + \ No newline at end of file diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java new file mode 100644 index 0000000000..c92d1f4b03 --- /dev/null +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -0,0 +1,100 @@ +package com.alibaba.csp.sentinel.datasource.zookeeper; + +import com.alibaba.csp.sentinel.datasource.AbstractDataSource; +import com.alibaba.csp.sentinel.datasource.ConfigParser; +import com.alibaba.csp.sentinel.log.RecordLog; +import com.alibaba.csp.sentinel.util.StringUtil; +import org.I0Itec.zkclient.IZkDataListener; +import org.I0Itec.zkclient.ZkClient; + +/** + * Zookeeper DataSource + * + * @author guonanjun + */ +public class ZookeeperDataSource extends AbstractDataSource { + + private static final int DEFAULT_TIMEOUT = 3000; + + private final IZkDataListener zkDataListener; + private final String groupId; + private final String dataId; + + private ZkClient zkClient = null; + + public ZookeeperDataSource(final String serverAddr, final String groupId, final String dataId, + ConfigParser parser) { + super(parser); + if (StringUtil.isBlank(serverAddr) || StringUtil.isBlank(groupId) || StringUtil.isBlank(dataId)) { + throw new IllegalArgumentException(String.format("Bad argument: serverAddr=[%s], groupId=[%s], dataId=[%s]", + serverAddr, groupId, dataId)); + } + this.groupId = groupId; + this.dataId = dataId; + this.zkDataListener = new IZkDataListener() { + @Override + public void handleDataChange(String path, Object configInfo) throws Exception { + RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s, %s): %s", + serverAddr, dataId, groupId, configInfo)); + System.out.println("===========" + configInfo); + T newValue = ZookeeperDataSource.this.parser.parse(String.valueOf(configInfo)); + // Update the new value to the property. + getProperty().updateValue(newValue); + } + + @Override + public void handleDataDeleted(String s) throws Exception { + RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s, %s): %s", + serverAddr, dataId, groupId, null)); + // Update the new value to the property. + getProperty().updateValue(null); + } + }; + initZookeeperListener(serverAddr); + loadInitialConfig(); + } + + private void loadInitialConfig() { + try { + T newValue = loadConfig(); + if (newValue == null) { + RecordLog.info("[ZookeeperDataSource] WARN: initial config is null, you may have to check your data source"); + } + getProperty().updateValue(newValue); + } catch (Exception ex) { + RecordLog.info("[ZookeeperDataSource] Error when loading initial config", ex); + } + } + + private void initZookeeperListener(String serverAddr) { + try { + this.zkClient = new ZkClient(serverAddr, DEFAULT_TIMEOUT); + String path = "/" + this.groupId + "/" + this.dataId; + if (!zkClient.exists(path)) { + zkClient.createPersistent(path, true); + } + zkClient.subscribeDataChanges(path, zkDataListener); + } catch (Exception e) { + RecordLog.info("[ZookeeperDataSource] Error occurred when initializing Zookeeper data source", e); + e.printStackTrace(); + } + } + + @Override + public String readSource() throws Exception { + if (zkClient == null) { + throw new IllegalStateException("Zookeeper has not been initialized or error occurred"); + } + String path = "/" + this.groupId + "/" + this.dataId; + return zkClient.readData(path); + } + + @Override + public void close() throws Exception { + if (zkClient != null) { + String path = "/" + this.groupId + "/" + this.dataId; + zkClient.unsubscribeDataChanges(path, zkDataListener); + } + zkClient.close(); + } +} From a3b5be34d4326c147b933206513eaaa74a83432d Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sat, 4 Aug 2018 08:35:51 +0800 Subject: [PATCH 2/8] =?UTF-8?q?=E5=8E=BB=E6=8E=89=E5=A4=9A=E4=BD=99?= =?UTF-8?q?=E7=9A=84=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../demo/datasource/zookeeper/ZookeeperConfigSender.java | 2 -- .../csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java | 1 - 2 files changed, 3 deletions(-) diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java index e6853f79b8..5d3565216f 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java @@ -11,8 +11,6 @@ public class ZookeeperConfigSender { public static void main(String[] args) throws Exception { - - final String remoteAddress = "127.0.0.1:2181"; final String groupId = "Sentinel-Demo"; final String dataId = "SYSTEM-CODE-DEMO-FLOW"; diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java index c92d1f4b03..d69de570aa 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -36,7 +36,6 @@ public ZookeeperDataSource(final String serverAddr, final String groupId, final public void handleDataChange(String path, Object configInfo) throws Exception { RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s, %s): %s", serverAddr, dataId, groupId, configInfo)); - System.out.println("===========" + configInfo); T newValue = ZookeeperDataSource.this.parser.parse(String.valueOf(configInfo)); // Update the new value to the property. getProperty().updateValue(newValue); From ed2932a52d0e4ed32b04fce7cdcd82ead2ddfe7b Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sat, 4 Aug 2018 08:43:04 +0800 Subject: [PATCH 3/8] =?UTF-8?q?=E8=A7=84=E8=8C=83=E5=91=BD=E5=90=8D?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java index d69de570aa..67058d00bc 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -42,7 +42,7 @@ public void handleDataChange(String path, Object configInfo) throws Exception { } @Override - public void handleDataDeleted(String s) throws Exception { + public void handleDataDeleted(String path) throws Exception { RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s, %s): %s", serverAddr, dataId, groupId, null)); // Update the new value to the property. From 5843e63217530e3a9a7d52ad011fc879d01dcd98 Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sat, 4 Aug 2018 15:27:04 +0800 Subject: [PATCH 4/8] =?UTF-8?q?=E5=B0=86zkclient=E6=9B=BF=E6=8D=A2?= =?UTF-8?q?=E4=B8=BAcurator?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../zookeeper/ZookeeperConfigSender.java | 17 ++-- .../sentinel-datasource-zookeeper/pom.xml | 20 ++++- .../zookeeper/ZookeeperDataSource.java | 80 +++++++++++++------ 3 files changed, 82 insertions(+), 35 deletions(-) diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java index 5d3565216f..f048de820d 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java @@ -1,6 +1,10 @@ package com.alibaba.csp.sentinel.demo.datasource.zookeeper; -import org.I0Itec.zkclient.ZkClient; +import org.apache.curator.framework.CuratorFramework; +import org.apache.curator.framework.CuratorFrameworkFactory; +import org.apache.curator.retry.RetryNTimes; +import org.apache.zookeeper.CreateMode; +import org.apache.zookeeper.data.Stat; /** * Zookeeper config sender for demo @@ -25,11 +29,14 @@ public static void main(String[] args) throws Exception { + " }\n" + "]"; - ZkClient zkClient = new ZkClient(remoteAddress, 5000); + CuratorFramework zkClient = CuratorFrameworkFactory.newClient(remoteAddress, new RetryNTimes(3, 5000)); + zkClient.start(); String path = "/" + groupId + "/" + dataId; - if (!zkClient.exists(path)) { - zkClient.createPersistent(path, true); + Stat stat = zkClient.checkExists().forPath(path); + if (stat == null) { + zkClient.create().creatingParentContainersIfNeeded().withMode(CreateMode.PERSISTENT).forPath(path, null); } - zkClient.writeData(path, rule); + zkClient.setData().forPath(path, rule.getBytes()); + // zkClient.delete().forPath(path); } } diff --git a/sentinel-extension/sentinel-datasource-zookeeper/pom.xml b/sentinel-extension/sentinel-datasource-zookeeper/pom.xml index 00fe649612..6c7b2e21f3 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/pom.xml +++ b/sentinel-extension/sentinel-datasource-zookeeper/pom.xml @@ -13,7 +13,8 @@ jar - 0.10 + 3.4.13 + 4.0.1 @@ -23,9 +24,20 @@ - com.101tec - zkclient - ${zkclient.version} + org.apache.zookeeper + zookeeper + ${zookeeper.version} + + + org.apache.curator + curator-recipes + ${curator.version} + + + org.apache.zookeeper + zookeeper + + \ No newline at end of file diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java index 67058d00bc..2a8af4fa1f 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -1,11 +1,23 @@ package com.alibaba.csp.sentinel.datasource.zookeeper; +import java.util.concurrent.ArrayBlockingQueue; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; + +import com.alibaba.csp.sentinel.concurrent.NamedThreadFactory; import com.alibaba.csp.sentinel.datasource.AbstractDataSource; import com.alibaba.csp.sentinel.datasource.ConfigParser; import com.alibaba.csp.sentinel.log.RecordLog; import com.alibaba.csp.sentinel.util.StringUtil; -import org.I0Itec.zkclient.IZkDataListener; -import org.I0Itec.zkclient.ZkClient; +import org.apache.curator.framework.CuratorFramework; +import org.apache.curator.framework.CuratorFrameworkFactory; +import org.apache.curator.framework.recipes.cache.ChildData; +import org.apache.curator.framework.recipes.cache.NodeCache; +import org.apache.curator.framework.recipes.cache.NodeCacheListener; +import org.apache.curator.retry.RetryNTimes; +import org.apache.zookeeper.CreateMode; +import org.apache.zookeeper.data.Stat; /** * Zookeeper DataSource @@ -14,13 +26,19 @@ */ public class ZookeeperDataSource extends AbstractDataSource { - private static final int DEFAULT_TIMEOUT = 3000; + private static final int RETRY_TIMES = 3; + private static final int SLEEP_TIME = 3000; + + private final ExecutorService pool = new ThreadPoolExecutor(1, 1, 0, TimeUnit.MILLISECONDS, + new ArrayBlockingQueue(1), new NamedThreadFactory("sentinel-zookeeper-ds-update"), + new ThreadPoolExecutor.DiscardOldestPolicy()); - private final IZkDataListener zkDataListener; + private final NodeCacheListener listener; private final String groupId; private final String dataId; - private ZkClient zkClient = null; + private CuratorFramework zkClient = null; + private NodeCache nodeCache = null; public ZookeeperDataSource(final String serverAddr, final String groupId, final String dataId, ConfigParser parser) { @@ -31,23 +49,21 @@ public ZookeeperDataSource(final String serverAddr, final String groupId, final } this.groupId = groupId; this.dataId = dataId; - this.zkDataListener = new IZkDataListener() { + this.listener = new NodeCacheListener() { @Override - public void handleDataChange(String path, Object configInfo) throws Exception { + public void nodeChanged() throws Exception { + String configInfo = null; + ChildData childData = nodeCache.getCurrentData(); + if (null != childData && childData.getData() != null) { + + configInfo = new String(childData.getData()); + } RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s, %s): %s", serverAddr, dataId, groupId, configInfo)); - T newValue = ZookeeperDataSource.this.parser.parse(String.valueOf(configInfo)); + T newValue = ZookeeperDataSource.this.parser.parse(configInfo); // Update the new value to the property. getProperty().updateValue(newValue); } - - @Override - public void handleDataDeleted(String path) throws Exception { - RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s, %s): %s", - serverAddr, dataId, groupId, null)); - // Update the new value to the property. - getProperty().updateValue(null); - } }; initZookeeperListener(serverAddr); loadInitialConfig(); @@ -67,12 +83,17 @@ private void loadInitialConfig() { private void initZookeeperListener(String serverAddr) { try { - this.zkClient = new ZkClient(serverAddr, DEFAULT_TIMEOUT); + this.zkClient = CuratorFrameworkFactory.newClient(serverAddr, new RetryNTimes(RETRY_TIMES, SLEEP_TIME)); + this.zkClient.start(); String path = "/" + this.groupId + "/" + this.dataId; - if (!zkClient.exists(path)) { - zkClient.createPersistent(path, true); + Stat stat = this.zkClient.checkExists().forPath(path); + if (stat == null) { + this.zkClient.create().creatingParentContainersIfNeeded().withMode(CreateMode.PERSISTENT).forPath(path, null); } - zkClient.subscribeDataChanges(path, zkDataListener); + + this.nodeCache = new NodeCache(this.zkClient, path); + this.nodeCache.getListenable().addListener(this.listener, this.pool); + this.nodeCache.start(); } catch (Exception e) { RecordLog.info("[ZookeeperDataSource] Error occurred when initializing Zookeeper data source", e); e.printStackTrace(); @@ -81,19 +102,26 @@ private void initZookeeperListener(String serverAddr) { @Override public String readSource() throws Exception { - if (zkClient == null) { + if (this.zkClient == null) { throw new IllegalStateException("Zookeeper has not been initialized or error occurred"); } String path = "/" + this.groupId + "/" + this.dataId; - return zkClient.readData(path); + byte[] data = this.zkClient.getData().forPath(path); + if (data != null) { + return new String(data); + } + return null; } @Override public void close() throws Exception { - if (zkClient != null) { - String path = "/" + this.groupId + "/" + this.dataId; - zkClient.unsubscribeDataChanges(path, zkDataListener); + if (this.nodeCache != null) { + this.nodeCache.getListenable().removeListener(listener); + this.nodeCache.close(); + } + if (this.zkClient != null) { + this.zkClient.close(); } - zkClient.close(); + pool.shutdown(); } } From e7f6e4cf89c9ba5cd85b2d50cab52c3380ec0bdb Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sat, 4 Aug 2018 17:35:07 +0800 Subject: [PATCH 5/8] =?UTF-8?q?=E5=A2=9E=E5=8A=A0curator-test?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../pom.xml | 24 ++++++++++++++++++ .../zookeeper/ZookeeperConfigSender.java | 21 +++++++++++++--- .../zookeeper/ZookeeperDataSourceDemo.java | 25 ++++++++----------- .../zookeeper/ZookeeperDataSource.java | 6 ++--- 4 files changed, 55 insertions(+), 21 deletions(-) diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml b/sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml index b320ecba59..2915c1d118 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/pom.xml @@ -11,6 +11,12 @@ sentinel-demo-zookeeper-datasource + + 3.4.13 + 4.0.1 + 2.12.0 + + com.alibaba.csp @@ -29,6 +35,24 @@ com.alibaba fastjson + + + org.apache.zookeeper + zookeeper + ${zookeeper.version} + + + + org.apache.curator + curator-test + ${curator-test.version} + + + org.apache.zookeeper + zookeeper + + + diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java index f048de820d..d82c8ac462 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java @@ -2,7 +2,8 @@ import org.apache.curator.framework.CuratorFramework; import org.apache.curator.framework.CuratorFrameworkFactory; -import org.apache.curator.retry.RetryNTimes; +import org.apache.curator.retry.ExponentialBackoffRetry; +import org.apache.curator.test.TestingServer; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.data.Stat; @@ -13,9 +14,14 @@ */ public class ZookeeperConfigSender { + private static final int RETRY_TIMES = 3; + private static final int SLEEP_TIME = 1000; + public static void main(String[] args) throws Exception { - final String remoteAddress = "127.0.0.1:2181"; + TestingServer server = new TestingServer(2181); + + final String remoteAddress = server.getConnectString(); final String groupId = "Sentinel-Demo"; final String dataId = "SYSTEM-CODE-DEMO-FLOW"; final String rule = "[\n" @@ -29,7 +35,7 @@ public static void main(String[] args) throws Exception { + " }\n" + "]"; - CuratorFramework zkClient = CuratorFrameworkFactory.newClient(remoteAddress, new RetryNTimes(3, 5000)); + CuratorFramework zkClient = CuratorFrameworkFactory.newClient(remoteAddress, new ExponentialBackoffRetry(SLEEP_TIME, RETRY_TIMES)); zkClient.start(); String path = "/" + groupId + "/" + dataId; Stat stat = zkClient.checkExists().forPath(path); @@ -38,5 +44,14 @@ public static void main(String[] args) throws Exception { } zkClient.setData().forPath(path, rule.getBytes()); // zkClient.delete().forPath(path); + + try { + Thread.sleep(30000L); + } catch (InterruptedException e) { + e.printStackTrace(); + } + + zkClient.close(); + server.stop(); } } diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java index be5e39f527..391ec5a30c 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java @@ -4,12 +4,8 @@ import com.alibaba.csp.sentinel.datasource.DataSource; import com.alibaba.csp.sentinel.datasource.zookeeper.ZookeeperDataSource; -import com.alibaba.csp.sentinel.slots.block.degrade.DegradeRule; -import com.alibaba.csp.sentinel.slots.block.degrade.DegradeRuleManager; import com.alibaba.csp.sentinel.slots.block.flow.FlowRule; import com.alibaba.csp.sentinel.slots.block.flow.FlowRuleManager; -import com.alibaba.csp.sentinel.slots.system.SystemRule; -import com.alibaba.csp.sentinel.slots.system.SystemRuleManager; import com.alibaba.fastjson.JSON; import com.alibaba.fastjson.TypeReference; @@ -20,31 +16,30 @@ */ public class ZookeeperDataSourceDemo { - private static final String KEY = "TestResource"; - public static void main(String[] args) { loadRules(); } private static void loadRules() { + final String remoteAddress = "127.0.0.1:2181"; final String groupId = "Sentinel-Demo"; final String flowDataId = "SYSTEM-CODE-DEMO-FLOW"; - final String degradeDataId = "SYSTEM-CODE-DEMO-DEGRADE"; - final String systemDataId = "SYSTEM-CODE-DEMO-SYSTEM"; + // final String degradeDataId = "SYSTEM-CODE-DEMO-DEGRADE"; + // final String systemDataId = "SYSTEM-CODE-DEMO-SYSTEM"; DataSource> flowRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, flowDataId, source -> JSON.parseObject(source, new TypeReference>() {})); FlowRuleManager.register2Property(flowRuleDataSource.getProperty()); - DataSource> degradeRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, degradeDataId, - source -> JSON.parseObject(source, new TypeReference>() {})); - DegradeRuleManager.register2Property(degradeRuleDataSource.getProperty()); - - DataSource> systemRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, systemDataId, - source -> JSON.parseObject(source, new TypeReference>() {})); - SystemRuleManager.register2Property(systemRuleDataSource.getProperty()); + // DataSource> degradeRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, degradeDataId, + // source -> JSON.parseObject(source, new TypeReference>() {})); + // DegradeRuleManager.register2Property(degradeRuleDataSource.getProperty()); + // + // DataSource> systemRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, systemDataId, + // source -> JSON.parseObject(source, new TypeReference>() {})); + // SystemRuleManager.register2Property(systemRuleDataSource.getProperty()); } } diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java index 2a8af4fa1f..6803bee10a 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -15,7 +15,7 @@ import org.apache.curator.framework.recipes.cache.ChildData; import org.apache.curator.framework.recipes.cache.NodeCache; import org.apache.curator.framework.recipes.cache.NodeCacheListener; -import org.apache.curator.retry.RetryNTimes; +import org.apache.curator.retry.ExponentialBackoffRetry; import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.data.Stat; @@ -27,7 +27,7 @@ public class ZookeeperDataSource extends AbstractDataSource { private static final int RETRY_TIMES = 3; - private static final int SLEEP_TIME = 3000; + private static final int SLEEP_TIME = 1000; private final ExecutorService pool = new ThreadPoolExecutor(1, 1, 0, TimeUnit.MILLISECONDS, new ArrayBlockingQueue(1), new NamedThreadFactory("sentinel-zookeeper-ds-update"), @@ -83,7 +83,7 @@ private void loadInitialConfig() { private void initZookeeperListener(String serverAddr) { try { - this.zkClient = CuratorFrameworkFactory.newClient(serverAddr, new RetryNTimes(RETRY_TIMES, SLEEP_TIME)); + this.zkClient = CuratorFrameworkFactory.newClient(serverAddr, new ExponentialBackoffRetry(SLEEP_TIME, RETRY_TIMES)); this.zkClient.start(); String path = "/" + this.groupId + "/" + this.dataId; Stat stat = this.zkClient.checkExists().forPath(path); From 2f8fafffb4b4d4f3e58f895cfd97a4bc0b259ba1 Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sun, 5 Aug 2018 11:25:57 +0800 Subject: [PATCH 6/8] =?UTF-8?q?groupId=E5=92=8CdataId=E6=94=AF=E6=8C=81/?= =?UTF-8?q?=E5=BC=80=E5=A4=B4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../zookeeper/ZookeeperConfigSender.java | 20 ++++++++++++++++++- .../zookeeper/ZookeeperDataSourceDemo.java | 3 +++ .../zookeeper/ZookeeperDataSource.java | 19 ++++++++++++++++-- 3 files changed, 39 insertions(+), 3 deletions(-) diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java index d82c8ac462..1b635e61fe 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperConfigSender.java @@ -19,6 +19,7 @@ public class ZookeeperConfigSender { public static void main(String[] args) throws Exception { + // 启动Zookeeper服务 TestingServer server = new TestingServer(2181); final String remoteAddress = server.getConnectString(); @@ -37,7 +38,7 @@ public static void main(String[] args) throws Exception { CuratorFramework zkClient = CuratorFrameworkFactory.newClient(remoteAddress, new ExponentialBackoffRetry(SLEEP_TIME, RETRY_TIMES)); zkClient.start(); - String path = "/" + groupId + "/" + dataId; + String path = getPath(groupId, dataId); Stat stat = zkClient.checkExists().forPath(path); if (stat == null) { zkClient.create().creatingParentContainersIfNeeded().withMode(CreateMode.PERSISTENT).forPath(path, null); @@ -52,6 +53,23 @@ public static void main(String[] args) throws Exception { } zkClient.close(); + + //停止zookeeper服务 server.stop(); } + + private static String getPath(String groupId, String dataId) { + String path = ""; + if (groupId.startsWith("/")) { + path += groupId; + } else { + path += "/" + groupId; + } + if (dataId.startsWith("/")) { + path += dataId; + } else { + path += "/" + dataId; + } + return path; + } } diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java index 391ec5a30c..3659785050 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java @@ -29,6 +29,9 @@ private static void loadRules() { // final String systemDataId = "SYSTEM-CODE-DEMO-SYSTEM"; + // 规则会持久化到zk的/groupId/flowDataId节点 + // groupId和和flowDataId可以用/开头也可以不用 + // 建议不用以/开头,目的是为了如果从Zookeeper切换到nacos的话,只需要改个数据源类名就可以 DataSource> flowRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, flowDataId, source -> JSON.parseObject(source, new TypeReference>() {})); FlowRuleManager.register2Property(flowRuleDataSource.getProperty()); diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java index 6803bee10a..b49ac82e38 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -85,7 +85,7 @@ private void initZookeeperListener(String serverAddr) { try { this.zkClient = CuratorFrameworkFactory.newClient(serverAddr, new ExponentialBackoffRetry(SLEEP_TIME, RETRY_TIMES)); this.zkClient.start(); - String path = "/" + this.groupId + "/" + this.dataId; + String path = getPath(this.groupId, this.dataId); Stat stat = this.zkClient.checkExists().forPath(path); if (stat == null) { this.zkClient.create().creatingParentContainersIfNeeded().withMode(CreateMode.PERSISTENT).forPath(path, null); @@ -105,7 +105,7 @@ public String readSource() throws Exception { if (this.zkClient == null) { throw new IllegalStateException("Zookeeper has not been initialized or error occurred"); } - String path = "/" + this.groupId + "/" + this.dataId; + String path = getPath(this.groupId, this.dataId); byte[] data = this.zkClient.getData().forPath(path); if (data != null) { return new String(data); @@ -124,4 +124,19 @@ public void close() throws Exception { } pool.shutdown(); } + + private String getPath(String groupId, String dataId) { + String path = ""; + if (groupId.startsWith("/")) { + path += groupId; + } else { + path += "/" + groupId; + } + if (dataId.startsWith("/")) { + path += dataId; + } else { + path += "/" + dataId; + } + return path; + } } From ab4933a6073cd9c0217779f1a2fdac525de14631 Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sun, 5 Aug 2018 12:41:07 +0800 Subject: [PATCH 7/8] =?UTF-8?q?=E6=94=AF=E6=8C=81zookeeper=E9=A3=8E?= =?UTF-8?q?=E6=A0=BC=E7=9A=84path?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../zookeeper/ZookeeperDataSourceDemo.java | 19 +++++++++- .../zookeeper/ZookeeperDataSource.java | 38 +++++++++++++------ 2 files changed, 45 insertions(+), 12 deletions(-) diff --git a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java index 3659785050..2077b7d834 100644 --- a/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java +++ b/sentinel-demo/sentinel-demo-zookeeper-datasource/src/main/java/com/alibaba/csp/sentinel/demo/datasource/zookeeper/ZookeeperDataSourceDemo.java @@ -17,12 +17,29 @@ public class ZookeeperDataSourceDemo { public static void main(String[] args) { + // 使用zookeeper的场景 loadRules(); + + // 方便扩展的场景 + //loadRules2(); } private static void loadRules() { final String remoteAddress = "127.0.0.1:2181"; + final String path = "/Sentinel-Demo/SYSTEM-CODE-DEMO-FLOW"; + + DataSource> flowRuleDataSource = new ZookeeperDataSource<>(remoteAddress, path, + source -> JSON.parseObject(source, new TypeReference>() {})); + FlowRuleManager.register2Property(flowRuleDataSource.getProperty()); + + + } + + private static void loadRules2() { + + final String remoteAddress = "127.0.0.1:2181"; + // 引入groupId和dataId的概念,是为了方便和Nacos进行切换 final String groupId = "Sentinel-Demo"; final String flowDataId = "SYSTEM-CODE-DEMO-FLOW"; // final String degradeDataId = "SYSTEM-CODE-DEMO-DEGRADE"; @@ -31,7 +48,7 @@ private static void loadRules() { // 规则会持久化到zk的/groupId/flowDataId节点 // groupId和和flowDataId可以用/开头也可以不用 - // 建议不用以/开头,目的是为了如果从Zookeeper切换到nacos的话,只需要改个数据源类名就可以 + // 建议不用以/开头,目的是为了如果从Zookeeper切换到Nacos的话,只需要改数据源类名就可以 DataSource> flowRuleDataSource = new ZookeeperDataSource<>(remoteAddress, groupId, flowDataId, source -> JSON.parseObject(source, new TypeReference>() {})); FlowRuleManager.register2Property(flowRuleDataSource.getProperty()); diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java index b49ac82e38..0683dff137 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -33,9 +33,10 @@ public class ZookeeperDataSource extends AbstractDataSource { new ArrayBlockingQueue(1), new NamedThreadFactory("sentinel-zookeeper-ds-update"), new ThreadPoolExecutor.DiscardOldestPolicy()); - private final NodeCacheListener listener; - private final String groupId; - private final String dataId; + private NodeCacheListener listener; + private String groupId; + private String dataId; + private String path; private CuratorFramework zkClient = null; private NodeCache nodeCache = null; @@ -49,6 +50,23 @@ public ZookeeperDataSource(final String serverAddr, final String groupId, final } this.groupId = groupId; this.dataId = dataId; + this.path = getPath(groupId, dataId); + + init(serverAddr); + } + + public ZookeeperDataSource(final String serverAddr, final String path, ConfigParser parser) { + super(parser); + if (StringUtil.isBlank(serverAddr) || StringUtil.isBlank(path)) { + throw new IllegalArgumentException(String.format("Bad argument: serverAddr=[%s], path=[%s]", + serverAddr, path)); + } + this.path = path; + + init(serverAddr); + } + + private void init(final String serverAddr) { this.listener = new NodeCacheListener() { @Override public void nodeChanged() throws Exception { @@ -58,8 +76,8 @@ public void nodeChanged() throws Exception { configInfo = new String(childData.getData()); } - RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s, %s): %s", - serverAddr, dataId, groupId, configInfo)); + RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s): %s", + serverAddr, path, configInfo)); T newValue = ZookeeperDataSource.this.parser.parse(configInfo); // Update the new value to the property. getProperty().updateValue(newValue); @@ -85,13 +103,12 @@ private void initZookeeperListener(String serverAddr) { try { this.zkClient = CuratorFrameworkFactory.newClient(serverAddr, new ExponentialBackoffRetry(SLEEP_TIME, RETRY_TIMES)); this.zkClient.start(); - String path = getPath(this.groupId, this.dataId); - Stat stat = this.zkClient.checkExists().forPath(path); + Stat stat = this.zkClient.checkExists().forPath(this.path); if (stat == null) { - this.zkClient.create().creatingParentContainersIfNeeded().withMode(CreateMode.PERSISTENT).forPath(path, null); + this.zkClient.create().creatingParentContainersIfNeeded().withMode(CreateMode.PERSISTENT).forPath(this.path, null); } - this.nodeCache = new NodeCache(this.zkClient, path); + this.nodeCache = new NodeCache(this.zkClient, this.path); this.nodeCache.getListenable().addListener(this.listener, this.pool); this.nodeCache.start(); } catch (Exception e) { @@ -105,8 +122,7 @@ public String readSource() throws Exception { if (this.zkClient == null) { throw new IllegalStateException("Zookeeper has not been initialized or error occurred"); } - String path = getPath(this.groupId, this.dataId); - byte[] data = this.zkClient.getData().forPath(path); + byte[] data = this.zkClient.getData().forPath(this.path); if (data != null) { return new String(data); } From 1a8bd285fb32828e70d5660759642dcd3aaa69d6 Mon Sep 17 00:00:00 2001 From: guonanjun Date: Sun, 5 Aug 2018 12:49:50 +0800 Subject: [PATCH 8/8] =?UTF-8?q?=E6=94=AF=E6=8C=81zookeeper=E9=A3=8E?= =?UTF-8?q?=E6=A0=BC=E7=9A=84path?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../zookeeper/ZookeeperDataSource.java | 44 ++++++++++--------- 1 file changed, 24 insertions(+), 20 deletions(-) diff --git a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java index 0683dff137..b8c2663d3f 100644 --- a/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java +++ b/sentinel-extension/sentinel-datasource-zookeeper/src/main/java/com/alibaba/csp/sentinel/datasource/zookeeper/ZookeeperDataSource.java @@ -34,9 +34,9 @@ public class ZookeeperDataSource extends AbstractDataSource { new ThreadPoolExecutor.DiscardOldestPolicy()); private NodeCacheListener listener; - private String groupId; - private String dataId; - private String path; + private final String groupId; + private final String dataId; + private final String path; private CuratorFramework zkClient = null; private NodeCache nodeCache = null; @@ -62,27 +62,13 @@ public ZookeeperDataSource(final String serverAddr, final String path, ConfigPar serverAddr, path)); } this.path = path; + this.groupId = null; + this.dataId = null; init(serverAddr); } private void init(final String serverAddr) { - this.listener = new NodeCacheListener() { - @Override - public void nodeChanged() throws Exception { - String configInfo = null; - ChildData childData = nodeCache.getCurrentData(); - if (null != childData && childData.getData() != null) { - - configInfo = new String(childData.getData()); - } - RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s): %s", - serverAddr, path, configInfo)); - T newValue = ZookeeperDataSource.this.parser.parse(configInfo); - // Update the new value to the property. - getProperty().updateValue(newValue); - } - }; initZookeeperListener(serverAddr); loadInitialConfig(); } @@ -99,8 +85,26 @@ private void loadInitialConfig() { } } - private void initZookeeperListener(String serverAddr) { + private void initZookeeperListener(final String serverAddr) { try { + + this.listener = new NodeCacheListener() { + @Override + public void nodeChanged() throws Exception { + String configInfo = null; + ChildData childData = nodeCache.getCurrentData(); + if (null != childData && childData.getData() != null) { + + configInfo = new String(childData.getData()); + } + RecordLog.info(String.format("[ZookeeperDataSource] New property value received for (%s, %s): %s", + serverAddr, path, configInfo)); + T newValue = ZookeeperDataSource.this.parser.parse(configInfo); + // Update the new value to the property. + getProperty().updateValue(newValue); + } + }; + this.zkClient = CuratorFrameworkFactory.newClient(serverAddr, new ExponentialBackoffRetry(SLEEP_TIME, RETRY_TIMES)); this.zkClient.start(); Stat stat = this.zkClient.checkExists().forPath(this.path);