From 743b36f08551850b30ce6dbefaa9e72cf8abd7c3 Mon Sep 17 00:00:00 2001 From: wusheng Date: Mon, 29 Feb 2016 11:14:58 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=B8=BASDK=E5=A2=9E=E5=8A=A0=E5=A4=A7?= =?UTF-8?q?=E9=87=8F=E7=9A=84=E5=81=A5=E5=BA=B7=E6=8A=A5=E5=91=8A=E6=A3=80?= =?UTF-8?q?=E6=9F=A5=E4=BB=A3=E7=A0=81=E3=80=82=E5=AE=9A=E6=9C=9F=E8=BE=93?= =?UTF-8?q?=E5=87=BA=E5=81=A5=E5=BA=B7=E6=97=A5=E5=BF=97=EF=BC=8C=E6=96=B9?= =?UTF-8?q?=E4=BE=BF=E9=94=99=E8=AF=AF=E6=8E=92=E6=9F=A5=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- skywalking-api/pom.xml | 8 -- .../cloud/skywalking/buffer/BufferGroup.java | 4 + .../ai/cloud/skywalking/conf/AuthDesc.java | 4 + .../com/ai/cloud/skywalking/conf/Config.java | 5 ++ .../selfexamination/HeathReading.java | 72 ++++++++++++++++++ .../selfexamination/SDKHealthCollector.java | 76 +++++++++++++++++++ .../cloud/skywalking/sender/DataSender.java | 5 ++ .../sender/DataSenderFactoryWithBalance.java | 9 ++- .../sender/DataSenderWithCopies.java | 4 + 9 files changed, 178 insertions(+), 9 deletions(-) create mode 100644 skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HeathReading.java create mode 100644 skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/SDKHealthCollector.java diff --git a/skywalking-api/pom.xml b/skywalking-api/pom.xml index 93965ccae..fcf5ce51e 100644 --- a/skywalking-api/pom.xml +++ b/skywalking-api/pom.xml @@ -134,12 +134,4 @@ - - - - company-private-nexus-library-snapshots - company-private-nexus-library-snapshots - http://10.1.228.199:18081/nexus/content/repositories/snapshots/ - - diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java index b74b3f352..258499838 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/BufferGroup.java @@ -11,6 +11,8 @@ import org.apache.logging.log4j.Logger; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.conf.Constants; import com.ai.cloud.skywalking.protocol.Span; +import com.ai.cloud.skywalking.selfexamination.HeathReading; +import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector; import com.ai.cloud.skywalking.sender.DataSenderFactoryWithBalance; import com.ai.cloud.skywalking.util.AtomicRangeInteger; @@ -42,8 +44,10 @@ public class BufferGroup { logger.warn( "Group[{}] index[{}] data collision, discard old data.", groupName, i); + SDKHealthCollector.getCurrentHeathReading("BufferGroup").updateData(HeathReading.WARNING, "BufferGroup index[" + i + "] data collision, data been coverd."); } dataBuffer[i] = span; + SDKHealthCollector.getCurrentHeathReading("BufferGroup").updateData(HeathReading.INFO, "save span"); } class ConsumerWorker extends Thread { diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/AuthDesc.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/AuthDesc.java index 546ae127f..f03a81914 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/AuthDesc.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/AuthDesc.java @@ -1,11 +1,15 @@ package com.ai.cloud.skywalking.conf; +import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector; + public class AuthDesc { static boolean isAuth = false; static { ConfigInitializer.initialize(); ConfigValidator.validate(); + + SDKHealthCollector.init(); } public static boolean isAuth() { diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java index caccd32cc..cc8a95f15 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Config.java @@ -74,4 +74,9 @@ public class Config { public static long RETRY_FIND_CONNECTION_SENDER = 1000; } + + public static class HealthCollector { + // 默认健康检查上报时间 + public static long REPORT_INTERVAL = 5 * 60 * 1000L; + } } \ No newline at end of file diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HeathReading.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HeathReading.java new file mode 100644 index 000000000..ae33d8865 --- /dev/null +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HeathReading.java @@ -0,0 +1,72 @@ +package com.ai.cloud.skywalking.selfexamination; + +import java.util.HashMap; +import java.util.Map; + +public class HeathReading { + public static final String ERROR = "[ERROR]"; + public static final String WARNING = "[WARNING]"; + public static final String INFO = "[INFO]"; + + private String id; + + private Map datas = new HashMap(); + + /** + * 健康读数,只应该在工作线程中创建 + * + */ + public HeathReading(String id) { + this.id = id; + } + + public void updateData(String key, String newData){ + if(datas.containsKey(key)){ + datas.get(key).updateData(newData); + }else{ + datas.put(key, new HeathDetailData(newData)); + } + } + + @Override + public String toString(){ + StringBuilder sb = new StringBuilder(); + sb.append("id<").append(this.id).append(">\n"); + for(Map.Entry data : datas.entrySet()){ + sb.append(data.getKey()).append(data.getValue().toString()).append("\n"); + } + + //reset data + datas = new HashMap(); + return sb.toString(); + } + + class HeathDetailData{ + private String data; + + private long statusTime; + + HeathDetailData(String initialData){ + data = initialData; + statusTime = System.currentTimeMillis(); + } + + void updateData(String newData){ + data = newData; + statusTime = System.currentTimeMillis(); + } + + String getData() { + return data; + } + + long getStatusTime() { + return statusTime; + } + + @Override + public String toString(){ + return data + "(t:" + statusTime + ")"; + } + } +} diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/SDKHealthCollector.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/SDKHealthCollector.java new file mode 100644 index 000000000..d9c6c770f --- /dev/null +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/SDKHealthCollector.java @@ -0,0 +1,76 @@ +package com.ai.cloud.skywalking.selfexamination; + +import java.util.Arrays; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; + +import com.ai.cloud.skywalking.conf.AuthDesc; +import com.ai.cloud.skywalking.conf.Config; +import com.ai.cloud.skywalking.util.BuriedPointMachineUtil; + +public class SDKHealthCollector extends Thread { + private Logger logger = LogManager.getLogger(SDKHealthCollector.class); + + private static Map heathReadings = new ConcurrentHashMap(); + + private SDKHealthCollector(){ + super("HealthCollector"); + } + + public static void init(){ + if(AuthDesc.isAuth()){ + new SDKHealthCollector().start(); + } + } + + public static HeathReading getCurrentHeathReading(String extraId) { + String id = getId(extraId); + if (!heathReadings.containsKey(id)) { + synchronized (heathReadings) { + if (!heathReadings.containsKey(id)) { + heathReadings.put(id, new HeathReading(id)); + } + } + } + return heathReadings.get(id); + } + + private static String getId(String extraId) { + return "SDK-API,M:" + BuriedPointMachineUtil.getHostDesc() + ",P:" + + BuriedPointMachineUtil.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---------SDK 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("SDKHealthCollector report error.", t); + } + } + } +} diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java index ca9a98b34..1f0529e03 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSender.java @@ -18,6 +18,8 @@ import com.ai.cloud.io.netty.handler.codec.LengthFieldBasedFrameDecoder; import com.ai.cloud.io.netty.handler.codec.LengthFieldPrepender; import com.ai.cloud.io.netty.handler.codec.bytes.ByteArrayDecoder; import com.ai.cloud.io.netty.handler.codec.bytes.ByteArrayEncoder; +import com.ai.cloud.skywalking.selfexamination.HeathReading; +import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector; public class DataSender implements IDataSender { private EventLoopGroup group; @@ -71,12 +73,15 @@ public class DataSender implements IDataSender { try { if (channel != null && channel.isActive()) { channel.writeAndFlush(data.getBytes()); + SDKHealthCollector.getCurrentHeathReading("sender").updateData(HeathReading.INFO, "DataSender send data successfully."); return true; }else{ DataSenderFactoryWithBalance.unRegister(this); + SDKHealthCollector.getCurrentHeathReading("sender").updateData(HeathReading.WARNING, "DataSender channel isn't active. unregister sender."); } } catch (Exception e) { DataSenderFactoryWithBalance.unRegister(this); + SDKHealthCollector.getCurrentHeathReading("sender").updateData(HeathReading.WARNING, "DataSender channel broken. unregister sender."); } return false; diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java index 8a021aef5..be4e11071 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactoryWithBalance.java @@ -20,6 +20,8 @@ import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import com.ai.cloud.skywalking.conf.Config; +import com.ai.cloud.skywalking.selfexamination.HeathReading; +import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector; import com.ai.cloud.skywalking.util.StringUtil; public class DataSenderFactoryWithBalance { @@ -131,6 +133,7 @@ public class DataSenderFactoryWithBalance { unusedServerAddresses.add(tmpDataSender .getServerIp()); senderIterator.remove(); + SDKHealthCollector.getCurrentHeathReading("remove").updateData(HeathReading.INFO, "remove disconnected sender."); } } @@ -141,7 +144,7 @@ public class DataSenderFactoryWithBalance { break; } usingDataSender.add(newSender); - + SDKHealthCollector.getCurrentHeathReading("add").updateData(HeathReading.INFO, "add new sender."); } // try to switch. @@ -175,11 +178,15 @@ public class DataSenderFactoryWithBalance { .getServerIp()); unusedServerAddresses.add(toBeSwitchSender .getServerIp()); + SDKHealthCollector.getCurrentHeathReading("switch").updateData(HeathReading.INFO, "switch existed sender."); } } sleepTime = 0; } + + SDKHealthCollector.getCurrentHeathReading(null).updateData(HeathReading.INFO, "using available DataSender size:" + usingDataSender); } catch (Throwable e) { + SDKHealthCollector.getCurrentHeathReading(null).updateData(HeathReading.ERROR, "DataSenderChecker running failed:" + e.getMessage()); logger.error("DataSenderChecker running failed", e); } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java index c9e762517..c9ec369cd 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderWithCopies.java @@ -5,6 +5,9 @@ import static com.ai.cloud.skywalking.conf.Config.Sender.MAX_COPY_NUM; import java.util.HashSet; import java.util.Set; +import com.ai.cloud.skywalking.selfexamination.HeathReading; +import com.ai.cloud.skywalking.selfexamination.SDKHealthCollector; + /** * 带副本的数据发送器 * @@ -52,6 +55,7 @@ public class DataSenderWithCopies implements IDataSender { successNum++; } } + SDKHealthCollector.getCurrentHeathReading("DataSenderWithCopies").updateData(HeathReading.INFO, "DataSender send data with copynum=" + successNum + " successfully."); if (senders.size() == 1 && successNum == 1) { return true; } else if (successNum >= 2) {