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);