From 72e430ef9a7e375b3bfa23d87c6e1dd1d7bc035f Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Thu, 3 Dec 2015 16:19:40 +0800 Subject: [PATCH] =?UTF-8?q?1.=20=E5=88=A0=E9=99=A4=E6=97=A0=E7=94=A8?= =?UTF-8?q?=E7=9A=84=E9=85=8D=E7=BD=AE=202.=20=E5=AE=9E=E7=8E=B0=E5=BD=93?= =?UTF-8?q?=E8=BF=9E=E6=8E=A5=E4=B8=8D=E4=B8=8ARedis=E7=9A=84=E6=97=B6?= =?UTF-8?q?=E5=80=99=EF=BC=8C=E5=BC=80=E5=90=AF=E5=B7=A1=E6=A3=80Redis?= =?UTF-8?q?=E7=9A=84=E7=BA=BF=E7=A8=8B=E3=80=82=203.=20=E5=A2=9E=E5=8A=A0?= =?UTF-8?q?=E5=AF=B9Redis=E7=9A=84Client=E4=B8=BA=E7=A9=BA=E6=A3=80?= =?UTF-8?q?=E9=AA=8C=EF=BC=8C=E5=8E=9F=E5=9B=A0=EF=BC=9A=E5=BD=93=E5=8F=91?= =?UTF-8?q?=E7=94=9F=E5=BC=82=E5=B8=B8=E6=97=B6=EF=BC=8CRedis=E7=9A=84Clie?= =?UTF-8?q?nt=E4=BC=9A=E8=BF=94=E5=9B=9Enull=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../skywalking/reciever/conf/Config.java | 2 - .../chain/alarm/AlarmMessageStorage.java | 41 ++++----- .../chain/alarm/RedisAccessController.java | 91 ++++++++++++++----- 3 files changed, 88 insertions(+), 46 deletions(-) diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java index cf2791f5b..3b8268a47 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java @@ -90,8 +90,6 @@ public class Config { public static String REDIS_SERVER_CONFIG = "127.0.0.1:6379"; - public static String BAK_REDIS_SERVER_CONFIG = "127.0.0.1:6379"; - public static int REDIS_MAX_IDLE = 10; public static int REDIS_MIN_IDLE = 1; 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 index 960e60d71..f873f679d 100644 --- 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 @@ -1,19 +1,24 @@ 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 java.util.Collection; - 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) { - return jedis.hexists(key, field); + if (jedis != null) { + return jedis.hexists(key, field); + } + return false; } }); } @@ -22,7 +27,10 @@ public class AlarmMessageStorage { return RedisAccessController.redis(new RedisAccessController.Executor() { @Override public Boolean exec(Jedis jedis) { - return jedis.exists(key); + if (jedis != null) { + return jedis.exists(key); + } + return false; } }); } @@ -31,32 +39,23 @@ public class AlarmMessageStorage { return RedisAccessController.redis(new RedisAccessController.Executor() { @Override public Long exec(Jedis jedis) { - if (exist(key, field)) { - return null; + if (jedis != null) { + if (exist(key, field)) { + return null; + } + jedis.hset(key, field, value); + return jedis.expire(key, ALARM_EXPIRE_SECONDS); } - jedis.hset(key, field, value); - return jedis.expire(key, ALARM_EXPIRE_SECONDS); + return null; } }); } - private static Collection get(final String key) { - return RedisAccessController.redis(new RedisAccessController.Executor>() { - @Override - public Collection exec(Jedis jedis) { - if (!exist(key)) { - return null; - } - return jedis.hgetAll(key).values(); - } - }); - } - - 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 index c41d62a01..9ae0058c8 100644 --- 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 @@ -12,6 +12,9 @@ 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(); @@ -19,26 +22,28 @@ public class RedisAccessController { if (redisServerConfig == null || redisServerConfig.length() <= 0) { logger.error("Redis server config is null."); Config.Alarm.ALARM_OFF_FLAG = true; - } - - - String[] config = redisServerConfig.split(":"); - if (config.length != 2) { - logger.error("Redis server config is illegal"); - Config.Alarm.ALARM_OFF_FLAG = true; - } - - - jedisPool = - new JedisPool(genericObjectPoolConfig, config[0], Integer.valueOf(config[1])); - - // Test connect redis. - RedisAccessController.redis(new Executor() { - @Override - public String exec(Jedis jedis) { - return jedis.get("ok"); + } 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; + } + }); } - }); + } } @@ -48,17 +53,26 @@ public class RedisAccessController { jedis = jedisPool.getResource(); return executor.exec(jedis); } catch (Exception e) { - logger.error("Failed to set data.", e); if (e instanceof JedisConnectionException) { - logger.error("Failed to connect redis. close alarm function.", e); - Config.Alarm.ALARM_OFF_FLAG = true; + // 发生连接不上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; } @@ -71,6 +85,37 @@ public class RedisAccessController { 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);