From dbec9c194ac9512a9d19e0658f2a51aa4699c480 Mon Sep 17 00:00:00 2001 From: zhangxin10 Date: Thu, 3 Dec 2015 14:34:45 +0800 Subject: [PATCH] =?UTF-8?q?=E5=AE=8C=E6=88=90Redis=E5=91=8A=E8=AD=A6?= =?UTF-8?q?=E9=97=AE=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../buriedpoint/ThreadBuriedPointSender.java | 6 ++- .../skywalking/util/ContextGenerator.java | 8 +-- .../skywalking/util/TraceIdGenerator.java | 2 +- .../example/web/SaveOrderServlet.java | 7 +++ .../src/main/webapp/WEB-INF/web.xml | 0 .../ai/cloud/skywalking/protocol/Span.java | 15 ++++-- .../cloud/skywalking/protocol/SpanData.java | 4 ++ .../skywalking/reciever/conf/Config.java | 2 + .../persistance/PersistenceThread.java | 12 +---- .../storage/StorageChainController.java | 2 + .../reciever/storage/chain/AlarmChain.java | 38 ++++++++++++++ .../chain/alarm/AlarmMessageStorage.java | 8 ++- .../chain/alarm/RedisAccessController.java | 50 ++++--------------- .../src/main/resources/config.properties | 10 +++- 14 files changed, 96 insertions(+), 68 deletions(-) create mode 100644 skywalking-example/web-application/src/main/java/com/ai/cloud/skywalking/example/web/SaveOrderServlet.java create mode 100644 skywalking-example/web-application/src/main/webapp/WEB-INF/web.xml create mode 100644 skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/AlarmChain.java diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadBuriedPointSender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadBuriedPointSender.java index 96747745c..29bfa9934 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadBuriedPointSender.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadBuriedPointSender.java @@ -28,10 +28,12 @@ public class ThreadBuriedPointSender implements IBuriedPointSender { // 从ThreadLocal中取出上下文 final Span parentSpanData = Context.getLastSpan(); if (parentSpanData == null) { - spanData = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID); + spanData = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID, + Config.SkyWalking.USER_ID); } else { // 如果不为空,则将当前的Context存放到上下文 - spanData = new Span(parentSpanData.getTraceId(), Config.SkyWalking.APPLICATION_ID); + spanData = new Span(parentSpanData.getTraceId(), Config.SkyWalking.APPLICATION_ID, + Config.SkyWalking.USER_ID); spanData.setParentLevel(parentSpanData.getParentLevel() + "." + parentSpanData.getLevelId()); spanData.setLevelId(threadSeqId); } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ContextGenerator.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ContextGenerator.java index 904dfa633..b945728f8 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ContextGenerator.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ContextGenerator.java @@ -31,10 +31,10 @@ public final class ContextGenerator { // 校验传入的参数是否为空,如果为空,则新创建一个 if (context == null || StringUtil.isEmpty(context.getTraceId())) { // 不存在,新创建一个Context - spanData = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID); + spanData = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID, Config.SkyWalking.USER_ID); } else { // 如果不为空,则将当前的Context存放到上下文 - spanData = new Span(context.getTraceId(), context.getParentLevel(), context.getLevelId(), Config.SkyWalking.APPLICATION_ID); + spanData = new Span(context.getTraceId(), context.getParentLevel(), context.getLevelId(), Config.SkyWalking.APPLICATION_ID, Config.SkyWalking.USER_ID); } initNewSpanData(spanData, id); @@ -58,12 +58,12 @@ public final class ContextGenerator { // 2 校验Context,Context是否存在 if (parentSpan == null) { // 不存在,新创建一个Context - span = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID); + span = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID, Config.SkyWalking.USER_ID); } else { // 根据ParentContextData的TraceId和RPCID // LevelId是由SpanNode类的nextSubSpanLevelId字段进行初始化的. // 所以在这里不需要初始化 - span = new Span(parentSpan.getTraceId(), Config.SkyWalking.APPLICATION_ID); + span = new Span(parentSpan.getTraceId(), Config.SkyWalking.APPLICATION_ID, Config.SkyWalking.USER_ID); if (!StringUtil.isEmpty(parentSpan.getParentLevel())) { span.setParentLevel(parentSpan.getParentLevel() + "." + parentSpan.getLevelId()); } else { diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/TraceIdGenerator.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/TraceIdGenerator.java index f27743362..f122516a8 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/TraceIdGenerator.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/TraceIdGenerator.java @@ -10,6 +10,6 @@ public final class TraceIdGenerator { } public static String generate() { - return UUID.randomUUID().toString().replaceAll("-", "") + Config.SkyWalking.USER_ID; + return UUID.randomUUID().toString().replaceAll("-", ""); } } diff --git a/skywalking-example/web-application/src/main/java/com/ai/cloud/skywalking/example/web/SaveOrderServlet.java b/skywalking-example/web-application/src/main/java/com/ai/cloud/skywalking/example/web/SaveOrderServlet.java new file mode 100644 index 000000000..a5e8060b9 --- /dev/null +++ b/skywalking-example/web-application/src/main/java/com/ai/cloud/skywalking/example/web/SaveOrderServlet.java @@ -0,0 +1,7 @@ +package com.ai.cloud.skywalking.example.web; + +/** + * Created by astraea on 2015/11/26. + */ +public class SaveOrderServlet { +} diff --git a/skywalking-example/web-application/src/main/webapp/WEB-INF/web.xml b/skywalking-example/web-application/src/main/webapp/WEB-INF/web.xml new file mode 100644 index 000000000..e69de29bb diff --git a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java index aeadf5912..24cdea067 100644 --- a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java +++ b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java @@ -13,15 +13,17 @@ public class Span extends SpanData { public Span() { } - public Span(String traceId, String applicationID) { + public Span(String traceId, String applicationID, String userId) { this.traceId = traceId; this.applicationId = applicationID; + this.userId = userId; } - public Span(String traceId, String parentLevelId, int levelId, String applicationID) { + public Span(String traceId, String parentLevelId, int levelId, String applicationID, String userId) { this.traceId = traceId; this.applicationId = applicationID; this.parentLevel = parentLevelId; + this.userId = userId; this.levelId = levelId; } @@ -47,6 +49,7 @@ public class Span extends SpanData { NEW_LINE_CHARACTER_PATTERN); processNo = fieldValues[12].trim(); applicationId = fieldValues[13].trim(); + userId = fieldValues[14].trim(); this.originData = originData; } @@ -114,7 +117,13 @@ public class Span extends SpanData { } if (isNonBlank(applicationId)) { - toStringValue.append(applicationId); + toStringValue.append(applicationId + SPAN_FIELD_SPILT_PATTERN); + } else { + toStringValue.append(" " + SPAN_FIELD_SPILT_PATTERN); + } + + if (isNonBlank(userId)) { + toStringValue.append(userId); } else { toStringValue.append(" " + SPAN_FIELD_SPILT_PATTERN); } diff --git a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java index 01dd054e8..556d3f4fb 100644 --- a/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java +++ b/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SpanData.java @@ -22,6 +22,7 @@ public abstract class SpanData { protected String processNo = ""; protected String applicationId = ""; protected String originData = ""; + protected String userId; public String getTraceId() { @@ -124,4 +125,7 @@ public abstract class SpanData { return applicationId; } + public String getUserId() { + return userId; + } } 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 a43c917ab..cf2791f5b 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 @@ -97,5 +97,7 @@ public class Config { public static int REDIS_MIN_IDLE = 1; public static int REDIS_MAX_TOTAL = 20; + + public static boolean ALARM_OFF_FLAG = false; } } \ No newline at end of file diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java index 1b0dafd11..87c29e986 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/persistance/PersistenceThread.java @@ -72,17 +72,7 @@ public class PersistenceThread extends Thread { offset += 1; continue; } -// -// if (tmpData == null || tmpData.length() <= 0) { -// MemoryRegister -// .instance() -// .doRegisterStatus( -// new FileRegisterEntry( -// file1.getName(), -// offset, -// FileRegisterEntry.FileRegisterEntryStatus.UNREGISTER)); -// break; -// } + ServerHealthCollector.getCurrentHeathReading(null) .updateData( ServerHeathReading.INFO, 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 d76e7c8e8..1142b9b35 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 @@ -3,6 +3,7 @@ package com.ai.cloud.skywalking.reciever.storage; import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.conf.Constants; +import com.ai.cloud.skywalking.reciever.storage.chain.AlarmChain; import com.ai.cloud.skywalking.reciever.storage.chain.SaveToHBaseChain; import com.ai.cloud.skywalking.reciever.storage.chain.SaveToMySQLChain; import org.apache.logging.log4j.LogManager; @@ -21,6 +22,7 @@ public class StorageChainController { static { if (STORAGE_TYPE.equalsIgnoreCase("hbase")) { + chainArray.add(new AlarmChain()); chainArray.add(new SaveToHBaseChain()); } else if (STORAGE_TYPE.equalsIgnoreCase("mysql")) { chainArray.add(new SaveToMySQLChain()); 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 new file mode 100644 index 000000000..50873a9ca --- /dev/null +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/AlarmChain.java @@ -0,0 +1,38 @@ +package com.ai.cloud.skywalking.reciever.storage.chain; + +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.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import java.util.List; + +public class AlarmChain implements IStorageChain { + + private Logger logger = LogManager.getLogger(AlarmChain.class); + + @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( + generateAlarmKey(span) + , span.getTraceId()); + } + chain.doChain(spans); + } + + private String generateAlarmKey(Span span) { + return span.getUserId() + "-" + + span.getApplicationId() + "-" + + (System.currentTimeMillis() / (10000 * 6)); + } +} 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 c03858a3a..960e60d71 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,5 +1,6 @@ package com.ai.cloud.skywalking.reciever.storage.chain.alarm; +import com.ai.cloud.skywalking.reciever.conf.Config; import redis.clients.jedis.Jedis; import java.util.Collection; @@ -54,11 +55,8 @@ public class AlarmMessageStorage { public static void saveAlarmMessage(String key, String traceId) { + if (Config.Alarm.ALARM_OFF_FLAG) + return; 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/storage/chain/alarm/RedisAccessController.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/alarm/RedisAccessController.java index 50e154580..c41d62a01 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 @@ -6,25 +6,29 @@ 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 Object lock = new Object(); 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; } + 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])); @@ -42,48 +46,12 @@ public class RedisAccessController { Jedis jedis = null; try { jedis = jedisPool.getResource(); - jedis.connect(); return executor.exec(jedis); } catch (Exception e) { - jedisPool = null; - logger.error("Failed to connect redis", e); - // 启动备用Redis - if (jedisPool == null) { - synchronized (lock) { - if (jedisPool == null) { - // 生成备份Redis的Redis Client Pool - GenericObjectPoolConfig genericObjectPoolConfig = buildGenericObjectPoolConfig(); - String bakRedisServerConfig = Config.Alarm.BAK_REDIS_SERVER_CONFIG; - if (bakRedisServerConfig == null || bakRedisServerConfig.length() <= 0) { - logger.error("Bak Redis server config is null."); - } - - String[] config = bakRedisServerConfig.split(":"); - if (config.length != 2) { - logger.error("Bak Redis server config is illegal"); - } - - jedisPool = - new JedisPool(genericObjectPoolConfig, config[0], Integer.valueOf(config[1])); - try { - jedis = jedisPool.getResource(); - jedis.connect(); - jedis.get("ok"); - } catch (Exception ex) { - logger.error("Failed to connect bak redis server.", ex); - // 备份Redis的都失败了,没有想好怎么提示 - //System.exit(-1); - } finally { - if (jedis != null) { - jedis.close(); - } - } - } - } - // 重新再获取Redis Client返回执行 - jedis = jedisPool.getResource(); - jedis.connect(); - return executor.exec(jedis); + 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; } } finally { if (jedis != null) { diff --git a/skywalking-server/src/main/resources/config.properties b/skywalking-server/src/main/resources/config.properties index d63507b9a..2e29564e5 100644 --- a/skywalking-server/src/main/resources/config.properties +++ b/skywalking-server/src/main/resources/config.properties @@ -45,4 +45,12 @@ hbaseconfig.client_port=29181 #告警失效时间 alarm.alarm_expire_seconds=3600000 #Redis配置 -alarm.redis_server_config=127.0.0.1:6379 +alarm.redis_server_config=127.0.0.1:16379 +#Redis最大空闲数量 +alarm.edis_max_idle=10 +#Redis最小空闲数量 +alarm.edis_min_idle=1 +#Redis最大个数 +alarm.edis_max_total=20 +#是否关闭告警 +alarm.larm_off_flag=false \ No newline at end of file