Browse Source

Merge remote-tracking branch 'upstream/edqs' into msa-edqs

pull/12701/head
dashevchenko 2 years ago
parent
commit
e360f56721
  1. 7
      common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java

7
common/queue/src/main/java/org/thingsboard/server/queue/discovery/ZkDiscoveryService.java

@ -143,6 +143,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
} else { } else {
log.info("Received application ready event. Starting current ZK node."); log.info("Received application ready event. Starting current ZK node.");
} }
subscribeToEvents();
if (client.getState() != CuratorFrameworkState.STARTED) { if (client.getState() != CuratorFrameworkState.STARTED) {
log.debug("Ignoring application ready event, ZK client is not started, ZK client state [{}]", client.getState()); log.debug("Ignoring application ready event, ZK client is not started, ZK client state [{}]", client.getState());
return; return;
@ -212,6 +213,7 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
try { try {
destroyZkClient(); destroyZkClient();
initZkClient(); initZkClient();
subscribeToEvents();
publishCurrentServer(); publishCurrentServer();
} catch (Exception e) { } catch (Exception e) {
log.error("Failed to reconnect to ZK: {}", e.getMessage(), e); log.error("Failed to reconnect to ZK: {}", e.getMessage(), e);
@ -227,7 +229,6 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
client.start(); client.start();
client.blockUntilConnected(); client.blockUntilConnected();
cache = new PathChildrenCache(client, zkNodesDir, true); cache = new PathChildrenCache(client, zkNodesDir, true);
cache.getListenable().addListener(this);
cache.start(); cache.start();
stopped = false; stopped = false;
log.info("ZK client connected"); log.info("ZK client connected");
@ -239,6 +240,10 @@ public class ZkDiscoveryService implements DiscoveryService, PathChildrenCacheLi
} }
} }
private void subscribeToEvents() {
cache.getListenable().addListener(this);
}
private void unpublishCurrentServer() { private void unpublishCurrentServer() {
try { try {
if (nodePath != null) { if (nodePath != null) {

Loading…
Cancel
Save