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));
+ }
+
+}