diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCClientInvokeMonitor.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCClientInvokeMonitor.java index 6dce553ad..3af476cf9 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCClientInvokeMonitor.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCClientInvokeMonitor.java @@ -30,7 +30,7 @@ public class RPCClientInvokeMonitor extends BaseInvokeMonitor { ContextBuffer.save(new RequestSpan(spanData)); CurrentThreadSpanStack.push(spanData); - return new ContextData(spanData.getTraceId(), generateSubParentLevelId(spanData), spanData.getCallType()); + return new ContextData(spanData.getTraceId(), generateSubParentLevelId(spanData)); } catch (Throwable t) { logger.error(t.getMessage(), t); return new EmptyContextData(); diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCServerInvokeMonitor.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCServerInvokeMonitor.java index 0393dc38e..c4a5f6a16 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCServerInvokeMonitor.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/invoke/monitor/RPCServerInvokeMonitor.java @@ -27,18 +27,18 @@ public class RPCServerInvokeMonitor extends BaseInvokeMonitor { invalidateAllSpanIfIsNotFirstSpan(spanData); - super.beforeInvoke(spanData); + super.beforeInvoke(spanData, id); } catch (Throwable t) { logger.error(t.getMessage(), t); } } - public void afterInvoke(){ + public void afterInvoke() { super.afterInvoke(); } - public void occurException(Throwable th){ + public void occurException(Throwable th) { super.occurException(th); } diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java index 27a6d4212..43c70a1bc 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/model/ContextData.java @@ -7,23 +7,20 @@ public class ContextData { private String traceId; private String parentLevel; private int levelId; - private String spanType; ContextData() { } - public ContextData(String traceId, String parentLevel, String spanType) { + public ContextData(String traceId, String parentLevel) { this.traceId = traceId; this.parentLevel = parentLevel; - this.spanType = spanType; } public ContextData(Span span) { this.traceId = span.getTraceId(); this.parentLevel = span.getParentLevel(); this.levelId = span.getLevelId(); - this.spanType = span.getSpanTypeDesc(); } public ContextData(String contextDataStr) { @@ -35,7 +32,6 @@ public class ContextData { this.traceId = value[0]; this.parentLevel = value[1].trim(); this.levelId = Integer.valueOf(value[2]); - this.spanType = value[3]; } @@ -51,10 +47,6 @@ public class ContextData { return levelId; } - public String getSpanType() { - return spanType; - } - @Override public String toString() { StringBuilder stringBuilder = new StringBuilder(); @@ -67,8 +59,6 @@ public class ContextData { } stringBuilder.append("-"); stringBuilder.append(levelId); - stringBuilder.append("-"); - stringBuilder.append(spanType); return stringBuilder.toString(); } } diff --git a/skywalking-collector/skywalking-api/src/test/java/test/ai/cloud/serialize/SerializeTest.java b/skywalking-collector/skywalking-api/src/test/java/test/ai/cloud/serialize/SerializeTest.java new file mode 100644 index 000000000..a52e01dbc --- /dev/null +++ b/skywalking-collector/skywalking-api/src/test/java/test/ai/cloud/serialize/SerializeTest.java @@ -0,0 +1,20 @@ +package test.ai.cloud.serialize; + +import com.ai.cloud.skywalking.buffer.ContextBuffer; +import com.ai.cloud.skywalking.protocol.AckSpan; +import com.ai.cloud.skywalking.protocol.RequestSpan; +import com.ai.cloud.skywalking.protocol.Span; +import com.ai.cloud.skywalking.protocol.common.CallType; +import com.ai.cloud.skywalking.protocol.common.SpanType; + +public class SerializeTest { + public static void main(String[] args) throws InterruptedException { + Span spandata = new Span("1.0b.1461060884539.7d6d06e.22489.1271.103", "", 0); + spandata.setSpanType(SpanType.LOCAL); + spandata.setStartDate(System.currentTimeMillis() - 1000 * 60); + AckSpan requestSpan = new AckSpan(spandata); + ContextBuffer.save(requestSpan); + + Thread.sleep(60 * 1000 * 10); + } +} diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java index 98bc46d4f..a6b12bbc0 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/AckSpan.java @@ -24,11 +24,11 @@ public class AckSpan extends AbstractDataSerializable { * 当前调用链的本机描述
* 如当前序号为:0.1.0时,levelId=0 */ - private int levelId = 0; + private int levelId = 0; /** * 节点调用花费时间 */ - private long cost = 0L; + private long cost = 0L; /** * 节点调用的状态
* 0:成功
@@ -40,7 +40,7 @@ public class AckSpan extends AbstractDataSerializable { * 节点调用的错误堆栈
* 堆栈以JAVA的exception为主要判断依据 */ - private String exceptionStack; + private String exceptionStack = ""; public AckSpan(Span spanData) { @@ -124,7 +124,8 @@ public class AckSpan extends AbstractDataSerializable { @Override public byte[] getData() { - return new byte[0]; + return TraceProtocol.AckSpan.newBuilder().setTraceId(traceId).setParentLevel(parentLevel). + setLevelId(levelId).setCost(cost).setStatusCode(statusCode).setExceptionStack(exceptionStack).build().toByteArray(); } @Override diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java index 38154bd3f..f17fbaeaf 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/Span.java @@ -43,7 +43,7 @@ public class Span { * 节点调用的错误堆栈
* 堆栈以JAVA的exception为主要判断依据 */ - protected String exceptionStack; + protected String exceptionStack = ""; /** * 节点的状态
@@ -60,7 +60,7 @@ public class Span { * 节点类型
* 如:RPC Client,RPC Server,Local */ - private SpanType spanType = SpanType.LOCAL; + private SpanType spanType = SpanType.LOCAL; public Span(String traceId) { this.traceId = traceId; diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java index 9c9174cb2..e900975a6 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/AbstractDataSerializable.java @@ -28,13 +28,12 @@ public abstract class AbstractDataSerializable implements ISerializable, Nullabl byte[] messageByteData = getData(); byte[] messagePackage = new byte[4 + messageByteData.length]; packData(messageByteData, messagePackage); - setPackageLength(messagePackage); + setProtocolType(messageByteData, messagePackage); return messagePackage; } - private void setPackageLength(byte[] messagePackage) { - byte[] type = IntegerAssist.intToBytes(this.getDataType()); - System.arraycopy(type, 0, messagePackage, 0, type.length); + private void setProtocolType(byte[] messageByteData, byte[] messagePackage) { + System.arraycopy(IntegerAssist.intToBytes(getDataType()), 0, messagePackage, 0, 4); } private void packData(byte[] messageByteData, byte[] messagePackage) { @@ -49,7 +48,7 @@ public abstract class AbstractDataSerializable implements ISerializable, Nullabl return new NullClass(); } - return this.convertData(Arrays.copyOfRange(data,4, data.length)); + return this.convertData(Arrays.copyOfRange(data, 4, data.length)); } diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/serialize/SerializedFactory.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/serialize/SerializedFactory.java index 1fae8df41..27469a423 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/serialize/SerializedFactory.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/serialize/SerializedFactory.java @@ -29,7 +29,7 @@ public class SerializedFactory { if (abstractDataSerializable != null) { NullableClass nullableClass = abstractDataSerializable.convert2Object(bytes); if (!nullableClass.isNull()) { - return abstractDataSerializable; + return (AbstractDataSerializable) nullableClass; } } return null; diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/TransportPackager.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/TransportPackager.java index f03d1d33e..bdbdcdc39 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/TransportPackager.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/util/TransportPackager.java @@ -17,18 +17,22 @@ public class TransportPackager { public static List unpack(byte[] dataPackage) { if (validateCheckSum(dataPackage)) { - return unpackDataText(dataPackage); + return unpackDataText(unpackCheckSum(dataPackage)); } else { return new ArrayList(); } } + private static byte[] unpackCheckSum(byte[] dataPackage) { + return Arrays.copyOfRange(dataPackage, 4, dataPackage.length); + } + private static List unpackDataText(byte[] dataPackage) { List serializeData = new ArrayList(); int currentLength = 0; while (true) { //读取长度 - int dataLength = IntegerAssist.bytesToInt(dataPackage, 0); + int dataLength = IntegerAssist.bytesToInt(dataPackage, currentLength); // 反序列化 byte[] data = new byte[dataLength]; System.arraycopy(dataPackage, currentLength + 4, data, 0, dataLength); @@ -56,8 +60,8 @@ public class TransportPackager { * @return */ private static byte[] generateChecksum(byte[] data, int offset) { - int result = data[0]; - for (int i = offset; i < data.length; i++) { + int result = data[offset]; + for (int i = offset + 1; i < data.length; i++) { result ^= data[i]; } @@ -73,12 +77,16 @@ public class TransportPackager { } private static byte[] packDataText(List beSendingData) { - byte[] dataText = new byte[1024 * 30]; + byte[] dataText = null; int currentIndex = 0; for (ISerializable sendingData : beSendingData) { byte[] dataElementText = appendingLength(sendingData.convert2Bytes()); - dataText = expansionCapacityIfNecessary(dataText, currentIndex, dataElementText); - appendBeSendingDataText(dataElementText, dataText, currentIndex); + if (dataText == null) { + dataText = new byte[dataElementText.length]; + } else { + dataText = Arrays.copyOf(dataText, dataText.length + dataElementText.length); + } + System.arraycopy(dataElementText, 0, dataText, currentIndex, dataElementText.length); currentIndex += dataElementText.length; } diff --git a/skywalking-server/src/main/resources/META-INF/services/com.ai.cloud.skywalking.reciever.processor.AbstractSpanProcessor b/skywalking-server/src/main/resources/META-INF/services/com.ai.cloud.skywalking.reciever.processor.IProcessor similarity index 100% rename from skywalking-server/src/main/resources/META-INF/services/com.ai.cloud.skywalking.reciever.processor.AbstractSpanProcessor rename to skywalking-server/src/main/resources/META-INF/services/com.ai.cloud.skywalking.reciever.processor.IProcessor diff --git a/skywalking-server/src/main/resources/config.properties b/skywalking-server/src/main/resources/config.properties index 2090f0042..a525da474 100644 --- a/skywalking-server/src/main/resources/config.properties +++ b/skywalking-server/src/main/resources/config.properties @@ -1,8 +1,8 @@ #采集服务器的端口 server.port=34000 -server.max_deal_data_thread_number=5 +server.max_deal_data_thread_number=1 server.failed_package_watching_time_windowss=300 -server.max_watching_failed_package_size=200; +server.max_watching_failed_package_size=200 #每个线程最大缓存数量 buffer.per_thread_max_buffer_number=1024 @@ -22,28 +22,8 @@ buffer.write_data_failure_retry_interval = 10000 #数据包的最大限制 datapackage.max_data_package=1048576 -#定位文件时,每次读取偏移量跳过大小 -persistence.step_size_for_location_file_offset=20480 -#切换数据文件,等待时间(单位:毫秒) -persistence.switch_file_wait_time=5000 -#追加EOF标志位的线程数量 -persistence.max_append_eof_flags_thread_number=2 -#当读取文件结束时最大等待时间 -persistence.read_ending_file_max_waite_time=50 -#每次存储的最大数量 -persistence.max_storage_size_per_time = 1048576 - -#偏移量注册文件的目录 -registerpersistence.register_file_parent_directory=d:/test-data/data/offset -#偏移量注册文件名 -registerpersistence.register_file_name=offset.txt -#偏移量注册备份文件名 -registerpersistence.register_bak_file_name=offset.txt.bak -#偏移量写入文件等待周期(单位:毫秒) -registerpersistence.offset_written_file_wait_cycle=5000 - #hbase表名 -hbaseconfig.table_name=sw-call-chain +hbaseconfig.table_name=sw-call-chain-new #hbase列簇名字 hbaseconfig.family_column_name=call-chain #hbase zk quorum