1. 移除AlarmMessageStorage的lamba写法,因为lamba在现在情况中会降低代码的可读性,

2. 重构代码,将AlarmMessageStorage下沉到AlarmChain类中
This commit is contained in:
zhangxin10 2015-12-03 18:36:01 +08:00
parent 72e430ef9a
commit 2d0f88ccac
3 changed files with 119 additions and 192 deletions

View File

@ -4,26 +4,61 @@ import com.ai.cloud.skywalking.protocol.Span;
import com.ai.cloud.skywalking.reciever.conf.Config;
import com.ai.cloud.skywalking.reciever.storage.Chain;
import com.ai.cloud.skywalking.reciever.storage.IStorageChain;
import com.ai.cloud.skywalking.reciever.storage.chain.alarm.AlarmMessageStorage;
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.List;
public class AlarmChain implements IStorageChain {
import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXPIRE_SECONDS;
private Logger logger = LogManager.getLogger(AlarmChain.class);
public class AlarmChain implements IStorageChain {
private static Logger logger = LogManager.getLogger(AlarmChain.class);
private static JedisPool jedisPool;
private static String[] config;
private static Object lock = new Object();
private static RedisInspector connector;
static {
GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig();
String redisServerConfig = Config.Alarm.REDIS_SERVER_CONFIG;
if (redisServerConfig == null || redisServerConfig.length() <= 0) {
logger.error("Redis server config is null.");
Config.Alarm.ALARM_OFF_FLAG = true;
} else {
config = redisServerConfig.split(":");
if (config.length != 2) {
logger.error("Redis server config is illegal");
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("Failed to set data.", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
}
}
@Override
public void doChain(List<Span> spans, Chain chain) {
if (Config.Alarm.ALARM_OFF_FLAG) {
chain.doChain(spans);
return;
}
for (Span span : spans) {
if (span.getStatusCode() != 1)
continue;
AlarmMessageStorage.saveAlarmMessage(
saveAlarmMessage(
generateAlarmKey(span)
, span.getTraceId());
}
@ -35,4 +70,80 @@ public class AlarmChain implements IStorageChain {
+ span.getApplicationId() + "-"
+ (System.currentTimeMillis() / (10000 * 6));
}
private void saveAlarmMessage(String key, String traceId) {
if (Config.Alarm.ALARM_OFF_FLAG) {
return;
}
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
jedis.hset(key, traceId, "");
jedis.expire(key, ALARM_EXPIRE_SECONDS);
} catch (Exception e) {
handleFailedToConnectRedisServerException(e);
logger.error("Failed to set data.", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
}
private static class RedisInspector extends Thread {
@Override
public void run() {
logger.info("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.info("Connected to redis success. Open alarm function.");
Config.Alarm.ALARM_OFF_FLAG = false;
// 清理当前线程
connector = null;
}
}
private static void handleFailedToConnectRedisServerException(Exception e) {
if (e instanceof JedisConnectionException) {
// 发生连接不上Redis
if (connector == null || !connector.isAlive()) {
synchronized (lock) {
if (connector == null || !connector.isAlive()) {
// 启动巡检线程
connector = new RedisInspector();
connector.start();
}
}
}
}
}
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;
}
}

View File

@ -1,61 +0,0 @@
package com.ai.cloud.skywalking.reciever.storage.chain.alarm;
import com.ai.cloud.skywalking.reciever.conf.Config;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import redis.clients.jedis.Jedis;
import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXPIRE_SECONDS;
public class AlarmMessageStorage {
private static Logger logger = LogManager.getLogger(AlarmMessageStorage.class);
private static boolean exist(final String key, final String field) {
return RedisAccessController.redis(new RedisAccessController.Executor<Boolean>() {
@Override
public Boolean exec(Jedis jedis) {
if (jedis != null) {
return jedis.hexists(key, field);
}
return false;
}
});
}
private static boolean exist(final String key) {
return RedisAccessController.redis(new RedisAccessController.Executor<Boolean>() {
@Override
public Boolean exec(Jedis jedis) {
if (jedis != null) {
return jedis.exists(key);
}
return false;
}
});
}
private static Long set(final String key, final String field, final String value) {
return RedisAccessController.redis(new RedisAccessController.Executor<Long>() {
@Override
public Long exec(Jedis jedis) {
if (jedis != null) {
if (exist(key, field)) {
return null;
}
jedis.hset(key, field, value);
return jedis.expire(key, ALARM_EXPIRE_SECONDS);
}
return null;
}
});
}
public static void saveAlarmMessage(String key, String traceId) {
if (Config.Alarm.ALARM_OFF_FLAG)
return;
set(key, traceId, "");
}
}

View File

@ -1,123 +0,0 @@
package com.ai.cloud.skywalking.reciever.storage.chain.alarm;
import com.ai.cloud.skywalking.reciever.conf.Config;
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;
public class RedisAccessController {
private static Logger logger = LogManager.getLogger(RedisAccessController.class);
private static JedisPool jedisPool;
private static String[] config;
private static Object lock = new Object();
private static RedisConnector connector;
static {
GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig();
String redisServerConfig = Config.Alarm.REDIS_SERVER_CONFIG;
if (redisServerConfig == null || redisServerConfig.length() <= 0) {
logger.error("Redis server config is null.");
Config.Alarm.ALARM_OFF_FLAG = true;
} else {
config = redisServerConfig.split(":");
if (config.length != 2) {
logger.error("Redis server config is illegal");
Config.Alarm.ALARM_OFF_FLAG = true;
} else {
jedisPool =
new JedisPool(genericObjectPoolConfig, config[0],
Integer.valueOf(config[1]));
// Test connect redis.
RedisAccessController.redis(new Executor<String>() {
@Override
public String exec(Jedis jedis) {
// 对Redis Client为空校验
if (jedis != null) {
return jedis.get("ok");
}
return null;
}
});
}
}
}
public static <T> T redis(Executor<T> executor) {
Jedis jedis = null;
try {
jedis = jedisPool.getResource();
return executor.exec(jedis);
} catch (Exception e) {
if (e instanceof JedisConnectionException) {
// 发生连接不上Redis
if (connector == null || !connector.isAlive()) {
synchronized (lock) {
if (connector == null || !connector.isAlive()) {
// 启动巡检线程
connector = new RedisConnector();
connector.start();
}
}
}
}
logger.error("Failed to set data.", e);
} finally {
if (jedis != null) {
jedis.close();
}
}
// 当发生异常的时候返回的Redis的Client会是null
// 需要对Redis Client为空校验
return null;
}
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;
}
static class RedisConnector extends Thread {
@Override
public void run() {
logger.info("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.info("Connected to redis success. Open alarm function.");
Config.Alarm.ALARM_OFF_FLAG = false;
// 清理当前线程
connector = null;
}
}
public interface Executor<R> {
R exec(Jedis jedis);
}
}