From b6d262c2faeaa73eb4a1f7c15b337a28eb815295 Mon Sep 17 00:00:00 2001 From: ascrutae Date: Thu, 4 Aug 2016 13:30:35 +0800 Subject: [PATCH] =?UTF-8?q?1.=20=E9=87=8D=E6=9E=84=E4=BB=A3=E7=A0=81=20=20?= =?UTF-8?q?2.=20=E4=BF=AE=E5=A4=8D=E5=BA=8F=E5=88=97=E5=8C=96=E9=97=AE?= =?UTF-8?q?=E9=A2=98?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../src/main/resources/sky-walking.auth | 2 +- .../resources/provider/dubbo-provider.xml | 4 +- .../resources/consumer/dubbo-consumer.xml | 4 +- .../src/main/resources/sky-walking.auth | 45 - .../context/CurrentThreadSpanStack.java | 13 - .../enhance/ClassEnhancePluginDefine.java | 5 +- .../skywalking-protocol/pom.xml | 5 +- .../protocol/proto/TraceProtocol.java | 1050 +++++++++++++---- .../skywalking/protocol/RequestSpan.java | 19 +- .../protocol/SerializedFactory.java | 2 +- .../ai/cloud/skywalking/protocol/Span.java | 26 +- .../protocol/TransportPackager.java | 69 +- .../skywalking/protocol/common/SpanType.java | 2 +- .../protocol/util/ByteDataUtil.java | 28 + .../src/main/proto/TraceProtocol.proto | 1 + .../dubbo-plugin/pom.xml | 5 +- .../dubbo/MonitorFilterInterceptor.java | 26 +- .../define/AbstractDatabasePluginDefine.java | 17 + .../define/DatabasePluginInterceptor.java | 36 + .../jdbc/plugin/H2DatabasePluginDefine.java | 10 + .../src/main/resources/skywalking-plugin.def | 2 +- .../tomcat78x/TomcatPluginInterceptor.java | 3 +- skywalking-server/pom.xml | 5 + .../reciever/buffer/AppendEOFFlagThread.java | 5 +- .../BufferDataAssist.java} | 28 +- .../reciever/buffer/DataBufferThread.java | 7 +- .../buffer/DataBufferThreadContainer.java | 5 +- .../skywalking/reciever/conf/Config.java | 8 +- .../handler/CollectionServerDataHandler.java | 25 +- .../peresistent/BufferFileReader.java | 51 +- .../PersistenceThreadLauncher.java | 8 +- .../src/main/resources/config.properties | 7 +- .../reciever/model/BufferDataAssistTest.java | 51 + 33 files changed, 1107 insertions(+), 467 deletions(-) delete mode 100644 samples/skywalking-example/example-web/src/main/resources/sky-walking.auth create mode 100644 skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/util/ByteDataUtil.java create mode 100644 skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/AbstractDatabasePluginDefine.java create mode 100644 skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/DatabasePluginInterceptor.java create mode 100644 skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/plugin/H2DatabasePluginDefine.java rename skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/{model/BufferDataPackagerGenerator.java => buffer/BufferDataAssist.java} (56%) create mode 100644 skywalking-server/src/test/java/com/ai/cloud/skywalking/reciever/model/BufferDataAssistTest.java diff --git a/samples/skywalking-auth/src/main/resources/sky-walking.auth b/samples/skywalking-auth/src/main/resources/sky-walking.auth index 90b8fff41..71f26d069 100644 --- a/samples/skywalking-auth/src/main/resources/sky-walking.auth +++ b/samples/skywalking-auth/src/main/resources/sky-walking.auth @@ -31,7 +31,7 @@ sender.max_send_length=20000 sender.retry_get_sender_wait_interval=2000 #最大消费线程数 -consumer.max_consumer=0 +consumer.max_consumer=1 #消费者最大等待时间 consumer.max_wait_time=5 #发送失败等待时间 diff --git a/samples/skywalking-example/example-dubbo/dubbo-impl/src/main/resources/provider/dubbo-provider.xml b/samples/skywalking-example/example-dubbo/dubbo-impl/src/main/resources/provider/dubbo-provider.xml index b195734d3..bbed60137 100644 --- a/samples/skywalking-example/example-dubbo/dubbo-impl/src/main/resources/provider/dubbo-provider.xml +++ b/samples/skywalking-example/example-dubbo/dubbo-impl/src/main/resources/provider/dubbo-provider.xml @@ -8,8 +8,8 @@ http://code.alibabatech.com/schema/dubbo/dubbo.xsd"> - + - \ No newline at end of file + diff --git a/samples/skywalking-example/example-web/src/main/resources/consumer/dubbo-consumer.xml b/samples/skywalking-example/example-web/src/main/resources/consumer/dubbo-consumer.xml index 3751f4815..cb9a67055 100644 --- a/samples/skywalking-example/example-web/src/main/resources/consumer/dubbo-consumer.xml +++ b/samples/skywalking-example/example-web/src/main/resources/consumer/dubbo-consumer.xml @@ -7,6 +7,6 @@ http://code.alibabatech.com/schema/dubbo http://code.alibabatech.com/schema/dubbo/dubbo.xsd"> - + - \ No newline at end of file + diff --git a/samples/skywalking-example/example-web/src/main/resources/sky-walking.auth b/samples/skywalking-example/example-web/src/main/resources/sky-walking.auth deleted file mode 100644 index 9be01f6e3..000000000 --- a/samples/skywalking-example/example-web/src/main/resources/sky-walking.auth +++ /dev/null @@ -1,45 +0,0 @@ -#skyWalking用户ID -skywalking.user_id=123 -#skyWalking应用编码 -skywalking.application_code=skywalking-sample-dubbo -#skywalking auth的环境变量名字 -skywalking.auth_system_env_name=SKYWALKING_RUN -#skywalking数据编码 -skywalking.charset=UTF-8 -skywalking.auth_override=true - -#是否打印数据 -buriedpoint.printf=true -#埋点异常的最大长度 -buriedpoint.max_exception_stack_length=4000 -#业务字段的最大长度 -buriedpoint.businesskey_max_length=300 -#过滤异常 -buriedpoint.exclusive_exceptions=java.lang.RuntimeException - -#最大发送者的连接数阀比例 -sender.connect_percent=100 -#发送服务端配置 -sender.servers_addr=127.0.0.1:34000 -#最大发送的副本数量 -sender.max_copy_num=2 -#发送的最大长度 -sender.max_send_length=20000 -#当没有Sender时,尝试获取sender的等待周期 -sender.retry_get_sender_wait_interval=2000 - -#最大消费线程数 -consumer.max_consumer=0 -#消费者最大等待时间 -consumer.max_wait_time=5 -#发送失败等待时间 -consumer.consumer_fail_retry_wait_interval=50 - -#每个Buffer的最大个数 -buffer.buffer_max_size=18000 -#Buffer池的最大长度 -buffer.pool_size=5 - -#发送检查线程检查周期 -senderchecker.check_polling_time=200 - diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java index 251c803d0..97d43ad82 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/context/CurrentThreadSpanStack.java @@ -33,14 +33,6 @@ public class CurrentThreadSpanStack { return nodes.get().pop(); } - public static void invalidatePresentSpans() { - if (nodes.get() == null) { - nodes.set(new SpanNodeStack()); - } - - nodes.get().invalidatePresentSpans(); - } - static class SpanNodeStack { /** * 单JVM的单线程,埋点数量一般不会超过20. @@ -84,11 +76,6 @@ public class CurrentThreadSpanStack { spans.add(spans.size(), spanNode); } - public void invalidatePresentSpans() { - for (SpanNode spanNode : spans) { - spanNode.getData().setValidate(true); - } - } } static class SpanNode { diff --git a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/plugin/interceptor/enhance/ClassEnhancePluginDefine.java b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/plugin/interceptor/enhance/ClassEnhancePluginDefine.java index 86ed27e67..d857146f4 100644 --- a/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/plugin/interceptor/enhance/ClassEnhancePluginDefine.java +++ b/skywalking-collector/skywalking-api/src/main/java/com/ai/cloud/skywalking/plugin/interceptor/enhance/ClassEnhancePluginDefine.java @@ -1,5 +1,6 @@ package com.ai.cloud.skywalking.plugin.interceptor.enhance; +import static net.bytebuddy.jar.asm.Opcodes.ACC_PRIVATE; import static net.bytebuddy.matcher.ElementMatchers.any; import static net.bytebuddy.matcher.ElementMatchers.not; @@ -20,6 +21,8 @@ import com.ai.cloud.skywalking.plugin.interceptor.EnhanceException; import com.ai.cloud.skywalking.plugin.interceptor.EnhancedClassInstanceContext; import com.ai.cloud.skywalking.plugin.interceptor.MethodMatcher; +import java.lang.reflect.Modifier; + public abstract class ClassEnhancePluginDefine extends AbstractClassEnhancePluginDefine { private static Logger logger = LogManager .getLogger(ClassEnhancePluginDefine.class); @@ -61,7 +64,7 @@ public abstract class ClassEnhancePluginDefine extends AbstractClassEnhancePlugi newClassBuilder = newClassBuilder .defineField(contextAttrName, - EnhancedClassInstanceContext.class) + EnhancedClassInstanceContext.class, ACC_PRIVATE) .constructor(any()) .intercept( SuperMethodCall.INSTANCE.andThen(MethodDelegation.to( diff --git a/skywalking-collector/skywalking-protocol/pom.xml b/skywalking-collector/skywalking-protocol/pom.xml index ed6c460a1..bba884047 100644 --- a/skywalking-collector/skywalking-protocol/pom.xml +++ b/skywalking-collector/skywalking-protocol/pom.xml @@ -21,15 +21,16 @@ com.google.protobuf protobuf-java - 2.6.1 + 3.0.0 junit junit - 3.8.1 + 4.12 test + diff --git a/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java b/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java index af9be17c9..77269571b 100644 --- a/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java +++ b/skywalking-collector/skywalking-protocol/src/main/gen-java/com/ai/cloud/skywalking/protocol/proto/TraceProtocol.java @@ -5,8 +5,14 @@ package com.ai.cloud.skywalking.protocol.proto; public final class TraceProtocol { private TraceProtocol() {} + public static void registerAllExtensions( + com.google.protobuf.ExtensionRegistryLite registry) { + } + public static void registerAllExtensions( com.google.protobuf.ExtensionRegistry registry) { + registerAllExtensions( + (com.google.protobuf.ExtensionRegistryLite) registry); } public interface AckSpanOrBuilder extends // @@protoc_insertion_point(interface_extends:AckSpan) @@ -84,37 +90,33 @@ public final class TraceProtocol { /** * Protobuf type {@code AckSpan} */ - public static final class AckSpan extends - com.google.protobuf.GeneratedMessage implements + public static final class AckSpan extends + com.google.protobuf.GeneratedMessageV3 implements // @@protoc_insertion_point(message_implements:AckSpan) AckSpanOrBuilder { // Use AckSpan.newBuilder() to construct. - private AckSpan(com.google.protobuf.GeneratedMessage.Builder builder) { + private AckSpan(com.google.protobuf.GeneratedMessageV3.Builder builder) { super(builder); - this.unknownFields = builder.getUnknownFields(); } - private AckSpan(boolean noInit) { this.unknownFields = com.google.protobuf.UnknownFieldSet.getDefaultInstance(); } - - private static final AckSpan defaultInstance; - public static AckSpan getDefaultInstance() { - return defaultInstance; + private AckSpan() { + traceId_ = ""; + parentLevel_ = ""; + levelId_ = 0; + cost_ = 0L; + statusCode_ = 0; + exceptionStack_ = ""; } - public AckSpan getDefaultInstanceForType() { - return defaultInstance; - } - - private final com.google.protobuf.UnknownFieldSet unknownFields; @java.lang.Override public final com.google.protobuf.UnknownFieldSet - getUnknownFields() { + getUnknownFields() { return this.unknownFields; } private AckSpan( com.google.protobuf.CodedInputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws com.google.protobuf.InvalidProtocolBufferException { - initFields(); + this(); int mutable_bitField0_ = 0; com.google.protobuf.UnknownFieldSet.Builder unknownFields = com.google.protobuf.UnknownFieldSet.newBuilder(); @@ -172,7 +174,7 @@ public final class TraceProtocol { throw e.setUnfinishedMessage(this); } catch (java.io.IOException e) { throw new com.google.protobuf.InvalidProtocolBufferException( - e.getMessage()).setUnfinishedMessage(this); + e).setUnfinishedMessage(this); } finally { this.unknownFields = unknownFields.build(); makeExtensionsImmutable(); @@ -183,31 +185,16 @@ public final class TraceProtocol { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_AckSpan_descriptor; } - protected com.google.protobuf.GeneratedMessage.FieldAccessorTable + protected com.google.protobuf.GeneratedMessageV3.FieldAccessorTable internalGetFieldAccessorTable() { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_AckSpan_fieldAccessorTable .ensureFieldAccessorsInitialized( com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan.class, com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan.Builder.class); } - public static com.google.protobuf.Parser PARSER = - new com.google.protobuf.AbstractParser() { - public AckSpan parsePartialFrom( - com.google.protobuf.CodedInputStream input, - com.google.protobuf.ExtensionRegistryLite extensionRegistry) - throws com.google.protobuf.InvalidProtocolBufferException { - return new AckSpan(input, extensionRegistry); - } - }; - - @java.lang.Override - public com.google.protobuf.Parser getParserForType() { - return PARSER; - } - private int bitField0_; public static final int TRACEID_FIELD_NUMBER = 1; - private java.lang.Object traceId_; + private volatile java.lang.Object traceId_; /** * required string traceId = 1; */ @@ -249,7 +236,7 @@ public final class TraceProtocol { } public static final int PARENTLEVEL_FIELD_NUMBER = 2; - private java.lang.Object parentLevel_; + private volatile java.lang.Object parentLevel_; /** * optional string parentLevel = 2; */ @@ -336,7 +323,7 @@ public final class TraceProtocol { } public static final int EXCEPTIONSTACK_FIELD_NUMBER = 6; - private java.lang.Object exceptionStack_; + private volatile java.lang.Object exceptionStack_; /** * optional string exceptionStack = 6; */ @@ -377,14 +364,6 @@ public final class TraceProtocol { } } - private void initFields() { - traceId_ = ""; - parentLevel_ = ""; - levelId_ = 0; - cost_ = 0L; - statusCode_ = 0; - exceptionStack_ = ""; - } private byte memoizedIsInitialized = -1; public final boolean isInitialized() { byte isInitialized = memoizedIsInitialized; @@ -413,12 +392,11 @@ public final class TraceProtocol { public void writeTo(com.google.protobuf.CodedOutputStream output) throws java.io.IOException { - getSerializedSize(); if (((bitField0_ & 0x00000001) == 0x00000001)) { - output.writeBytes(1, getTraceIdBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 1, traceId_); } if (((bitField0_ & 0x00000002) == 0x00000002)) { - output.writeBytes(2, getParentLevelBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 2, parentLevel_); } if (((bitField0_ & 0x00000004) == 0x00000004)) { output.writeInt32(3, levelId_); @@ -430,24 +408,21 @@ public final class TraceProtocol { output.writeInt32(5, statusCode_); } if (((bitField0_ & 0x00000020) == 0x00000020)) { - output.writeBytes(6, getExceptionStackBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 6, exceptionStack_); } - getUnknownFields().writeTo(output); + unknownFields.writeTo(output); } - private int memoizedSerializedSize = -1; public int getSerializedSize() { - int size = memoizedSerializedSize; + int size = memoizedSize; if (size != -1) return size; size = 0; if (((bitField0_ & 0x00000001) == 0x00000001)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(1, getTraceIdBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(1, traceId_); } if (((bitField0_ & 0x00000002) == 0x00000002)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(2, getParentLevelBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(2, parentLevel_); } if (((bitField0_ & 0x00000004) == 0x00000004)) { size += com.google.protobuf.CodedOutputStream @@ -462,19 +437,94 @@ public final class TraceProtocol { .computeInt32Size(5, statusCode_); } if (((bitField0_ & 0x00000020) == 0x00000020)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(6, getExceptionStackBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(6, exceptionStack_); } - size += getUnknownFields().getSerializedSize(); - memoizedSerializedSize = size; + size += unknownFields.getSerializedSize(); + memoizedSize = size; return size; } private static final long serialVersionUID = 0L; @java.lang.Override - protected java.lang.Object writeReplace() - throws java.io.ObjectStreamException { - return super.writeReplace(); + public boolean equals(final java.lang.Object obj) { + if (obj == this) { + return true; + } + if (!(obj instanceof com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan)) { + return super.equals(obj); + } + com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan other = (com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan) obj; + + boolean result = true; + result = result && (hasTraceId() == other.hasTraceId()); + if (hasTraceId()) { + result = result && getTraceId() + .equals(other.getTraceId()); + } + result = result && (hasParentLevel() == other.hasParentLevel()); + if (hasParentLevel()) { + result = result && getParentLevel() + .equals(other.getParentLevel()); + } + result = result && (hasLevelId() == other.hasLevelId()); + if (hasLevelId()) { + result = result && (getLevelId() + == other.getLevelId()); + } + result = result && (hasCost() == other.hasCost()); + if (hasCost()) { + result = result && (getCost() + == other.getCost()); + } + result = result && (hasStatusCode() == other.hasStatusCode()); + if (hasStatusCode()) { + result = result && (getStatusCode() + == other.getStatusCode()); + } + result = result && (hasExceptionStack() == other.hasExceptionStack()); + if (hasExceptionStack()) { + result = result && getExceptionStack() + .equals(other.getExceptionStack()); + } + result = result && unknownFields.equals(other.unknownFields); + return result; + } + + @java.lang.Override + public int hashCode() { + if (memoizedHashCode != 0) { + return memoizedHashCode; + } + int hash = 41; + hash = (19 * hash) + getDescriptorForType().hashCode(); + if (hasTraceId()) { + hash = (37 * hash) + TRACEID_FIELD_NUMBER; + hash = (53 * hash) + getTraceId().hashCode(); + } + if (hasParentLevel()) { + hash = (37 * hash) + PARENTLEVEL_FIELD_NUMBER; + hash = (53 * hash) + getParentLevel().hashCode(); + } + if (hasLevelId()) { + hash = (37 * hash) + LEVELID_FIELD_NUMBER; + hash = (53 * hash) + getLevelId(); + } + if (hasCost()) { + hash = (37 * hash) + COST_FIELD_NUMBER; + hash = (53 * hash) + com.google.protobuf.Internal.hashLong( + getCost()); + } + if (hasStatusCode()) { + hash = (37 * hash) + STATUSCODE_FIELD_NUMBER; + hash = (53 * hash) + getStatusCode(); + } + if (hasExceptionStack()) { + hash = (37 * hash) + EXCEPTIONSTACK_FIELD_NUMBER; + hash = (53 * hash) + getExceptionStack().hashCode(); + } + hash = (29 * hash) + unknownFields.hashCode(); + memoizedHashCode = hash; + return hash; } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan parseFrom( @@ -500,46 +550,57 @@ public final class TraceProtocol { } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan parseFrom(java.io.InputStream input) throws java.io.IOException { - return PARSER.parseFrom(input); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan parseFrom( java.io.InputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws java.io.IOException { - return PARSER.parseFrom(input, extensionRegistry); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input, extensionRegistry); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan parseDelimitedFrom(java.io.InputStream input) throws java.io.IOException { - return PARSER.parseDelimitedFrom(input); + return com.google.protobuf.GeneratedMessageV3 + .parseDelimitedWithIOException(PARSER, input); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan parseDelimitedFrom( java.io.InputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws java.io.IOException { - return PARSER.parseDelimitedFrom(input, extensionRegistry); + return com.google.protobuf.GeneratedMessageV3 + .parseDelimitedWithIOException(PARSER, input, extensionRegistry); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan parseFrom( com.google.protobuf.CodedInputStream input) throws java.io.IOException { - return PARSER.parseFrom(input); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan parseFrom( com.google.protobuf.CodedInputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws java.io.IOException { - return PARSER.parseFrom(input, extensionRegistry); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input, extensionRegistry); } - public static Builder newBuilder() { return Builder.create(); } public Builder newBuilderForType() { return newBuilder(); } - public static Builder newBuilder(com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan prototype) { - return newBuilder().mergeFrom(prototype); + public static Builder newBuilder() { + return DEFAULT_INSTANCE.toBuilder(); + } + public static Builder newBuilder(com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan prototype) { + return DEFAULT_INSTANCE.toBuilder().mergeFrom(prototype); + } + public Builder toBuilder() { + return this == DEFAULT_INSTANCE + ? new Builder() : new Builder().mergeFrom(this); } - public Builder toBuilder() { return newBuilder(this); } @java.lang.Override protected Builder newBuilderForType( - com.google.protobuf.GeneratedMessage.BuilderParent parent) { + com.google.protobuf.GeneratedMessageV3.BuilderParent parent) { Builder builder = new Builder(parent); return builder; } @@ -547,7 +608,7 @@ public final class TraceProtocol { * Protobuf type {@code AckSpan} */ public static final class Builder extends - com.google.protobuf.GeneratedMessage.Builder implements + com.google.protobuf.GeneratedMessageV3.Builder implements // @@protoc_insertion_point(builder_implements:AckSpan) com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpanOrBuilder { public static final com.google.protobuf.Descriptors.Descriptor @@ -555,7 +616,7 @@ public final class TraceProtocol { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_AckSpan_descriptor; } - protected com.google.protobuf.GeneratedMessage.FieldAccessorTable + protected com.google.protobuf.GeneratedMessageV3.FieldAccessorTable internalGetFieldAccessorTable() { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_AckSpan_fieldAccessorTable .ensureFieldAccessorsInitialized( @@ -568,18 +629,15 @@ public final class TraceProtocol { } private Builder( - com.google.protobuf.GeneratedMessage.BuilderParent parent) { + com.google.protobuf.GeneratedMessageV3.BuilderParent parent) { super(parent); maybeForceBuilderInitialization(); } private void maybeForceBuilderInitialization() { - if (com.google.protobuf.GeneratedMessage.alwaysUseFieldBuilders) { + if (com.google.protobuf.GeneratedMessageV3 + .alwaysUseFieldBuilders) { } } - private static Builder create() { - return new Builder(); - } - public Builder clear() { super.clear(); traceId_ = ""; @@ -597,10 +655,6 @@ public final class TraceProtocol { return this; } - public Builder clone() { - return create().mergeFrom(buildPartial()); - } - public com.google.protobuf.Descriptors.Descriptor getDescriptorForType() { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_AckSpan_descriptor; @@ -651,6 +705,32 @@ public final class TraceProtocol { return result; } + public Builder clone() { + return (Builder) super.clone(); + } + public Builder setField( + com.google.protobuf.Descriptors.FieldDescriptor field, + Object value) { + return (Builder) super.setField(field, value); + } + public Builder clearField( + com.google.protobuf.Descriptors.FieldDescriptor field) { + return (Builder) super.clearField(field); + } + public Builder clearOneof( + com.google.protobuf.Descriptors.OneofDescriptor oneof) { + return (Builder) super.clearOneof(oneof); + } + public Builder setRepeatedField( + com.google.protobuf.Descriptors.FieldDescriptor field, + int index, Object value) { + return (Builder) super.setRepeatedField(field, index, value); + } + public Builder addRepeatedField( + com.google.protobuf.Descriptors.FieldDescriptor field, + Object value) { + return (Builder) super.addRepeatedField(field, value); + } public Builder mergeFrom(com.google.protobuf.Message other) { if (other instanceof com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan) { return mergeFrom((com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan)other); @@ -686,25 +766,22 @@ public final class TraceProtocol { exceptionStack_ = other.exceptionStack_; onChanged(); } - this.mergeUnknownFields(other.getUnknownFields()); + this.mergeUnknownFields(other.unknownFields); + onChanged(); return this; } public final boolean isInitialized() { if (!hasTraceId()) { - return false; } if (!hasLevelId()) { - return false; } if (!hasCost()) { - return false; } if (!hasStatusCode()) { - return false; } return true; @@ -719,7 +796,7 @@ public final class TraceProtocol { parsedMessage = PARSER.parsePartialFrom(input, extensionRegistry); } catch (com.google.protobuf.InvalidProtocolBufferException e) { parsedMessage = (com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan) e.getUnfinishedMessage(); - throw e; + throw e.unwrapIOException(); } finally { if (parsedMessage != null) { mergeFrom(parsedMessage); @@ -1052,16 +1129,53 @@ public final class TraceProtocol { onChanged(); return this; } + public final Builder setUnknownFields( + final com.google.protobuf.UnknownFieldSet unknownFields) { + return super.setUnknownFields(unknownFields); + } + + public final Builder mergeUnknownFields( + final com.google.protobuf.UnknownFieldSet unknownFields) { + return super.mergeUnknownFields(unknownFields); + } + // @@protoc_insertion_point(builder_scope:AckSpan) } + // @@protoc_insertion_point(class_scope:AckSpan) + private static final com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan DEFAULT_INSTANCE; static { - defaultInstance = new AckSpan(true); - defaultInstance.initFields(); + DEFAULT_INSTANCE = new com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan(); + } + + public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan getDefaultInstance() { + return DEFAULT_INSTANCE; + } + + @java.lang.Deprecated public static final com.google.protobuf.Parser + PARSER = new com.google.protobuf.AbstractParser() { + public AckSpan parsePartialFrom( + com.google.protobuf.CodedInputStream input, + com.google.protobuf.ExtensionRegistryLite extensionRegistry) + throws com.google.protobuf.InvalidProtocolBufferException { + return new AckSpan(input, extensionRegistry); + } + }; + + public static com.google.protobuf.Parser parser() { + return PARSER; + } + + @java.lang.Override + public com.google.protobuf.Parser getParserForType() { + return PARSER; + } + + public com.ai.cloud.skywalking.protocol.proto.TraceProtocol.AckSpan getDefaultInstanceForType() { + return DEFAULT_INSTANCE; } - // @@protoc_insertion_point(class_scope:AckSpan) } public interface RequestSpanOrBuilder extends @@ -1220,41 +1334,77 @@ public final class TraceProtocol { */ com.google.protobuf.ByteString getAgentIdBytes(); + + /** + * map<string, string> parameters = 13; + */ + int getParametersCount(); + /** + * map<string, string> parameters = 13; + */ + boolean containsParameters( + java.lang.String key); + /** + * Use {@link #getParametersMap()} instead. + */ + @java.lang.Deprecated + java.util.Map + getParameters(); + /** + * map<string, string> parameters = 13; + */ + java.util.Map + getParametersMap(); + /** + * map<string, string> parameters = 13; + */ + + java.lang.String getParametersOrDefault( + java.lang.String key, + java.lang.String defaultValue); + /** + * map<string, string> parameters = 13; + */ + + java.lang.String getParametersOrThrow( + java.lang.String key); } /** * Protobuf type {@code RequestSpan} */ - public static final class RequestSpan extends - com.google.protobuf.GeneratedMessage implements + public static final class RequestSpan extends + com.google.protobuf.GeneratedMessageV3 implements // @@protoc_insertion_point(message_implements:RequestSpan) RequestSpanOrBuilder { // Use RequestSpan.newBuilder() to construct. - private RequestSpan(com.google.protobuf.GeneratedMessage.Builder builder) { + private RequestSpan(com.google.protobuf.GeneratedMessageV3.Builder builder) { super(builder); - this.unknownFields = builder.getUnknownFields(); } - private RequestSpan(boolean noInit) { this.unknownFields = com.google.protobuf.UnknownFieldSet.getDefaultInstance(); } - - private static final RequestSpan defaultInstance; - public static RequestSpan getDefaultInstance() { - return defaultInstance; + private RequestSpan() { + traceId_ = ""; + parentLevel_ = ""; + levelId_ = 0; + viewPointId_ = ""; + startDate_ = 0L; + spanTypeDesc_ = ""; + callType_ = ""; + spanType_ = 0; + applicationId_ = ""; + userId_ = ""; + bussinessKey_ = ""; + agentId_ = ""; } - public RequestSpan getDefaultInstanceForType() { - return defaultInstance; - } - - private final com.google.protobuf.UnknownFieldSet unknownFields; @java.lang.Override public final com.google.protobuf.UnknownFieldSet - getUnknownFields() { + getUnknownFields() { return this.unknownFields; } private RequestSpan( com.google.protobuf.CodedInputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws com.google.protobuf.InvalidProtocolBufferException { - initFields(); + this(); int mutable_bitField0_ = 0; com.google.protobuf.UnknownFieldSet.Builder unknownFields = com.google.protobuf.UnknownFieldSet.newBuilder(); @@ -1342,13 +1492,25 @@ public final class TraceProtocol { agentId_ = bs; break; } + case 106: { + if (!((mutable_bitField0_ & 0x00001000) == 0x00001000)) { + parameters_ = com.google.protobuf.MapField.newMapField( + ParametersDefaultEntryHolder.defaultEntry); + mutable_bitField0_ |= 0x00001000; + } + com.google.protobuf.MapEntry + parameters = input.readMessage( + ParametersDefaultEntryHolder.defaultEntry.getParserForType(), extensionRegistry); + parameters_.getMutableMap().put(parameters.getKey(), parameters.getValue()); + break; + } } } } catch (com.google.protobuf.InvalidProtocolBufferException e) { throw e.setUnfinishedMessage(this); } catch (java.io.IOException e) { throw new com.google.protobuf.InvalidProtocolBufferException( - e.getMessage()).setUnfinishedMessage(this); + e).setUnfinishedMessage(this); } finally { this.unknownFields = unknownFields.build(); makeExtensionsImmutable(); @@ -1359,31 +1521,27 @@ public final class TraceProtocol { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_RequestSpan_descriptor; } - protected com.google.protobuf.GeneratedMessage.FieldAccessorTable + @SuppressWarnings({"rawtypes"}) + protected com.google.protobuf.MapField internalGetMapField( + int number) { + switch (number) { + case 13: + return internalGetParameters(); + default: + throw new RuntimeException( + "Invalid map field number: " + number); + } + } + protected com.google.protobuf.GeneratedMessageV3.FieldAccessorTable internalGetFieldAccessorTable() { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_RequestSpan_fieldAccessorTable .ensureFieldAccessorsInitialized( com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan.class, com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan.Builder.class); } - public static com.google.protobuf.Parser PARSER = - new com.google.protobuf.AbstractParser() { - public RequestSpan parsePartialFrom( - com.google.protobuf.CodedInputStream input, - com.google.protobuf.ExtensionRegistryLite extensionRegistry) - throws com.google.protobuf.InvalidProtocolBufferException { - return new RequestSpan(input, extensionRegistry); - } - }; - - @java.lang.Override - public com.google.protobuf.Parser getParserForType() { - return PARSER; - } - private int bitField0_; public static final int TRACEID_FIELD_NUMBER = 1; - private java.lang.Object traceId_; + private volatile java.lang.Object traceId_; /** * required string traceId = 1; */ @@ -1425,7 +1583,7 @@ public final class TraceProtocol { } public static final int PARENTLEVEL_FIELD_NUMBER = 2; - private java.lang.Object parentLevel_; + private volatile java.lang.Object parentLevel_; /** * optional string parentLevel = 2; */ @@ -1482,7 +1640,7 @@ public final class TraceProtocol { } public static final int VIEWPOINTID_FIELD_NUMBER = 4; - private java.lang.Object viewPointId_; + private volatile java.lang.Object viewPointId_; /** * required string viewPointId = 4; */ @@ -1539,7 +1697,7 @@ public final class TraceProtocol { } public static final int SPANTYPEDESC_FIELD_NUMBER = 6; - private java.lang.Object spanTypeDesc_; + private volatile java.lang.Object spanTypeDesc_; /** * required string spanTypeDesc = 6; */ @@ -1581,7 +1739,7 @@ public final class TraceProtocol { } public static final int CALLTYPE_FIELD_NUMBER = 7; - private java.lang.Object callType_; + private volatile java.lang.Object callType_; /** * required string callType = 7; */ @@ -1638,7 +1796,7 @@ public final class TraceProtocol { } public static final int APPLICATIONID_FIELD_NUMBER = 9; - private java.lang.Object applicationId_; + private volatile java.lang.Object applicationId_; /** * required string applicationId = 9; */ @@ -1680,7 +1838,7 @@ public final class TraceProtocol { } public static final int USERID_FIELD_NUMBER = 10; - private java.lang.Object userId_; + private volatile java.lang.Object userId_; /** * required string userId = 10; */ @@ -1722,7 +1880,7 @@ public final class TraceProtocol { } public static final int BUSSINESSKEY_FIELD_NUMBER = 11; - private java.lang.Object bussinessKey_; + private volatile java.lang.Object bussinessKey_; /** * optional string bussinessKey = 11; */ @@ -1764,7 +1922,7 @@ public final class TraceProtocol { } public static final int AGENTID_FIELD_NUMBER = 12; - private java.lang.Object agentId_; + private volatile java.lang.Object agentId_; /** * required string agentId = 12; */ @@ -1805,20 +1963,82 @@ public final class TraceProtocol { } } - private void initFields() { - traceId_ = ""; - parentLevel_ = ""; - levelId_ = 0; - viewPointId_ = ""; - startDate_ = 0L; - spanTypeDesc_ = ""; - callType_ = ""; - spanType_ = 0; - applicationId_ = ""; - userId_ = ""; - bussinessKey_ = ""; - agentId_ = ""; + public static final int PARAMETERS_FIELD_NUMBER = 13; + private static final class ParametersDefaultEntryHolder { + static final com.google.protobuf.MapEntry< + java.lang.String, java.lang.String> defaultEntry = + com.google.protobuf.MapEntry + .newDefaultInstance( + com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_RequestSpan_ParametersEntry_descriptor, + com.google.protobuf.WireFormat.FieldType.STRING, + "", + com.google.protobuf.WireFormat.FieldType.STRING, + ""); } + private com.google.protobuf.MapField< + java.lang.String, java.lang.String> parameters_; + private com.google.protobuf.MapField + internalGetParameters() { + if (parameters_ == null) { + return com.google.protobuf.MapField.emptyMapField( + ParametersDefaultEntryHolder.defaultEntry); + } + return parameters_; + } + + public int getParametersCount() { + return internalGetParameters().getMap().size(); + } + /** + * map<string, string> parameters = 13; + */ + + public boolean containsParameters( + java.lang.String key) { + if (key == null) { throw new java.lang.NullPointerException(); } + return internalGetParameters().getMap().containsKey(key); + } + /** + * Use {@link #getParametersMap()} instead. + */ + @java.lang.Deprecated + public java.util.Map getParameters() { + return getParametersMap(); + } + /** + * map<string, string> parameters = 13; + */ + + public java.util.Map getParametersMap() { + return internalGetParameters().getMap(); + } + /** + * map<string, string> parameters = 13; + */ + + public java.lang.String getParametersOrDefault( + java.lang.String key, + java.lang.String defaultValue) { + if (key == null) { throw new java.lang.NullPointerException(); } + java.util.Map map = + internalGetParameters().getMap(); + return map.containsKey(key) ? map.get(key) : defaultValue; + } + /** + * map<string, string> parameters = 13; + */ + + public java.lang.String getParametersOrThrow( + java.lang.String key) { + if (key == null) { throw new java.lang.NullPointerException(); } + java.util.Map map = + internalGetParameters().getMap(); + if (!map.containsKey(key)) { + throw new java.lang.IllegalArgumentException(); + } + return map.get(key); + } + private byte memoizedIsInitialized = -1; public final boolean isInitialized() { byte isInitialized = memoizedIsInitialized; @@ -1871,110 +2091,254 @@ public final class TraceProtocol { public void writeTo(com.google.protobuf.CodedOutputStream output) throws java.io.IOException { - getSerializedSize(); if (((bitField0_ & 0x00000001) == 0x00000001)) { - output.writeBytes(1, getTraceIdBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 1, traceId_); } if (((bitField0_ & 0x00000002) == 0x00000002)) { - output.writeBytes(2, getParentLevelBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 2, parentLevel_); } if (((bitField0_ & 0x00000004) == 0x00000004)) { output.writeInt32(3, levelId_); } if (((bitField0_ & 0x00000008) == 0x00000008)) { - output.writeBytes(4, getViewPointIdBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 4, viewPointId_); } if (((bitField0_ & 0x00000010) == 0x00000010)) { output.writeInt64(5, startDate_); } if (((bitField0_ & 0x00000020) == 0x00000020)) { - output.writeBytes(6, getSpanTypeDescBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 6, spanTypeDesc_); } if (((bitField0_ & 0x00000040) == 0x00000040)) { - output.writeBytes(7, getCallTypeBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 7, callType_); } if (((bitField0_ & 0x00000080) == 0x00000080)) { output.writeUInt32(8, spanType_); } if (((bitField0_ & 0x00000100) == 0x00000100)) { - output.writeBytes(9, getApplicationIdBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 9, applicationId_); } if (((bitField0_ & 0x00000200) == 0x00000200)) { - output.writeBytes(10, getUserIdBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 10, userId_); } if (((bitField0_ & 0x00000400) == 0x00000400)) { - output.writeBytes(11, getBussinessKeyBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 11, bussinessKey_); } if (((bitField0_ & 0x00000800) == 0x00000800)) { - output.writeBytes(12, getAgentIdBytes()); + com.google.protobuf.GeneratedMessageV3.writeString(output, 12, agentId_); } - getUnknownFields().writeTo(output); + for (java.util.Map.Entry entry + : internalGetParameters().getMap().entrySet()) { + com.google.protobuf.MapEntry + parameters = ParametersDefaultEntryHolder.defaultEntry.newBuilderForType() + .setKey(entry.getKey()) + .setValue(entry.getValue()) + .build(); + output.writeMessage(13, parameters); + } + unknownFields.writeTo(output); } - private int memoizedSerializedSize = -1; public int getSerializedSize() { - int size = memoizedSerializedSize; + int size = memoizedSize; if (size != -1) return size; size = 0; if (((bitField0_ & 0x00000001) == 0x00000001)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(1, getTraceIdBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(1, traceId_); } if (((bitField0_ & 0x00000002) == 0x00000002)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(2, getParentLevelBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(2, parentLevel_); } if (((bitField0_ & 0x00000004) == 0x00000004)) { size += com.google.protobuf.CodedOutputStream .computeInt32Size(3, levelId_); } if (((bitField0_ & 0x00000008) == 0x00000008)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(4, getViewPointIdBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(4, viewPointId_); } if (((bitField0_ & 0x00000010) == 0x00000010)) { size += com.google.protobuf.CodedOutputStream .computeInt64Size(5, startDate_); } if (((bitField0_ & 0x00000020) == 0x00000020)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(6, getSpanTypeDescBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(6, spanTypeDesc_); } if (((bitField0_ & 0x00000040) == 0x00000040)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(7, getCallTypeBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(7, callType_); } if (((bitField0_ & 0x00000080) == 0x00000080)) { size += com.google.protobuf.CodedOutputStream .computeUInt32Size(8, spanType_); } if (((bitField0_ & 0x00000100) == 0x00000100)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(9, getApplicationIdBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(9, applicationId_); } if (((bitField0_ & 0x00000200) == 0x00000200)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(10, getUserIdBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(10, userId_); } if (((bitField0_ & 0x00000400) == 0x00000400)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(11, getBussinessKeyBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(11, bussinessKey_); } if (((bitField0_ & 0x00000800) == 0x00000800)) { - size += com.google.protobuf.CodedOutputStream - .computeBytesSize(12, getAgentIdBytes()); + size += com.google.protobuf.GeneratedMessageV3.computeStringSize(12, agentId_); } - size += getUnknownFields().getSerializedSize(); - memoizedSerializedSize = size; + for (java.util.Map.Entry entry + : internalGetParameters().getMap().entrySet()) { + com.google.protobuf.MapEntry + parameters = ParametersDefaultEntryHolder.defaultEntry.newBuilderForType() + .setKey(entry.getKey()) + .setValue(entry.getValue()) + .build(); + size += com.google.protobuf.CodedOutputStream + .computeMessageSize(13, parameters); + } + size += unknownFields.getSerializedSize(); + memoizedSize = size; return size; } private static final long serialVersionUID = 0L; @java.lang.Override - protected java.lang.Object writeReplace() - throws java.io.ObjectStreamException { - return super.writeReplace(); + public boolean equals(final java.lang.Object obj) { + if (obj == this) { + return true; + } + if (!(obj instanceof com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan)) { + return super.equals(obj); + } + com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan other = (com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan) obj; + + boolean result = true; + result = result && (hasTraceId() == other.hasTraceId()); + if (hasTraceId()) { + result = result && getTraceId() + .equals(other.getTraceId()); + } + result = result && (hasParentLevel() == other.hasParentLevel()); + if (hasParentLevel()) { + result = result && getParentLevel() + .equals(other.getParentLevel()); + } + result = result && (hasLevelId() == other.hasLevelId()); + if (hasLevelId()) { + result = result && (getLevelId() + == other.getLevelId()); + } + result = result && (hasViewPointId() == other.hasViewPointId()); + if (hasViewPointId()) { + result = result && getViewPointId() + .equals(other.getViewPointId()); + } + result = result && (hasStartDate() == other.hasStartDate()); + if (hasStartDate()) { + result = result && (getStartDate() + == other.getStartDate()); + } + result = result && (hasSpanTypeDesc() == other.hasSpanTypeDesc()); + if (hasSpanTypeDesc()) { + result = result && getSpanTypeDesc() + .equals(other.getSpanTypeDesc()); + } + result = result && (hasCallType() == other.hasCallType()); + if (hasCallType()) { + result = result && getCallType() + .equals(other.getCallType()); + } + result = result && (hasSpanType() == other.hasSpanType()); + if (hasSpanType()) { + result = result && (getSpanType() + == other.getSpanType()); + } + result = result && (hasApplicationId() == other.hasApplicationId()); + if (hasApplicationId()) { + result = result && getApplicationId() + .equals(other.getApplicationId()); + } + result = result && (hasUserId() == other.hasUserId()); + if (hasUserId()) { + result = result && getUserId() + .equals(other.getUserId()); + } + result = result && (hasBussinessKey() == other.hasBussinessKey()); + if (hasBussinessKey()) { + result = result && getBussinessKey() + .equals(other.getBussinessKey()); + } + result = result && (hasAgentId() == other.hasAgentId()); + if (hasAgentId()) { + result = result && getAgentId() + .equals(other.getAgentId()); + } + result = result && internalGetParameters().equals( + other.internalGetParameters()); + result = result && unknownFields.equals(other.unknownFields); + return result; + } + + @java.lang.Override + public int hashCode() { + if (memoizedHashCode != 0) { + return memoizedHashCode; + } + int hash = 41; + hash = (19 * hash) + getDescriptorForType().hashCode(); + if (hasTraceId()) { + hash = (37 * hash) + TRACEID_FIELD_NUMBER; + hash = (53 * hash) + getTraceId().hashCode(); + } + if (hasParentLevel()) { + hash = (37 * hash) + PARENTLEVEL_FIELD_NUMBER; + hash = (53 * hash) + getParentLevel().hashCode(); + } + if (hasLevelId()) { + hash = (37 * hash) + LEVELID_FIELD_NUMBER; + hash = (53 * hash) + getLevelId(); + } + if (hasViewPointId()) { + hash = (37 * hash) + VIEWPOINTID_FIELD_NUMBER; + hash = (53 * hash) + getViewPointId().hashCode(); + } + if (hasStartDate()) { + hash = (37 * hash) + STARTDATE_FIELD_NUMBER; + hash = (53 * hash) + com.google.protobuf.Internal.hashLong( + getStartDate()); + } + if (hasSpanTypeDesc()) { + hash = (37 * hash) + SPANTYPEDESC_FIELD_NUMBER; + hash = (53 * hash) + getSpanTypeDesc().hashCode(); + } + if (hasCallType()) { + hash = (37 * hash) + CALLTYPE_FIELD_NUMBER; + hash = (53 * hash) + getCallType().hashCode(); + } + if (hasSpanType()) { + hash = (37 * hash) + SPANTYPE_FIELD_NUMBER; + hash = (53 * hash) + getSpanType(); + } + if (hasApplicationId()) { + hash = (37 * hash) + APPLICATIONID_FIELD_NUMBER; + hash = (53 * hash) + getApplicationId().hashCode(); + } + if (hasUserId()) { + hash = (37 * hash) + USERID_FIELD_NUMBER; + hash = (53 * hash) + getUserId().hashCode(); + } + if (hasBussinessKey()) { + hash = (37 * hash) + BUSSINESSKEY_FIELD_NUMBER; + hash = (53 * hash) + getBussinessKey().hashCode(); + } + if (hasAgentId()) { + hash = (37 * hash) + AGENTID_FIELD_NUMBER; + hash = (53 * hash) + getAgentId().hashCode(); + } + if (!internalGetParameters().getMap().isEmpty()) { + hash = (37 * hash) + PARAMETERS_FIELD_NUMBER; + hash = (53 * hash) + internalGetParameters().hashCode(); + } + hash = (29 * hash) + unknownFields.hashCode(); + memoizedHashCode = hash; + return hash; } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan parseFrom( @@ -2000,46 +2364,57 @@ public final class TraceProtocol { } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan parseFrom(java.io.InputStream input) throws java.io.IOException { - return PARSER.parseFrom(input); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan parseFrom( java.io.InputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws java.io.IOException { - return PARSER.parseFrom(input, extensionRegistry); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input, extensionRegistry); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan parseDelimitedFrom(java.io.InputStream input) throws java.io.IOException { - return PARSER.parseDelimitedFrom(input); + return com.google.protobuf.GeneratedMessageV3 + .parseDelimitedWithIOException(PARSER, input); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan parseDelimitedFrom( java.io.InputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws java.io.IOException { - return PARSER.parseDelimitedFrom(input, extensionRegistry); + return com.google.protobuf.GeneratedMessageV3 + .parseDelimitedWithIOException(PARSER, input, extensionRegistry); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan parseFrom( com.google.protobuf.CodedInputStream input) throws java.io.IOException { - return PARSER.parseFrom(input); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input); } public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan parseFrom( com.google.protobuf.CodedInputStream input, com.google.protobuf.ExtensionRegistryLite extensionRegistry) throws java.io.IOException { - return PARSER.parseFrom(input, extensionRegistry); + return com.google.protobuf.GeneratedMessageV3 + .parseWithIOException(PARSER, input, extensionRegistry); } - public static Builder newBuilder() { return Builder.create(); } public Builder newBuilderForType() { return newBuilder(); } - public static Builder newBuilder(com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan prototype) { - return newBuilder().mergeFrom(prototype); + public static Builder newBuilder() { + return DEFAULT_INSTANCE.toBuilder(); + } + public static Builder newBuilder(com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan prototype) { + return DEFAULT_INSTANCE.toBuilder().mergeFrom(prototype); + } + public Builder toBuilder() { + return this == DEFAULT_INSTANCE + ? new Builder() : new Builder().mergeFrom(this); } - public Builder toBuilder() { return newBuilder(this); } @java.lang.Override protected Builder newBuilderForType( - com.google.protobuf.GeneratedMessage.BuilderParent parent) { + com.google.protobuf.GeneratedMessageV3.BuilderParent parent) { Builder builder = new Builder(parent); return builder; } @@ -2047,7 +2422,7 @@ public final class TraceProtocol { * Protobuf type {@code RequestSpan} */ public static final class Builder extends - com.google.protobuf.GeneratedMessage.Builder implements + com.google.protobuf.GeneratedMessageV3.Builder implements // @@protoc_insertion_point(builder_implements:RequestSpan) com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpanOrBuilder { public static final com.google.protobuf.Descriptors.Descriptor @@ -2055,7 +2430,29 @@ public final class TraceProtocol { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_RequestSpan_descriptor; } - protected com.google.protobuf.GeneratedMessage.FieldAccessorTable + @SuppressWarnings({"rawtypes"}) + protected com.google.protobuf.MapField internalGetMapField( + int number) { + switch (number) { + case 13: + return internalGetParameters(); + default: + throw new RuntimeException( + "Invalid map field number: " + number); + } + } + @SuppressWarnings({"rawtypes"}) + protected com.google.protobuf.MapField internalGetMutableMapField( + int number) { + switch (number) { + case 13: + return internalGetMutableParameters(); + default: + throw new RuntimeException( + "Invalid map field number: " + number); + } + } + protected com.google.protobuf.GeneratedMessageV3.FieldAccessorTable internalGetFieldAccessorTable() { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_RequestSpan_fieldAccessorTable .ensureFieldAccessorsInitialized( @@ -2068,18 +2465,15 @@ public final class TraceProtocol { } private Builder( - com.google.protobuf.GeneratedMessage.BuilderParent parent) { + com.google.protobuf.GeneratedMessageV3.BuilderParent parent) { super(parent); maybeForceBuilderInitialization(); } private void maybeForceBuilderInitialization() { - if (com.google.protobuf.GeneratedMessage.alwaysUseFieldBuilders) { + if (com.google.protobuf.GeneratedMessageV3 + .alwaysUseFieldBuilders) { } } - private static Builder create() { - return new Builder(); - } - public Builder clear() { super.clear(); traceId_ = ""; @@ -2106,13 +2500,10 @@ public final class TraceProtocol { bitField0_ = (bitField0_ & ~0x00000400); agentId_ = ""; bitField0_ = (bitField0_ & ~0x00000800); + internalGetMutableParameters().clear(); return this; } - public Builder clone() { - return create().mergeFrom(buildPartial()); - } - public com.google.protobuf.Descriptors.Descriptor getDescriptorForType() { return com.ai.cloud.skywalking.protocol.proto.TraceProtocol.internal_static_RequestSpan_descriptor; @@ -2182,11 +2573,39 @@ public final class TraceProtocol { to_bitField0_ |= 0x00000800; } result.agentId_ = agentId_; + result.parameters_ = internalGetParameters(); + result.parameters_.makeImmutable(); result.bitField0_ = to_bitField0_; onBuilt(); return result; } + public Builder clone() { + return (Builder) super.clone(); + } + public Builder setField( + com.google.protobuf.Descriptors.FieldDescriptor field, + Object value) { + return (Builder) super.setField(field, value); + } + public Builder clearField( + com.google.protobuf.Descriptors.FieldDescriptor field) { + return (Builder) super.clearField(field); + } + public Builder clearOneof( + com.google.protobuf.Descriptors.OneofDescriptor oneof) { + return (Builder) super.clearOneof(oneof); + } + public Builder setRepeatedField( + com.google.protobuf.Descriptors.FieldDescriptor field, + int index, Object value) { + return (Builder) super.setRepeatedField(field, index, value); + } + public Builder addRepeatedField( + com.google.protobuf.Descriptors.FieldDescriptor field, + Object value) { + return (Builder) super.addRepeatedField(field, value); + } public Builder mergeFrom(com.google.protobuf.Message other) { if (other instanceof com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan) { return mergeFrom((com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan)other); @@ -2252,49 +2671,42 @@ public final class TraceProtocol { agentId_ = other.agentId_; onChanged(); } - this.mergeUnknownFields(other.getUnknownFields()); + internalGetMutableParameters().mergeFrom( + other.internalGetParameters()); + this.mergeUnknownFields(other.unknownFields); + onChanged(); return this; } public final boolean isInitialized() { if (!hasTraceId()) { - return false; } if (!hasLevelId()) { - return false; } if (!hasViewPointId()) { - return false; } if (!hasStartDate()) { - return false; } if (!hasSpanTypeDesc()) { - return false; } if (!hasCallType()) { - return false; } if (!hasSpanType()) { - return false; } if (!hasApplicationId()) { - return false; } if (!hasUserId()) { - return false; } if (!hasAgentId()) { - return false; } return true; @@ -2309,7 +2721,7 @@ public final class TraceProtocol { parsedMessage = PARSER.parsePartialFrom(input, extensionRegistry); } catch (com.google.protobuf.InvalidProtocolBufferException e) { parsedMessage = (com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan) e.getUnfinishedMessage(); - throw e; + throw e.unwrapIOException(); } finally { if (parsedMessage != null) { mergeFrom(parsedMessage); @@ -3099,47 +3511,211 @@ public final class TraceProtocol { return this; } + private com.google.protobuf.MapField< + java.lang.String, java.lang.String> parameters_; + private com.google.protobuf.MapField + internalGetParameters() { + if (parameters_ == null) { + return com.google.protobuf.MapField.emptyMapField( + ParametersDefaultEntryHolder.defaultEntry); + } + return parameters_; + } + private com.google.protobuf.MapField + internalGetMutableParameters() { + onChanged();; + if (parameters_ == null) { + parameters_ = com.google.protobuf.MapField.newMapField( + ParametersDefaultEntryHolder.defaultEntry); + } + if (!parameters_.isMutable()) { + parameters_ = parameters_.copy(); + } + return parameters_; + } + + public int getParametersCount() { + return internalGetParameters().getMap().size(); + } + /** + * map<string, string> parameters = 13; + */ + + public boolean containsParameters( + java.lang.String key) { + if (key == null) { throw new java.lang.NullPointerException(); } + return internalGetParameters().getMap().containsKey(key); + } + /** + * Use {@link #getParametersMap()} instead. + */ + @java.lang.Deprecated + public java.util.Map getParameters() { + return getParametersMap(); + } + /** + * map<string, string> parameters = 13; + */ + + public java.util.Map getParametersMap() { + return internalGetParameters().getMap(); + } + /** + * map<string, string> parameters = 13; + */ + + public java.lang.String getParametersOrDefault( + java.lang.String key, + java.lang.String defaultValue) { + if (key == null) { throw new java.lang.NullPointerException(); } + java.util.Map map = + internalGetParameters().getMap(); + return map.containsKey(key) ? map.get(key) : defaultValue; + } + /** + * map<string, string> parameters = 13; + */ + + public java.lang.String getParametersOrThrow( + java.lang.String key) { + if (key == null) { throw new java.lang.NullPointerException(); } + java.util.Map map = + internalGetParameters().getMap(); + if (!map.containsKey(key)) { + throw new java.lang.IllegalArgumentException(); + } + return map.get(key); + } + + public Builder clearParameters() { + getMutableParameters().clear(); + return this; + } + /** + * map<string, string> parameters = 13; + */ + + public Builder removeParameters( + java.lang.String key) { + if (key == null) { throw new java.lang.NullPointerException(); } + getMutableParameters().remove(key); + return this; + } + /** + * Use alternate mutation accessors instead. + */ + @java.lang.Deprecated + public java.util.Map + getMutableParameters() { + return internalGetMutableParameters().getMutableMap(); + } + /** + * map<string, string> parameters = 13; + */ + public Builder putParameters( + java.lang.String key, + java.lang.String value) { + if (key == null) { throw new java.lang.NullPointerException(); } + if (value == null) { throw new java.lang.NullPointerException(); } + getMutableParameters().put(key, value); + return this; + } + /** + * map<string, string> parameters = 13; + */ + + public Builder putAllParameters( + java.util.Map values) { + getMutableParameters().putAll(values); + return this; + } + public final Builder setUnknownFields( + final com.google.protobuf.UnknownFieldSet unknownFields) { + return super.setUnknownFields(unknownFields); + } + + public final Builder mergeUnknownFields( + final com.google.protobuf.UnknownFieldSet unknownFields) { + return super.mergeUnknownFields(unknownFields); + } + + // @@protoc_insertion_point(builder_scope:RequestSpan) } + // @@protoc_insertion_point(class_scope:RequestSpan) + private static final com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan DEFAULT_INSTANCE; static { - defaultInstance = new RequestSpan(true); - defaultInstance.initFields(); + DEFAULT_INSTANCE = new com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan(); + } + + public static com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan getDefaultInstance() { + return DEFAULT_INSTANCE; + } + + @java.lang.Deprecated public static final com.google.protobuf.Parser + PARSER = new com.google.protobuf.AbstractParser() { + public RequestSpan parsePartialFrom( + com.google.protobuf.CodedInputStream input, + com.google.protobuf.ExtensionRegistryLite extensionRegistry) + throws com.google.protobuf.InvalidProtocolBufferException { + return new RequestSpan(input, extensionRegistry); + } + }; + + public static com.google.protobuf.Parser parser() { + return PARSER; + } + + @java.lang.Override + public com.google.protobuf.Parser getParserForType() { + return PARSER; + } + + public com.ai.cloud.skywalking.protocol.proto.TraceProtocol.RequestSpan getDefaultInstanceForType() { + return DEFAULT_INSTANCE; } - // @@protoc_insertion_point(class_scope:RequestSpan) } private static final com.google.protobuf.Descriptors.Descriptor internal_static_AckSpan_descriptor; - private static - com.google.protobuf.GeneratedMessage.FieldAccessorTable + private static final + com.google.protobuf.GeneratedMessageV3.FieldAccessorTable internal_static_AckSpan_fieldAccessorTable; private static final com.google.protobuf.Descriptors.Descriptor internal_static_RequestSpan_descriptor; - private static - com.google.protobuf.GeneratedMessage.FieldAccessorTable + private static final + com.google.protobuf.GeneratedMessageV3.FieldAccessorTable internal_static_RequestSpan_fieldAccessorTable; + private static final com.google.protobuf.Descriptors.Descriptor + internal_static_RequestSpan_ParametersEntry_descriptor; + private static final + com.google.protobuf.GeneratedMessageV3.FieldAccessorTable + internal_static_RequestSpan_ParametersEntry_fieldAccessorTable; public static com.google.protobuf.Descriptors.FileDescriptor getDescriptor() { return descriptor; } - private static com.google.protobuf.Descriptors.FileDescriptor + private static com.google.protobuf.Descriptors.FileDescriptor descriptor; static { java.lang.String[] descriptorData = { "\n\023TraceProtocol.proto\"z\n\007AckSpan\022\017\n\007trac" + "eId\030\001 \002(\t\022\023\n\013parentLevel\030\002 \001(\t\022\017\n\007levelI" + "d\030\003 \002(\005\022\014\n\004cost\030\004 \002(\003\022\022\n\nstatusCode\030\005 \002(" + - "\005\022\026\n\016exceptionStack\030\006 \001(\t\"\364\001\n\013RequestSpa" + + "\005\022\026\n\016exceptionStack\030\006 \001(\t\"\331\002\n\013RequestSpa" + "n\022\017\n\007traceId\030\001 \002(\t\022\023\n\013parentLevel\030\002 \001(\t\022" + "\017\n\007levelId\030\003 \002(\005\022\023\n\013viewPointId\030\004 \002(\t\022\021\n" + "\tstartDate\030\005 \002(\003\022\024\n\014spanTypeDesc\030\006 \002(\t\022\020" + "\n\010callType\030\007 \002(\t\022\020\n\010spanType\030\010 \002(\r\022\025\n\rap" + "plicationId\030\t \002(\t\022\016\n\006userId\030\n \002(\t\022\024\n\014bus" + - "sinessKey\030\013 \001(\t\022\017\n\007agentId\030\014 \002(\tB(\n&com.", - "ai.cloud.skywalking.protocol.proto" + "sinessKey\030\013 \001(\t\022\017\n\007agentId\030\014 \002(\t\0220\n\npara", + "meters\030\r \003(\0132\034.RequestSpan.ParametersEnt" + + "ry\0321\n\017ParametersEntry\022\013\n\003key\030\001 \001(\t\022\r\n\005va" + + "lue\030\002 \001(\t:\0028\001B(\n&com.ai.cloud.skywalking" + + ".protocol.proto" }; com.google.protobuf.Descriptors.FileDescriptor.InternalDescriptorAssigner assigner = new com.google.protobuf.Descriptors.FileDescriptor. InternalDescriptorAssigner() { @@ -3156,15 +3732,21 @@ public final class TraceProtocol { internal_static_AckSpan_descriptor = getDescriptor().getMessageTypes().get(0); internal_static_AckSpan_fieldAccessorTable = new - com.google.protobuf.GeneratedMessage.FieldAccessorTable( + com.google.protobuf.GeneratedMessageV3.FieldAccessorTable( internal_static_AckSpan_descriptor, new java.lang.String[] { "TraceId", "ParentLevel", "LevelId", "Cost", "StatusCode", "ExceptionStack", }); internal_static_RequestSpan_descriptor = getDescriptor().getMessageTypes().get(1); internal_static_RequestSpan_fieldAccessorTable = new - com.google.protobuf.GeneratedMessage.FieldAccessorTable( + com.google.protobuf.GeneratedMessageV3.FieldAccessorTable( internal_static_RequestSpan_descriptor, - new java.lang.String[] { "TraceId", "ParentLevel", "LevelId", "ViewPointId", "StartDate", "SpanTypeDesc", "CallType", "SpanType", "ApplicationId", "UserId", "BussinessKey", "AgentId", }); + new java.lang.String[] { "TraceId", "ParentLevel", "LevelId", "ViewPointId", "StartDate", "SpanTypeDesc", "CallType", "SpanType", "ApplicationId", "UserId", "BussinessKey", "AgentId", "Parameters", }); + internal_static_RequestSpan_ParametersEntry_descriptor = + internal_static_RequestSpan_descriptor.getNestedTypes().get(0); + internal_static_RequestSpan_ParametersEntry_fieldAccessorTable = new + com.google.protobuf.GeneratedMessageV3.FieldAccessorTable( + internal_static_RequestSpan_ParametersEntry_descriptor, + new java.lang.String[] { "Key", "Value", }); } // @@protoc_insertion_point(outer_class_scope) diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java index 820ee2ecc..7525f5c22 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/RequestSpan.java @@ -213,11 +213,20 @@ public class RequestSpan extends AbstractDataSerializable { @Override public byte[] getData() { - return TraceProtocol.RequestSpan.newBuilder().setTraceId(traceId).setParentLevel(parentLevel) - .setLevelId(levelId).setViewPointId(viewPointId).setStartDate(startDate) - .setSpanType(spanType.getValue()).setSpanTypeDesc(spanTypeDesc).setBussinessKey(businessKey) - .setCallType(callType).setApplicationId(applicationId).setUserId(userId).setBussinessKey(businessKey) - .setAgentId(agentId).build().toByteArray(); + TraceProtocol.RequestSpan.Builder builder = + TraceProtocol.RequestSpan.newBuilder().setTraceId(traceId).setParentLevel(parentLevel) + .setLevelId(levelId).setViewPointId(viewPointId).setStartDate(startDate) + .setSpanType(spanType.getValue()).setSpanTypeDesc(spanTypeDesc); + if (businessKey != null && businessKey.length() > 0) { + builder.setBussinessKey(businessKey); + } + + if (parameters != null && parameters.size() > 0) { + builder.getParametersMap().putAll(parameters); + } + + return builder.setCallType(callType).setApplicationId(applicationId).setUserId(userId).setAgentId(agentId) + .build().toByteArray(); } @Override diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java index 8fd2d7f25..9cbcd6ac5 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/SerializedFactory.java @@ -22,7 +22,7 @@ public class SerializedFactory { } } - public static AbstractDataSerializable unSerialize(byte[] bytes) throws ConvertFailedException { + public static AbstractDataSerializable deserialize(byte[] bytes) throws ConvertFailedException { try { AbstractDataSerializable abstractDataSerializable = serializableMap.get(IntegerAssist.bytesToInt(bytes, 0)); if (abstractDataSerializable != null) { 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 c9a1c9051..58c9746ec 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 @@ -47,12 +47,6 @@ public class Span { */ protected String exceptionStack = ""; - /** - * 节点的状态
- * 不参与序列化 - */ - protected boolean isValidate = true; - /** * 节点调用过程中的业务字段
* 如:业务系统设置的订单号,SQL语句等 @@ -82,6 +76,7 @@ public class Span { this.traceId = traceId; this.applicationId = applicationId; this.userId = userId; + this.parentLevel = ""; } public Span(String traceId, String parentLevel, int levelId, String applicationId, String userId) { @@ -120,10 +115,6 @@ public class Span { this.startDate = startDate; } - public boolean isValidate() { - return isValidate; - } - public byte getStatusCode() { return statusCode; } @@ -144,17 +135,6 @@ public class Span { this.parameters = parameters; } - public void setValidate(boolean validate) { - isValidate = validate; - } - - public boolean isRPCClientSpan() { - if (spanType == SpanType.RPC_CLIENT) { - return true; - } - return false; - } - public void setSpanType(SpanType spanType) { this.spanType = spanType; } @@ -228,10 +208,6 @@ public class Span { this.viewPointId = viewPointId; } - public void appendParameter(String key, String value) { - this.parameters.put(key, value); - } - public void setInvokeResult(String result){ this.parameters.put(INVOKE_RESULT_PARAMETER_KEY, result); } diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java index a4eb659fd..46fa4ec78 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/TransportPackager.java @@ -9,58 +9,21 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.List; -public class TransportPackager { +import static com.ai.cloud.skywalking.protocol.util.ByteDataUtil.generateChecksum; +public class TransportPackager { + public static byte[] pack(List beSendingData) { // 对协议格式进行修改 // | check sum(4 byte) | data - byte[] dataText = packSerializableObjects(beSendingData); - byte[] dataPackage = packCheckSum(dataText); + byte[] data = serializeObjects(beSendingData); + byte[] dataPackage = appendCheckSum(data); return dataPackage; } - public static List unpackSerializableObjects(byte[] dataPackage) { - List serializeData = new ArrayList(); - int currentLength = 0; - while (true) { - //读取长度 - int dataLength = IntegerAssist.bytesToInt(dataPackage, currentLength); - // 反序列化 - byte[] data = new byte[dataLength]; - System.arraycopy(dataPackage, currentLength + 4, data, 0, dataLength); - try { - AbstractDataSerializable abstractDataSerializable = SerializedFactory.unSerialize(data); - serializeData.add(abstractDataSerializable); - } catch (ConvertFailedException e) { - e.printStackTrace(); - } - currentLength += 4 + dataLength; - if (currentLength >= dataPackage.length) { - break; - } - } - - return serializeData; - } - - /** - * 生成校验和参数 - * - * @param data - * @return - */ - private static byte[] generateChecksum(byte[] data, int offset) { - int result = data[offset]; - for (int i = offset + 1; i < data.length; i++) { - result ^= data[i]; - } - - return IntegerAssist.intToBytes(result); - } - - private static byte[] packCheckSum(byte[] dataText) { + private static byte[] appendCheckSum(byte[] dataText) { byte[] dataPackage = new byte[dataText.length + 4]; byte[] checkSum = generateChecksum(dataText, 0); System.arraycopy(checkSum, 0, dataPackage, 0, 4); @@ -68,19 +31,20 @@ public class TransportPackager { return dataPackage; } - private static byte[] packSerializableObjects(List beSendingData) { - byte[] dataText = null; + + private static byte[] serializeObjects(List beSendingData) { + byte[] data = null; int currentIndex = 0; for (ISerializable sendingData : beSendingData) { - byte[] dataElementText = packSerializableObject(sendingData); - dataText = appendingDataBytes(dataText, currentIndex, dataElementText); - currentIndex += dataElementText.length; + byte[] elementData = serialize(sendingData); + data = appendData(data, currentIndex, elementData); + currentIndex += elementData.length; } - return dataText; + return data; } - private static byte[] appendingDataBytes(byte[] dataText, int currentIndex, byte[] dataElementText) { + private static byte[] appendData(byte[] dataText, int currentIndex, byte[] dataElementText) { if (dataText == null) { dataText = new byte[dataElementText.length]; } else { @@ -91,9 +55,10 @@ public class TransportPackager { } - public static byte[] packSerializableObject(ISerializable serializable) { + public static byte[] serialize(ISerializable serializable) { byte[] serializableBytes = serializable.convert2Bytes(); - byte[] dataText = Arrays.copyOf(serializableBytes, serializableBytes.length + 4); + byte[] dataText = new byte[serializableBytes.length + 4]; + System.arraycopy(serializableBytes, 0, dataText, 4, serializableBytes.length); byte[] length = IntegerAssist.intToBytes(serializableBytes.length); System.arraycopy(length, 0, dataText, 0, 4); return dataText; diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/SpanType.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/SpanType.java index 3086de9e1..d864e4df3 100644 --- a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/SpanType.java +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/common/SpanType.java @@ -23,7 +23,7 @@ public enum SpanType { return LOCAL; case 2: return RPC_CLIENT; - case 3: + case 4: return RPC_SERVER; default: throw new SpanTypeCannotConvertException(spanTypeValue + ""); diff --git a/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/util/ByteDataUtil.java b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/util/ByteDataUtil.java new file mode 100644 index 000000000..a00b49067 --- /dev/null +++ b/skywalking-collector/skywalking-protocol/src/main/java/com/ai/cloud/skywalking/protocol/util/ByteDataUtil.java @@ -0,0 +1,28 @@ +package com.ai.cloud.skywalking.protocol.util; + + +import java.util.Arrays; + +public class ByteDataUtil { + public static byte[] unpackCheckSum(byte[] msg) { + return Arrays.copyOfRange(msg, 4, msg.length); + } + + + public static boolean validateCheckSum(byte[] dataPackage) { + byte[] checkSum = generateChecksum(dataPackage, 4); + byte[] originCheckSum = new byte[4]; + System.arraycopy(dataPackage, 0, originCheckSum, 0, 4); + return Arrays.equals(checkSum, originCheckSum); + } + + + public static byte[] generateChecksum(byte[] data, int offset) { + int result = data[offset]; + for (int i = offset + 1; i < data.length; i++) { + result ^= data[i]; + } + + return IntegerAssist.intToBytes(result); + } +} diff --git a/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto b/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto index 77b490abd..f5f4ec143 100644 --- a/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto +++ b/skywalking-collector/skywalking-protocol/src/main/proto/TraceProtocol.proto @@ -25,4 +25,5 @@ message RequestSpan { required string userId = 10; optional string bussinessKey = 11; required string agentId = 12; + map parameters = 13; } diff --git a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml index abfe850df..b1adec626 100644 --- a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml +++ b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/pom.xml @@ -97,20 +97,21 @@ compile --> + com.alibaba dubbo 2.5.3 provided - --> + org.apache.httpcomponents httpclient diff --git a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/MonitorFilterInterceptor.java b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/MonitorFilterInterceptor.java index 88ba1c0c3..753c0e17c 100644 --- a/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/MonitorFilterInterceptor.java +++ b/skywalking-collector/skywalking-sdk-plugin/dubbo-plugin/src/main/java/com/ai/cloud/skywalking/plugin/dubbo/MonitorFilterInterceptor.java @@ -33,9 +33,7 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto boolean isConsumer = rpcContext.isConsumerSide(); context.set("isConsumer", isConsumer); if (isConsumer) { - RPCClientInvokeMonitor rpcClientInvokeMonitor = new RPCClientInvokeMonitor(); - context.set("rpcClientInvokeMonitor", rpcClientInvokeMonitor); - ContextData contextData = rpcClientInvokeMonitor.beforeInvoke(createIdentification(invoker, invocation)); + ContextData contextData = new RPCClientInvokeMonitor().beforeInvoke(createIdentification(invoker, invocation)); String contextDataStr = contextData.toString(); //追加参数 @@ -60,8 +58,6 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto } } else { // 读取参数 - RPCServerInvokeMonitor rpcServerInvokeMonitor = new RPCServerInvokeMonitor(); - context.set("rpcServerInvokeMonitor", rpcServerInvokeMonitor); String contextDataStr; if (!BugFixAcitve.isActive) { @@ -75,7 +71,7 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto contextData = new ContextData(contextDataStr); } - rpcServerInvokeMonitor.beforeInvoke(contextData, createIdentification(invoker, invocation)); + new RPCServerInvokeMonitor().beforeInvoke(contextData, createIdentification(invoker, invocation)); } } @@ -88,6 +84,12 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto dealException(result.getException(), context); } + if (isConsumer(context)){ + new RPCClientInvokeMonitor().afterInvoke(); + }else{ + new RPCServerInvokeMonitor().afterInvoke(); + } + return ret; } @@ -97,12 +99,16 @@ public class MonitorFilterInterceptor implements InstanceMethodsAroundIntercepto dealException(t, context); } + + private boolean isConsumer(EnhancedClassInstanceContext context){ + return (boolean) context.get("isConsumer"); + } + private void dealException(Throwable t, EnhancedClassInstanceContext context) { - boolean isConsumer = (boolean) context.get("isConsumer"); - if (isConsumer) { - ((RPCClientInvokeMonitor) context.get("rpcClientInvokeMonitor")).occurException(t); + if (isConsumer(context)) { + new RPCClientInvokeMonitor().occurException(t); } else { - ((RPCServerInvokeMonitor) context.get("rpcServerInvokeMonitor")).occurException(t); + new RPCServerInvokeMonitor().occurException(t); } } diff --git a/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/AbstractDatabasePluginDefine.java b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/AbstractDatabasePluginDefine.java new file mode 100644 index 000000000..37f3b854a --- /dev/null +++ b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/AbstractDatabasePluginDefine.java @@ -0,0 +1,17 @@ +package com.ai.cloud.skywalking.plugin.jdbc.define; + +import com.ai.cloud.skywalking.plugin.interceptor.MethodMatcher; +import com.ai.cloud.skywalking.plugin.interceptor.enhance.ClassInstanceMethodsEnhancePluginDefine; +import com.ai.cloud.skywalking.plugin.interceptor.matcher.SimpleMethodMatcher; + +public abstract class AbstractDatabasePluginDefine extends ClassInstanceMethodsEnhancePluginDefine { + @Override + protected MethodMatcher[] getInstanceMethodsMatchers() { + return new MethodMatcher[]{new SimpleMethodMatcher("connect")}; + } + + @Override + protected String getInstanceMethodsInterceptor() { + return "com.ai.cloud.skywalking.plugin.jdbc.define.DatabasePluginInterceptor"; + } +} diff --git a/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/DatabasePluginInterceptor.java b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/DatabasePluginInterceptor.java new file mode 100644 index 000000000..d39eddec1 --- /dev/null +++ b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/define/DatabasePluginInterceptor.java @@ -0,0 +1,36 @@ +package com.ai.cloud.skywalking.plugin.jdbc.define; + +import com.ai.cloud.skywalking.plugin.interceptor.EnhancedClassInstanceContext; +import com.ai.cloud.skywalking.plugin.interceptor.enhance.ConstructorInvokeContext; +import com.ai.cloud.skywalking.plugin.interceptor.enhance.InstanceMethodInvokeContext; +import com.ai.cloud.skywalking.plugin.interceptor.enhance.InstanceMethodsAroundInterceptor; +import com.ai.cloud.skywalking.plugin.interceptor.enhance.MethodInterceptResult; +import com.ai.cloud.skywalking.plugin.jdbc.SWConnection; + +import java.sql.Connection; +import java.util.Properties; + +public class DatabasePluginInterceptor implements InstanceMethodsAroundInterceptor { + @Override + public void onConstruct(EnhancedClassInstanceContext context, ConstructorInvokeContext interceptorContext) { + } + + @Override + public void beforeMethod(EnhancedClassInstanceContext context, InstanceMethodInvokeContext interceptorContext, + MethodInterceptResult result) { + + } + + @Override + public Object afterMethod(EnhancedClassInstanceContext context, InstanceMethodInvokeContext interceptorContext, + Object ret) { + return new SWConnection((String) interceptorContext.allArguments()[0], + (Properties) interceptorContext.allArguments()[1], (Connection) ret); + } + + @Override + public void handleMethodException(Throwable t, EnhancedClassInstanceContext context, + InstanceMethodInvokeContext interceptorContext, Object ret) { + + } +} diff --git a/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/plugin/H2DatabasePluginDefine.java b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/plugin/H2DatabasePluginDefine.java new file mode 100644 index 000000000..09dc09134 --- /dev/null +++ b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/java/com/ai/cloud/skywalking/plugin/jdbc/plugin/H2DatabasePluginDefine.java @@ -0,0 +1,10 @@ +package com.ai.cloud.skywalking.plugin.jdbc.plugin; + +import com.ai.cloud.skywalking.plugin.jdbc.define.AbstractDatabasePluginDefine; + +public class H2DatabasePluginDefine extends AbstractDatabasePluginDefine { + @Override + protected String enhanceClassName() { + return "org.h2.Driver"; + } +} diff --git a/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/resources/skywalking-plugin.def b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/resources/skywalking-plugin.def index c371f19ab..0d21475f6 100644 --- a/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/resources/skywalking-plugin.def +++ b/skywalking-collector/skywalking-sdk-plugin/jdbc-plugin/src/main/resources/skywalking-plugin.def @@ -1 +1 @@ -com.ai.cloud.skywalking.plugin.jdbc.JDBCPluginDefine +com.ai.cloud.skywalking.plugin.jdbc.plugin.H2DatabasePluginDefine diff --git a/skywalking-collector/skywalking-sdk-plugin/tomcat-7.x-8.x-plugin/src/main/java/com/ai/cloud/skywalking/plugin/tomcat78x/TomcatPluginInterceptor.java b/skywalking-collector/skywalking-sdk-plugin/tomcat-7.x-8.x-plugin/src/main/java/com/ai/cloud/skywalking/plugin/tomcat78x/TomcatPluginInterceptor.java index e1f709ae4..40d86607d 100644 --- a/skywalking-collector/skywalking-sdk-plugin/tomcat-7.x-8.x-plugin/src/main/java/com/ai/cloud/skywalking/plugin/tomcat78x/TomcatPluginInterceptor.java +++ b/skywalking-collector/skywalking-sdk-plugin/tomcat-7.x-8.x-plugin/src/main/java/com/ai/cloud/skywalking/plugin/tomcat78x/TomcatPluginInterceptor.java @@ -62,11 +62,12 @@ public class TomcatPluginInterceptor implements InstanceMethodsAroundInterceptor Object[] args = interceptorContext.allArguments(); HttpServletResponse httpServletResponse = (HttpServletResponse) args[1]; httpServletResponse.addHeader(TRACE_ID_HEADER_NAME, Tracing.getTraceId()); + new RPCServerInvokeMonitor().afterInvoke(); return ret; } @Override public void handleMethodException(Throwable t, EnhancedClassInstanceContext context, InstanceMethodInvokeContext interceptorContext, Object ret) { - // DO Nothing + new RPCServerInvokeMonitor().occurException(t); } } diff --git a/skywalking-server/pom.xml b/skywalking-server/pom.xml index f7089f9a0..aa998daa1 100644 --- a/skywalking-server/pom.xml +++ b/skywalking-server/pom.xml @@ -56,6 +56,11 @@ gson 2.2.2 + + com.google.protobuf + protobuf-java + 3.0.0 + diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java index 6a8201419..4ecaa6937 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/AppendEOFFlagThread.java @@ -2,7 +2,6 @@ package com.ai.cloud.skywalking.reciever.buffer; import com.ai.cloud.skywalking.protocol.BufferFileEOFProtocol; import com.ai.cloud.skywalking.protocol.TransportPackager; -import com.ai.cloud.skywalking.reciever.model.BufferDataPackagerGenerator; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -29,8 +28,8 @@ class AppendEOFFlagThread extends Thread { try { logger.info("Add EOF flags to unprocessed data file[{}]", file.getName()); fileOutputStream = new FileOutputStream(new File(file.getParent(), file.getName()), true); - fileOutputStream.write(BufferDataPackagerGenerator - .pack(TransportPackager.packSerializableObject(new BufferFileEOFProtocol()))); + fileOutputStream.write(BufferDataAssist + .appendLengthAndSplit(TransportPackager.serialize(new BufferFileEOFProtocol()))); } catch (IOException e) { logger.info("Add EOF flags to the unprocessed data file failed.", e); } finally { diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BufferDataPackagerGenerator.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/BufferDataAssist.java similarity index 56% rename from skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BufferDataPackagerGenerator.java rename to skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/BufferDataAssist.java index a2ff30dd0..b19006a47 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/model/BufferDataPackagerGenerator.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/BufferDataAssist.java @@ -1,31 +1,23 @@ -package com.ai.cloud.skywalking.reciever.model; +package com.ai.cloud.skywalking.reciever.buffer; +import com.ai.cloud.skywalking.protocol.SerializedFactory; +import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; +import com.ai.cloud.skywalking.protocol.exception.ConvertFailedException; import com.ai.cloud.skywalking.protocol.util.IntegerAssist; -public class BufferDataPackagerGenerator { +import java.util.ArrayList; +import java.util.List; + +public class BufferDataAssist { private static byte[] SPILT = new byte[] {127, 127, 127, 127}; private static byte[] EOF = null; - private BufferDataPackagerGenerator() { + private BufferDataAssist() { //DO Nothing } - public static byte[] generateEOFPackage() { - if (EOF != null) { - return EOF; - } - - EOF = generatePackage("EOF".getBytes()); - return EOF; - } - - public static byte[] pack(byte[] msg) { - return generatePackage(msg); - } - - - private static byte[] generatePackage(byte[] msg) { + public static byte[] appendLengthAndSplit(byte[] msg) { byte[] dataPackage = new byte[msg.length + 8]; // 前四位长度 System.arraycopy(IntegerAssist.intToBytes(msg.length), 0, dataPackage, 0, 4); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java index 2f0745a03..215bb5c56 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThread.java @@ -4,7 +4,6 @@ import com.ai.cloud.skywalking.protocol.BufferFileEOFProtocol; import com.ai.cloud.skywalking.protocol.TransportPackager; import com.ai.cloud.skywalking.protocol.util.AtomicRangeInteger; import com.ai.cloud.skywalking.reciever.conf.Config; -import com.ai.cloud.skywalking.reciever.model.BufferDataPackagerGenerator; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHealthCollector; import com.ai.cloud.skywalking.reciever.selfexamination.ServerHeathReading; import org.apache.logging.log4j.LogManager; @@ -44,7 +43,7 @@ public class DataBufferThread extends Thread { } try { - fileOutputStream.write(BufferDataPackagerGenerator.pack(data[i])); + fileOutputStream.write(BufferDataAssist.appendLengthAndSplit(data[i])); length += data[i].length; data[i] = null; } catch (IOException e) { @@ -69,8 +68,8 @@ public class DataBufferThread extends Thread { private void closeCurrentBufferFile(FileOutputStream fileOutputStream) { try { fileOutputStream.flush(); - fileOutputStream.write(BufferDataPackagerGenerator - .pack(TransportPackager.packSerializableObject(new BufferFileEOFProtocol()))); + fileOutputStream.write(BufferDataAssist + .appendLengthAndSplit(TransportPackager.serialize(new BufferFileEOFProtocol()))); } catch (IOException e) { logger.error("Failed to write msg.", e); } finally { diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java index f40d4dd94..a0a0c8455 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/buffer/DataBufferThreadContainer.java @@ -58,10 +58,9 @@ public class DataBufferThreadContainer { countDownLatch.await(); } - logger.info("Data buffer thread size {} begin to init ", Config.Server. - MAX_DEAL_DATA_THREAD_NUMBER); + logger.info("Data buffer thread size {} begin to init ", Config.Buffer.BUFFER_DEAL_THREAD_NUMBER); - for (int i = 0; i < Config.Server.MAX_DEAL_DATA_THREAD_NUMBER; i++) { + for (int i = 0; i < Config.Buffer.BUFFER_DEAL_THREAD_NUMBER; i++) { DataBufferThread dataBufferThread = new DataBufferThread(i); dataBufferThread.start(); buffers.add(dataBufferThread); diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java index 5ca805b6e..f550a75e5 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/conf/Config.java @@ -6,8 +6,7 @@ public class Config { public static class Server { // 采集服务器的端口 public static int PORT = 34000; - // 最大数据处理线程数量 - public static int MAX_DEAL_DATA_THREAD_NUMBER = 3; + // 异常数据的时间间隔 public static int FAILED_PACKAGE_WATCHING_TIME_WINDOW = 5 * 60; // 时间间隔内最大异常数据次数 @@ -31,6 +30,8 @@ public class Config { public static long BUFFER_FILE_MAX_LENGTH = 30 * 1024 * 1024; + public static int BUFFER_DEAL_THREAD_NUMBER = 0; + } @@ -50,6 +51,9 @@ public class Config { public static class Persistence { + // 最大数据处理线程数量 + public static int MAX_DEAL_DATA_THREAD_NUMBER = 3; + // 切换文件,等待时间 public static long SWITCH_FILE_WAIT_TIME = 5000L; diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java index 958b5b095..f0907c157 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/handler/CollectionServerDataHandler.java @@ -1,6 +1,5 @@ package com.ai.cloud.skywalking.reciever.handler; -import com.ai.cloud.skywalking.protocol.util.IntegerAssist; import com.ai.cloud.skywalking.reciever.buffer.DataBufferThreadContainer; import com.ai.cloud.skywalking.reciever.conf.Config; import com.ai.cloud.skywalking.reciever.util.RedisConnector; @@ -9,7 +8,9 @@ import io.netty.channel.SimpleChannelInboundHandler; import redis.clients.jedis.Jedis; import java.net.InetSocketAddress; -import java.util.Arrays; + +import static com.ai.cloud.skywalking.protocol.util.ByteDataUtil.unpackCheckSum; +import static com.ai.cloud.skywalking.protocol.util.ByteDataUtil.validateCheckSum; public class CollectionServerDataHandler extends SimpleChannelInboundHandler { @@ -27,26 +28,6 @@ public class CollectionServerDataHandler extends SimpleChannelInboundHandler serializables; public BufferFileReader(File bufferFile, int currentOffset) { @@ -58,13 +59,12 @@ public class BufferFileReader { int length = unpackLength(); byte[] dataContext = readByte(length); // 转换对象 - serializables = new ArrayList<>(); - serializables = TransportPackager.unpackSerializableObjects(dataContext); + serializables = deserializableObjects(dataContext); byte[] skip = new byte[4]; bufferInputStream.read(skip); - if (!Arrays.equals(SPILT_BALE_ARRAY, skip)) { - skipToNextBufferBale(); + if (!Arrays.equals(DATA_SPILT, skip)) { + skipToNext(); } MemoryRegister.instance().updateOffSet(bufferFile.getName(), currentOffset); } catch (IOException e) { @@ -122,14 +122,14 @@ public class BufferFileReader { return lengthByte; } - public void skipToNextBufferBale() throws IOException { + public void skipToNext() throws IOException { byte[] previousDataByte = new byte[4]; byte[] currentDataByte = new byte[4]; byte[] compactDataByte = new byte[8]; while (true) { currentDataByte = readByte(4); - if (Arrays.equals(currentDataByte, SPILT_BALE_ARRAY)) { + if (Arrays.equals(currentDataByte, DATA_SPILT)) { remainderLength = 0; break; } @@ -137,7 +137,7 @@ public class BufferFileReader { System.arraycopy(previousDataByte, 0, compactDataByte, 0, 4); System.arraycopy(currentDataByte, 0, compactDataByte, 4, 4); - int index = bytesIndexOf(compactDataByte, SPILT_BALE_ARRAY, 0, 8); + int index = bytesIndexOf(compactDataByte, DATA_SPILT, 0, 8); if (index != -1) { recodeRemainderByteAndLength(compactDataByte, index); break; @@ -185,4 +185,31 @@ public class BufferFileReader { bufferFile = null; } + + public static List deserializableObjects(byte[] dataPackage) { + List serializeData = new ArrayList(); + int currentLength = 0; + while (true) { + //读取长度 + int dataLength = IntegerAssist.bytesToInt(dataPackage, currentLength); + // 反序列化 + byte[] data = new byte[dataLength]; + System.arraycopy(dataPackage, currentLength + 4, data, 0, dataLength); + + try { + AbstractDataSerializable abstractDataSerializable = SerializedFactory.deserialize(data); + serializeData.add(abstractDataSerializable); + } catch (ConvertFailedException e) { + // FIXME: 16/8/4 logger日志输出 + e.printStackTrace(); + } + + currentLength += 4 + dataLength; + if (currentLength >= dataPackage.length) { + break; + } + } + + return serializeData; + } } diff --git a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java index 8f3d2e526..33e7348dd 100644 --- a/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java +++ b/skywalking-server/src/main/java/com/ai/cloud/skywalking/reciever/peresistent/PersistenceThreadLauncher.java @@ -4,9 +4,11 @@ import com.ai.cloud.skywalking.reciever.conf.Config; public class PersistenceThreadLauncher { public static void doLaunch() { - new RegisterPersistenceThread().start(); - for (int i = 0; i < Config.Server.MAX_DEAL_DATA_THREAD_NUMBER; i++) { - new PersistenceThread(i).start(); + if (Config.Persistence.MAX_DEAL_DATA_THREAD_NUMBER > 0) { + new RegisterPersistenceThread().start(); + for (int i = 0; i < Config.Persistence.MAX_DEAL_DATA_THREAD_NUMBER; i++) { + new PersistenceThread(i).start(); + } } } } diff --git a/skywalking-server/src/main/resources/config.properties b/skywalking-server/src/main/resources/config.properties index 76995a6e2..836c0cfe2 100644 --- a/skywalking-server/src/main/resources/config.properties +++ b/skywalking-server/src/main/resources/config.properties @@ -1,9 +1,10 @@ #采集服务器的端口 server.port=34000 -server.max_deal_data_thread_number=1 + server.failed_package_watching_time_windowss=300 server.max_watching_failed_package_size=200 - +# +buffer.buffer_deal_thread_number=1 #每个线程最大缓存数量 buffer.per_thread_max_buffer_number=1024 #无数据处理时轮询等待时间(单位:毫秒) @@ -21,6 +22,8 @@ buffer.write_data_failure_retry_interval = 10000 persistence.switch_file_wait_time=5000 #追加EOF标志位的线程数量 persistence.max_append_eof_flags_thread_number=1 +#持久化线程个数 +persistence.max_deal_data_thread_number=0 #偏移量注册文件的目录 registerpersistence.register_file_parent_directory=/tmp/skywalking/data/offset diff --git a/skywalking-server/src/test/java/com/ai/cloud/skywalking/reciever/model/BufferDataAssistTest.java b/skywalking-server/src/test/java/com/ai/cloud/skywalking/reciever/model/BufferDataAssistTest.java new file mode 100644 index 000000000..c1b7a6ca1 --- /dev/null +++ b/skywalking-server/src/test/java/com/ai/cloud/skywalking/reciever/model/BufferDataAssistTest.java @@ -0,0 +1,51 @@ +package com.ai.cloud.skywalking.reciever.model; + +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.TransportPackager; +import com.ai.cloud.skywalking.protocol.common.AbstractDataSerializable; +import com.ai.cloud.skywalking.protocol.common.ISerializable; +import com.ai.cloud.skywalking.protocol.common.SpanType; +import com.ai.cloud.skywalking.protocol.util.IntegerAssist; +import com.ai.cloud.skywalking.reciever.buffer.BufferDataAssist; +import com.ai.cloud.skywalking.reciever.peresistent.BufferFileReader; +import org.junit.Before; +import org.junit.Test; + +import java.util.ArrayList; +import java.util.Arrays; +import java.util.List; + +import static com.ai.cloud.skywalking.protocol.util.ByteDataUtil.unpackCheckSum; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +public class BufferDataAssistTest { + private List serializableList; + + @Test + public void pack() throws Exception { + byte[] byteData = + BufferDataAssist.appendLengthAndSplit(unpackCheckSum(TransportPackager.pack(serializableList))); + + int length = IntegerAssist.bytesToInt(Arrays.copyOfRange(byteData, 0, 4), 0); + byte[] serializableByteData = Arrays.copyOfRange(byteData, 4, length + 4); + List serializables = BufferFileReader.deserializableObjects(serializableByteData); + assertEquals(2, serializables.size()); + assertNotNull(serializables.get(0)); + + } + + @Before + public void initData() { + serializableList = new ArrayList(); + Span span = new Span("test-traceID", "test", "10"); + span.setStartDate(System.currentTimeMillis() - 1000 * 10); + span.setViewPointId("test-viewpoint"); + span.setSpanType(SpanType.LOCAL); + serializableList.add(new RequestSpan(span)); + serializableList.add(new AckSpan(span)); + } + +}