diff --git a/skywalking-alarm/pom.xml b/skywalking-alarm/pom.xml
new file mode 100644
index 000000000..8eb07bd75
--- /dev/null
+++ b/skywalking-alarm/pom.xml
@@ -0,0 +1,88 @@
+
+ 4.0.0
+
+ com.ai.cloud
+ skywalking-alarm
+ 1.0-SNAPSHOT
+ jar
+
+ skywalking-alarm
+ http://maven.apache.org
+
+
+ UTF-8
+
+
+
+ Company
+ http://10.1.228.199:18081/nexus/content/groups/public/
+
+
+
+
+
+ junit
+ junit
+ 4.12
+ test
+
+
+ org.apache.logging.log4j
+ log4j-core
+ 2.4.1
+
+
+ redis.clients
+ jedis
+ 2.8.0
+
+
+ com.sun.mail
+ javax.mail
+ 1.5.4
+
+
+ org.freemarker
+ freemarker
+ 2.3.23
+
+
+ org.apache.curator
+ curator-framework
+ 2.8.0
+
+
+ org.apache.zookeeper
+ zookeeper
+ 3.4.0
+
+
+ org.apache.curator
+ curator-recipes
+ 2.8.0
+
+
+ com.google.code.gson
+ gson
+ 2.2.2
+
+
+ mysql
+ mysql-connector-java
+ 5.1.37
+
+
+
+
+
+ org.apache.maven.plugins
+ maven-compiler-plugin
+
+ 1.7
+ 1.7
+
+
+
+
+
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
new file mode 100644
index 000000000..d285c3380
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmMessageProcessThread.java
@@ -0,0 +1,135 @@
+package com.ai.cloud.skywalking.alarm;
+
+import com.ai.cloud.skywalking.alarm.dao.AlarmMessageDao;
+import com.ai.cloud.skywalking.alarm.model.ApplicationInfo;
+import com.ai.cloud.skywalking.alarm.model.UserInfo;
+import com.ai.cloud.skywalking.alarm.zk.ZKUtil;
+import org.apache.curator.framework.recipes.locks.InterProcessMutex;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.ArrayList;
+import java.util.List;
+import java.util.UUID;
+
+
+public class AlarmMessageProcessThread extends Thread {
+ private int rank;
+ private String threadId;
+ private Logger logger = LogManager.getLogger(AlarmMessageProcessThread.class);
+ private boolean isChanged = false;
+ private List locks;
+ private List toBeProcessUsers;
+
+ public AlarmMessageProcessThread() {
+ threadId = UUID.randomUUID().toString();
+ registerServer(threadId);
+ rank = ballot(threadId);
+ ZKUtil.watch(AlarmServerRegisterWatcher.getInstance());
+ toBeProcessUsers = getToBeProcessUsers();
+ }
+
+
+ @Override
+ public void run() {
+ lockAllUsers();
+ while (true) {
+ for (UserInfo userInfo : toBeProcessUsers) {
+ for (ApplicationInfo applicationInfo : userInfo.getApplicationInfos()) {
+ System.out.println(threadId + applicationInfo.getAppId());
+ }
+ }
+
+ if (isChanged) {
+ unlockAllUsers();
+ rank = ballot(threadId);
+ toBeProcessUsers = getToBeProcessUsers();
+ lockAllUsers();
+ isChanged = false;
+ }
+
+ try {
+ Thread.sleep(1000L);
+ } catch (InterruptedException e) {
+ logger.error("Sleep failed", e);
+ }
+ }
+ }
+
+ private void unlockAllUsers() {
+ for (InterProcessMutex lock : locks) {
+ try {
+ lock.release();
+ } catch (Exception e) {
+ logger.error("Failed to release lock[{}]", e);
+ }
+ }
+ }
+
+ private void lockAllUsers() {
+ locks = new ArrayList(toBeProcessUsers.size());
+ InterProcessMutex tmpLock;
+ for (UserInfo userInfo : toBeProcessUsers) {
+ tmpLock = ZKUtil.getProcessUserLock(userInfo.getUserId());
+ locks.add(tmpLock);
+ try {
+ tmpLock.acquire();
+ } catch (Exception e) {
+ logger.error("Failed to lock ");
+ }
+
+ }
+ }
+
+ private List getToBeProcessUsers() {
+ List userInfos = AlarmMessageDao.selectAllUserInfo();
+ List allServerIds = ZKUtil.selectAllServerIds();
+ int step = (int) Math.ceil(userInfos.size() * 1.0 / allServerIds.size());
+ int start = rank * step;
+ int end = (rank + 1) * step;
+ if (end > userInfos.size()) {
+ return new ArrayList();
+ }
+
+ List toBeProcessUsers = userInfos.subList(start, end);
+ for (UserInfo userInfo : toBeProcessUsers) {
+ userInfo.setApplicationInfos(AlarmMessageDao.selectAllApplicationsByUserId(userInfo.getUserId()));
+ }
+ return toBeProcessUsers;
+ }
+
+
+ private void registerServer(String serverId) {
+ InterProcessMutex registerLock = ZKUtil.getRegisterLock();
+ try {
+ registerLock.acquire();
+ ZKUtil.register(serverId);
+ } catch (Exception e) {
+ logger.error("Failed to lock.", e);
+ } finally {
+ if (registerLock != null) {
+ try {
+ registerLock.release();
+ } catch (Exception e) {
+ logger.error("Failed to release lock.", e);
+ }
+ }
+ }
+ }
+
+ private int ballot(String threadId) {
+ List serverIds = ZKUtil.selectAllServerIds();
+ int rank = 0;
+ for (String tmpServerId : serverIds) {
+ if (tmpServerId.hashCode() < threadId.hashCode()) {
+ rank++;
+ }
+ }
+ return rank;
+ }
+
+
+ public void setChanged(boolean changed) {
+ isChanged = changed;
+ }
+}
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
new file mode 100644
index 000000000..32b481868
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmProcessServer.java
@@ -0,0 +1,44 @@
+package com.ai.cloud.skywalking.alarm;
+
+import com.ai.cloud.skywalking.alarm.conf.Config;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.util.ArrayList;
+import java.util.List;
+
+public class AlarmProcessServer {
+
+ private static Logger logger = LogManager.getLogger(AlarmProcessServer.class);
+
+ private static List processThreads =
+ new ArrayList();
+
+ public static void main(String[] main) {
+ logger.info("Begin to start alarm process server....");
+ logger.info("Begin to start process thread...");
+ AlarmMessageProcessThread tmpThread;
+ for (int i = 0; i < Config.Server.PROCESS_THREAD_SIZE; i++) {
+ tmpThread = new AlarmMessageProcessThread();
+ tmpThread.start();
+ 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 {
+ Thread.sleep(Config.Server.DAEMON_THREAD_WAITE_INTERVAL);
+ } catch (InterruptedException e) {
+ logger.error("Sleep failed", e);
+ }
+ }
+ }
+
+
+ public static List getProcessThreads() {
+ return processThreads;
+ }
+}
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmServerRegisterWatcher.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmServerRegisterWatcher.java
new file mode 100644
index 000000000..ab8499e14
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/AlarmServerRegisterWatcher.java
@@ -0,0 +1,26 @@
+package com.ai.cloud.skywalking.alarm;
+
+import com.ai.cloud.skywalking.alarm.zk.ZKUtil;
+import org.apache.curator.framework.api.CuratorWatcher;
+import org.apache.zookeeper.WatchedEvent;
+import org.apache.zookeeper.Watcher;
+
+public class AlarmServerRegisterWatcher implements CuratorWatcher {
+
+ private static AlarmServerRegisterWatcher watcher = new AlarmServerRegisterWatcher();
+
+ @Override
+ public void process(WatchedEvent watchedEvent) throws Exception {
+ if (watchedEvent.getType() == Watcher.Event.EventType.NodeChildrenChanged) {
+ for (AlarmMessageProcessThread thread : AlarmProcessServer.getProcessThreads()) {
+ thread.setChanged(true);
+ }
+ }
+
+ ZKUtil.watch(getInstance());
+ }
+
+ public static AlarmServerRegisterWatcher getInstance() {
+ return watcher;
+ }
+}
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
new file mode 100644
index 000000000..14f46ead0
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/UserInfoInspector.java
@@ -0,0 +1,37 @@
+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
new file mode 100644
index 000000000..20192beb1
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/Config.java
@@ -0,0 +1,58 @@
+package com.ai.cloud.skywalking.alarm.conf;
+
+public class Config {
+
+ public static class Server {
+ public static int PROCESS_THREAD_SIZE = 2;
+
+ public static long DAEMON_THREAD_WAITE_INTERVAL = 50000L;
+
+ }
+
+ public static class ZKPath {
+
+ public static String CONNECT_STR = "127.0.0.1:2181";
+ public static int CONNECT_TIMEOUT = 1000;
+ public static int RETRY_TIMEOUT = 1000;
+ public static int RETRY_TIMES = 3;
+
+ public static String NODE_PREFIX = "/skywalking";
+
+ public static String SERVER_REGISTER_LOCK_PATH = "/alarm-server/register";
+ public static String REGISTER_SERVER_PATH = "/alarm-server/servers";
+ public static String USER_REGISTER_LOCK_PATH = "/alarm-server/users";
+ }
+
+ public static class DB {
+ public static String PASSWORD = "devrdbusr13";
+ public static String USER_NAME = "devrdbusr13";
+ public static String DRIVER_CLASS = "com.mysql.jdbc.Driver";
+ public static String URL = "jdbc:mysql://10.1.228.200:31306/test";
+ }
+
+ public static class Alarm {
+
+ public static String REDIS_SERVER = "127.0.0.1:6379";
+
+ public static int REDIS_MAX_IDLE = 10;
+
+ public static int REDIS_MIN_IDLE = 1;
+
+ public static int REDIS_MAX_TOTAL = 20;
+
+ public static boolean ALARM_OFF_FLAG = false;
+ }
+
+ public static class MailSender {
+
+ public static String HOST = "";
+
+ public static String TRANSPORT_PROTOCOL = "";
+
+ public static boolean SMTP_AUTH = true;
+
+ public static String USER_NAME = "";
+
+ public static String PASSWORD = "";
+ }
+}
\ No newline at end of file
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/ConfigInitializer.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/ConfigInitializer.java
new file mode 100644
index 000000000..b4b41d78d
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/ConfigInitializer.java
@@ -0,0 +1,85 @@
+package com.ai.cloud.skywalking.alarm.conf;
+
+import java.io.IOException;
+import java.io.InputStream;
+import java.lang.reflect.Field;
+import java.lang.reflect.Modifier;
+import java.util.LinkedList;
+import java.util.Properties;
+import java.util.logging.Level;
+import java.util.logging.Logger;
+
+public class ConfigInitializer {
+ private static Logger logger = Logger.getLogger(ConfigInitializer.class.getName());
+
+ public static void initialize() {
+ InputStream inputStream = ConfigInitializer.class.getResourceAsStream("/sky-walking.auth");
+ if (inputStream == null) {
+ logger.log(Level.ALL, "No provider sky-walking certification documents, buried point won't work");
+ } else {
+ try {
+ Properties properties = new Properties();
+ properties.load(inputStream);
+ initNextLevel(properties, Config.class, new ConfigDesc());
+ } catch (IllegalAccessException e) {
+ logger.log(Level.ALL, "Parsing certification file failed, buried won't work");
+ } catch (IOException e) {
+ logger.log(Level.ALL, "Failed to read the certification file, buried won't work");
+ }
+ }
+ }
+
+ private static void initNextLevel(Properties properties, Class> recentConfigType, ConfigDesc parentDesc) throws NumberFormatException, IllegalArgumentException, IllegalAccessException {
+ for (Field field : recentConfigType.getFields()) {
+ if (Modifier.isPublic(field.getModifiers()) && Modifier.isStatic(field.getModifiers())) {
+ String configKey = (parentDesc + "." +
+ field.getName()).toLowerCase();
+ String value = properties.getProperty(configKey);
+ if (value != null) {
+ if (field.getType().equals(int.class))
+ field.set(null, Integer.valueOf(value));
+ if (field.getType().equals(String.class))
+ field.set(null, value);
+ if (field.getType().equals(long.class))
+ field.set(null, Long.valueOf(value));
+ if (field.getType().equals(boolean.class))
+ field.set(null, Boolean.valueOf(value));
+ }
+ }
+ }
+ for (Class> innerConfiguration : recentConfigType.getClasses()) {
+ parentDesc.append(innerConfiguration.getSimpleName());
+ initNextLevel(properties, innerConfiguration, parentDesc);
+ parentDesc.removeLastDesc();
+ }
+ }
+}
+
+class ConfigDesc {
+ private LinkedList descs = new LinkedList();
+
+ void append(String currentDesc) {
+ descs.addLast(currentDesc);
+ }
+
+ void removeLastDesc() {
+ descs.removeLast();
+ }
+
+ @Override
+ public String toString() {
+ if (descs.size() == 0) {
+ return "";
+ }
+ StringBuilder ret = new StringBuilder(descs.getFirst());
+ boolean first = true;
+ for (String desc : descs) {
+ if (first) {
+ first = false;
+ continue;
+ }
+ ret.append(".").append(desc);
+ }
+ return ret.toString();
+ }
+}
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
new file mode 100644
index 000000000..7256362b8
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/dao/AlarmMessageDao.java
@@ -0,0 +1,121 @@
+package com.ai.cloud.skywalking.alarm.dao;
+
+import com.ai.cloud.skywalking.alarm.conf.Config;
+import com.ai.cloud.skywalking.alarm.model.ApplicationInfo;
+import com.ai.cloud.skywalking.alarm.model.UserInfo;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+
+import java.sql.*;
+import java.util.ArrayList;
+import java.util.List;
+
+public class AlarmMessageDao {
+
+ private static Logger logger = LogManager.getLogger(AlarmMessageDao.class);
+ private static Connection con;
+
+ static {
+ try {
+ Class.forName(Config.DB.DRIVER_CLASS);
+ } catch (ClassNotFoundException e) {
+ logger.error("Failed to found DB driver class.", e);
+ System.exit(-1);
+ }
+
+ try {
+ con = DriverManager.getConnection(Config.DB.URL, Config.DB.USER_NAME, Config.DB.PASSWORD);
+ } catch (SQLException e) {
+ logger.error("Failed to connect DB", e);
+ System.exit(-1);
+ }
+ }
+
+ public static List selectAllUserInfo() {
+ List result = new ArrayList();
+ try {
+ PreparedStatement ps = con.prepareStatement("SELECT user_info.uid FROM user_info WHERE sts = ?");
+ ps.setString(1, "A");
+ ResultSet rs = ps.executeQuery();
+ UserInfo userInfo;
+ while (rs.next()) {
+ userInfo = new UserInfo(rs.getString("uid"));
+ result.add(userInfo);
+ }
+ } catch (SQLException e) {
+ logger.error("Failed to select all user info", e);
+ }
+
+ return result;
+ }
+
+ public static int selectUserCount() {
+ try {
+ PreparedStatement ps = con.prepareStatement("SELECT count(user_info.uid) as totalNumber FROM user_info WHERE sts = ?");
+ ps.setString(1, "A");
+ ResultSet rs = ps.executeQuery();
+ rs.next();
+ return rs.getInt("totalNumber");
+ } catch (SQLException e) {
+ logger.error("Failed to select all user info", e);
+ }
+
+ return 0;
+ }
+
+ public static List selectAllApplicationsByUserId(String userId) {
+ List result = new ArrayList();
+ try {
+ PreparedStatement ps = con.prepareStatement("SELECT alarm_rule.app_id, alarm_rule.uid,alarm_rule.is_global, alarm_rule.todo_type," +
+ " alarm_rule.config_args FROM alarm_rule WHERE uid = ? AND sts = ?");
+ ps.setString(1, userId);
+ ps.setString(2, "A");
+ ResultSet rs = ps.executeQuery();
+ ApplicationInfo globalConfig = null;
+ ApplicationInfo tmpApplication;
+ while (rs.next()) {
+ if ("1".equals(rs.getString("is_global"))) {
+ globalConfig = new ApplicationInfo();
+ globalConfig.setConfigArgs(rs.getString("config_args"));
+ continue;
+ }
+ tmpApplication = new ApplicationInfo();
+ tmpApplication.setAppId(rs.getString("app_id"));
+ tmpApplication.setUId(rs.getString("uid"));
+ tmpApplication.setConfigArgs(rs.getString("config_args"));
+ tmpApplication.setToDoType(rs.getString("todo_type"));
+ result.add(tmpApplication);
+ }
+
+ if (globalConfig == null) {
+ //
+ throw new IllegalArgumentException("Can not found the global config");
+ }
+
+ List allApplication = new ArrayList();
+ ps = con.prepareStatement("SELECT application_info.app_id, application_info.uid FROM application_info WHERE uid = ? AND sts = ?");
+ ps.setString(1, userId);
+ ps.setString(2, "A");
+ rs = ps.executeQuery();
+ ApplicationInfo applicationInfo;
+ while (rs.next()) {
+ applicationInfo = new ApplicationInfo();
+ applicationInfo.setAppId(rs.getString("app_id"));
+ applicationInfo.setUId(rs.getString("uid"));
+ allApplication.add(applicationInfo);
+ }
+
+ allApplication.removeAll(result);
+
+ for (ApplicationInfo app : allApplication) {
+ app.setConfigArgs(globalConfig.getConfigArgs());
+ app.setToDoType(globalConfig.getToDoType());
+ result.add(app);
+ }
+
+ } catch (SQLException e) {
+ logger.error("Failed to query applications.", e);
+ }
+ return result;
+ }
+}
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ApplicationInfo.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ApplicationInfo.java
new file mode 100644
index 000000000..85559361c
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ApplicationInfo.java
@@ -0,0 +1,70 @@
+package com.ai.cloud.skywalking.alarm.model;
+
+import com.google.gson.Gson;
+
+public class ApplicationInfo {
+ private String appId;
+ private String configArgs;
+ private String UId;
+ private String toDoType;
+ private ConfigArgsDescriber configArgsDescriber;
+
+ public ApplicationInfo() {
+ }
+
+ public String getAppId() {
+ return appId;
+ }
+
+ public void setAppId(String appId) {
+ this.appId = appId;
+ }
+
+ public void setConfigArgs(String configArgs) {
+ this.configArgs = configArgs;
+ configArgsDescriber = new Gson().fromJson(configArgs, ConfigArgsDescriber.class);
+ }
+
+ public ConfigArgsDescriber getConfigArgsDescriber() {
+ return configArgsDescriber;
+ }
+
+ public String getConfigArgs() {
+ return configArgs;
+ }
+
+ public void setUId(String UId) {
+ this.UId = UId;
+ }
+
+ public String getUId() {
+ return UId;
+ }
+
+ public void setToDoType(String toDoType) {
+ this.toDoType = toDoType;
+ }
+
+ public String getToDoType() {
+ return toDoType;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) return true;
+ if (!(o instanceof ApplicationInfo)) return false;
+
+ ApplicationInfo that = (ApplicationInfo) o;
+
+ if (getAppId() != null ? !getAppId().equals(that.getAppId()) : that.getAppId() != null) return false;
+ return !(getUId() != null ? !getUId().equals(that.getUId()) : that.getUId() != null);
+
+ }
+
+ @Override
+ public int hashCode() {
+ int result = getAppId() != null ? getAppId().hashCode() : 0;
+ result = 31 * result + (getUId() != null ? getUId().hashCode() : 0);
+ return result;
+ }
+}
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ConfigArgsDescriber.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ConfigArgsDescriber.java
new file mode 100644
index 000000000..0af77200f
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/ConfigArgsDescriber.java
@@ -0,0 +1,13 @@
+package com.ai.cloud.skywalking.alarm.model;
+
+public class ConfigArgsDescriber {
+ private int period;
+
+ public int getPeriod() {
+ return period;
+ }
+
+ public void setPeriod(int period) {
+ this.period = period;
+ }
+}
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/UserInfo.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/UserInfo.java
new file mode 100644
index 000000000..74398d7b6
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/model/UserInfo.java
@@ -0,0 +1,25 @@
+package com.ai.cloud.skywalking.alarm.model;
+
+import java.util.List;
+
+public class UserInfo {
+ private String userId;
+
+ private List applicationInfos;
+
+ public UserInfo(String userId) {
+ this.userId = userId;
+ }
+
+ public String getUserId() {
+ return userId;
+ }
+
+ public List getApplicationInfos() {
+ return applicationInfos;
+ }
+
+ public void setApplicationInfos(List applicationInfos) {
+ this.applicationInfos = applicationInfos;
+ }
+}
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/redis/RedisUtil.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/redis/RedisUtil.java
new file mode 100644
index 000000000..ec16c8d26
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/redis/RedisUtil.java
@@ -0,0 +1,128 @@
+package com.ai.cloud.skywalking.alarm.redis;
+
+import com.ai.cloud.skywalking.alarm.conf.Config;
+import com.ai.cloud.skywalking.alarm.model.ApplicationInfo;
+import org.apache.commons.pool2.impl.GenericObjectPoolConfig;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import redis.clients.jedis.Jedis;
+import redis.clients.jedis.JedisPool;
+import redis.clients.jedis.exceptions.JedisConnectionException;
+
+import java.util.ArrayList;
+import java.util.List;
+
+public class RedisUtil {
+ private static Logger logger = LogManager.getLogger(RedisUtil.class);
+ private static JedisPool jedisPool;
+ private static String[] config;
+ private static RedisInspector connector = new RedisInspector();
+ private static Object lock = new Object();
+
+ static {
+ GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig();
+ String redisServerConfig = Config.Alarm.REDIS_SERVER;
+ if (redisServerConfig == null || redisServerConfig.length() <= 0) {
+ logger.error("Redis server is not setting. Switch off alarm module. ");
+ } else {
+ config = redisServerConfig.split(":");
+ if (config.length != 2) {
+ logger.error("Redis server address is illegal setting, need to be 'ip:port'. Switch off alarm module. ");
+ Config.Alarm.ALARM_OFF_FLAG = true;
+ } else {
+ jedisPool = new JedisPool(genericObjectPoolConfig, config[0],
+ Integer.valueOf(config[1]));
+ // Test connect redis.
+ Jedis jedis = null;
+ try {
+ jedis = jedisPool.getResource();
+ } catch (Exception e) {
+ handleFailedToConnectRedisServerException(e);
+ logger.error("can't connect to redis["
+ + Config.Alarm.REDIS_SERVER + "]", e);
+ } finally {
+ if (jedis != null) {
+ jedis.close();
+ }
+ }
+ }
+ }
+ }
+
+ public static List getAlarmMessage(ApplicationInfo applicationInfo) {
+ Jedis jedis = null;
+ try {
+ jedis = jedisPool.getResource();
+ return new ArrayList(jedis.hgetAll(generateAlarmKey(applicationInfo)).values());
+ } catch (Exception e) {
+ handleFailedToConnectRedisServerException(e);
+ logger.error("Failed to set data.", e);
+ } finally {
+ if (jedis != null) {
+ jedis.close();
+ }
+ }
+
+ return new ArrayList();
+ }
+
+ private static String generateAlarmKey(ApplicationInfo applicationInfo) {
+ return applicationInfo.getUId() + "-" + applicationInfo.getAppId() + "-"
+ + ((System.currentTimeMillis() / (10000 * 6))
+ - applicationInfo.getConfigArgsDescriber().getPeriod());
+ }
+
+ private static GenericObjectPoolConfig buildGenericObjectPoolConfig() {
+ GenericObjectPoolConfig genericObjectPoolConfig = new GenericObjectPoolConfig();
+ genericObjectPoolConfig.setTestOnBorrow(true);
+ genericObjectPoolConfig.setMaxIdle(Config.Alarm.REDIS_MAX_IDLE);
+ genericObjectPoolConfig.setMinIdle(Config.Alarm.REDIS_MIN_IDLE);
+ genericObjectPoolConfig.setMaxTotal(Config.Alarm.REDIS_MAX_TOTAL);
+ return genericObjectPoolConfig;
+ }
+
+ private static void handleFailedToConnectRedisServerException(Exception e) {
+ if (e instanceof JedisConnectionException) {
+ // 发生连接不上Redis
+ if (connector == null || !connector.isAlive()) {
+ synchronized (lock) {
+ if (!connector.isAlive()) {
+ // 启动巡检线程
+ connector.start();
+ }
+ }
+ }
+ }
+ }
+
+ private static class RedisInspector extends Thread {
+ @Override
+ public void run() {
+ logger.debug("Connecting to redis....");
+ Jedis jedis;
+ while (true) {
+ try {
+ jedisPool = new JedisPool(buildGenericObjectPoolConfig(),
+ config[0], Integer.valueOf(config[1]));
+ jedis = jedisPool.getResource();
+ jedis.get("ok");
+ break;
+ } catch (Exception e) {
+ if (e instanceof JedisConnectionException) {
+ try {
+ Thread.sleep(5000L);
+ } catch (InterruptedException e1) {
+ logger.error("Sleep failed", e);
+ }
+ continue;
+ }
+ }
+ }
+ logger.debug("Connected to redis success. Open alarm function.");
+ Config.Alarm.ALARM_OFF_FLAG = false;
+ // 清理当前线程
+ connector = null;
+ }
+ }
+
+}
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/zk/ZKUtil.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/zk/ZKUtil.java
new file mode 100644
index 000000000..31ee439de
--- /dev/null
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/zk/ZKUtil.java
@@ -0,0 +1,96 @@
+package com.ai.cloud.skywalking.alarm.zk;
+
+import com.ai.cloud.skywalking.alarm.AlarmServerRegisterWatcher;
+import com.ai.cloud.skywalking.alarm.conf.Config;
+import org.apache.curator.RetryPolicy;
+import org.apache.curator.framework.CuratorFramework;
+import org.apache.curator.framework.CuratorFrameworkFactory;
+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 java.util.ArrayList;
+import java.util.List;
+
+public class ZKUtil {
+
+ private static Logger logger = LogManager.getLogger(ZKUtil.class);
+
+ private static CuratorFramework client;
+
+ static {
+ try {
+ RetryPolicy retryPolicy = new ExponentialBackoffRetry(Config.ZKPath.RETRY_TIMEOUT,
+ Config.ZKPath.RETRY_TIMES);
+ CuratorFrameworkFactory.Builder builder = CuratorFrameworkFactory.builder().
+ connectString(Config.ZKPath.CONNECT_STR)
+ .connectionTimeoutMs(Config.ZKPath.CONNECT_TIMEOUT).retryPolicy(retryPolicy);
+ client = builder.build();
+ client.start();
+
+ createIfNotExists(getRegisterLockPath());
+ createIfNotExists(getUserLockPathPrefix());
+ } catch (Exception e) {
+ logger.error("Failed to connect zookeeper.", e);
+ System.exit(-1);
+ }
+ }
+
+ public static InterProcessMutex getRegisterLock() {
+ return new InterProcessMutex(client, getRegisterLockPath());
+ }
+
+ public static InterProcessMutex getProcessUserLock(String uid) {
+ return new InterProcessMutex(client, getUserLockPath(uid));
+ }
+
+ public static void register(String serverId) {
+ try {
+ client.create().creatingParentsIfNeeded().forPath(getRegisterServerPathPrefix() + "/" + serverId);
+ } catch (Exception e) {
+ logger.error("Failed to register server", e);
+ }
+ }
+
+ public static List selectAllServerIds() {
+ try {
+ return client.getChildren().forPath(getRegisterServerPathPrefix());
+ } catch (Exception e) {
+ logger.error("Failed to get children of Path[" + getRegisterServerPathPrefix() + "].", e);
+ }
+ return new ArrayList();
+ }
+
+ public static void watch(AlarmServerRegisterWatcher instance) {
+ try {
+ client.getChildren().usingWatcher(instance).forPath(getRegisterServerPathPrefix());
+ } catch (Exception e) {
+ logger.error("Failed to get children for path[" + getRegisterServerPathPrefix() + "].", e);
+ }
+ }
+
+
+ private static String getRegisterLockPath() {
+ return Config.ZKPath.NODE_PREFIX + Config.ZKPath.SERVER_REGISTER_LOCK_PATH;
+ }
+
+ private static String getUserLockPath(String uid) {
+ return getUserLockPathPrefix() + "/" + uid;
+ }
+
+ private static String getUserLockPathPrefix() {
+ return Config.ZKPath.NODE_PREFIX + Config.ZKPath.USER_REGISTER_LOCK_PATH;
+ }
+
+
+ private static void createIfNotExists(String path) throws Exception {
+ if (client.checkExists().forPath(path) == null) {
+ client.create().creatingParentsIfNeeded().forPath(path);
+ }
+ }
+
+ private static String getRegisterServerPathPrefix() {
+ return Config.ZKPath.NODE_PREFIX + Config.ZKPath.REGISTER_SERVER_PATH;
+ }
+}
diff --git a/skywalking-alarm/src/main/resources/log4j.xml b/skywalking-alarm/src/main/resources/log4j.xml
new file mode 100644
index 000000000..0da28a6d5
--- /dev/null
+++ b/skywalking-alarm/src/main/resources/log4j.xml
@@ -0,0 +1,19 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
+
diff --git a/skywalking-alarm/src/main/resources/log4j2.xml b/skywalking-alarm/src/main/resources/log4j2.xml
new file mode 100644
index 000000000..025d5b2fb
--- /dev/null
+++ b/skywalking-alarm/src/main/resources/log4j2.xml
@@ -0,0 +1,13 @@
+
+
+
+
+
+
+
+
+
+
+
+
+
\ No newline at end of file