From 4c806457e6259acbba9df88157f6b809b9c239f5 Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Fri, 11 Dec 2015 20:04:18 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=8C=E6=88=90=E5=8D=8F=E8=B0=83=E5=99=A8?= =?UTF-8?q?=E5=8A=9F=E8=83=BD=E3=80=82=20=E5=8D=8F=E8=B0=83=E5=99=A8?= =?UTF-8?q?=E8=B0=83=E7=94=A8=E6=B5=81=E7=A8=8B=E5=A6=82=E4=B8=8B=EF=BC=9A?= =?UTF-8?q?=20=E5=A4=84=E7=90=86=E7=BA=BF=E7=A8=8B=EF=BC=9A=20=E5=88=9D?= =?UTF-8?q?=E5=A7=8B=E5=8C=96=EF=BC=9A=20=E8=BF=90=E8=A1=8C=EF=BC=9A=20=20?= =?UTF-8?q?=20=20=20=E6=B3=A8=E5=86=8C=E6=9C=8D=E5=8A=A1(=E9=BB=98?= =?UTF-8?q?=E8=AE=A4=E4=B8=BA=E7=A9=BA=E9=97=B2=E7=8A=B6=E6=80=81)=20=20?= =?UTF-8?q?=20=20=20=E6=A3=80=E6=9F=A5=E6=98=AF=E5=90=A6=E4=B8=BA=E5=BF=99?= =?UTF-8?q?=E7=A2=8C=E7=8A=B6=E6=80=81=20=20=20=20=20=E5=A4=84=E7=90=86?= =?UTF-8?q?=E5=91=8A=E8=AD=A6=E4=BF=A1=E6=81=AF=20=20=20=20=20=E6=A3=80?= =?UTF-8?q?=E6=9F=A5=E6=98=AF=E5=90=A6=E5=88=86=E9=85=8D=E7=BA=BF=E7=A8=8B?= =?UTF-8?q?=E7=9A=84=E7=8A=B6=E6=80=81(=E9=87=8D=E6=96=B0=E5=88=86?= =?UTF-8?q?=E9=85=8D=E7=8A=B6=E6=80=81)=20=20=20=20=20=E9=87=8A=E6=94=BE?= =?UTF-8?q?=E7=94=A8=E6=88=B7=E9=94=81=20=20=20=20=20=E4=BF=AE=E6=94=B9?= =?UTF-8?q?=E8=87=AA=E8=BA=AB=E7=8A=B6=E6=80=81=EF=BC=9A(=E7=A9=BA?= =?UTF-8?q?=E9=97=B2=E7=8A=B6=E6=80=81)=20=20=20=20=20=E6=A3=80=E6=9F=A5?= =?UTF-8?q?=E5=88=86=E9=85=8D=E7=BA=BF=E7=A8=8B=E7=9A=84=E7=8A=B6=E6=80=81?= =?UTF-8?q?(=E5=88=86=E9=85=8D=E5=AE=8C=E6=88=90=E7=8A=B6=E6=80=81)=20=20?= =?UTF-8?q?=20=20=20=E9=87=8D=E6=96=B0=E8=8E=B7=E5=8F=96=E5=BE=85=E5=A4=84?= =?UTF-8?q?=E7=90=86=E7=9A=84=E7=94=A8=E6=88=B7=20=20=20=20=20=E7=BB=99?= =?UTF-8?q?=E7=94=A8=E6=88=B7=E5=8A=A0=E9=94=81=20=20=20=20=20=E4=BF=AE?= =?UTF-8?q?=E6=94=B9=E8=87=AA=E8=BA=AB=E7=8A=B6=E6=80=81=20=EF=BC=9A(?= =?UTF-8?q?=E5=BF=99=E7=A2=8C=E7=8A=B6=E6=80=81)?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 协调线程: 初始化: 初始化状态:未分配 运行: 检查是否有新服务注册 修改状态 (开始重新分配状态) 检查所有的服务是否都未待分配状态, 将用户重新分配给服务 修改状态(分配完成) --- skywalking-alarm/pom.xml | 14 +- .../alarm/AlarmMessageProcessThread.java | 158 ++++++++-------- .../skywalking/alarm/AlarmProcessServer.java | 9 +- .../skywalking/alarm/UserInfoCoordinator.java | 173 ++++++++++++++---- .../cloud/skywalking/alarm/conf/Config.java | 20 +- .../skywalking/alarm/dao/AlarmMessageDao.java | 3 +- .../alarm/model/ProcessThreadStatus.java | 20 +- .../alarm/model/ProcessThreadValue.java | 6 +- .../cloud/skywalking/alarm/util/ZKUtil.java | 32 +++- skywalking-alarm/src/main/resources/log4j.xml | 3 +- .../src/main/resources/log4j2.xml | 5 +- .../src/main/resources/mail/mail.config | 6 - 12 files changed, 295 insertions(+), 154 deletions(-) delete mode 100644 skywalking-alarm/src/main/resources/mail/mail.config diff --git a/skywalking-alarm/pom.xml b/skywalking-alarm/pom.xml index 8eb07bd75..549781f4e 100644 --- a/skywalking-alarm/pom.xml +++ b/skywalking-alarm/pom.xml @@ -76,11 +76,19 @@ - org.apache.maven.plugins maven-compiler-plugin - 1.7 - 1.7 + 1.6 + 1.6 + ${project.build.sourceEncoding} + + + + org.apache.maven.plugins + maven-resources-plugin + 2.4.3 + + ${project.build.sourceEncoding} 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 bdb719426..0b1b32375 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 @@ -4,19 +4,22 @@ import com.ai.cloud.skywalking.alarm.conf.Config; import com.ai.cloud.skywalking.alarm.dao.AlarmMessageDao; import com.ai.cloud.skywalking.alarm.model.AlarmRule; import com.ai.cloud.skywalking.alarm.model.ProcessThreadStatus; +import com.ai.cloud.skywalking.alarm.model.ProcessThreadValue; import com.ai.cloud.skywalking.alarm.model.UserInfo; import com.ai.cloud.skywalking.alarm.procesor.AlarmMessageProcessor; import com.ai.cloud.skywalking.alarm.util.ProcessUtil; import com.ai.cloud.skywalking.alarm.util.ZKUtil; +import com.google.gson.Gson; import org.apache.curator.framework.api.CuratorWatcher; -import org.apache.curator.framework.recipes.locks.InterProcessMutex; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; +import org.apache.zookeeper.CreateMode; import org.apache.zookeeper.WatchedEvent; import org.apache.zookeeper.Watcher; -import java.util.ArrayList; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.UUID; public class AlarmMessageProcessThread extends Thread { @@ -25,9 +28,10 @@ public class AlarmMessageProcessThread extends Thread { private String threadId; private ProcessThreadStatus status; - private String[] processUserIds; - private List usersLocks; + private List processUserIds; private CoordinatorStatusWatcher watcher = new CoordinatorStatusWatcher(); + private Map> cacheRules = new HashMap>(); + private static AlarmMessageProcessor processor = new AlarmMessageProcessor(); public AlarmMessageProcessThread() { // 初始化生成ThreadId @@ -39,96 +43,85 @@ public class AlarmMessageProcessThread extends Thread { //注册服务(默认为空闲状态) registerProcessThread(threadId, ProcessThreadStatus.FREE); while (true) { - //检查是否为忙碌状态 - if (status == ProcessThreadStatus.BUSY) { - //处理告警信息 - for (String userId : processUserIds) { - List rules = AlarmMessageDao.selectAlarmRulesByUserId(userId); - UserInfo userInfo = AlarmMessageDao.selectUser(userId); - for (AlarmRule rule : rules) { - new AlarmMessageProcessor().process(userInfo, rule); + try { + //检查是否为忙碌状态 + if (status == ProcessThreadStatus.BUSY) { + //处理告警信息 + 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()); + } } } - } - //检查是否分配线程的状态(重新分配状态) - if (status == ProcessThreadStatus.REDISTRIBUTING) { - // 释放用户锁 - releaseUserLock(); - // 修改自身状态:(空闲状态) - status = ProcessThreadStatus.FREE; - ProcessUtil.changeProcessThreadStatus(threadId, ProcessThreadStatus.FREE); - } + //检查是否分配线程的状态(重新分配状态) + if (status == ProcessThreadStatus.REDISTRIBUTING) { + // 修改自身状态:(空闲状态) + status = ProcessThreadStatus.FREE; + ProcessUtil.changeProcessThreadStatus(threadId, ProcessThreadStatus.FREE); - //检查分配线程的状态(分配完成状态) - if (status == ProcessThreadStatus.REDISTRIBUTE_SUCCESS) { - // 获取待处理的用户 - processUserIds = acquireProcessedUsers(); - // 给用户加锁 - lockUser(processUserIds); - // 修改自身状态 :(忙碌状态) - status = ProcessThreadStatus.BUSY; - ProcessUtil.changeProcessThreadStatus(threadId, ProcessThreadStatus.BUSY); - } + //清空缓存数据 + clearCacheProcessUser(cacheRules); + } - try { - Thread.sleep(Config.ProcessThread.THREAD_WAIT_INTERVAL); - } catch (InterruptedException e) { - logger.error("Sleep failed.", e); + //检查分配线程的状态(分配完成状态) + if (status == ProcessThreadStatus.REDISTRIBUTE_SUCCESS) { + // 获取待处理的用户 + processUserIds = acquireProcessedUsers(); + + // 缓存数据 + cacheProcessUser(processUserIds); + + // 修改自身状态 :(忙碌状态) + status = ProcessThreadStatus.BUSY; + ProcessUtil.changeProcessThreadStatus(threadId, ProcessThreadStatus.BUSY); + } + + try { + Thread.sleep(Config.ProcessThread.THREAD_WAIT_INTERVAL); + } catch (InterruptedException e) { + logger.error("Sleep failed.", e); + } + } catch (Exception e) { + logger.error("Failed to process data.", e); } } } - private String[] acquireProcessedUsers() { + private void clearCacheProcessUser(Map> cacheRules) { + cacheRules.clear(); + } + + private void cacheProcessUser(List processUserIds) { + // TODO 需要重新获取 + UserInfo tmpUserInfo; + List alarmRules; + for (String userId : processUserIds) { + tmpUserInfo = AlarmMessageDao.selectUser(userId); + if (tmpUserInfo == null) { + continue; + } + alarmRules = AlarmMessageDao.selectAlarmRulesByUserId(userId); + cacheRules.put(tmpUserInfo, alarmRules); + } + } + + private List acquireProcessedUsers() { String path = Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId; String value = ZKUtil.getPathData(path); - String[] valueArrays = value.split("@"); - String toBeProcessUserId; - if (valueArrays.length < 2) { - toBeProcessUserId = ""; - } else { - toBeProcessUserId = valueArrays[1]; - } - - return toBeProcessUserId.split(";"); - } - - private void lockUser(String[] userIds) { - usersLocks = new ArrayList(); - String userLockPath = Config.ZKPath.USER_REGISTER_LOCK_PATH + "/"; - InterProcessMutex tmpLock; - for (String userId : userIds) { - tmpLock = ZKUtil.getLock(userLockPath + userId); - try { - tmpLock.acquire(); - } catch (Exception e) { - logger.error("Failed to lock user[{}]", userId, e); - //TODO 锁失败,该怎么处理? - } - usersLocks.add(tmpLock); - } - } - - private void releaseUserLock() { - for (InterProcessMutex lock : usersLocks) { - try { - lock.release(); - } catch (Exception e) { - // - logger.error("Failed to release lock user.", e); - //TODO 释放锁,该怎么处理。 - } - } - - usersLocks.clear(); + ProcessThreadValue processThreadValue = new Gson().fromJson(value, ProcessThreadValue.class); + return processThreadValue.getDealUserIds(); } private void registerProcessThread(String threadId, ProcessThreadStatus status) { try { String registerPath = Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId; - String registerValue = status + "@ "; + ProcessThreadValue initValue = new ProcessThreadValue(); + initValue.setStatus(status.getValue()); ZKUtil.getZkClient().create().creatingParentsIfNeeded() - .forPath(registerPath, registerValue.getBytes()); + .withMode(CreateMode.EPHEMERAL).forPath + (registerPath, new Gson().toJson(initValue).getBytes()); this.status = status; @@ -144,11 +137,16 @@ public class AlarmMessageProcessThread extends Thread { @Override public void process(WatchedEvent watchedEvent) { if (watchedEvent.getType() == Watcher.Event.EventType.NodeDataChanged) { - String value = ZKUtil.getPathData(Config.ZKPath.COORDINATOR_STATUS_PATH); - status = ProcessThreadStatus.convert(value); + String value = ZKUtil.getPathData(Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId); + ProcessThreadValue processThreadValue = new Gson().fromJson(value, ProcessThreadValue.class); + status = ProcessThreadStatus.convert(processThreadValue.getStatus()); } - ZKUtil.getPathDataWithWatch(Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId, watcher); + try { + ZKUtil.getPathDataWithWatch(Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId, watcher); + } catch (Exception e) { + e.printStackTrace(); + } } } diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmProcessServer.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmProcessServer.java index d845ab50c..4a434eaa3 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmProcessServer.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmProcessServer.java @@ -1,6 +1,7 @@ package com.ai.cloud.skywalking.alarm; import com.ai.cloud.skywalking.alarm.conf.Config; +import com.ai.cloud.skywalking.alarm.util.ZKUtil; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -16,6 +17,11 @@ public class AlarmProcessServer { public static void main(String[] main) { logger.info("Begin to start alarm process server...."); + if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH)) { + ZKUtil.createPath(Config.ZKPath.REGISTER_SERVER_PATH); + } + new UserInfoCoordinator().start(); + logger.info("Begin to start process thread..."); AlarmMessageProcessThread tmpThread; for (int i = 0; i < Config.Server.PROCESS_THREAD_SIZE; i++) { @@ -24,9 +30,6 @@ public class AlarmProcessServer { processThreads.add(tmpThread); } logger.info("Successfully launched {} processing threads.", Config.Server.PROCESS_THREAD_SIZE); - logger.info("Begin to start user inspector thread...."); - new UserInfoInspector().start(); - logger.info("Start user inspector thread success...."); logger.info("Alarm process server successfully started."); while (true) { try { 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 7a8aa0de2..248e8532b 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 @@ -8,14 +8,17 @@ import com.ai.cloud.skywalking.alarm.util.ProcessUtil; import com.ai.cloud.skywalking.alarm.util.ZKUtil; import com.google.gson.Gson; import org.apache.curator.framework.api.CuratorWatcher; +import org.apache.curator.framework.recipes.locks.InterProcessMutex; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.apache.zookeeper.WatchedEvent; import org.apache.zookeeper.Watcher; +import java.util.ArrayList; import java.util.HashSet; import java.util.List; import java.util.Set; +import java.util.concurrent.TimeUnit; public class UserInfoCoordinator extends Thread { @@ -24,54 +27,122 @@ public class UserInfoCoordinator extends Thread { 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() { 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; + } + + if (coordinatorFlag) { + watcherRegisterServerPath(); + } } @Override public void run() { - while (true) { - //检查是否有新服务注册或者在重分配过程做有新处理线程启动了 - if (!redistributing || !newServerComingFlag) { + + if (!coordinatorFlag) { + // + while (!retryBecomeCoordinator()) { try { - Thread.sleep(1000L); - } catch (InterruptedException e) { - logger.error("Sleep error", e); + Thread.sleep(Config.Coordinator.RETRY_BECOME_COORDINATOR_WAIT_TIME); + } catch (Exception e) { + logger.error("Sleep Failed.", e); } - - continue; } - - // 设置正在处理的标志位 + // 新官上任三把火,重新分配 + watcherRegisterServerPath(); redistributing = true; - newServerComingFlag = false; + } - //获取当前所有的注册的处理线程 - List registeredThreads = acquireAllRegisteredThread(); - //修改状态 (开始重新分配状态) - changeStatus(registeredThreads, ProcessThreadStatus.REDISTRIBUTING); - //检查所有的服务是否都处于空闲状态 - while (!isAllProcessThreadFree(registeredThreads)) { - try { - Thread.sleep(100L); - } catch (InterruptedException e) { - logger.error("Sleep failed", e); + while (true) { + try { + //检查是否有新服务注册或者在重分配过程做有新处理线程启动了 + if (!redistributing && !newServerComingFlag) { + try { + Thread.sleep(Config.Coordinator.CHECK_REDISTRIBUTE_INTERVAL); + } catch (InterruptedException e) { + logger.error("Sleep error", e); + } + + continue; } + + newServerComingFlag = false; + + //获取当前所有的注册的处理线程 + 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); + } + } + + //查询当前有多少用户 + List users = AlarmMessageDao.selectAllUserIds(); + + //将用户重新分配给服务 + List realRedistributeThread = allocationUser(registeredThreads, users); + + //修改状态(分配完成) + changeStatus(realRedistributeThread, ProcessThreadStatus.REDISTRIBUTE_SUCCESS); + + //检查所有的服务是否都处于忙碌状态 + while (!checkAllProcessStatus(realRedistributeThread, ProcessThreadStatus.BUSY)) { + try { + Thread.sleep(Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL); + } catch (InterruptedException e) { + logger.error("Sleep failed", e); + } + } + + redistributing = false; + + } catch (Exception e) { + logger.error("Failed to redistribute User ", e); } - - //查询当前有多少用户 - List users = AlarmMessageDao.selectAllUserIds(); - - //将用户重新分配给服务 - allocationUser(registeredThreads, users); - - //修改状态(分配完成) - changeStatus(registeredThreads, ProcessThreadStatus.REDISTRIBUTE_SUCCESS); } } - private void allocationUser(List registeredThreads, List userIds) { + 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); + } + } + } + } + + private boolean retryBecomeCoordinator() { + try { + return lock.acquire(5, TimeUnit.SECONDS); + } catch (Exception e) { + logger.error("Failed to acquire lock .", e); + return 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; @@ -82,30 +153,41 @@ public class UserInfoCoordinator extends Thread { } 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); start = end; end += step; - - if (end > userIds.size()) { + if (start >= userIds.size()) { break; } + if (end > userIds.size()) { + end = userIds.size(); + } } + return realRedistributeThread; } private List acquireAllRegisteredThread() { return ZKUtil.getChildren(Config.ZKPath.REGISTER_SERVER_PATH); } - private boolean isAllProcessThreadFree(List registeredThreadIds) { + private boolean checkAllProcessStatus(List registeredThreadIds, ProcessThreadStatus status) { String registerPathPrefix = Config.ZKPath.REGISTER_SERVER_PATH + "/"; for (String threadId : registeredThreadIds) { + + if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/" + threadId)) + continue; + if (getProcessThreadStatus(registerPathPrefix, threadId) - != ProcessThreadStatus.FREE) { + != status) { return false; } } @@ -113,7 +195,11 @@ public class UserInfoCoordinator extends Thread { } 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()); } @@ -128,14 +214,29 @@ public class UserInfoCoordinator extends Thread { public class RegisterServerWatcher implements CuratorWatcher { @Override - public void process(WatchedEvent watchedEvent) throws Exception { + 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) { + // + e.printStackTrace(); + if (lock != null && lock.isAcquiredInThisProcess()) { + try { + lock.release(); + } catch (Exception e1) { + logger.error("Failed to release master locker.", e1); + } + } + } } } } 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 37f622d83..743182a3a 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 @@ -4,7 +4,7 @@ public class Config { public static class Server { - public static int PROCESS_THREAD_SIZE = 2; + public static int PROCESS_THREAD_SIZE = 1; public static long DAEMON_THREAD_WAIT_INTERVAL = 50000L; @@ -12,6 +12,7 @@ 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 { @@ -26,11 +27,22 @@ public class Config { public static String NODE_PREFIX = "/skywalking"; - public static String COORDINATOR_STATUS_PATH = "/alarm-server/coordinator/status"; + public static String REGISTER_SERVER_PATH = NODE_PREFIX + "/alarm-server/register-servers"; - public static String REGISTER_SERVER_PATH = "/alarm-server/register-servers"; + public static String COORDINATOR_PATH = NODE_PREFIX + "/alarm-server/coordinator/lock"; - public static String USER_REGISTER_LOCK_PATH = "/alarm-server/users"; + } + + + public static class Coordinator { + // 单位:(秒) + public static long RETRY_GET_COORDINATOR_LOCK_INTERVAL = 5; + + public static long RETRY_BECOME_COORDINATOR_WAIT_TIME = 10 * 1000L; + // 单位:(毫秒) + public static long CHECK_REDISTRIBUTE_INTERVAL = 5 * 1000; + // 单位:(毫秒) + public static long CHECK_ALL_PROCESS_THREAD_INTERVAL = 100L; } public static class DB { diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/dao/AlarmMessageDao.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/dao/AlarmMessageDao.java index ca01cd07c..4fafab1a9 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/dao/AlarmMessageDao.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/dao/AlarmMessageDao.java @@ -51,8 +51,9 @@ public class AlarmMessageDao { public static UserInfo selectUser(String userId) { UserInfo userInfo = null; try { - PreparedStatement ps = DBConnectUtil.getConnection().prepareStatement("SELECT user_info.uid,user_info.user_name FROM user_info WHERE sts = ?"); + PreparedStatement ps = DBConnectUtil.getConnection().prepareStatement("SELECT user_info.uid,user_info.user_name FROM user_info WHERE sts = ? AND uid = ?"); ps.setString(1, "A"); + ps.setString(2, userId); ResultSet rs = ps.executeQuery(); rs.next(); userInfo = new UserInfo(rs.getString("uid")); diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadStatus.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadStatus.java index b3d3fae42..3079123cd 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadStatus.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadStatus.java @@ -1,27 +1,33 @@ package com.ai.cloud.skywalking.alarm.model; public enum ProcessThreadStatus { - REDISTRIBUTING("1"), REDISTRIBUTE_SUCCESS("2"), FREE("0"), BUSY("3"); + REDISTRIBUTING(1), REDISTRIBUTE_SUCCESS(2), FREE(0), BUSY(3); - private String value; + private int value; - ProcessThreadStatus(String value) { + ProcessThreadStatus(int value) { this.value = value; } - public String getValue() { + public int getValue() { return value; } - public static ProcessThreadStatus convert(String value) { + public static ProcessThreadStatus convert(int value) { ProcessThreadStatus status; switch (value) { - case "0": + case 0: + status = FREE; + break; + case 1: status = REDISTRIBUTING; break; - case "1": + case 2: status = REDISTRIBUTE_SUCCESS; break; + case 3: + status = BUSY; + break; default: throw new IllegalArgumentException("Coordinator status illegal"); } diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadValue.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadValue.java index ae515272e..3b0dd81c0 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadValue.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ProcessThreadValue.java @@ -3,14 +3,14 @@ package com.ai.cloud.skywalking.alarm.model; import java.util.List; public class ProcessThreadValue { - private String status; + private int status; private List dealUserIds; - public String getStatus() { + public int getStatus() { return status; } - public void setStatus(String status) { + public void setStatus(int status) { this.status = status; } diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/ZKUtil.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/ZKUtil.java index 89cc6a03b..24f4a491a 100644 --- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/ZKUtil.java +++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/ZKUtil.java @@ -9,6 +9,7 @@ import org.apache.curator.framework.recipes.locks.InterProcessMutex; import org.apache.curator.retry.ExponentialBackoffRetry; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; +import org.apache.zookeeper.CreateMode; import java.util.ArrayList; import java.util.List; @@ -50,13 +51,8 @@ public class ZKUtil { return ""; } - public static String getPathDataWithWatch(String path, CuratorWatcher watcher) { - try { - return new String(client.getData().usingWatcher(watcher).forPath(path)); - } catch (Exception e) { - logger.error("Failed to get the value of path[{}]", path, e); - } - return ""; + public static String getPathDataWithWatch(String path, CuratorWatcher watcher) throws Exception { + return new String(client.getData().usingWatcher(watcher).forPath(path)); } public static void setPathData(String path, String value) { @@ -75,4 +71,26 @@ public class ZKUtil { } return new ArrayList(); } + + public static List getChildrenWithWatcher(String registerServerPath, CuratorWatcher watcher) throws Exception { + return client.getChildren().usingWatcher(watcher).forPath(registerServerPath); + } + + public static void createPath(String path) { + try { + client.create().creatingParentsIfNeeded().withMode(CreateMode.PERSISTENT).forPath(path); + } catch (Exception e) { + logger.error("Failed to create path."); + } + } + + public static boolean exists(String registerServerPath) { + try { + return client.checkExists().forPath(registerServerPath) != null; + } catch (Exception e) { + logger.error("Failed check exists for path"); + } + + return false; + } } diff --git a/skywalking-alarm/src/main/resources/log4j.xml b/skywalking-alarm/src/main/resources/log4j.xml index 0da28a6d5..a93df5db7 100644 --- a/skywalking-alarm/src/main/resources/log4j.xml +++ b/skywalking-alarm/src/main/resources/log4j.xml @@ -11,9 +11,8 @@ - + - diff --git a/skywalking-alarm/src/main/resources/log4j2.xml b/skywalking-alarm/src/main/resources/log4j2.xml index 025d5b2fb..a497568d1 100644 --- a/skywalking-alarm/src/main/resources/log4j2.xml +++ b/skywalking-alarm/src/main/resources/log4j2.xml @@ -1,13 +1,14 @@ - + - + + \ No newline at end of file diff --git a/skywalking-alarm/src/main/resources/mail/mail.config b/skywalking-alarm/src/main/resources/mail/mail.config deleted file mode 100644 index 6b968bf65..000000000 --- a/skywalking-alarm/src/main/resources/mail/mail.config +++ /dev/null @@ -1,6 +0,0 @@ -mail.host=mail.asiainfo.com -mail.transport.protocol=smtp -mail.smtp.auth=true -mail.username=zhangxin10 -mail.password=!qaz2wsx -mail.account.prefix=@asiainfo.com \ No newline at end of file