diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterNodeExistException.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterNodeExistException.java new file mode 100644 index 000000000..d449920ea --- /dev/null +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterNodeExistException.java @@ -0,0 +1,17 @@ +package org.skywalking.apm.collector.cluster; + +import org.skywalking.apm.collector.core.client.ClientException; + +/** + * @author pengys5 + */ +public class ClusterNodeExistException extends ClientException { + + public ClusterNodeExistException(String message) { + super(message); + } + + public ClusterNodeExistException(String message, Throwable cause) { + super(message, cause); + } +} diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java index 7e973716a..a2e6c862f 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java @@ -10,6 +10,7 @@ import org.apache.zookeeper.ZooDefs; import org.skywalking.apm.collector.client.zookeeper.ZookeeperClient; import org.skywalking.apm.collector.client.zookeeper.ZookeeperClientException; import org.skywalking.apm.collector.client.zookeeper.util.PathUtils; +import org.skywalking.apm.collector.cluster.ClusterNodeExistException; import org.skywalking.apm.collector.core.client.Client; import org.skywalking.apm.collector.core.client.ClientException; import org.skywalking.apm.collector.core.client.DataMonitor; @@ -74,7 +75,11 @@ public class ClusterZKDataMonitor implements DataMonitor, Watcher { String serverPath = path + "/" + value.getHostPort(); listener.addAddress(value.getHostPort() + contextPath); - setData(serverPath, contextPath); + if (client.exists(serverPath, false) == null) { + setData(serverPath, contextPath); + } else { + throw new ClusterNodeExistException("current address: " + value.getHostPort() + " has been registered, check the host and port configuration or wait a moment."); + } } @Override public ClusterDataListener getListener(String path) {