diff --git a/skywalking-server/pom.xml b/skywalking-server/pom.xml index 3cef3b490..366f7f342 100644 --- a/skywalking-server/pom.xml +++ b/skywalking-server/pom.xml @@ -46,6 +46,16 @@ skywalking-protocol 1.0-SNAPSHOT + + redis.clients + jedis + 2.8.0 + + + com.google.code.gson + gson + 2.2.2 + diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/alarm/AlarmOperator.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/alarm/AlarmOperator.java new file mode 100644 index 000000000..40848e74d --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/alarm/AlarmOperator.java @@ -0,0 +1,65 @@ +package com.ai.cloud.skywalking.reciever.alarm; + +import com.ai.cloud.skywalking.reciever.alarm.redis.RedisInitializer; +import redis.clients.jedis.Jedis; + +import java.util.Collection; + +import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXPIRE_SECONDS; + +public class AlarmOperator { + + private static boolean exist(final String key, final String field) { + return RedisInitializer.redis(new RedisInitializer.Executor() { + @Override + public Boolean exec(Jedis jedis) { + return jedis.hexists(key, field); + } + }); + } + + private static boolean exist(final String key) { + return RedisInitializer.redis(new RedisInitializer.Executor() { + @Override + public Boolean exec(Jedis jedis) { + return jedis.exists(key); + } + }); + } + + private static Long set(final String key, final String field, final String value) { + return RedisInitializer.redis(new RedisInitializer.Executor() { + @Override + public Long exec(Jedis jedis) { + if (exist(key, field)) { + return null; + } + jedis.hset(key, field, value); + return jedis.expire(key, ALARM_EXPIRE_SECONDS); + } + }); + + } + + private static Collection get(final String key) { + return RedisInitializer.redis(new RedisInitializer.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) { + set(key, traceId, ""); + } + + public static Collection getAlarmMessage(String key) { + return get(key); + } + +} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/alarm/redis/RedisInitializer.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/alarm/redis/RedisInitializer.java new file mode 100644 index 000000000..bda7e226c --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/alarm/redis/RedisInitializer.java @@ -0,0 +1,69 @@ +package com.ai.cloud.skywalking.reciever.alarm.redis; + +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; + +public class RedisInitializer { + + private static Logger logger = LogManager.getLogger(RedisInitializer.class); + private static JedisPool jedisPool; + private static final int REDIS_MAX_IDLE = 10; + private static final int REDIS_MIN_IDLE = 1; + private static final int REDIS_MAX_TOTAL = 20; + + static { + GenericObjectPoolConfig genericObjectPoolConfig = new GenericObjectPoolConfig(); + genericObjectPoolConfig.setTestOnBorrow(true); + genericObjectPoolConfig.setMaxIdle(REDIS_MAX_IDLE); + genericObjectPoolConfig.setMinIdle(REDIS_MIN_IDLE); + genericObjectPoolConfig.setMaxTotal(REDIS_MAX_TOTAL); + + String redisServerConfig = Config.Alarm.REDIS_SERVER_CONFIG; + if (redisServerConfig == null || redisServerConfig.length() <= 0) { + logger.error("Redis server config is null."); + } + + String[] config = redisServerConfig.split(":"); + if (config.length != 2) { + logger.error("Redis server config is illegal"); + } + + jedisPool = + new JedisPool(genericObjectPoolConfig, config[0], Integer.valueOf(config[1])); + + // Test connect redis. + RedisInitializer.redis(new Executor() { + @Override + public String exec(Jedis jedis) { + return jedis.get("ok"); + } + }); + + } + + public static T redis(Executor executor) { + Jedis jedis = null; + try { + jedis = jedisPool.getResource(); + return executor.exec(jedis); + } catch (Exception e) { + logger.error("Failed to connect redis", e); + //TODO 启动备用Redis + } finally { + if (jedis != null) { + jedis.close(); + } + } + + return null; + } + + + public interface Executor { + R exec(Jedis jedis); + } +} 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 25094507a..27141583c 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 @@ -83,4 +83,11 @@ public class Config { public static String STORAGE_TYPE = "hbase"; } + + public static class Alarm { + + public static int ALARM_EXPIRE_SECONDS = 1000 * 60 * 90; + + public static String REDIS_SERVER_CONFIG = "127.0.0.1:6379"; + } } \ No newline at end of file diff --git a/skywalking-server/src/main/resources/config.properties b/skywalking-server/src/main/resources/config.properties index 893ef4e5f..d63507b9a 100644 --- a/skywalking-server/src/main/resources/config.properties +++ b/skywalking-server/src/main/resources/config.properties @@ -40,4 +40,9 @@ hbaseconfig.family_column_name=call-chain #hbase zk quorum hbaseconfig.zk_hostname=10.1.235.197,10.1.235.198,10.1.235.199 #hbase zk port -hbaseconfig.client_port=29181 \ No newline at end of file +hbaseconfig.client_port=29181 + +#告警失效时间 +alarm.alarm_expire_seconds=3600000 +#Redis配置 +alarm.redis_server_config=127.0.0.1:6379 diff --git a/skywalking-server/src/test/java/com/ai/cloud/skywalking/reciever/alarm/AlarmOperatorTest.java b/skywalking-server/src/test/java/com/ai/cloud/skywalking/reciever/alarm/AlarmOperatorTest.java new file mode 100644 index 000000000..32328362c --- /dev/null +++ b/skywalking-server/src/test/java/com/ai/cloud/skywalking/reciever/alarm/AlarmOperatorTest.java @@ -0,0 +1,26 @@ +package com.ai.cloud.skywalking.reciever.alarm; + +import com.ai.cloud.skywalking.reciever.conf.Config; +import com.ai.cloud.skywalking.reciever.conf.ConfigInitializer; +import org.junit.Test; + +import java.util.Properties; +import java.util.UUID; + +import static org.junit.Assert.assertEquals; + +public class AlarmOperatorTest { + + @Test + public void testSaveAlarmMessage() throws Exception { + Properties properties = new Properties(); + properties.load(AlarmOperatorTest.class.getResourceAsStream("/config.properties")); + ConfigInitializer.initialize(properties, Config.class); + + String key = UUID.randomUUID().toString(); + AlarmOperator.saveAlarmMessage(key, UUID.randomUUID().toString()); + assertEquals(1, AlarmOperator.getAlarmMessage(key).size()); + } + + +} \ No newline at end of file