From d5b85dfb6177f4fe0b69d69f59ae431ed0c87161 Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Mon, 14 Dec 2015 15:57:49 +0800 Subject: [PATCH] =?UTF-8?q?1.=20=E4=BF=AE=E6=94=B9=E6=95=B0=E6=8D=AE?= =?UTF-8?q?=E5=BA=93=E8=BF=9E=E6=8E=A5=E6=B1=A0=202.=20=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=E6=89=93=E5=8D=B0=E6=97=A5=E5=BF=97?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- skywalking-alarm/pom.xml | 5 + .../skywalking/alarm/UserInfoCoordinator.java | 379 +++++++++--------- .../cloud/skywalking/alarm/conf/Config.java | 4 + .../skywalking/alarm/util/DBConnectUtil.java | 42 +- 4 files changed, 222 insertions(+), 208 deletions(-) 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."); + } + } }