parent
aa059fbbeb
commit
d5b85dfb61
|
|
@ -72,6 +72,11 @@
|
|||
<artifactId>mysql-connector-java</artifactId>
|
||||
<version>5.1.37</version>
|
||||
</dependency>
|
||||
<dependency>
|
||||
<groupId>com.zaxxer</groupId>
|
||||
<artifactId>HikariCP</artifactId>
|
||||
<version>2.4.3</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
<build>
|
||||
<plugins>
|
||||
|
|
|
|||
|
|
@ -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<String> 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<String> users = AlarmMessageDao.selectAllUserIds();
|
||||
// 获取当前所有的注册的处理线程
|
||||
List<String> 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<String> 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<String> users = AlarmMessageDao.selectAllUserIds();
|
||||
logger.info("Query to the {} pending user from the database", users.size());
|
||||
// 将用户重新分配给服务
|
||||
List<String> 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<String> allocationUser(List<String> registeredThreads,
|
||||
List<String> userIds) throws Exception {
|
||||
List<String> realRedistributeThread = new ArrayList<String>();
|
||||
Set<String> sortThreadIds = new HashSet<String>(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<String> allocationUser(List<String> registeredThreads,
|
||||
List<String> userIds) throws Exception {
|
||||
List<String> realRedistributeThread = new ArrayList<String>();
|
||||
Set<String> sortThreadIds = new HashSet<String>(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<String> acquireAllRegisteredThread() throws Exception {
|
||||
return ZKUtil.getChildren(Config.ZKPath.REGISTER_SERVER_PATH);
|
||||
}
|
||||
}
|
||||
return realRedistributeThread;
|
||||
}
|
||||
|
||||
private boolean checkAllProcessStatus(List<String> registeredThreadIds,
|
||||
ProcessThreadStatus status) throws Exception {
|
||||
String registerPathPrefix = Config.ZKPath.REGISTER_SERVER_PATH + "/";
|
||||
for (String threadId : registeredThreadIds) {
|
||||
private List<String> 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<String> 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<String> 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<String> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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.");
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue