From ddc9885103350a60f8782bb60f4effa27649c2ec Mon Sep 17 00:00:00 2001 From: ascrutae Date: Wed, 17 Aug 2016 22:21:38 +0800 Subject: [PATCH 1/3] =?UTF-8?q?=E4=BF=AE=E5=A4=8DDubbo=E7=9A=84viewpoint?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../analysis/config/HBaseTableMetaData.java | 2 +- .../dubbo/MonitorFilterInterceptor.java | 24 ++++++++++++------- .../resources/spring/springmvc-servlet.xml | 4 ++-- 3 files changed, 18 insertions(+), 12 deletions(-) diff --git a/skywalking-analysis/src/main/java/com/a/eye/skywalking/analysis/config/HBaseTableMetaData.java b/skywalking-analysis/src/main/java/com/a/eye/skywalking/analysis/config/HBaseTableMetaData.java index 77939840c..bfd6bc884 100644 --- a/skywalking-analysis/src/main/java/com/a/eye/skywalking/analysis/config/HBaseTableMetaData.java +++ b/skywalking-analysis/src/main/java/com/a/eye/skywalking/analysis/config/HBaseTableMetaData.java @@ -7,7 +7,7 @@ public class HBaseTableMetaData { * @author wusheng */ public final static class TABLE_CALL_CHAIN { - public static final String TABLE_NAME = "sw-call-chain"; + public static final String TABLE_NAME = "trace-data"; public static final String FAMILY_NAME = "call-chain"; } diff --git a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/a/eye/skywalking/plugin/dubbo/MonitorFilterInterceptor.java b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/a/eye/skywalking/plugin/dubbo/MonitorFilterInterceptor.java index ac879a831..aee8314bc 100644 --- a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/a/eye/skywalking/plugin/dubbo/MonitorFilterInterceptor.java +++ b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/a/eye/skywalking/plugin/dubbo/MonitorFilterInterceptor.java @@ -33,8 +33,8 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto boolean isConsumer = rpcContext.isConsumerSide(); context.set("isConsumer", isConsumer); if (isConsumer) { - ContextData - contextData = new RPCClientInvokeMonitor().beforeInvoke(createIdentification(invoker, invocation)); + ContextData contextData = + new RPCClientInvokeMonitor().beforeInvoke(createIdentification(invoker, invocation, true)); String contextDataStr = contextData.toString(); //追加参数 @@ -72,7 +72,7 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto contextData = new ContextData(contextDataStr); } - new RPCServerInvokeMonitor().beforeInvoke(contextData, createIdentification(invoker, invocation)); + new RPCServerInvokeMonitor().beforeInvoke(contextData, createIdentification(invoker, invocation, false)); } } @@ -85,9 +85,9 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto dealException(result.getException(), context); } - if (isConsumer(context)){ + if (isConsumer(context)) { new RPCClientInvokeMonitor().afterInvoke(); - }else{ + } else { new RPCServerInvokeMonitor().afterInvoke(); } @@ -95,25 +95,31 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto } @Override - public void handleMethodException(Throwable t, EnhancedClassInstanceContext context, InstanceMethodInvokeContext interceptorContext) { + public void handleMethodException(Throwable t, EnhancedClassInstanceContext context, + InstanceMethodInvokeContext interceptorContext) { dealException(t, context); } - private boolean isConsumer(EnhancedClassInstanceContext context){ + private boolean isConsumer(EnhancedClassInstanceContext context) { return (boolean) context.get("isConsumer"); } private void dealException(Throwable t, EnhancedClassInstanceContext context) { if (isConsumer(context)) { - new RPCClientInvokeMonitor().occurException(t); + new RPCClientInvokeMonitor().occurException(t); } else { new RPCServerInvokeMonitor().occurException(t); } } - private static Identification createIdentification(Invoker invoker, Invocation invocation) { + private static Identification createIdentification(Invoker invoker, Invocation invocation, boolean isConsumer) { StringBuilder viewPoint = new StringBuilder(); + if (isConsumer) { + viewPoint.append("comsumer:"); + } else { + viewPoint.append("provider:"); + } viewPoint.append(invoker.getUrl().getProtocol() + "://"); viewPoint.append(invoker.getUrl().getHost()); viewPoint.append(":" + invoker.getUrl().getPort()); diff --git a/skywalking-webui/src/main/resources/spring/springmvc-servlet.xml b/skywalking-webui/src/main/resources/spring/springmvc-servlet.xml index 70fa9d303..94a8b4d6a 100644 --- a/skywalking-webui/src/main/resources/spring/springmvc-servlet.xml +++ b/skywalking-webui/src/main/resources/spring/springmvc-servlet.xml @@ -12,7 +12,7 @@ - + @@ -59,4 +59,4 @@ - \ No newline at end of file + From b39f00a58071d4f14234420193e3673ad8e9f012 Mon Sep 17 00:00:00 2001 From: ascrutae Date: Thu, 18 Aug 2016 07:56:05 +0800 Subject: [PATCH 2/3] =?UTF-8?q?=E5=B0=86=E6=89=80=E6=9C=89=E7=9A=84?= =?UTF-8?q?=E7=BA=BF=E7=A8=8B=E6=94=B9=E4=B8=BA=E5=AE=88=E6=8A=A4=E7=BA=BF?= =?UTF-8?q?=E7=A8=8B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../alarm/AlarmMessageProcessThread.java | 1 + .../skywalking/alarm/UserInfoCoordinator.java | 1 + .../alarm/UsersChangedDetectionThread.java | 4 + .../a/eye/skywalking/buffer/BufferGroup.java | 1 + .../selfexamination/SDKHealthCollector.java | 1 + .../sender/DataSenderFactoryWithBalance.java | 1 + .../reciever/buffer/AppendEOFFlagThread.java | 1 + .../reciever/buffer/DataBufferThread.java | 1 + .../peresistent/PersistenceThread.java | 1 + .../RegisterPersistenceThread.java | 1 + .../ackspan/alarm/AlarmRedisConnector.java | 1 + .../ServerHealthCollector.java | 117 +++++++++--------- 12 files changed, 73 insertions(+), 58 deletions(-) diff --git a/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/AlarmMessageProcessThread.java b/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/AlarmMessageProcessThread.java index 449d441bc..8761d9835 100644 --- a/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/AlarmMessageProcessThread.java +++ b/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/AlarmMessageProcessThread.java @@ -37,6 +37,7 @@ public class AlarmMessageProcessThread extends Thread { public AlarmMessageProcessThread() { // 初始化生成ThreadId threadId = UUID.randomUUID().toString(); + this.setDaemon(true); } @Override diff --git a/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UserInfoCoordinator.java b/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UserInfoCoordinator.java index 7b658ae78..03dd491d9 100644 --- a/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UserInfoCoordinator.java +++ b/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UserInfoCoordinator.java @@ -31,6 +31,7 @@ public class UserInfoCoordinator extends Thread { private boolean isCoordinator = false; public UserInfoCoordinator() { + this.setDaemon(true); } @Override diff --git a/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UsersChangedDetectionThread.java b/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UsersChangedDetectionThread.java index 6fb1a6d28..e2f71f123 100644 --- a/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UsersChangedDetectionThread.java +++ b/skywalking-alarm/src/main/java/com/a/eye/skywalking/alarm/UsersChangedDetectionThread.java @@ -19,6 +19,10 @@ public class UsersChangedDetectionThread extends Thread { private String userIdsEncryptedStr; private Logger logger = LogManager.getLogger(UsersChangedDetectionThread.class); + public UsersChangedDetectionThread() { + this.setDaemon(true); + } + public void run() { while (true) { try { diff --git a/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/buffer/BufferGroup.java b/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/buffer/BufferGroup.java index aef25a12e..2f286450a 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/buffer/BufferGroup.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/buffer/BufferGroup.java @@ -61,6 +61,7 @@ public class BufferGroup { super("ConsumerWorker"); this.start = start; this.end = end; + this.setDaemon(true); } @Override diff --git a/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/selfexamination/SDKHealthCollector.java b/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/selfexamination/SDKHealthCollector.java index 6e38bf66f..ac91bc648 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/selfexamination/SDKHealthCollector.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/selfexamination/SDKHealthCollector.java @@ -18,6 +18,7 @@ public class SDKHealthCollector extends Thread { private SDKHealthCollector() { super("HealthCollector"); + this.setDaemon(true); } public static void init() { diff --git a/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/sender/DataSenderFactoryWithBalance.java b/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/sender/DataSenderFactoryWithBalance.java index ba09daf2b..c35001e71 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/sender/DataSenderFactoryWithBalance.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/a/eye/skywalking/sender/DataSenderFactoryWithBalance.java @@ -102,6 +102,7 @@ public class DataSenderFactoryWithBalance { public static class DataSenderChecker extends Thread { public DataSenderChecker() { super("Data-Sender-Checker"); + this.setDaemon(true); } @Override diff --git a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/AppendEOFFlagThread.java b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/AppendEOFFlagThread.java index 7f87645eb..738da1052 100644 --- a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/AppendEOFFlagThread.java +++ b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/AppendEOFFlagThread.java @@ -19,6 +19,7 @@ class AppendEOFFlagThread extends Thread { super("AppendEOFFlagThread"); this.dataBufferFiles = dataBufferFiles; this.countDownLatch = countDownLatch; + this.setDaemon(true); } @Override diff --git a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/DataBufferThread.java b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/DataBufferThread.java index 9e8a54825..9cab70f4f 100644 --- a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/DataBufferThread.java +++ b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/buffer/DataBufferThread.java @@ -26,6 +26,7 @@ public class DataBufferThread extends Thread { public DataBufferThread(int threadIdx) { super("DataBufferThread_" + threadIdx); + this.setDaemon(true); } @Override diff --git a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/PersistenceThread.java b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/PersistenceThread.java index 1eb4db18f..37bd733ee 100644 --- a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/PersistenceThread.java +++ b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/PersistenceThread.java @@ -24,6 +24,7 @@ public class PersistenceThread extends Thread { public PersistenceThread(int trdIndex) { super("PersistentThread" + trdIndex); + this.setDaemon(true); } @Override diff --git a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/RegisterPersistenceThread.java b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/RegisterPersistenceThread.java index 61a186fb0..b630b78c0 100644 --- a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/RegisterPersistenceThread.java +++ b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/peresistent/RegisterPersistenceThread.java @@ -30,6 +30,7 @@ public class RegisterPersistenceThread extends Thread { Config.RegisterPersistence.REGISTER_FILE_PARENT_DIRECTORY, Config.RegisterPersistence.REGISTER_FILE_NAME); bakOffsetFile = new File( Config.RegisterPersistence.REGISTER_FILE_PARENT_DIRECTORY, Config.RegisterPersistence.REGISTER_BAK_FILE_NAME); + this.setDaemon(true); } @Override diff --git a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/processor/ackspan/alarm/AlarmRedisConnector.java b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/processor/ackspan/alarm/AlarmRedisConnector.java index ebf237038..a08512d23 100644 --- a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/processor/ackspan/alarm/AlarmRedisConnector.java +++ b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/processor/ackspan/alarm/AlarmRedisConnector.java @@ -57,6 +57,7 @@ public class AlarmRedisConnector { Config.Alarm.ALARM_OFF_FLAG = true; } } + this.setDaemon(true); } private RedisInspector connect() { diff --git a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/selfexamination/ServerHealthCollector.java b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/selfexamination/ServerHealthCollector.java index c9ba49ecf..f73f8c95e 100644 --- a/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/selfexamination/ServerHealthCollector.java +++ b/skywalking-server/src/main/java/com/a/eye/skywalking/reciever/selfexamination/ServerHealthCollector.java @@ -10,66 +10,67 @@ import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; public class ServerHealthCollector extends Thread { - private Logger logger = LogManager.getLogger(ServerHealthCollector.class); + private Logger logger = LogManager.getLogger(ServerHealthCollector.class); - private static Map heathReadings = new ConcurrentHashMap(); + private static Map heathReadings = new ConcurrentHashMap(); - private ServerHealthCollector(){ - super("ServerHealthCollector"); - } - - public static void init(){ - new ServerHealthCollector().start(); - } - - public static ServerHeathReading getCurrentHeathReading(String extraId) { - String id = getId(extraId); - if (!heathReadings.containsKey(id)) { - synchronized (heathReadings) { - if (!heathReadings.containsKey(id)) { - if(heathReadings.keySet().size() > 5000){ - throw new RuntimeException("use ServerHealthCollector illegal. There is an overflow trend of Server Health Collector Report Data."); - } - heathReadings.put(id, new ServerHeathReading(id)); - } - } - } - return heathReadings.get(id); - } + private ServerHealthCollector() { + super("ServerHealthCollector"); + this.setDaemon(true); + } - private static String getId(String extraId) { - return "SkyWalkingServer,M:" + MachineUtil.getHostDesc() + ",P:" - + MachineUtil.getProcessNo() + ",T:" - + Thread.currentThread().getName() + "(" - + Thread.currentThread().getId() + ")" - + (extraId == null ? "" : ",extra:" + extraId); - } + public static void init() { + new ServerHealthCollector().start(); + } - @Override - public void run() { - while (true) { - try { - Map heathReadingsSnapshot = heathReadings; - heathReadings = new ConcurrentHashMap(); - String[] keyList = heathReadingsSnapshot.keySet().toArray(new String[0]); - Arrays.sort(keyList); - StringBuilder log = new StringBuilder(); - log.append("\n---------Server Health Collector Report---------\n"); - for(String key : keyList){ - log.append(heathReadingsSnapshot.get(key)).append("\n"); - } - log.append("------------------------------------------------\n"); - - logger.info(log); - - try { - Thread.sleep(Config.HealthCollector.REPORT_INTERVAL); - } catch (InterruptedException e) { - logger.warn("sleep error.", e); - } - } catch (Throwable t) { - logger.error("ServerHealthCollector report error.", t); - } - } - } + public static ServerHeathReading getCurrentHeathReading(String extraId) { + String id = getId(extraId); + if (!heathReadings.containsKey(id)) { + synchronized (heathReadings) { + if (!heathReadings.containsKey(id)) { + if (heathReadings.keySet().size() > 5000) { + throw new RuntimeException( + "use ServerHealthCollector illegal. There is an overflow trend of Server Health Collector Report Data."); + } + heathReadings.put(id, new ServerHeathReading(id)); + } + } + } + return heathReadings.get(id); + } + + private static String getId(String extraId) { + return "SkyWalkingServer,M:" + MachineUtil.getHostDesc() + ",P:" + MachineUtil.getProcessNo() + ",T:" + Thread + .currentThread().getName() + "(" + Thread.currentThread().getId() + ")" + (extraId == null ? + "" : + ",extra:" + extraId); + } + + @Override + public void run() { + while (true) { + try { + Map heathReadingsSnapshot = heathReadings; + heathReadings = new ConcurrentHashMap(); + String[] keyList = heathReadingsSnapshot.keySet().toArray(new String[0]); + Arrays.sort(keyList); + StringBuilder log = new StringBuilder(); + log.append("\n---------Server Health Collector Report---------\n"); + for (String key : keyList) { + log.append(heathReadingsSnapshot.get(key)).append("\n"); + } + log.append("------------------------------------------------\n"); + + logger.info(log); + + try { + Thread.sleep(Config.HealthCollector.REPORT_INTERVAL); + } catch (InterruptedException e) { + logger.warn("sleep error.", e); + } + } catch (Throwable t) { + logger.error("ServerHealthCollector report error.", t); + } + } + } } From cd7fc294a4ac486c0726a327fa3a31e178b27da6 Mon Sep 17 00:00:00 2001 From: ascrutae Date: Thu, 18 Aug 2016 10:31:15 +0800 Subject: [PATCH 3/3] =?UTF-8?q?=E5=A2=9E=E5=8A=A0=E6=8E=88=E6=9D=83?= =?UTF-8?q?=E6=96=87=E4=BB=B6=E7=9A=84=E5=88=9D=E5=A7=8B=E5=8C=96=E8=84=9A?= =?UTF-8?q?=E6=9C=AC?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- skywalking-webui/src/main/sql/table.mysql | 4 ++++ 1 file changed, 4 insertions(+) diff --git a/skywalking-webui/src/main/sql/table.mysql b/skywalking-webui/src/main/sql/table.mysql index 92e42504e..d6a54a75c 100644 --- a/skywalking-webui/src/main/sql/table.mysql +++ b/skywalking-webui/src/main/sql/table.mysql @@ -243,6 +243,10 @@ INSERT INTO `auth_file_config` (`config_id`, `key`, `value0`, `value1`, `key_des INSERT INTO `auth_file_config` (`config_id`, `key`, `value0`, `value1`, `key_desc`, `sts`) VALUES ('14', 'buffer.pool_size', '5', '5', 'Buffer池的最大长度', 'A'); INSERT INTO `auth_file_config` (`config_id`, `key`, `value0`, `value1`, `key_desc`, `sts`) VALUES ('15', 'senderchecker.check_polling_time', '200', '200', '发送检查线程检查周期', 'A'); INSERT INTO `auth_file_config` (`config_id`, `key`, `value0`, `value1`, `key_desc`, `sts`) VALUES ('16', 'skywalking.charset', 'UTF-8', 'UTF-8', 'skywalking数据编码', 'A'); +INSERT INTO `auth_file_config` (`config_id`,`key`,`value0`,`value1`,`value2`,`value3`,`value4`,`key_desc`,`sts`) VALUES ('17','plugin.customlocalmethodinterceptorplugin.is_enable','false','false',NULL,NULL,NULL,'自定义本地方法插件是否开启','A'); +INSERT INTO `auth_file_config` (`config_id`,`key`,`value0`,`value1`,`value2`,`value3`,`value4`,`key_desc`,`sts`) VALUES ('18','plugin.customlocalmethodinterceptorplugin.package_prefix','','',NULL,NULL,NULL,'自定义插件拦截的包前缀','A'); +INSERT INTO `auth_file_config` (`config_id`,`key`,`value0`,`value1`,`value2`,`value3`,`value4`,`key_desc`,`sts`) VALUES ('19','plugin.customlocalmethodinterceptorplugin.record_param_enable','false','false',NULL,NULL,NULL,'自定义插件是否记录入参',NULL); + # alter table since 2016-4-8 ALTER TABLE `application_info`