完成协调器功能。
协调器调用流程如下:
处理线程:
初始化:
运行:
注册服务(默认为空闲状态)
检查是否为忙碌状态
处理告警信息
检查是否分配线程的状态(重新分配状态)
释放用户锁
修改自身状态:(空闲状态)
检查分配线程的状态(分配完成状态)
重新获取待处理的用户
给用户加锁
修改自身状态 :(忙碌状态)
协调线程:
初始化:
初始化状态:未分配
运行:
检查是否有新服务注册
修改状态 (开始重新分配状态)
检查所有的服务是否都未待分配状态,
将用户重新分配给服务
修改状态(分配完成)
This commit is contained in:
parent
2e4d45d206
commit
4c806457e6
|
|
@ -76,11 +76,19 @@
|
|||
<build>
|
||||
<plugins>
|
||||
<plugin>
|
||||
<groupId>org.apache.maven.plugins</groupId>
|
||||
<artifactId>maven-compiler-plugin</artifactId>
|
||||
<configuration>
|
||||
<source>1.7</source>
|
||||
<target>1.7</target>
|
||||
<source>1.6</source>
|
||||
<target>1.6</target>
|
||||
<encoding>${project.build.sourceEncoding}</encoding>
|
||||
</configuration>
|
||||
</plugin>
|
||||
<plugin>
|
||||
<groupId>org.apache.maven.plugins</groupId>
|
||||
<artifactId>maven-resources-plugin</artifactId>
|
||||
<version>2.4.3</version>
|
||||
<configuration>
|
||||
<encoding>${project.build.sourceEncoding}</encoding>
|
||||
</configuration>
|
||||
</plugin>
|
||||
</plugins>
|
||||
|
|
|
|||
|
|
@ -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<InterProcessMutex> usersLocks;
|
||||
private List<String> processUserIds;
|
||||
private CoordinatorStatusWatcher watcher = new CoordinatorStatusWatcher();
|
||||
private Map<UserInfo, List<AlarmRule>> cacheRules = new HashMap<UserInfo, List<AlarmRule>>();
|
||||
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<AlarmRule> 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<UserInfo, List<AlarmRule>> 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<UserInfo, List<AlarmRule>> cacheRules) {
|
||||
cacheRules.clear();
|
||||
}
|
||||
|
||||
private void cacheProcessUser(List<String> processUserIds) {
|
||||
// TODO 需要重新获取
|
||||
UserInfo tmpUserInfo;
|
||||
List<AlarmRule> alarmRules;
|
||||
for (String userId : processUserIds) {
|
||||
tmpUserInfo = AlarmMessageDao.selectUser(userId);
|
||||
if (tmpUserInfo == null) {
|
||||
continue;
|
||||
}
|
||||
alarmRules = AlarmMessageDao.selectAlarmRulesByUserId(userId);
|
||||
cacheRules.put(tmpUserInfo, alarmRules);
|
||||
}
|
||||
}
|
||||
|
||||
private List<String> 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<InterProcessMutex>();
|
||||
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();
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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<String> 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<String> 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<String> users = AlarmMessageDao.selectAllUserIds();
|
||||
|
||||
//将用户重新分配给服务
|
||||
List<String> 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<String> users = AlarmMessageDao.selectAllUserIds();
|
||||
|
||||
//将用户重新分配给服务
|
||||
allocationUser(registeredThreads, users);
|
||||
|
||||
//修改状态(分配完成)
|
||||
changeStatus(registeredThreads, ProcessThreadStatus.REDISTRIBUTE_SUCCESS);
|
||||
}
|
||||
}
|
||||
|
||||
private void allocationUser(List<String> registeredThreads, List<String> 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<String> allocationUser(List<String> registeredThreads, List<String> userIds) {
|
||||
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;
|
||||
|
|
@ -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<String> acquireAllRegisteredThread() {
|
||||
return ZKUtil.getChildren(Config.ZKPath.REGISTER_SERVER_PATH);
|
||||
}
|
||||
|
||||
private boolean isAllProcessThreadFree(List<String> registeredThreadIds) {
|
||||
private boolean checkAllProcessStatus(List<String> 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);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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 {
|
||||
|
|
|
|||
|
|
@ -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"));
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<String> dealUserIds;
|
||||
|
||||
public String getStatus() {
|
||||
public int getStatus() {
|
||||
return status;
|
||||
}
|
||||
|
||||
public void setStatus(String status) {
|
||||
public void setStatus(int status) {
|
||||
this.status = status;
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -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<String>();
|
||||
}
|
||||
|
||||
public static List<String> 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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,9 +11,8 @@
|
|||
</appender>
|
||||
|
||||
<root>
|
||||
<priority value="debug" />
|
||||
<priority value="INFO" />
|
||||
<appender-ref ref="CONSOLE" />
|
||||
</root>
|
||||
|
||||
|
||||
</log4j:configuration>
|
||||
|
|
|
|||
|
|
@ -1,13 +1,14 @@
|
|||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<Configuration status="debug">
|
||||
<Configuration status="INFO">
|
||||
<Appenders>
|
||||
<Console name="Console" target="SYSTEM_OUT">
|
||||
<PatternLayout pattern="%d{HH:mm:ss.SSS} [%t] %-5level %logger{36} - %msg%n"/>
|
||||
</Console>
|
||||
</Appenders>
|
||||
<Loggers>
|
||||
<Root level="debug">
|
||||
<Root level="INFO">
|
||||
<AppenderRef ref="Console"/>
|
||||
</Root>
|
||||
<Logger name="org.apache.zookeeper" level="OFF"/>
|
||||
</Loggers>
|
||||
</Configuration>
|
||||
|
|
@ -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
|
||||
Loading…
Reference in New Issue