From 5b23fd09c3cd37b60f853c16cb6ce77b32346fe2 Mon Sep 17 00:00:00 2001 From: ascrutae Date: Thu, 25 Feb 2016 15:54:13 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BF=AE=E5=A4=8DSpan=E5=9C=A8StorageChain?= =?UTF-8?q?=E4=B8=AD=E5=8F=91=E7=94=9F=E5=BC=82=E5=B8=B8=E6=97=B6=EF=BC=8C?= =?UTF-8?q?=E7=A8=8B=E5=BA=8F=E8=BF=9B=E5=85=A5=E6=AD=BB=E5=BE=AA=E7=8E=AF?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../skywalking/reciever/conf/Config.java | 2 + .../storage/StorageChainController.java | 5 +- .../reciever/storage/chain/AlarmChain.java | 247 +++++++++--------- .../storage/chain/SaveToHBaseChain.java | 36 +-- 4 files changed, 150 insertions(+), 140 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 a3f357e9a..2a7d7df2e 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 @@ -86,6 +86,8 @@ public class Config { public static long RETRY_STORAGE_WAIT_TIME = 50L; public static String STORAGE_TYPE = "hbase"; + + public static int RETRY_STORAGE_TIMES = 3; } public static class Alarm { diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java index 1142b9b35..889a9ee7e 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/StorageChainController.java @@ -49,12 +49,13 @@ public class StorageChainController { } } - while (true) { + int retryTimes = 0; + while(retryTimes++ < Config.StorageChain.RETRY_STORAGE_TIMES) { try { Chain chain = new Chain(chainArray); chain.doChain(spans); - break; } catch (Throwable e) { + // 主要的异常可能跟环境有关系,比如Redis,HBase,将会重试N次 try { Thread.sleep(Config.StorageChain.RETRY_STORAGE_WAIT_TIME); } catch (InterruptedException e1) { 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 605d20269..425aa1ae7 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 @@ -13,138 +13,139 @@ import redis.clients.jedis.exceptions.JedisConnectionException; import java.util.List; -import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXPIRE_SECONDS; import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXCEPTION_STACK_LENGTH; +import static com.ai.cloud.skywalking.reciever.conf.Config.Alarm.ALARM_EXPIRE_SECONDS; 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 = new RedisInspector();; + 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 = new RedisInspector(); + ; - static { - GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig(); - String redisServerConfig = Config.Alarm.REDIS_SERVER; - if (redisServerConfig == null || redisServerConfig.length() <= 0) { - logger.error("Redis server is not setting. Switch off alarm module. "); - Config.Alarm.ALARM_OFF_FLAG = true; - } else { - config = redisServerConfig.split(":"); - if (config.length != 2) { - logger.error("Redis server address is illegal setting, need to be 'ip:port'. Switch off alarm module. "); - 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("can't connect to redis[" - + Config.Alarm.REDIS_SERVER + "]", e); - } finally { - if (jedis != null) { - jedis.close(); - } - } - } - } - } + static { + GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig(); + String redisServerConfig = Config.Alarm.REDIS_SERVER; + if (redisServerConfig == null || redisServerConfig.length() <= 0) { + logger.error("Redis server is not setting. Switch off alarm module. "); + Config.Alarm.ALARM_OFF_FLAG = true; + } else { + config = redisServerConfig.split(":"); + if (config.length != 2) { + logger.error("Redis server address is illegal setting, need to be 'ip:port'. Switch off alarm module. "); + 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("can't connect to redis[" + + Config.Alarm.REDIS_SERVER + "]", e); + } finally { + if (jedis != null) { + jedis.close(); + } + } + } + } + } - @Override - public void doChain(List spans, Chain chain) { - for (Span span : spans) { - if (span.getStatusCode() != 1) - continue; - String exceptionStack = span.getExceptionStack(); - if(exceptionStack == null){ - exceptionStack = ""; - }else if(exceptionStack.length() > ALARM_EXCEPTION_STACK_LENGTH){ - exceptionStack = exceptionStack.substring(0, ALARM_EXCEPTION_STACK_LENGTH); - } - saveAlarmMessage(generateAlarmKey(span), span.getTraceId(), exceptionStack); - } - chain.doChain(spans); - } + @Override + public void doChain(List spans, Chain chain) { + for (Span span : spans) { + if (span.getStatusCode() != 1) + continue; + String exceptionStack = span.getExceptionStack(); + if (exceptionStack == null) { + exceptionStack = ""; + } else if (exceptionStack.length() > ALARM_EXCEPTION_STACK_LENGTH) { + exceptionStack = exceptionStack.substring(0, ALARM_EXCEPTION_STACK_LENGTH); + } + saveAlarmMessage(generateAlarmKey(span), span.getTraceId(), exceptionStack); + } + chain.doChain(spans); + } - private String generateAlarmKey(Span span) { - return span.getUserId() + "-" + span.getApplicationId() + "-" - + (System.currentTimeMillis() / (10000 * 6)); - } + private String generateAlarmKey(Span span) { + return span.getUserId() + "-" + span.getApplicationId() + "-" + + (System.currentTimeMillis() / (10000 * 6)); + } - private void saveAlarmMessage(String key, String traceId, String exceptionMsgOutline) { - if (Config.Alarm.ALARM_OFF_FLAG) { - return; - } + private void saveAlarmMessage(String key, String traceId, String exceptionMsgOutline) { + if (Config.Alarm.ALARM_OFF_FLAG) { + return; + } - Jedis jedis = null; - try { - jedis = jedisPool.getResource(); - jedis.hsetnx(key, traceId, exceptionMsgOutline); - jedis.expire(key, ALARM_EXPIRE_SECONDS); - } catch (Exception e) { - handleFailedToConnectRedisServerException(e); - logger.error("Failed to set data.", e); - } finally { - if (jedis != null) { - jedis.close(); - } - } - } + Jedis jedis = null; + try { + jedis = jedisPool.getResource(); + jedis.hsetnx(key, traceId, exceptionMsgOutline); + 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.debug("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.debug("Connected to redis success. Open alarm function."); - Config.Alarm.ALARM_OFF_FLAG = false; - // 清理当前线程 - connector = null; - } - } + private static class RedisInspector extends Thread { + @Override + public void run() { + logger.debug("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.debug("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.isAlive()) { - // 启动巡检线程 - connector.start(); - } - } - } - } - } + private static void handleFailedToConnectRedisServerException(Exception e) { + if (e instanceof JedisConnectionException) { + // 发生连接不上Redis + if (connector == null || !connector.isAlive()) { + synchronized (lock) { + if (!connector.isAlive()) { + // 启动巡检线程 + 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; - } + 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/SaveToHBaseChain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java index 05bcb6d62..6b228ea57 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToHBaseChain.java @@ -77,24 +77,30 @@ public class SaveToHBaseChain implements IStorageChain { if (spans == null || spans.size() <= 0) return; List puts = new ArrayList(); - Put put; + Put put = null; String columnName; for (Span span : spans) { - put = new Put(Bytes.toBytes(span.getTraceId()), getTSBySpanTraceId(span)); - if (StringUtils.isEmpty(span.getParentLevel().trim())) { - columnName = span.getLevelId() + ""; - if (span.isReceiver()) { - columnName = span.getLevelId() + "-S"; + try { + put = new Put(Bytes.toBytes(span.getTraceId()), getTSBySpanTraceId(span)); + if (StringUtils.isEmpty(span.getParentLevel().trim())) { + columnName = span.getLevelId() + ""; + if (span.isReceiver()) { + columnName = span.getLevelId() + "-S"; + } + put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName), + Bytes.toBytes(span.getOriginData())); + } else { + columnName = span.getParentLevel() + "." + span.getLevelId(); + if (span.isReceiver()) { + columnName = span.getParentLevel() + "." + span.getLevelId() + "-S"; + } + put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName), + Bytes.toBytes(span.getOriginData())); } - put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName), - Bytes.toBytes(span.getOriginData())); - } else { - columnName = span.getParentLevel() + "." + span.getLevelId(); - if (span.isReceiver()) { - columnName = span.getParentLevel() + "." + span.getLevelId() + "-S"; - } - put.addColumn(Bytes.toBytes(Config.HBaseConfig.FAMILY_COLUMN_NAME), Bytes.toBytes(columnName), - Bytes.toBytes(span.getOriginData())); + } catch (Exception e) { + // 不合规范的数据 + logger.error("Failed to convert Span[" + span.getTraceId() + "] to put Object", e); + continue; } puts.add(put); }