diff --git a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java index 6b30a1a5e2..075c573415 100644 --- a/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java +++ b/dolphinscheduler-server/src/main/java/org/apache/dolphinscheduler/server/master/registry/ServerNodeManager.java @@ -189,26 +189,30 @@ public class ServerNodeManager implements InitializingBean { @Override public void run() { - // sync worker node info - Map newWorkerNodeInfo = registryClient.getServerMaps(NodeType.WORKER, true); - syncAllWorkerNodeInfo(newWorkerNodeInfo); + try { + // sync worker node info + Map newWorkerNodeInfo = registryClient.getServerMaps(NodeType.WORKER, true); + syncAllWorkerNodeInfo(newWorkerNodeInfo); - // sync worker group nodes from database - List workerGroupList = workerGroupMapper.queryAllWorkerGroup(); - if (CollectionUtils.isNotEmpty(workerGroupList)) { - for (WorkerGroup wg : workerGroupList) { - String workerGroup = wg.getName(); - Set nodes = new HashSet<>(); - String[] addrs = wg.getAddrList().split(Constants.COMMA); - for (String addr : addrs) { - if (newWorkerNodeInfo.containsKey(addr)) { - nodes.add(addr); + // sync worker group nodes from database + List workerGroupList = workerGroupMapper.queryAllWorkerGroup(); + if (CollectionUtils.isNotEmpty(workerGroupList)) { + for (WorkerGroup wg : workerGroupList) { + String workerGroup = wg.getName(); + Set nodes = new HashSet<>(); + String[] addrs = wg.getAddrList().split(Constants.COMMA); + for (String addr : addrs) { + if (newWorkerNodeInfo.containsKey(addr)) { + nodes.add(addr); + } + } + if (!nodes.isEmpty()) { + syncWorkerGroupNodes(workerGroup, nodes); } } - if (!nodes.isEmpty()) { - syncWorkerGroupNodes(workerGroup, nodes); - } } + } catch (Exception e) { + logger.error("WorkerNodeInfoAndGroupDbSyncTask error:", e); } } } @@ -359,7 +363,6 @@ public class ServerNodeManager implements InitializingBean { private void syncWorkerGroupNodes(String workerGroup, Collection nodes) { workerGroupLock.lock(); try { - workerGroup = workerGroup.toLowerCase(); Set workerNodes = workerGroupNodes.getOrDefault(workerGroup, new HashSet<>()); workerNodes.clear(); workerNodes.addAll(nodes); @@ -385,7 +388,6 @@ public class ServerNodeManager implements InitializingBean { if (StringUtils.isEmpty(workerGroup)) { workerGroup = Constants.DEFAULT_WORKER_GROUP; } - workerGroup = workerGroup.toLowerCase(); Set nodes = workerGroupNodes.get(workerGroup); if (CollectionUtils.isNotEmpty(nodes)) { // avoid ConcurrentModificationException