From d37dbc8d76db60034bf8aa9105d5bd111e046642 Mon Sep 17 00:00:00 2001 From: wusheng Date: Mon, 3 Jul 2017 22:16:20 +0800 Subject: [PATCH] FInish codes about TraceSegment to UpstreamSegment --- .../context/trace/AbstractTracingSpan.java | 36 +++++++++++ .../agent/core/context/trace/EntrySpan.java | 3 + .../agent/core/context/trace/ExitSpan.java | 11 ++++ .../agent/core/context/trace/LocalSpan.java | 6 ++ .../core/context/trace/LogDataEntity.java | 10 +++ .../agent/core/context/trace/SpanLayer.java | 18 ++++-- .../core/context/trace/TraceSegment.java | 62 ++++++++++++------- .../core/context/trace/TraceSegmentRef.java | 27 +++++--- .../agent/core/context/util/KeyValuePair.java | 9 +++ .../remote/TraceSegmentServiceClient.java | 22 ++++--- 10 files changed, 161 insertions(+), 43 deletions(-) diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java index f1c486234..0d951bec0 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/AbstractTracingSpan.java @@ -5,6 +5,8 @@ import java.util.List; import org.skywalking.apm.agent.core.context.util.KeyValuePair; import org.skywalking.apm.agent.core.context.util.ThrowableTransformer; import org.skywalking.apm.agent.core.dictionary.DictionaryUtil; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.SpanType; import org.skywalking.apm.network.trace.component.Component; /** @@ -152,4 +154,38 @@ public abstract class AbstractTracingSpan implements AbstractSpan { this.componentName = componentName; return this; } + + public SpanObject.Builder transform(){ + SpanObject.Builder spanBuilder = SpanObject.newBuilder(); + + spanBuilder.setSpanId(this.spanId); + spanBuilder.setParentSpanId(parentSpanId); + spanBuilder.setStartTime(startTime); + spanBuilder.setEndTime(endTime); + if (operationId == DictionaryUtil.nullValue()) { + spanBuilder.setOperationNameId(operationId); + } else { + spanBuilder.setOperationName(operationName); + } + spanBuilder.setSpanType(SpanType.Entry); + spanBuilder.setSpanLayerValue(this.layer.getCode()); + if (componentId == DictionaryUtil.nullValue()) { + spanBuilder.setComponentId(componentId); + } else { + spanBuilder.setComponent(componentName); + } + spanBuilder.setIsError(errorOccurred); + if (this.tags != null) { + for (KeyValuePair tag : this.tags) { + spanBuilder.addTags(tag.transform()); + } + } + if (this.logs != null) { + for (LogDataEntity log : this.logs) { + spanBuilder.addLogs(log.transform()); + } + } + + return spanBuilder; + } } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/EntrySpan.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/EntrySpan.java index aae12a8a3..f6bc563b9 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/EntrySpan.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/EntrySpan.java @@ -1,6 +1,9 @@ package org.skywalking.apm.agent.core.context.trace; +import org.skywalking.apm.agent.core.context.util.KeyValuePair; import org.skywalking.apm.agent.core.dictionary.DictionaryUtil; +import org.skywalking.apm.network.proto.SpanObject; +import org.skywalking.apm.network.proto.SpanType; import org.skywalking.apm.network.trace.component.Component; /** diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/ExitSpan.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/ExitSpan.java index 640e71842..95fee05a9 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/ExitSpan.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/ExitSpan.java @@ -1,6 +1,7 @@ package org.skywalking.apm.agent.core.context.trace; import org.skywalking.apm.agent.core.dictionary.DictionaryUtil; +import org.skywalking.apm.network.proto.SpanObject; import org.skywalking.apm.network.trace.component.Component; /** @@ -102,6 +103,16 @@ public class ExitSpan extends AbstractTracingSpan { return this; } + @Override public SpanObject.Builder transform() { + SpanObject.Builder spanBuilder = super.transform(); + if (peerId == DictionaryUtil.nullValue()) { + spanBuilder.setPeerId(peerId); + } else { + spanBuilder.setPeer(peer); + } + return spanBuilder; + } + public int getPeerId() { return peerId; } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LocalSpan.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LocalSpan.java index e28edce34..9f76419e1 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LocalSpan.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LocalSpan.java @@ -1,5 +1,7 @@ package org.skywalking.apm.agent.core.context.trace; +import org.skywalking.apm.network.proto.SpanObject; + /** * The LocalSpan represents a normal tracing point, such as a local method. * @@ -27,6 +29,10 @@ public class LocalSpan extends AbstractTracingSpan { return this; } + @Override public SpanObject transform() { + return null; + } + @Override public boolean isEntry() { return false; } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LogDataEntity.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LogDataEntity.java index 1c6156aeb..49aa98e64 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LogDataEntity.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/LogDataEntity.java @@ -3,6 +3,8 @@ package org.skywalking.apm.agent.core.context.trace; import java.util.LinkedList; import java.util.List; import org.skywalking.apm.agent.core.context.util.KeyValuePair; +import org.skywalking.apm.network.proto.KeyWithStringValue; +import org.skywalking.apm.network.proto.LogMessage; /** * The LogDataEntity represents a collection of {@link KeyValuePair}, @@ -39,4 +41,12 @@ public class LogDataEntity { return new LogDataEntity(logs); } } + + public LogMessage transform() { + LogMessage.Builder logMessageBuilder = LogMessage.newBuilder(); + for (KeyValuePair log : logs) { + logMessageBuilder.addData(log.transform()); + } + return logMessageBuilder.build(); + } } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/SpanLayer.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/SpanLayer.java index 25b4a67a0..446fa752f 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/SpanLayer.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/SpanLayer.java @@ -4,10 +4,20 @@ package org.skywalking.apm.agent.core.context.trace; * @author wusheng */ public enum SpanLayer { - DB, - RPC_FRAMEWORK, - HTTP, - MQ; + DB(0), + RPC_FRAMEWORK(1), + HTTP(2), + MQ(3); + + private int code; + + SpanLayer(int code) { + this.code = code; + } + + public int getCode() { + return code; + } public static void asDB(AbstractSpan span) { span.setLayer(SpanLayer.DB); diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegment.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegment.java index b7ca41b99..7e34dffad 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegment.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegment.java @@ -3,6 +3,7 @@ package org.skywalking.apm.agent.core.context.trace; import java.util.LinkedList; import java.util.List; import org.skywalking.apm.agent.core.conf.Config; +import org.skywalking.apm.agent.core.conf.RemoteDownstreamConfig; import org.skywalking.apm.agent.core.context.ids.DistributedTraceId; import org.skywalking.apm.agent.core.context.ids.DistributedTraceIds; import org.skywalking.apm.agent.core.context.ids.GlobalIdGenerator; @@ -11,6 +12,9 @@ import org.skywalking.apm.agent.core.dictionary.DictionaryManager; import org.skywalking.apm.agent.core.dictionary.PossibleFound; import org.skywalking.apm.logging.ILog; import org.skywalking.apm.logging.LogManager; +import org.skywalking.apm.network.proto.TraceSegmentObject; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.skywalking.apm.network.proto.UpstreamSegment; /** * {@link TraceSegment} is a segment or fragment of the distributed trace. @@ -33,16 +37,6 @@ public class TraceSegment { */ private String traceSegmentId; - /** - * The start time of this trace segment. - */ - private long startTime; - - /** - * The end time of this trace segment. - */ - private long endTime; - /** * The refs of parent trace segments, except the primary one. * For most RPC call, {@link #refs} contains only one element, @@ -103,7 +97,6 @@ public class TraceSegment { } } ); - this.startTime = System.currentTimeMillis(); this.traceSegmentId = GlobalIdGenerator.generate(ID_TYPE); this.spans = new LinkedList(); this.relatedGlobalTraces = new DistributedTraceIds(); @@ -154,7 +147,6 @@ public class TraceSegment { * return this, for chaining */ public TraceSegment finish() { - this.endTime = System.currentTimeMillis(); return this; } @@ -162,14 +154,6 @@ public class TraceSegment { return traceSegmentId; } - public long getStartTime() { - return startTime; - } - - public long getEndTime() { - return endTime; - } - public int getApplicationId() { return applicationId; } @@ -198,12 +182,46 @@ public class TraceSegment { this.ignore = ignore; } + /** + * This is a high CPU cost method, only called when sending to collector or test cases. + * + * @return the segment as GRPC service parameter + */ + public UpstreamSegment transform() { + UpstreamSegment.Builder upstreamBuilder = UpstreamSegment.newBuilder(); + for (DistributedTraceId distributedTraceId : getRelatedGlobalTraces()) { + upstreamBuilder = upstreamBuilder.addGlobalTraceIds(distributedTraceId.get()); + } + TraceSegmentObject.Builder traceSegmentBuilder = TraceSegmentObject.newBuilder(); + /** + * Trace Segment + */ + traceSegmentBuilder.setTraceSegmentId(this.traceSegmentId); + // TraceSegmentReference + if (this.refs != null) { + for (TraceSegmentRef ref : this.refs) { + traceSegmentBuilder.addRefs(ref.transform()); + } + } + // globalTraceIds + for (DistributedTraceId distributedTraceId : getRelatedGlobalTraces()) { + traceSegmentBuilder.addGlobalTraceIds(distributedTraceId.get()); + } + // SpanObject + for (AbstractTracingSpan span : this.spans) { + traceSegmentBuilder.addSpans(span.transform()); + } + traceSegmentBuilder.setApplicationId(this.applicationId); + traceSegmentBuilder.setApplicationInstanceId(RemoteDownstreamConfig.Agent.APPLICATION_ID); + + upstreamBuilder.setSegment(traceSegmentBuilder.build().toByteString()); + return upstreamBuilder.build(); + } + @Override public String toString() { return "TraceSegment{" + "traceSegmentId='" + traceSegmentId + '\'' + - ", startTime=" + startTime + - ", endTime=" + endTime + ", refs=" + refs + ", spans=" + spans + ", applicationId='" + applicationId + '\'' + diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegmentRef.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegmentRef.java index 8502111dd..342c9702b 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegmentRef.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/trace/TraceSegmentRef.java @@ -4,6 +4,8 @@ import java.util.List; import org.skywalking.apm.agent.core.context.ContextCarrier; import org.skywalking.apm.agent.core.context.ids.DistributedTraceId; import org.skywalking.apm.agent.core.dictionary.DictionaryUtil; +import org.skywalking.apm.network.proto.TraceSegmentReference; +import org.skywalking.apm.network.proto.UpstreamSegment; /** * {@link TraceSegmentRef} is like a pointer, which ref to another {@link TraceSegment}, @@ -26,11 +28,6 @@ public class TraceSegmentRef { private int operationId = DictionaryUtil.nullValue(); - /** - * {@link DistributedTraceId} - */ - private List distributedTraceIds; - /** * Transform a {@link ContextCarrier} to the TraceSegmentRef * @@ -52,8 +49,6 @@ public class TraceSegmentRef { } else { this.operationId = Integer.parseInt(entryOperationName); } - - this.distributedTraceIds = carrier.getDistributedTraceIds(); } public String getOperationName() { @@ -64,6 +59,24 @@ public class TraceSegmentRef { return operationId; } + public TraceSegmentReference transform() { + TraceSegmentReference.Builder refBuilder = TraceSegmentReference.newBuilder(); + refBuilder.setParentTraceSegmentId(traceSegmentId); + refBuilder.setParentSpanId(spanId); + refBuilder.setParentApplicationId(applicationId); + if (peerId == DictionaryUtil.nullValue()) { + refBuilder.setNetworkAddress(peerHost); + } else { + refBuilder.setNetworkAddressId(peerId); + } + if (operationId == DictionaryUtil.nullValue()) { + refBuilder.setEntryServiceName(operationName); + } else { + refBuilder.setEntryServiceId(operationId); + } + return refBuilder.build(); + } + @Override public boolean equals(Object o) { if (this == o) diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/util/KeyValuePair.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/util/KeyValuePair.java index 6ae5ebcae..76f4227df 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/util/KeyValuePair.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/context/util/KeyValuePair.java @@ -1,5 +1,7 @@ package org.skywalking.apm.agent.core.context.util; +import org.skywalking.apm.network.proto.KeyWithStringValue; + /** * The KeyValuePair represents a object which contains a string key and a string value. * @@ -21,4 +23,11 @@ public class KeyValuePair { public String getValue() { return value; } + + public KeyWithStringValue transform() { + KeyWithStringValue.Builder keyValueBuilder = KeyWithStringValue.newBuilder(); + keyValueBuilder.setKey(key); + keyValueBuilder.setValue(value); + return keyValueBuilder.build(); + } } diff --git a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java index afb43c351..f40611e47 100644 --- a/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java +++ b/apm-sniffer/apm-agent-core/src/main/java/org/skywalking/apm/agent/core/remote/TraceSegmentServiceClient.java @@ -13,9 +13,9 @@ import org.skywalking.apm.agent.core.datacarrier.buffer.BufferStrategy; import org.skywalking.apm.agent.core.datacarrier.consumer.IConsumer; import org.skywalking.apm.logging.ILog; import org.skywalking.apm.logging.LogManager; -import org.skywalking.apm.network.collecor.proto.Downstream; -import org.skywalking.apm.network.trace.proto.TraceSegmentServiceGrpc; -import org.skywalking.apm.network.trace.proto.UpstreamSegment; +import org.skywalking.apm.network.proto.Downstream; +import org.skywalking.apm.network.proto.TraceSegmentServiceGrpc; +import org.skywalking.apm.network.proto.UpstreamSegment; import static org.skywalking.apm.agent.core.conf.Config.Buffer.BUFFER_SIZE; import static org.skywalking.apm.agent.core.conf.Config.Buffer.CHANNEL_SIZE; @@ -78,14 +78,13 @@ public class TraceSegmentServiceClient implements BootService, IConsumer