From 406d4e112cdc764cb228d3db42bb0b1eba3cb470 Mon Sep 17 00:00:00 2001 From: wusheng Date: Mon, 23 Nov 2015 16:37:26 +0800 Subject: [PATCH] =?UTF-8?q?1.=E4=BF=AE=E6=94=B9BusinessKeyAppender?= =?UTF-8?q?=E7=9A=84setBusinessKey2Trace=E6=96=B9=E6=B3=95=202.=E5=A2=9E?= =?UTF-8?q?=E5=8A=A0=E5=AE=A2=E6=88=B7=E7=AB=AF=E7=9A=84=E6=97=A5=E5=BF=97?= =?UTF-8?q?=E6=94=B6=E9=9B=86=EF=BC=88=E8=BF=98=E7=BC=BA=E5=B0=91=E4=B8=8A?= =?UTF-8?q?=E6=8A=A5=E7=A8=8B=E5=BA=8F=EF=BC=89=EF=BC=8C=E5=8C=85=E4=B8=BA?= =?UTF-8?q?=EF=BC=9Acom.ai.cloud.skywalking.selfexamination=203.=E7=A7=BB?= =?UTF-8?q?=E9=99=A4=E6=97=A0=E7=94=A8=E7=9A=84constant=E5=8C=85=EF=BC=8C?= =?UTF-8?q?=E5=8E=9F=E6=9C=89=E7=B1=BB=E6=9B=B4=E6=8D=A2=E5=8C=85=E5=90=8D?= =?UTF-8?q?=E3=80=82?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../skywalking/api/BusinessKeyAppender.java | 2 +- .../cloud/skywalking/buffer/BufferGroup.java | 17 ++--- .../{constants => conf}/Constants.java | 4 +- .../com/ai/cloud/skywalking/context/Span.java | 2 +- .../selfexamination/HealthCollector.java | 34 +++++++++ .../selfexamination/HeathReading.java | 74 +++++++++++++++++++ .../skywalking/sender/DataSenderFactory.java | 48 +++++++----- .../util/BuriedPointMachineUtil.java | 4 + .../skywalking/util/ContextGenerator.java | 2 +- .../skywalking/util/ExceptionHandleUtil.java | 2 +- .../plugin/spring/common/CallChainE.java | 2 +- .../plugin/spring/common/CallChainG.java | 2 +- .../plugin/spring/common/CallChainH.java | 2 +- .../plugin/spring/common/CallChainI.java | 2 +- .../plugin/spring/common/CallChainJ.java | 2 +- .../plugin/spring/common/CallChainK.java | 2 +- 16 files changed, 161 insertions(+), 40 deletions(-) rename skywalking-api/src/main/java/com/ai/cloud/skywalking/{constants => conf}/Constants.java (72%) create mode 100644 skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HealthCollector.java create mode 100644 skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HeathReading.java diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/BusinessKeyAppender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/BusinessKeyAppender.java index f135c9e451..2b9d7f4145 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/BusinessKeyAppender.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/BusinessKeyAppender.java @@ -16,7 +16,7 @@ public final class BusinessKeyAppender { * * @param businessKey */ - public static void trace(String businessKey) { + public static void setBusinessKey2Trace(String businessKey) { if (!AuthDesc.isAuth()) return; 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 6a8c55a1ee..c849e12d5e 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 @@ -3,6 +3,8 @@ package com.ai.cloud.skywalking.buffer; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Span; +import com.ai.cloud.skywalking.selfexamination.HealthCollector; +import com.ai.cloud.skywalking.selfexamination.HeathReading; import com.ai.cloud.skywalking.sender.DataSenderFactory; import java.util.concurrent.atomic.AtomicInteger; @@ -23,10 +25,9 @@ public class BufferGroup { int step = (int) Math.ceil(BUFFER_MAX_SIZE * 1.0 / MAX_CONSUMER); int start = 0, end = 0; - int i = 1; while (true) { if (end + step >= BUFFER_MAX_SIZE) { - new ConsumerWorker(groupName + "-consumer-" + (i++), start, BUFFER_MAX_SIZE).start(); + new ConsumerWorker(start, BUFFER_MAX_SIZE).start(); break; } end += step; @@ -38,8 +39,7 @@ public class BufferGroup { public void save(Span span) { int i = Math.abs(index.getAndIncrement() % BUFFER_MAX_SIZE); if (dataBuffer[i] != null) { - // TODO 需要上报 - System.out.println(span.getLevelId() + "在Group[" + groupName + "]的第" + i + "位冲突"); + HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.WARNING, span.getLevelId() + "在Group[" + groupName + "]的第" + i + "位冲突"); } dataBuffer[i] = span; } @@ -49,17 +49,11 @@ public class BufferGroup { private int end = BUFFER_MAX_SIZE; private ConsumerWorker(int start, int end) { + super("ConsumerWorker"); this.start = start; this.end = end; } - private ConsumerWorker(String threadName, int start, int end) { - super(threadName); - this.start = start; - this.end = end; - } - - @Override public void run() { StringBuilder data = new StringBuilder(); @@ -78,6 +72,7 @@ public class BufferGroup { logger.log(Level.ALL, "Sleep Failure"); } } + HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.INFO, "send buried-point data."); data = new StringBuilder(); } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/constants/Constants.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Constants.java similarity index 72% rename from skywalking-api/src/main/java/com/ai/cloud/skywalking/constants/Constants.java rename to skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Constants.java index 5418df2941..5b609dd8d3 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/constants/Constants.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Constants.java @@ -1,4 +1,4 @@ -package com.ai.cloud.skywalking.constants; +package com.ai.cloud.skywalking.conf; public class Constants { @@ -9,4 +9,6 @@ public class Constants { public static final String NEW_LINE_CHARACTER_PATTERN = "\\n"; public static final String EXCEPTION_SPILT_PATTERN = "^"; + + public static final String HEALTH_DATA_SPILT_PATTERN = "^~"; } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Span.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Span.java index 3f81667a63..ddf9333cfe 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Span.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Span.java @@ -1,6 +1,6 @@ package com.ai.cloud.skywalking.context; -import com.ai.cloud.skywalking.constants.Constants; +import com.ai.cloud.skywalking.conf.Constants; import com.ai.cloud.skywalking.util.StringUtil; public class Span { diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HealthCollector.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HealthCollector.java new file mode 100644 index 0000000000..9655c01925 --- /dev/null +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HealthCollector.java @@ -0,0 +1,34 @@ +package com.ai.cloud.skywalking.selfexamination; + +import java.util.HashMap; +import java.util.Map; + +import static com.ai.cloud.skywalking.conf.Config.SkyWalking.USER_ID; +import com.ai.cloud.skywalking.util.BuriedPointMachineUtil; + +public class HealthCollector extends Thread{ + private static Map heathReadings = new HashMap(); + + 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,U:" + USER_ID + ",M:" + BuriedPointMachineUtil.getHostDesc() +",P:" + BuriedPointMachineUtil.getProcessNo() + ",T:" + + Thread.currentThread().getName() + "(" + + Thread.currentThread().getId() + ")" + (extraId == null? "" : ",extra:" + extraId); + } + + @Override + public void run(){ + //TODO: 使用专有的端口,将数据上报给服务端,定时上报,默认应为分钟级别,降低服务端压力 + } +} 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 0000000000..1d7a894760 --- /dev/null +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/selfexamination/HeathReading.java @@ -0,0 +1,74 @@ +package com.ai.cloud.skywalking.selfexamination; + +import java.util.HashMap; +import java.util.Map; + +import com.ai.cloud.skywalking.conf.Constants; + +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(this.id); + sb.append(Constants.HEALTH_DATA_SPILT_PATTERN); + for(Map.Entry data : datas.entrySet()){ + sb.append(data.getKey()).append(Constants.HEALTH_DATA_SPILT_PATTERN).append(data.getValue().toString()).append(Constants.HEALTH_DATA_SPILT_PATTERN); + } + + //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 "d:" + data + Constants.HEALTH_DATA_SPILT_PATTERN + "t:" + statusTime; + } + } +} diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactory.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactory.java index d6ac3172f3..be15596278 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactory.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/sender/DataSenderFactory.java @@ -1,23 +1,29 @@ package com.ai.cloud.skywalking.sender; -import com.ai.cloud.skywalking.conf.Config; -import com.ai.cloud.skywalking.util.StringUtil; - -import java.io.IOException; -import java.net.InetSocketAddress; -import java.net.SocketAddress; -import java.util.*; -import java.util.concurrent.ThreadLocalRandom; -import java.util.logging.Level; -import java.util.logging.Logger; - import static com.ai.cloud.skywalking.conf.Config.Sender.CONNECT_PERCENT; import static com.ai.cloud.skywalking.conf.Config.Sender.RETRY_GET_SENDER_WAIT_INTERVAL; import static com.ai.cloud.skywalking.conf.Config.SenderChecker.CHECK_POLLING_TIME; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.net.SocketAddress; +import java.util.ArrayList; +import java.util.HashSet; +import java.util.Iterator; +import java.util.List; +import java.util.Set; +import java.util.concurrent.ThreadLocalRandom; +import java.util.logging.Level; +import java.util.logging.Logger; + +import com.ai.cloud.skywalking.conf.Config; +import com.ai.cloud.skywalking.selfexamination.HealthCollector; +import com.ai.cloud.skywalking.selfexamination.HeathReading; +import com.ai.cloud.skywalking.util.StringUtil; + public class DataSenderFactory { - //private static Logger logger = Logger.getLogger(DataSenderFactory.getSender().toString()); + private static Logger logger = Logger.getLogger(DataSenderFactory.getSender().toString()); private static Set socketAddresses = new HashSet(); private static Set unUsedSocketAddresses = new HashSet(); @@ -37,7 +43,7 @@ public class DataSenderFactory { socketAddresses.add(new InetSocketAddress(server[0], Integer.valueOf(server[1]))); } } catch (Exception e) { - // logger.log(Level.ALL, "Collection service configuration error."); + logger.log(Level.ALL, "Collection service configuration error.", e); System.exit(-1); } @@ -49,7 +55,7 @@ public class DataSenderFactory { try { Thread.sleep(RETRY_GET_SENDER_WAIT_INTERVAL); } catch (InterruptedException e) { - // logger.log(Level.ALL, "Sleep failure"); + logger.log(Level.ALL, "Sleep failure", e); } } return availableSenders.get(ThreadLocalRandom.current().nextInt(0, availableSenders.size())); @@ -60,8 +66,10 @@ public class DataSenderFactory { private int availableSize; public DataSenderChecker() { + super("DataSenderChecker"); + if (CONNECT_PERCENT <= 0 || CONNECT_PERCENT > 100) { - // logger.log(Level.ALL, "CONNECT_PERCENT must between 1 and 100"); + logger.log(Level.ALL, "CONNECT_PERCENT must between 1 and 100"); System.exit(-1); } availableSize = (int) Math.ceil(socketAddresses.size() * ((1.0 * CONNECT_PERCENT / 100) % 100)); @@ -92,25 +100,29 @@ public class DataSenderFactory { tmpScoketAddress = unUsedSocketAddressIterator.next(); if (availableSenders.size() >= availableSize) { + HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.INFO, "the num of available senders is enough."); break; } synchronized (lock) { try { - + HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.INFO, "increasing available senders."); availableSenders.add(new DataSender(tmpScoketAddress)); unUsedSocketAddresses.remove(tmpScoketAddress); } catch (IOException e) { } } - + } + + if (availableSenders.size() >= availableSize) { + HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.WARNING, "the num of available senders is not enough (" + availableSenders.size() + ")."); } try { Thread.sleep(CHECK_POLLING_TIME); } catch (InterruptedException e) { - //logger.log(Level.ALL, "Sleep Failure"); + logger.log(Level.ALL, "Sleep Failure"); } } } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/BuriedPointMachineUtil.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/BuriedPointMachineUtil.java index 71537327d4..bfd9fc661b 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/BuriedPointMachineUtil.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/BuriedPointMachineUtil.java @@ -52,6 +52,10 @@ public final class BuriedPointMachineUtil { } return hostName; } + + public static String getHostDesc(){ + return getHostName() + "/" + getHostIp(); + } private BuriedPointMachineUtil() { // Non 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 e36c61cedf..83b04950e1 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 @@ -48,7 +48,7 @@ public final class ContextGenerator { // 设置基本信息 spanData.setStartDate(System.currentTimeMillis()); spanData.setProcessNo(BuriedPointMachineUtil.getProcessNo()); - spanData.setAddress(BuriedPointMachineUtil.getHostName() + "/" + BuriedPointMachineUtil.getHostIp()); + spanData.setAddress(BuriedPointMachineUtil.getHostDesc()); } private static Span getSpanFromThreadLocal() { diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ExceptionHandleUtil.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ExceptionHandleUtil.java index 21026ca0f1..5df383102a 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ExceptionHandleUtil.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ExceptionHandleUtil.java @@ -1,6 +1,6 @@ package com.ai.cloud.skywalking.util; -import com.ai.cloud.skywalking.constants.Constants; +import com.ai.cloud.skywalking.conf.Constants; import com.ai.cloud.skywalking.context.Context; import com.ai.cloud.skywalking.context.Span; diff --git a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainE.java b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainE.java index 0a7125e754..06aad71fb9 100644 --- a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainE.java +++ b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainE.java @@ -12,7 +12,7 @@ public class CallChainE { @Tracing public void doBusiness() throws InterruptedException { Thread.sleep(ThreadLocalRandom.current().nextInt(10)); - BusinessKeyAppender.trace("key-value"); + BusinessKeyAppender.setBusinessKey2Trace("key-value"); callChainG.doBusiness(); Thread.sleep(ThreadLocalRandom.current().nextInt(10)); } diff --git a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainG.java b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainG.java index 24186d8370..eca0406a19 100644 --- a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainG.java +++ b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainG.java @@ -12,7 +12,7 @@ public class CallChainG { @Tracing public void doBusiness() throws InterruptedException { Thread.sleep(ThreadLocalRandom.current().nextInt(10)); - BusinessKeyAppender.trace("key-value"); + BusinessKeyAppender.setBusinessKey2Trace("key-value"); callChainH.doBusiness(); Thread.sleep(ThreadLocalRandom.current().nextInt(10)); diff --git a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainH.java b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainH.java index 04de8bbcf7..e905defad8 100644 --- a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainH.java +++ b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainH.java @@ -12,7 +12,7 @@ public class CallChainH { @Tracing public void doBusiness() throws InterruptedException { Thread.sleep(ThreadLocalRandom.current().nextInt(10)); - BusinessKeyAppender.trace("key-value"); + BusinessKeyAppender.setBusinessKey2Trace("key-value"); callChainI.doBusiness(); Thread.sleep(ThreadLocalRandom.current().nextInt(10)); diff --git a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainI.java b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainI.java index 37074d455a..4d825cc38b 100644 --- a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainI.java +++ b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainI.java @@ -11,7 +11,7 @@ public class CallChainI { @Tracing public void doBusiness() throws InterruptedException { Thread.sleep(ThreadLocalRandom.current().nextInt(10)); - BusinessKeyAppender.trace("key-value"); + BusinessKeyAppender.setBusinessKey2Trace("key-value"); callChainJ.doBusiness(); Thread.sleep(ThreadLocalRandom.current().nextInt(10)); diff --git a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainJ.java b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainJ.java index 8f1828f807..a33e3fc738 100644 --- a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainJ.java +++ b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainJ.java @@ -10,7 +10,7 @@ public class CallChainJ { @Tracing public void doBusiness() throws InterruptedException { Thread.sleep(ThreadLocalRandom.current().nextInt(10)); - BusinessKeyAppender.trace("key-value"); + BusinessKeyAppender.setBusinessKey2Trace("key-value"); callChainK.doBusiness(); Thread.sleep(ThreadLocalRandom.current().nextInt(10)); diff --git a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainK.java b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainK.java index 8516ac318e..d440c1b6e4 100644 --- a/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainK.java +++ b/skywalking-sdk-plugin/spring-plugin/src/test/java/com/ai/cloud/skywalking/plugin/spring/common/CallChainK.java @@ -9,7 +9,7 @@ public class CallChainK { @Tracing public void doBusiness() throws InterruptedException { Thread.sleep(ThreadLocalRandom.current().nextInt(10)); - BusinessKeyAppender.trace("key-value"); + BusinessKeyAppender.setBusinessKey2Trace("key-value"); Thread.sleep(ThreadLocalRandom.current().nextInt(10)); } }