From 2c1fb7913ae17f0caca66651d10c5ae08d91fcc2 Mon Sep 17 00:00:00 2001 From: wusheng Date: Fri, 11 Dec 2015 21:14:59 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BF=AE=E6=95=B4UserInfoCoordinator?= =?UTF-8?q?=E4=BB=A3=E7=A0=81=EF=BC=8C=E7=A1=AE=E4=BF=9D=E9=80=BB=E8=BE=91?= =?UTF-8?q?=E7=9A=84=E6=9C=89=E6=95=88=E6=80=A7=E5=92=8C=E5=81=A5=E5=A3=AE?= =?UTF-8?q?=E6=80=A7=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../alarm/AlarmMessageProcessThread.java | 1 - .../skywalking/alarm/UserInfoCoordinator.java | 375 +++++++++--------- .../skywalking/alarm/UserInfoInspector.java | 37 -- .../cloud/skywalking/alarm/conf/Config.java | 3 +- 4 files changed, 190 insertions(+), 226 deletions(-) delete mode 100644 skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoInspector.java diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmMessageProcessThread.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmMessageProcessThread.java index 0b1b32375..9c45e3d44 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmMessageProcessThread.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmMessageProcessThread.java @@ -50,7 +50,6 @@ public class AlarmMessageProcessThread extends Thread { for (Map.Entry> entry : cacheRules.entrySet()) { for (AlarmRule rule : entry.getValue()) { processor.process(entry.getKey(), rule); -// System.out.println(currentThread().getName() + " @~ " + entry.getKey().getUserId() + " @~ " + rule.getRuleId()); } } } diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoCoordinator.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoCoordinator.java index accd8f36c..f3858bbd2 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoCoordinator.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoCoordinator.java @@ -22,219 +22,222 @@ import java.util.concurrent.TimeUnit; public class UserInfoCoordinator extends Thread { + private Logger logger = LogManager.getLogger(UserInfoCoordinator.class); - private Logger logger = LogManager.getLogger(UserInfoInspector.class); + private boolean redistributing; + private RegisterServerWatcher watcher = new RegisterServerWatcher(); + private InterProcessMutex lock = new InterProcessMutex( + ZKUtil.getZkClient(), Config.ZKPath.COORDINATOR_PATH); + private boolean isCoordinator = false; - private boolean redistributing; - private boolean newServerComingFlag = false; - private RegisterServerWatcher watcher = new RegisterServerWatcher(); - private InterProcessMutex lock = new InterProcessMutex(ZKUtil.getZkClient(), Config.ZKPath.COORDINATOR_PATH); - private boolean coordinatorFlag; + public UserInfoCoordinator() { + } - public UserInfoCoordinator() { - redistributing = false; - try { - coordinatorFlag = lock.acquire(Config.Coordinator.RETRY_GET_COORDINATOR_LOCK_INTERVAL - , TimeUnit.SECONDS); - } catch (Exception e) { - logger.error("Failed to", e); - coordinatorFlag = false; - } + @Override + public void run() { + while (true) { + try { + if (!isCoordinator) { + while (!retryBecomeCoordinator()) { + try { + Thread.sleep(Config.Coordinator.RETRY_BECOME_COORDINATOR_WAIT_TIME); + } catch (Exception e) { + logger.error("Sleep Failed.", e); + } + } + + isCoordinator = true; + watcherRegisterServerPath(); + redistributing = true; + } - if (coordinatorFlag) { - watcherRegisterServerPath(); - } - } + // 检查是否有新服务注册或者在重分配过程做有新处理线程启动了 + if (!redistributing) { + try { + Thread.sleep(Config.Coordinator.CHECK_REDISTRIBUTE_INTERVAL); + } catch (InterruptedException e) { + logger.error("Sleep error", e); + } - @Override - public void run() { + continue; + } - if (!coordinatorFlag) { - // - while (!retryBecomeCoordinator()) { - try { - Thread.sleep(Config.Coordinator.RETRY_BECOME_COORDINATOR_WAIT_TIME); - } catch (Exception e) { - logger.error("Sleep Failed.", e); - } - } - // 新官上任三把火,重新分配 - watcherRegisterServerPath(); - redistributing = true; - } + redistributing = false; - while (true) { - try { - //检查是否有新服务注册或者在重分配过程做有新处理线程启动了 - if (!redistributing && !newServerComingFlag) { - try { - Thread.sleep(Config.Coordinator.CHECK_REDISTRIBUTE_INTERVAL); - } catch (InterruptedException e) { - logger.error("Sleep error", e); - } + // 获取当前所有的注册的处理线程 + List registeredThreads = acquireAllRegisteredThread(); + // 修改状态 (开始重新分配状态) + changeStatus(registeredThreads, + ProcessThreadStatus.REDISTRIBUTING); + // 检查所有的服务是否都处于空闲状态 + int retryTimes = 0; + while (!checkAllProcessStatus(registeredThreads, + ProcessThreadStatus.FREE)) { + try { + Thread.sleep(Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL); + retryTimes++; + } catch (InterruptedException e) { + logger.error("Sleep failed", e); + } + + if(retryTimes > 1000){ + logger.warn("checking all processors are free, waiting {}ms", Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL * retryTimes); + retryTimes = 0; + } + } - continue; - } + // 查询当前有多少用户 + List users = AlarmMessageDao.selectAllUserIds(); - newServerComingFlag = false; + // 将用户重新分配给服务 + List realRedistributeThread = allocationUser( + registeredThreads, users); - //获取当前所有的注册的处理线程 - List registeredThreads = acquireAllRegisteredThread(); - //修改状态 (开始重新分配状态) - changeStatus(registeredThreads, ProcessThreadStatus.REDISTRIBUTING); - //检查所有的服务是否都处于空闲状态 - while (!checkAllProcessStatus(registeredThreads, ProcessThreadStatus.FREE)) { - try { - Thread.sleep(Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL); - } catch (InterruptedException e) { - logger.error("Sleep failed", e); - } - } + // 修改状态(分配完成) + changeStatus(realRedistributeThread, + ProcessThreadStatus.REDISTRIBUTE_SUCCESS); - //查询当前有多少用户 - List users = AlarmMessageDao.selectAllUserIds(); + // 检查所有的服务是否都处于忙碌状态 + while (!checkAllProcessStatus(realRedistributeThread, + ProcessThreadStatus.BUSY)) { + try { + Thread.sleep(Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL); + } catch (InterruptedException e) { + logger.error("Sleep failed", e); + } + + if(retryTimes > 1000){ + logger.warn("checking all processors are busy, waiting {}ms", Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL * retryTimes); + retryTimes = 0; + } + } - //将用户重新分配给服务 - List realRedistributeThread = allocationUser(registeredThreads, users); + } catch (Exception e) { + logger.error("Failed to coordinate, retry. ", e); + releaseCoordinator(); + isCoordinator = false; + } + } + } - //修改状态(分配完成) - changeStatus(realRedistributeThread, ProcessThreadStatus.REDISTRIBUTE_SUCCESS); + private boolean retryBecomeCoordinator() { + try { + return lock.acquire( + Config.Coordinator.RETRY_GET_COORDINATOR_LOCK_INTERVAL, + TimeUnit.SECONDS); + } catch (Exception e) { + logger.error("Failed to acquire lock .", e); + return false; + } + } - //检查所有的服务是否都处于忙碌状态 - while (!checkAllProcessStatus(realRedistributeThread, ProcessThreadStatus.BUSY)) { - try { - Thread.sleep(Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL); - } catch (InterruptedException e) { - logger.error("Sleep failed", e); - } - } + private void releaseCoordinator() { + if (lock != null && lock.isAcquiredInThisProcess()) { + try { + lock.release(); + } catch (Exception e1) { + logger.error("Failed to release lock.", e1); + } + } + } - redistributing = false; + private List allocationUser(List registeredThreads, + List userIds) { + List realRedistributeThread = new ArrayList(); + Set sortThreadIds = new HashSet(registeredThreads); + int step = (int) Math.ceil(userIds.size() * 1.0 / sortThreadIds.size()); + int start = 0; + int end = step; - } catch (Exception e) { - logger.error("Failed to redistribute User ", e); - } - } - } + if (end > userIds.size()) { + end = userIds.size(); + } - private void watcherRegisterServerPath() { - try { - ZKUtil.getChildrenWithWatcher(Config.ZKPath.REGISTER_SERVER_PATH, watcher); - } catch (Exception e) { - logger.error("Failed to set watcher for get children", e); - if (lock != null && lock.isAcquiredInThisProcess()) { - try { - lock.release(); - } catch (Exception e1) { - logger.error("Failed to release lock.", e1); - } - } - } - } + for (String thread : sortThreadIds) { + if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/" + + thread)) + continue; + String value = ZKUtil + .getPathData(Config.ZKPath.REGISTER_SERVER_PATH + "/" + + thread); + ProcessThreadValue value1 = new Gson().fromJson(value, + ProcessThreadValue.class); + value1.setDealUserIds(userIds.subList(start, end)); + ZKUtil.setPathData(Config.ZKPath.REGISTER_SERVER_PATH + "/" + + thread, new Gson().toJson(value1)); + // 实际重新分配的线程Id + realRedistributeThread.add(thread); - private boolean retryBecomeCoordinator() { - try { - return lock.acquire(5, TimeUnit.SECONDS); - } catch (Exception e) { - logger.error("Failed to acquire lock .", e); - return false; - } - } + start = end; + end += step; + if (start >= userIds.size()) { + break; + } + if (end > userIds.size()) { + end = userIds.size(); + } - private List allocationUser(List registeredThreads, List userIds) { - List realRedistributeThread = new ArrayList(); - Set sortThreadIds = new HashSet(registeredThreads); - int step = (int) Math.ceil(userIds.size() * 1.0 / sortThreadIds.size()); - int start = 0; - int end = step; + } + return realRedistributeThread; + } - if (end > userIds.size()) { - end = userIds.size(); - } + private List acquireAllRegisteredThread() { + return ZKUtil.getChildren(Config.ZKPath.REGISTER_SERVER_PATH); + } - for (String thread : sortThreadIds) { - if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/" + thread)) - continue; - String value = ZKUtil.getPathData(Config.ZKPath.REGISTER_SERVER_PATH + "/" + thread); - ProcessThreadValue value1 = new Gson().fromJson(value, ProcessThreadValue.class); - value1.setDealUserIds(userIds.subList(start, end)); - ZKUtil.setPathData(Config.ZKPath.REGISTER_SERVER_PATH + "/" + thread, new Gson().toJson(value1)); - // 实际重新分配的线程Id - realRedistributeThread.add(thread); + private boolean checkAllProcessStatus(List registeredThreadIds, + ProcessThreadStatus status) { + String registerPathPrefix = Config.ZKPath.REGISTER_SERVER_PATH + "/"; + for (String threadId : registeredThreadIds) { - start = end; - end += step; - if (start >= userIds.size()) { - break; - } - if (end > userIds.size()) { - end = userIds.size(); - } + if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/" + + threadId)) + continue; - } - return realRedistributeThread; - } + if (getProcessThreadStatus(registerPathPrefix, threadId) != status) { + return false; + } + } + return true; + } - private List acquireAllRegisteredThread() { - return ZKUtil.getChildren(Config.ZKPath.REGISTER_SERVER_PATH); - } + private ProcessThreadStatus getProcessThreadStatus( + String registerPathPrefix, String threadId) { + if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId)) + return ProcessThreadStatus.FREE; + String value = ZKUtil.getPathData(registerPathPrefix + threadId); + if (value == null || value.length() == 0) + return ProcessThreadStatus.FREE; + ProcessThreadValue value1 = new Gson().fromJson(value, + ProcessThreadValue.class); + return ProcessThreadStatus.convert(value1.getStatus()); + } - private boolean checkAllProcessStatus(List registeredThreadIds, ProcessThreadStatus status) { - String registerPathPrefix = Config.ZKPath.REGISTER_SERVER_PATH + "/"; - for (String threadId : registeredThreadIds) { + private void changeStatus(List registeredThreadIds, + ProcessThreadStatus status) { + for (String threadId : registeredThreadIds) { + ProcessUtil.changeProcessThreadStatus(threadId, status); + } + } - if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId)) - continue; + public class RegisterServerWatcher implements CuratorWatcher { - if (getProcessThreadStatus(registerPathPrefix, threadId) - != status) { - return false; - } - } - return true; - } + @Override + public void process(WatchedEvent watchedEvent) { + if (watchedEvent.getType() == Watcher.Event.EventType.NodeChildrenChanged) { + redistributing = true; + } - private ProcessThreadStatus getProcessThreadStatus(String registerPathPrefix, String threadId) { - if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId)) - return ProcessThreadStatus.FREE; - String value = ZKUtil.getPathData(registerPathPrefix + threadId); - if (value == null || value.length() == 0) - return ProcessThreadStatus.FREE; - ProcessThreadValue value1 = new Gson().fromJson(value, ProcessThreadValue.class); - return ProcessThreadStatus.convert(value1.getStatus()); - } + watcherRegisterServerPath(); + } + } - - private void changeStatus(List registeredThreadIds, ProcessThreadStatus status) { - for (String threadId : registeredThreadIds) { - ProcessUtil.changeProcessThreadStatus(threadId, status); - } - } - - public class RegisterServerWatcher implements CuratorWatcher { - - @Override - public void process(WatchedEvent watchedEvent) { - if (watchedEvent.getType() == Watcher.Event.EventType.NodeChildrenChanged) { - if (redistributing) { - redistributing = false; - newServerComingFlag = true; - } else { - redistributing = true; - } - } - - try { - ZKUtil.getChildrenWithWatcher(Config.ZKPath.REGISTER_SERVER_PATH, watcher); - } catch (Exception e) { - if (lock != null && lock.isAcquiredInThisProcess()) { - try { - lock.release(); - } catch (Exception e1) { - logger.error("Failed to release master locker.", e1); - } - } - } - } - } + private void watcherRegisterServerPath() { + try { + ZKUtil.getChildrenWithWatcher(Config.ZKPath.REGISTER_SERVER_PATH, + watcher); + } catch (Exception e) { + logger.error("Failed to set watcher for get children", e); + } + } } diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoInspector.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoInspector.java deleted file mode 100644 index 955b58d0a..000000000 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoInspector.java +++ /dev/null @@ -1,37 +0,0 @@ -package com.ai.cloud.skywalking.alarm; - -import com.ai.cloud.skywalking.alarm.dao.AlarmMessageDao; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - -public class UserInfoInspector extends Thread { - - private Logger logger = LogManager.getLogger(UserInfoInspector.class); - - private int preUserSize; - - public UserInfoInspector() { - preUserSize = AlarmMessageDao.selectUserCount(); - } - - @Override - public void run() { - int currentUserSize; - while (true) { - try { - Thread.sleep(10 * 1000L); - } catch (InterruptedException e) { - logger.error("Sleep Failed", e); - } - - currentUserSize = AlarmMessageDao.selectUserCount(); - - if (currentUserSize != preUserSize) { - logger.info("Total user has been changed. Notice all process thread to change process date."); - for (AlarmMessageProcessThread thread : AlarmProcessServer.getProcessThreads()) { - //thread.setChanged(true); - } - } - } - } -} diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/Config.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/Config.java index 743182a3a..68ae891a2 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/Config.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/Config.java @@ -12,7 +12,6 @@ public class Config { public static class ProcessThread { public static long THREAD_WAIT_INTERVAL = 60 * 1000L; - //public static long THREAD_WAIT_INTERVAL = 1 * 1000L; } public static class ZKPath { @@ -42,7 +41,7 @@ public class Config { // 单位:(毫秒) public static long CHECK_REDISTRIBUTE_INTERVAL = 5 * 1000; // 单位:(毫秒) - public static long CHECK_ALL_PROCESS_THREAD_INTERVAL = 100L; + public static long CHECK_ALL_PROCESS_THREAD_INTERVAL = 500L; } public static class DB {