diff --git a/skywalking-api/pom.xml b/skywalking-api/pom.xml index e40baa6d7..d6b3974ad 100644 --- a/skywalking-api/pom.xml +++ b/skywalking-api/pom.xml @@ -27,5 +27,15 @@ 3.13 test + + com.ai.cloud + skywalking-protocol + 1.0-SNAPSHOT + + + org.apache.htrace + htrace-core + 3.1.0-incubating + 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 2b9d7f414..68eb7701a 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 @@ -3,7 +3,7 @@ package com.ai.cloud.skywalking.api; import com.ai.cloud.skywalking.conf.AuthDesc; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; +import com.ai.cloud.skywalking.protocol.Span; public final class BusinessKeyAppender { diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/Tracing.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/Tracing.java index 29ee97b05..ec6f50d6f 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/Tracing.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/api/Tracing.java @@ -2,7 +2,7 @@ package com.ai.cloud.skywalking.api; import com.ai.cloud.skywalking.conf.AuthDesc; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; +import com.ai.cloud.skywalking.protocol.Span; public class Tracing { /** 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 31fcb8090..070695cbc 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 @@ -2,7 +2,7 @@ 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.protocol.Span; import com.ai.cloud.skywalking.selfexamination.HealthCollector; import com.ai.cloud.skywalking.selfexamination.HeathReading; import com.ai.cloud.skywalking.sender.DataSenderFactory; diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/ContextBuffer.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/ContextBuffer.java index b21cce885..fd2321af6 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/ContextBuffer.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buffer/ContextBuffer.java @@ -1,6 +1,6 @@ package com.ai.cloud.skywalking.buffer; -import com.ai.cloud.skywalking.context.Span; +import com.ai.cloud.skywalking.protocol.Span; import java.util.concurrent.ThreadLocalRandom; diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/LocalBuriedPointSender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/LocalBuriedPointSender.java index f1b0b2109..5160a002a 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/LocalBuriedPointSender.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/LocalBuriedPointSender.java @@ -5,12 +5,11 @@ import com.ai.cloud.skywalking.buffer.ContextBuffer; import com.ai.cloud.skywalking.conf.AuthDesc; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; import com.ai.cloud.skywalking.model.ContextData; import com.ai.cloud.skywalking.model.EmptyContextData; import com.ai.cloud.skywalking.model.Identification; +import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.util.ContextGenerator; -import com.ai.cloud.skywalking.util.ExceptionHandleUtil; import java.util.logging.Level; import java.util.logging.Logger; @@ -55,6 +54,7 @@ public class LocalBuriedPointSender implements IBuriedPointSender { } public void handleException(Throwable e) { - ExceptionHandleUtil.handleException(e); + Span span = Context.getLastSpan(); + span.handleException(e, Config.BuriedPoint.MAX_EXCEPTION_STACK_LENGTH); } } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointReceiver.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointReceiver.java index 42c249d41..698fa342e 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointReceiver.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointReceiver.java @@ -5,11 +5,10 @@ import com.ai.cloud.skywalking.buffer.ContextBuffer; import com.ai.cloud.skywalking.conf.AuthDesc; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; import com.ai.cloud.skywalking.model.ContextData; import com.ai.cloud.skywalking.model.Identification; +import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.util.ContextGenerator; -import com.ai.cloud.skywalking.util.ExceptionHandleUtil; import java.util.logging.Level; import java.util.logging.Logger; @@ -47,6 +46,7 @@ public class RPCBuriedPointReceiver implements IBuriedPointReceiver { } public void handleException(Throwable e) { - ExceptionHandleUtil.handleException(e); + Span span = Context.getLastSpan(); + span.handleException(e, Config.BuriedPoint.MAX_EXCEPTION_STACK_LENGTH); } } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointSender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointSender.java index adfa6f9d6..9036f3f10 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointSender.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/RPCBuriedPointSender.java @@ -5,12 +5,11 @@ import com.ai.cloud.skywalking.buffer.ContextBuffer; import com.ai.cloud.skywalking.conf.AuthDesc; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; import com.ai.cloud.skywalking.model.ContextData; import com.ai.cloud.skywalking.model.EmptyContextData; import com.ai.cloud.skywalking.model.Identification; +import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.util.ContextGenerator; -import com.ai.cloud.skywalking.util.ExceptionHandleUtil; import java.util.logging.Level; import java.util.logging.Logger; @@ -54,6 +53,7 @@ public class RPCBuriedPointSender implements IBuriedPointSender { } public void handleException(Throwable e) { - ExceptionHandleUtil.handleException(e); + Span span = Context.getLastSpan(); + span.handleException(e, Config.BuriedPoint.MAX_EXCEPTION_STACK_LENGTH); } } 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 2d5025214..96747745c 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 @@ -5,12 +5,11 @@ import com.ai.cloud.skywalking.buffer.ContextBuffer; import com.ai.cloud.skywalking.conf.AuthDesc; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; import com.ai.cloud.skywalking.model.ContextData; import com.ai.cloud.skywalking.model.EmptyContextData; import com.ai.cloud.skywalking.model.Identification; +import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.util.BuriedPointMachineUtil; -import com.ai.cloud.skywalking.util.ExceptionHandleUtil; import com.ai.cloud.skywalking.util.TraceIdGenerator; import java.util.logging.Level; @@ -24,15 +23,15 @@ public class ThreadBuriedPointSender implements IBuriedPointSender { public ThreadBuriedPointSender(int threadSeqId) { if (!AuthDesc.isAuth()) - return ; + return; Span spanData; // 从ThreadLocal中取出上下文 final Span parentSpanData = Context.getLastSpan(); if (parentSpanData == null) { - spanData = new Span(TraceIdGenerator.generate()); + spanData = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID); } else { // 如果不为空,则将当前的Context存放到上下文 - spanData = new Span(parentSpanData.getTraceId()); + spanData = new Span(parentSpanData.getTraceId(), Config.SkyWalking.APPLICATION_ID); spanData.setParentLevel(parentSpanData.getParentLevel() + "." + parentSpanData.getLevelId()); spanData.setLevelId(threadSeqId); } @@ -70,7 +69,8 @@ public class ThreadBuriedPointSender implements IBuriedPointSender { } } - public void handleException(Throwable th) { - ExceptionHandleUtil.handleException(th); + public void handleException(Throwable e) { + Span span = Context.getLastSpan(); + span.handleException(e, Config.BuriedPoint.MAX_EXCEPTION_STACK_LENGTH); } } diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadFactoryBuriedPointSender.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadFactoryBuriedPointSender.java index db864108a..8d22ed8e9 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadFactoryBuriedPointSender.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/buriedpoint/ThreadFactoryBuriedPointSender.java @@ -5,12 +5,11 @@ import com.ai.cloud.skywalking.buffer.ContextBuffer; import com.ai.cloud.skywalking.conf.AuthDesc; import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; import com.ai.cloud.skywalking.model.ContextData; import com.ai.cloud.skywalking.model.EmptyContextData; import com.ai.cloud.skywalking.model.Identification; +import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.util.ContextGenerator; -import com.ai.cloud.skywalking.util.ExceptionHandleUtil; import java.util.logging.Level; import java.util.logging.Logger; @@ -51,7 +50,9 @@ public class ThreadFactoryBuriedPointSender implements IBuriedPointSender { } } - public void handleException(Throwable th) { - ExceptionHandleUtil.handleException(th); + public void handleException(Throwable e) { + Span span = Context.getLastSpan(); + span.handleException(e, Config.BuriedPoint.MAX_EXCEPTION_STACK_LENGTH); + } } 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 76be48395..8eb6b850a 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 @@ -4,6 +4,8 @@ public class Config { public static class SkyWalking { public static String USER_ID; + + public static String APPLICATION_ID; } public static class BuriedPoint { diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Constants.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Constants.java index 5b609dd8d..c48e4a324 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Constants.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/conf/Constants.java @@ -2,13 +2,5 @@ package com.ai.cloud.skywalking.conf; public class Constants { - public static final String SPAN_FIELD_SPILT_PATTERN = "^~"; - - public static final String BUSINESSKEY_SPILT_PATTERN = "~^"; - - 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/Context.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Context.java index b6ad2f93d..a5aebfacb 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Context.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Context.java @@ -1,5 +1,7 @@ package com.ai.cloud.skywalking.context; +import com.ai.cloud.skywalking.protocol.Span; + import java.util.Stack; public class Context { 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 deleted file mode 100644 index ddf9333cf..000000000 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/Span.java +++ /dev/null @@ -1,189 +0,0 @@ -package com.ai.cloud.skywalking.context; - -import com.ai.cloud.skywalking.conf.Constants; -import com.ai.cloud.skywalking.util.StringUtil; - -public class Span { - - private String traceId; - private String parentLevel; - private int levelId; - private String viewPointId; - private long startDate; - private long cost; - private String address; - private byte statusCode = 0; - private String exceptionStack; - private char spanType; - private boolean isReceiver = false; - private String businessKey; - private String processNo; - - public Span(String traceId) { - this.traceId = traceId; - } - - public String getProcessNo() { - return processNo; - } - - public void setProcessNo(String processNo) { - this.processNo = processNo; - } - - - public long getCost() { - return cost; - } - - public void setCost(long cost) { - this.cost = cost; - } - - public long getStartDate() { - return startDate; - } - - public void setStartDate(long startDate) { - this.startDate = startDate; - } - - public String getViewPointId() { - return viewPointId; - } - - public void setViewPointId(String viewPointId) { - this.viewPointId = viewPointId; - } - - public String getParentLevel() { - if (!StringUtil.isEmpty(parentLevel)) { - return parentLevel; - } - return ""; - } - - public void setParentLevel(String parentLevel) { - this.parentLevel = parentLevel; - } - - public int getLevelId() { - return levelId; - } - - public void setLevelId(int levelId) { - this.levelId = levelId; - } - - public String getTraceId() { - return traceId; - } - - public void setTraceId(String traceId) { - this.traceId = traceId; - } - - public String getAddress() { - return address; - } - - public void setAddress(String address) { - this.address = address; - } - - public byte getStatusCode() { - return statusCode; - } - - public void setStatusCode(byte statusCode) { - this.statusCode = statusCode; - } - - public String getExceptionStack() { - return exceptionStack; - } - - public void setExceptionStack(String exceptionStack) { - this.exceptionStack = exceptionStack; - } - - public char getSpanType() { - return spanType; - } - - public void setSpanType(char spanType) { - this.spanType = spanType; - } - - public boolean isReceiver() { - return isReceiver; - } - - public void setReceiver(boolean receiver) { - isReceiver = receiver; - } - - public String getBusinessKey() { - return businessKey; - } - - public void setBusinessKey(String businessKey) { - this.businessKey = businessKey; - } - - @Override - public String toString() { - StringBuilder toStringValue = new StringBuilder(); - toStringValue.append(traceId + Constants.SPAN_FIELD_SPILT_PATTERN); - - if (!StringUtil.isEmpty(parentLevel)) { - toStringValue.append(parentLevel + Constants.SPAN_FIELD_SPILT_PATTERN); - } else { - toStringValue.append(" " + Constants.SPAN_FIELD_SPILT_PATTERN); - } - - toStringValue.append(levelId + Constants.SPAN_FIELD_SPILT_PATTERN); - - if (!StringUtil.isEmpty(viewPointId)) { - toStringValue.append(viewPointId + Constants.SPAN_FIELD_SPILT_PATTERN); - } else { - toStringValue.append(" " + Constants.SPAN_FIELD_SPILT_PATTERN); - } - - toStringValue.append(startDate + Constants.SPAN_FIELD_SPILT_PATTERN); - toStringValue.append(cost + Constants.SPAN_FIELD_SPILT_PATTERN); - - if (!StringUtil.isEmpty(address)) { - toStringValue.append(address + Constants.SPAN_FIELD_SPILT_PATTERN); - } else { - toStringValue.append(" " + Constants.SPAN_FIELD_SPILT_PATTERN); - } - - toStringValue.append(statusCode + Constants.SPAN_FIELD_SPILT_PATTERN); - - if (!StringUtil.isEmpty(exceptionStack)) { - toStringValue.append(exceptionStack + Constants.SPAN_FIELD_SPILT_PATTERN); - } else { - toStringValue.append(" " + Constants.SPAN_FIELD_SPILT_PATTERN); - } - - toStringValue.append(spanType + Constants.SPAN_FIELD_SPILT_PATTERN); - toStringValue.append(isReceiver + Constants.SPAN_FIELD_SPILT_PATTERN); - - - if (!StringUtil.isEmpty(businessKey)) { - toStringValue.append(businessKey.replaceAll(Constants.NEW_LINE_CHARACTER_PATTERN, - Constants.BUSINESSKEY_SPILT_PATTERN) + Constants.SPAN_FIELD_SPILT_PATTERN); - } else { - toStringValue.append(" " + Constants.SPAN_FIELD_SPILT_PATTERN); - } - - if (!StringUtil.isEmpty(processNo)) { - toStringValue.append(processNo); - } else { - toStringValue.append(" " + Constants.SPAN_FIELD_SPILT_PATTERN); - } - - return toStringValue.toString(); - } -} diff --git a/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java index ddbc95a84..ac450c752 100644 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java +++ b/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java @@ -1,6 +1,7 @@ package com.ai.cloud.skywalking.model; -import com.ai.cloud.skywalking.context.Span; + +import com.ai.cloud.skywalking.protocol.Span; public class ContextData { private String traceId; 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 be1559627..af82457c4 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,29 +1,25 @@ package com.ai.cloud.skywalking.sender; -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; +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; + public class DataSenderFactory { - private static Logger logger = Logger.getLogger(DataSenderFactory.getSender().toString()); + private static Logger logger = Logger.getLogger(DataSenderFactory.class.getName()); private static Set socketAddresses = new HashSet(); private static Set unUsedSocketAddresses = new HashSet(); @@ -43,7 +39,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.", e); + logger.log(Level.ALL, "Collection service configuration error.", e); System.exit(-1); } @@ -55,7 +51,7 @@ public class DataSenderFactory { try { Thread.sleep(RETRY_GET_SENDER_WAIT_INTERVAL); } catch (InterruptedException e) { - logger.log(Level.ALL, "Sleep failure", e); + logger.log(Level.ALL, "Sleep failure", e); } } return availableSenders.get(ThreadLocalRandom.current().nextInt(0, availableSenders.size())); @@ -66,10 +62,10 @@ public class DataSenderFactory { private int availableSize; public DataSenderChecker() { - super("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)); @@ -93,30 +89,30 @@ public class DataSenderFactory { public void run() { Iterator unUsedSocketAddressIterator; - SocketAddress tmpScoketAddress; + SocketAddress tmpSocketAddress; while (true) { unUsedSocketAddressIterator = unUsedSocketAddresses.iterator(); while (unUsedSocketAddressIterator.hasNext()) { - tmpScoketAddress = unUsedSocketAddressIterator.next(); + tmpSocketAddress = unUsedSocketAddressIterator.next(); if (availableSenders.size() >= availableSize) { - HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.INFO, "the num of available senders is enough."); + 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); + HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.INFO, "increasing available senders."); + availableSenders.add(new DataSender(tmpSocketAddress)); + unUsedSocketAddresses.remove(tmpSocketAddress); } catch (IOException e) { } } } - + if (availableSenders.size() >= availableSize) { - HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.WARNING, "the num of available senders is not enough (" + availableSenders.size() + ")."); + HealthCollector.getCurrentHeathReading(null).updateData(HeathReading.WARNING, "the num of available senders is not enough (" + availableSenders.size() + ")."); } try { 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 83b04950e..e30b31940 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 @@ -1,9 +1,10 @@ package com.ai.cloud.skywalking.util; +import com.ai.cloud.skywalking.conf.Config; import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; import com.ai.cloud.skywalking.model.ContextData; import com.ai.cloud.skywalking.model.Identification; +import com.ai.cloud.skywalking.protocol.Span; public final class ContextGenerator { /** @@ -30,7 +31,7 @@ public final class ContextGenerator { // 校验传入的参数是否为空,如果为空,则新创建一个 if (context == null || StringUtil.isEmpty(context.getTraceId())) { // 不存在,新创建一个Context - spanData = new Span(TraceIdGenerator.generate()); + spanData = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID); } else { // 如果不为空,则将当前的Context存放到上下文 spanData = new Span(context.getTraceId()); @@ -58,10 +59,10 @@ public final class ContextGenerator { // 2 校验Context,Context是否存在 if (parentSpan == null) { // 不存在,新创建一个Context - span = new Span(TraceIdGenerator.generate()); + span = new Span(TraceIdGenerator.generate(), Config.SkyWalking.APPLICATION_ID); } else { // 根据ParentContextData的TraceId和RPCID - span = new Span(parentSpan.getTraceId()); + span = new Span(parentSpan.getTraceId(), Config.SkyWalking.APPLICATION_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/ExceptionHandleUtil.java b/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ExceptionHandleUtil.java deleted file mode 100644 index 5df383102..000000000 --- a/skywalking-api/src/main/java/com/ai/cloud/skywalking/util/ExceptionHandleUtil.java +++ /dev/null @@ -1,43 +0,0 @@ -package com.ai.cloud.skywalking.util; - -import com.ai.cloud.skywalking.conf.Constants; -import com.ai.cloud.skywalking.context.Context; -import com.ai.cloud.skywalking.context.Span; - -import java.io.ByteArrayOutputStream; -import java.io.IOException; - -import static com.ai.cloud.skywalking.conf.Config.BuriedPoint.MAX_EXCEPTION_STACK_LENGTH; - -public final class ExceptionHandleUtil { - private static String extractExceptionStackMessage(final Throwable e) { - ByteArrayOutputStream buf = new ByteArrayOutputStream(); - StringBuilder expMessage = new StringBuilder(); - Throwable causeException = e; - while (causeException != null && (causeException.getCause() != null || expMessage.length() < MAX_EXCEPTION_STACK_LENGTH)) { - causeException.printStackTrace(new java.io.PrintWriter(buf, true)); - expMessage.append(buf.toString()); - causeException = causeException.getCause(); - } - try { - buf.close(); - } catch (IOException e1) { - expMessage.append("\nClose exception stack input stream failed:\n"); - expMessage.append(e1.getCause().getMessage()); - } - if (expMessage.length() <= MAX_EXCEPTION_STACK_LENGTH) { - return expMessage.toString().replaceAll(Constants.NEW_LINE_CHARACTER_PATTERN, Constants.EXCEPTION_SPILT_PATTERN); - } else { - return expMessage.toString().replaceAll(Constants.NEW_LINE_CHARACTER_PATTERN, Constants.EXCEPTION_SPILT_PATTERN) - .substring(0, MAX_EXCEPTION_STACK_LENGTH); - } - } - - public static void handleException(Throwable e) { - Span spanData = Context.getLastSpan(); - // 设置错误信息 - byte errorCode = 1; - spanData.setStatusCode(errorCode); - spanData.setExceptionStack(extractExceptionStackMessage(e)); - } -} diff --git a/skywalking-server/pom.xml b/skywalking-server/pom.xml index 07a230619..3cef3b490 100644 --- a/skywalking-server/pom.xml +++ b/skywalking-server/pom.xml @@ -41,6 +41,11 @@ hbase-client 1.1.2 + + com.ai.cloud + skywalking-protocol + 1.0-SNAPSHOT + diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BuriedPointEntry.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BuriedPointEntry.java deleted file mode 100644 index 4c849f785..000000000 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BuriedPointEntry.java +++ /dev/null @@ -1,151 +0,0 @@ -package com.ai.cloud.skywalking.reciever.model; - -import com.ai.cloud.skywalking.reciever.conf.Constants; - -import java.util.ArrayList; -import java.util.Date; -import java.util.List; - -public class BuriedPointEntry { - private String traceId; - private String parentLevel; - private int levelId; - private String viewPointId; - private Date startDate; - private long cost; - private String address; - private byte statusCode = 0; - private String exceptionStack; - private char spanType; - private boolean isReceiver = false; - private String businessKey; - private String processNo; - // 元数据 - private String originData; - - private BuriedPointEntry() { - - } - - private BuriedPointEntry(String originData) { - - } - - public String getTraceId() { - return traceId; - } - - public String getParentLevel() { - return parentLevel; - } - - public int getLevelId() { - return levelId; - } - - public String getViewPointId() { - return viewPointId; - } - - public Date getStartDate() { - return startDate; - } - - public long getCost() { - return cost; - } - - public String getAddress() { - return address; - } - - public byte getStatusCode() { - return statusCode; - } - - public String getExceptionStack() { - return exceptionStack; - } - - public char getSpanType() { - return spanType; - } - - public boolean isReceiver() { - return isReceiver; - } - - public String getBusinessKey() { - return businessKey; - } - - public String getProcessNo() { - return processNo; - } - - public String getOriginData() { - return originData; - } - - public static BuriedPointEntry convert(String str) { - BuriedPointEntry result = new BuriedPointEntry(); - String[] fieldValues = str.split(Constants.spiltRegx); - result.traceId = fieldValues[0].trim(); - result.parentLevel = fieldValues[1].trim(); - result.levelId = Integer.valueOf(fieldValues[2]); - result.viewPointId = fieldValues[3].trim(); - result.startDate = new Date(Long.valueOf(fieldValues[4])); - result.cost = Long.parseLong(fieldValues[5]); - result.address = fieldValues[6].trim(); - result.statusCode = Byte.valueOf(fieldValues[7].trim()); - //异常情况才会存在exceptionStack - if (result.statusCode == 1) { - result.exceptionStack = fieldValues[8].trim().replaceAll("^","\n"); - } - result.spanType = fieldValues[9].charAt(0); - result.isReceiver = Boolean.valueOf(fieldValues[10]); - result.businessKey = fieldValues[11].trim(); - result.processNo = fieldValues[12].trim(); - result.originData = str; - return result; - } - - private static String[] spilt(String spilt, String value) { - List list = new ArrayList(); - int resultSize = 0; - int start = 0, offset = 0; - for (int index = 0; index < value.length() - 1; index++, offset++) { - if (spilt.charAt(0) == value.charAt(index) && spilt.charAt(1) == value.charAt(index + 1)) { - list.add(value.substring(start, offset)); - index++;offset++; - start = offset + 1; - resultSize++; - } - } - if (start < value.length()) { - list.add(value.substring(start, offset)); - resultSize++; - } - String[] result = new String[resultSize]; - return list.subList(0, resultSize).toArray(result); - } - - @Override - public String toString() { - return "BuriedPointEntry{" + - "traceId='" + traceId + '\'' + - ", parentLevel='" + parentLevel + '\'' + - ", levelId=" + levelId + - ", viewPointId='" + viewPointId + '\'' + - ", startDate=" + startDate + - ", cost=" + cost + - ", address='" + address + '\'' + - ", statusCode=" + statusCode + - ", exceptionStack='" + exceptionStack + '\'' + - ", spanType=" + spanType + - ", isReceiver=" + isReceiver + - ", businessKey='" + businessKey + '\'' + - ", processNo='" + processNo + '\'' + - '}'; - } -} diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/Chain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/Chain.java index 30d673f1c..0141e85e0 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/Chain.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/Chain.java @@ -1,25 +1,25 @@ package com.ai.cloud.skywalking.reciever.storage; +import com.ai.cloud.skywalking.protocol.Span; + import java.util.List; -import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry; - public class Chain { - private List chains; - - private int index = 0; - - public Chain(List chains){ - this.chains = chains; - } - - public void doChain(List enties){ - if(index < chains.size()){ - chains.get(index++).doChain(enties,this); - } - } - - void addChain(IStorageChain chain){ - chains.add(chain); - } + private List chains; + + private int index = 0; + + public Chain(List chains) { + this.chains = chains; + } + + public void doChain(List spans) { + if (index < chains.size()) { + chains.get(index++).doChain(spans, this); + } + } + + void addChain(IStorageChain chain) { + chains.add(chain); + } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/IStorageChain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/IStorageChain.java index e63ac5068..3e5189275 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/IStorageChain.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/IStorageChain.java @@ -1,9 +1,10 @@ package com.ai.cloud.skywalking.reciever.storage; -import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry; + +import com.ai.cloud.skywalking.protocol.Span; import java.util.List; public interface IStorageChain { - void doChain(List entry, Chain chain); + void doChain(List spans, Chain chain); } 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 f5953513a..c57bbb189 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 @@ -1,17 +1,16 @@ package com.ai.cloud.skywalking.reciever.storage; -import static com.ai.cloud.skywalking.reciever.conf.Config.StorageChain.STORAGE_TYPE; +import com.ai.cloud.skywalking.protocol.Span; +import com.ai.cloud.skywalking.reciever.conf.Config; +import com.ai.cloud.skywalking.reciever.storage.chain.SaveToHBaseChain; +import com.ai.cloud.skywalking.reciever.storage.chain.SaveToMySQLChain; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; import java.util.ArrayList; import java.util.List; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; - -import com.ai.cloud.skywalking.reciever.conf.Config; -import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry; -import com.ai.cloud.skywalking.reciever.storage.chain.SaveToHBaseChain; -import com.ai.cloud.skywalking.reciever.storage.chain.SaveToMySQLChain; +import static com.ai.cloud.skywalking.reciever.conf.Config.StorageChain.STORAGE_TYPE; public class StorageChainController { private static Logger logger = LogManager @@ -20,13 +19,13 @@ public class StorageChainController { private static List chainArray = new ArrayList(); static { - if(STORAGE_TYPE.equalsIgnoreCase("hbase")){ - chainArray.add(new SaveToHBaseChain()); - }else if(STORAGE_TYPE.equalsIgnoreCase("mysql")){ - chainArray.add(new SaveToMySQLChain()); - }else{ - throw new RuntimeException("illegal storage type."); - } + if (STORAGE_TYPE.equalsIgnoreCase("hbase")) { + chainArray.add(new SaveToHBaseChain()); + } else if (STORAGE_TYPE.equalsIgnoreCase("mysql")) { + chainArray.add(new SaveToMySQLChain()); + } else { + throw new RuntimeException("illegal storage type."); + } } public static void doStorage(String buriedPointDatas) { @@ -34,13 +33,13 @@ public class StorageChainController { if (buriedPointData == null || buriedPointData.length == 0) { return; } - List entries = new ArrayList(); + List spans = new ArrayList(); for (String buriedPoint : buriedPointData) { try { if (buriedPoint == null || buriedPoint.trim().length() == 0) { continue; } - entries.add(BuriedPointEntry.convert(buriedPoint)); + spans.add(new Span(buriedPoint)); } catch (Throwable e) { logger.error("ready to save buriedPoint error, choose to ignore. data=" + buriedPoint, e); @@ -50,7 +49,7 @@ public class StorageChainController { while (true) { try { Chain chain = new Chain(chainArray); - chain.doChain(entries); + chain.doChain(spans); break; } catch (Throwable e) { try { 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 0cdbb113c..66fb9389c 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 @@ -1,13 +1,12 @@ 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.model.BuriedPointEntry; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; import com.ai.cloud.skywalking.reciever.storage.Chain; import com.ai.cloud.skywalking.reciever.storage.ChainException; import com.ai.cloud.skywalking.reciever.storage.IStorageChain; - import org.apache.commons.lang.StringUtils; import org.apache.hadoop.conf.Configuration; import org.apache.hadoop.hbase.*; @@ -26,9 +25,9 @@ public class SaveToHBaseChain implements IStorageChain { private static Connection connection; @Override - public void doChain(List entry, Chain chain) { - bulkInsertBuriedPointData(entry); - chain.doChain(entry); + public void doChain(List spans, Chain chain) { + bulkInsertBuriedPointData(spans); + chain.doChain(spans); } private static void initHBaseClient() throws IOException { @@ -74,35 +73,35 @@ public class SaveToHBaseChain implements IStorageChain { } - private static void bulkInsertBuriedPointData(List entries) { - if (entries == null || entries.size() <= 0) + private static void bulkInsertBuriedPointData(List spans) { + if (spans == null || spans.size() <= 0) return; List puts = new ArrayList(); Put put; String columnName; - for (BuriedPointEntry buriedPointEntry : entries) { - put = new Put(Bytes.toBytes(buriedPointEntry.getTraceId())); - if (StringUtils.isEmpty(buriedPointEntry.getParentLevel().trim())) { - columnName = buriedPointEntry.getLevelId() + ""; - if (buriedPointEntry.isReceiver()) { - columnName = buriedPointEntry.getLevelId() + "-S"; + for (Span span : spans) { + put = new Put(Bytes.toBytes(span.getTraceId())); + 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(buriedPointEntry.getOriginData())); + Bytes.toBytes(span.getOriginData())); } else { - columnName = buriedPointEntry.getParentLevel() + "." + buriedPointEntry.getLevelId(); - if (buriedPointEntry.isReceiver()) { - columnName = buriedPointEntry.getParentLevel() + "." + buriedPointEntry.getLevelId() + "-S"; + 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(buriedPointEntry.getOriginData())); + Bytes.toBytes(span.getOriginData())); } puts.add(put); } bulkInsertBuriedPointData(Config.HBaseConfig.TABLE_NAME, puts); - - ServerHealthCollector.getCurrentHeathReading("hbase").updateData(ServerHeathReading.INFO, "save " + entries.size() + " BuriedPointEntries." ); + + ServerHealthCollector.getCurrentHeathReading("hbase").updateData(ServerHeathReading.INFO, "save " + spans.size() + " BuriedPointEntries."); } private static void bulkInsertBuriedPointData(String tableName, List data) { @@ -127,14 +126,14 @@ public class SaveToHBaseChain implements IStorageChain { } - public static List selectByTraceId(String traceId) throws IOException { - List entries = new ArrayList(); + public static List selectByTraceId(String traceId) throws IOException { + List entries = new ArrayList(); Table table = connection.getTable(TableName.valueOf(Config.HBaseConfig.TABLE_NAME)); Get g = new Get(Bytes.toBytes(traceId)); Result r = table.get(g); for (Cell cell : r.rawCells()) { if (cell.getValueArray().length > 0) - entries.add(BuriedPointEntry.convert(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength()))); + entries.add(new Span(Bytes.toString(cell.getValueArray(), cell.getValueOffset(), cell.getValueLength()))); } return entries; } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToMySQLChain.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToMySQLChain.java index b79877ac2..8bd832bbf 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToMySQLChain.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/storage/chain/SaveToMySQLChain.java @@ -1,17 +1,17 @@ package com.ai.cloud.skywalking.reciever.storage.chain; -import java.util.List; - -import com.ai.cloud.skywalking.reciever.model.BuriedPointEntry; +import com.ai.cloud.skywalking.protocol.Span; import com.ai.cloud.skywalking.reciever.storage.Chain; import com.ai.cloud.skywalking.reciever.storage.IStorageChain; +import java.util.List; + public class SaveToMySQLChain implements IStorageChain { - @Override - public void doChain(List entry, Chain chain) { - // TODO Auto-generated method stub + @Override + public void doChain(List entry, Chain chain) { + // TODO Auto-generated method stub - } + } }