diff --git a/skywalking-alarm/pom.xml b/skywalking-alarm/pom.xml
index 1a0056a3f..8c0b1dd5d 100644
--- a/skywalking-alarm/pom.xml
+++ b/skywalking-alarm/pom.xml
@@ -72,6 +72,11 @@
mysql-connector-java
5.1.37
+
+ com.zaxxer
+ HikariCP
+ 2.4.3
+
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 0ae85d09c..9eebcde80 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,222 +22,223 @@ import java.util.concurrent.TimeUnit;
public class UserInfoCoordinator extends Thread {
- private Logger logger = LogManager.getLogger(UserInfoCoordinator.class);
+ private Logger logger = LogManager.getLogger(UserInfoCoordinator.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 RegisterServerWatcher watcher = new RegisterServerWatcher();
+ private InterProcessMutex lock = new InterProcessMutex(
+ ZKUtil.getZkClient(), Config.ZKPath.COORDINATOR_PATH);
+ private boolean isCoordinator = false;
- public UserInfoCoordinator() {
- }
+ public UserInfoCoordinator() {
+ }
- @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;
- }
+ @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);
+ }
+ }
- // 检查是否有新服务注册或者在重分配过程做有新处理线程启动了
- if (!redistributing) {
- try {
- Thread.sleep(Config.Coordinator.CHECK_REDISTRIBUTE_INTERVAL);
- } catch (InterruptedException e) {
- logger.error("Sleep error", e);
- }
+ isCoordinator = true;
+ watcherRegisterServerPath();
+ redistributing = true;
+ }
- continue;
- }
+ // 检查是否有新服务注册或者在重分配过程做有新处理线程启动了
+ if (!redistributing) {
+ try {
+ Thread.sleep(Config.Coordinator.CHECK_REDISTRIBUTE_INTERVAL);
+ } catch (InterruptedException e) {
+ logger.error("Sleep error", e);
+ }
- redistributing = false;
+ continue;
+ }
- // 获取当前所有的注册的处理线程
- 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;
- }
- }
+ redistributing = false;
- // 查询当前有多少用户
- List users = AlarmMessageDao.selectAllUserIds();
+ // 获取当前所有的注册的处理线程
+ 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);
+ }
- // 将用户重新分配给服务
- List realRedistributeThread = allocationUser(
- registeredThreads, users);
+ if (retryTimes > 1000) {
+ logger.warn("checking all processors are free, waiting {}ms", Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL * retryTimes);
+ retryTimes = 0;
+ }
+ }
- // 修改状态(分配完成)
- 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);
- }
-
- if(retryTimes > 1000){
- logger.warn("checking all processors are busy, waiting {}ms", Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL * retryTimes);
- retryTimes = 0;
- }
- }
+ List users = AlarmMessageDao.selectAllUserIds();
+ logger.info("Query to the {} pending user from the database", users.size());
+ // 将用户重新分配给服务
+ List realRedistributeThread = allocationUser(
+ registeredThreads, users);
+ logger.info("All users are assigned to {} processing threads", realRedistributeThread.size());
+ // 修改状态(分配完成)
+ changeStatus(realRedistributeThread,
+ ProcessThreadStatus.REDISTRIBUTE_SUCCESS);
+ logger.info("Change state of {} processing threads to idle state", realRedistributeThread.size());
+ // 检查所有的服务是否都处于忙碌状态
+ while (!checkAllProcessStatus(realRedistributeThread,
+ ProcessThreadStatus.BUSY)) {
+ try {
+ Thread.sleep(Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL);
+ } catch (InterruptedException e) {
+ logger.error("Sleep failed", e);
+ }
- } catch (Exception e) {
- logger.error("Failed to coordinate, retry. ", e);
- releaseCoordinator();
- isCoordinator = false;
- }
- }
- }
+ if (retryTimes > 1000) {
+ logger.warn("checking all processors are busy, waiting {}ms", Config.Coordinator.CHECK_ALL_PROCESS_THREAD_INTERVAL * retryTimes);
+ retryTimes = 0;
+ }
+ }
- 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;
- }
- }
+ } catch (Exception e) {
+ logger.error("Failed to coordinate, retry. ", e);
+ releaseCoordinator();
+ isCoordinator = false;
+ }
+ }
+ }
- private void releaseCoordinator() {
- 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(
+ Config.Coordinator.RETRY_GET_COORDINATOR_LOCK_INTERVAL,
+ TimeUnit.SECONDS);
+ } catch (Exception e) {
+ logger.error("Failed to acquire lock .", e);
+ return false;
+ }
+ }
- private List allocationUser(List registeredThreads,
- List userIds) throws Exception {
- 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;
+ private void releaseCoordinator() {
+ if (lock != null && lock.isAcquiredInThisProcess()) {
+ try {
+ lock.release();
+ } catch (Exception e1) {
+ logger.error("Failed to release lock.", e1);
+ }
+ }
+ }
- if (end > userIds.size()) {
- end = userIds.size();
- }
+ private List allocationUser(List registeredThreads,
+ List userIds) throws Exception {
+ 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;
- 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);
+ if (end > userIds.size()) {
+ end = userIds.size();
+ }
- start = end;
- end += step;
- if (start >= userIds.size()) {
- break;
- }
- if (end > userIds.size()) {
- end = userIds.size();
- }
+ 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);
- }
- return realRedistributeThread;
- }
+ start = end;
+ end += step;
+ if (start >= userIds.size()) {
+ break;
+ }
+ if (end > userIds.size()) {
+ end = userIds.size();
+ }
- private List acquireAllRegisteredThread() throws Exception {
- return ZKUtil.getChildren(Config.ZKPath.REGISTER_SERVER_PATH);
- }
+ }
+ return realRedistributeThread;
+ }
- private boolean checkAllProcessStatus(List registeredThreadIds,
- ProcessThreadStatus status) throws Exception {
- String registerPathPrefix = Config.ZKPath.REGISTER_SERVER_PATH + "/";
- for (String threadId : registeredThreadIds) {
+ private List acquireAllRegisteredThread() throws Exception {
+ return ZKUtil.getChildren(Config.ZKPath.REGISTER_SERVER_PATH);
+ }
- if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/"
- + threadId))
- continue;
+ private boolean checkAllProcessStatus(List registeredThreadIds,
+ ProcessThreadStatus status) throws Exception {
+ String registerPathPrefix = Config.ZKPath.REGISTER_SERVER_PATH + "/";
+ for (String threadId : registeredThreadIds) {
- if (getProcessThreadStatus(registerPathPrefix, threadId) != status) {
- return false;
- }
- }
- return true;
- }
+ if (!ZKUtil.exists(Config.ZKPath.REGISTER_SERVER_PATH + "/"
+ + threadId))
+ continue;
- private ProcessThreadStatus getProcessThreadStatus(
- String registerPathPrefix, String threadId) throws Exception {
- 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());
- }
+ if (getProcessThreadStatus(registerPathPrefix, threadId) != status) {
+ return false;
+ }
+ }
+ return true;
+ }
- private void changeStatus(List registeredThreadIds,
- ProcessThreadStatus status) throws Exception {
- for (String threadId : registeredThreadIds) {
- ProcessUtil.changeProcessThreadStatus(threadId, status);
- }
- }
+ private ProcessThreadStatus getProcessThreadStatus(
+ String registerPathPrefix, String threadId) throws Exception {
+ 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());
+ }
- public class RegisterServerWatcher implements CuratorWatcher {
+ private void changeStatus(List registeredThreadIds,
+ ProcessThreadStatus status) throws Exception {
+ for (String threadId : registeredThreadIds) {
+ ProcessUtil.changeProcessThreadStatus(threadId, status);
+ }
+ }
- @Override
- public void process(WatchedEvent watchedEvent) {
- if (watchedEvent.getType() == Watcher.Event.EventType.NodeChildrenChanged) {
- redistributing = true;
- }
+ public class RegisterServerWatcher implements CuratorWatcher {
- watcherRegisterServerPath();
- }
- }
+ @Override
+ public void process(WatchedEvent watchedEvent) {
+ if (watchedEvent.getType() == Watcher.Event.EventType.NodeChildrenChanged) {
+ redistributing = true;
+ }
- 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);
- }
- }
+ watcherRegisterServerPath();
+ }
+ }
+
+ 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/conf/Config.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/conf/Config.java
index 7907cb3df..4ffb8594c 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
@@ -55,6 +55,10 @@ public class Config {
public static String URL = "jdbc:mysql://10.1.228.200:31306/test";
+ public static int MAX_IDLE = 1;
+
+ public static int MAX_POOL_SIZE = 20;
+
}
public static class Alarm {
diff --git a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/DBConnectUtil.java b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/DBConnectUtil.java
index 7cb587e84..836376e39 100644
--- a/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/DBConnectUtil.java
+++ b/skywalking-alarm/src/main/java/com/ai/cloud/skywalking/alarm/util/DBConnectUtil.java
@@ -1,36 +1,40 @@
package com.ai.cloud.skywalking.alarm.util;
import com.ai.cloud.skywalking.alarm.conf.Config;
+import com.zaxxer.hikari.HikariConfig;
+import com.zaxxer.hikari.HikariDataSource;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.sql.Connection;
-import java.sql.DriverManager;
import java.sql.SQLException;
public class DBConnectUtil {
private static Logger logger = LogManager.getLogger(DBConnectUtil.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);
- }
- }
+ private static HikariDataSource hikariDataSource;
public static Connection getConnection() {
- return con;
+ if (hikariDataSource == null) {
+ HikariConfig config = new HikariConfig();
+ config.setJdbcUrl(Config.DB.URL);
+ config.setUsername(Config.DB.USER_NAME);
+ config.setPassword(Config.DB.PASSWORD);
+ config.setDriverClassName(Config.DB.DRIVER_CLASS);
+ config.addDataSourceProperty("cachePrepStmts", "true");
+ config.addDataSourceProperty("prepStmtCacheSize", "250");
+ config.setMinimumIdle(Config.DB.MAX_IDLE);
+ config.setMaximumPoolSize(Config.DB.MAX_POOL_SIZE);
+ config.addDataSourceProperty("prepStmtCacheSqlLimit", "2048");
+ hikariDataSource = new HikariDataSource(config);
+ }
+ try {
+ return hikariDataSource.getConnection();
+ } catch (SQLException e) {
+ logger.error("Failed to get connection", e);
+ throw new RuntimeException("Cannot get connection.");
+ }
+
}
}