From 2d0f88ccac08e1eb399965d53b975c4add66f387 Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Thu, 3 Dec 2015 18:36:01 +0800 Subject: [PATCH] =?UTF-8?q?1.=20=E7=A7=BB=E9=99=A4AlarmMessageStorage?= =?UTF-8?q?=E7=9A=84lamba=E5=86=99=E6=B3=95=EF=BC=8C=E5=9B=A0=E4=B8=BAlamb?= =?UTF-8?q?a=E5=9C=A8=E7=8E=B0=E5=9C=A8=E6=83=85=E5=86=B5=E4=B8=AD?= =?UTF-8?q?=E4=BC=9A=E9=99=8D=E4=BD=8E=E4=BB=A3=E7=A0=81=E7=9A=84=E5=8F=AF?= =?UTF-8?q?=E8=AF=BB=E6=80=A7=EF=BC=8C=202.=20=E9=87=8D=E6=9E=84=E4=BB=A3?= =?UTF-8?q?=E7=A0=81=EF=BC=8C=E5=B0=86AlarmMessageStorage=E4=B8=8B?= =?UTF-8?q?=E6=B2=89=E5=88=B0AlarmChain=E7=B1=BB=E4=B8=AD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../reciever/storage/chain/AlarmChain.java | 127 ++++++++++++++++-- .../chain/alarm/AlarmMessageStorage.java | 61 --------- .../chain/alarm/RedisAccessController.java | 123 ----------------- 3 files changed, 119 insertions(+), 192 deletions(-) delete mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/AlarmMessageStorage.java delete mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/RedisAccessController.java diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/AlarmChain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/AlarmChain.java index 50873a9ca..914045289 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/AlarmChain.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/AlarmChain.java @@ -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 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; + } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/AlarmMessageStorage.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/AlarmMessageStorage.java deleted file mode 100644 index f873f679d..000000000 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/AlarmMessageStorage.java +++ /dev/null @@ -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() { - @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() { - @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() { - @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, ""); - } - -} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/RedisAccessController.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/RedisAccessController.java deleted file mode 100644 index 9ae0058c8..000000000 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/RedisAccessController.java +++ /dev/null @@ -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() { - @Override - public String exec(Jedis jedis) { - // 对Redis Client为空校验 - if (jedis != null) { - return jedis.get("ok"); - } - return null; - } - }); - } - } - - } - - public static T redis(Executor 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 exec(Jedis jedis); - } -}