From 6284fcbf51408041c151b1f9abb2d6d26b71cbcd Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Tue, 6 Jun 2017 22:50:17 +0800 Subject: [PATCH] temporary storage --- .../analysis/GlobalTraceAnalysis.java | 7 +- .../worker/httpserver/AbstractPost.java | 67 ++++++---- .../collector/worker/segment/SegmentPost.java | 8 +- .../segment/analysis/SegmentAnalysis.java | 5 +- .../segment/entity/DeserializeObject.java | 16 --- .../worker/segment/entity/GlobalTraceId.java | 49 ++++++-- .../worker/segment/entity/LogData.java | 47 +------ .../worker/segment/entity/Segment.java | 104 +++------------ .../worker/segment/entity/SegmentAndJson.java | 23 ++++ .../segment/entity/SegmentDeserialize.java | 11 +- .../collector/worker/segment/entity/Span.java | 119 +++--------------- .../segment/entity/TraceSegmentRef.java | 50 +------- .../segment/persistence/SegmentSave.java | 29 ++--- .../SegmentTopSearchWithGlobalTraceId.java | 5 +- .../SegmentTopSearchWithTimeSlice.java | 5 +- .../span/persistence/SpanSearchWithId.java | 1 - .../worker/tools/JsonFileReader.java | 13 +- .../httpserver/AbstractPostTestCase.java | 37 +++++- .../segment/entity/LogDataTestCase.java | 30 ----- .../entity/TraceSegmentRefTestCase.java | 26 ---- .../worker/segment/mock/SegmentMock.java | 4 + .../segment/post/normal/cache-service.json | 2 +- 22 files changed, 244 insertions(+), 414 deletions(-) delete mode 100644 apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/DeserializeObject.java create mode 100644 apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndJson.java delete mode 100644 apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/LogDataTestCase.java delete mode 100644 apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRefTestCase.java diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/globaltrace/analysis/GlobalTraceAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/globaltrace/analysis/GlobalTraceAnalysis.java index df53c9540..54e59963c 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/globaltrace/analysis/GlobalTraceAnalysis.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/globaltrace/analysis/GlobalTraceAnalysis.java @@ -33,11 +33,12 @@ public class GlobalTraceAnalysis extends JoinAndSplitAnalysisMember { SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice) message; Segment segment = segmentWithTimeSlice.getSegment(); String subSegmentId = segment.getTraceSegmentId(); - List globalTraceIdList = segment.getRelatedGlobalTraces(); + List globalTraceIdList = null; +// List globalTraceIdList = segment.getRelatedGlobalTraces(); if (CollectionTools.isNotEmpty(globalTraceIdList)) { for (GlobalTraceId disTraceId : globalTraceIdList) { - String traceId = disTraceId.get(); - set(traceId, GlobalTraceIndex.SUB_SEG_IDS, subSegmentId); +// String traceId = disTraceId.get(); +// set(traceId, GlobalTraceIndex.SUB_SEG_IDS, subSegmentId); } } } else { diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractPost.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractPost.java index 366c3a253..c48e520f7 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractPost.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractPost.java @@ -1,17 +1,21 @@ package org.skywalking.apm.collector.worker.httpserver; +import com.google.gson.Gson; import com.google.gson.JsonObject; -import com.google.gson.stream.JsonReader; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; -import org.skywalking.apm.collector.actor.*; -import org.skywalking.apm.collector.worker.segment.entity.Segment; - +import java.io.BufferedReader; +import java.io.IOException; import javax.servlet.ServletException; import javax.servlet.http.HttpServletRequest; import javax.servlet.http.HttpServletResponse; -import java.io.BufferedReader; -import java.io.IOException; +import org.apache.logging.log4j.LogManager; +import org.apache.logging.log4j.Logger; +import org.skywalking.apm.collector.actor.AbstractLocalAsyncWorker; +import org.skywalking.apm.collector.actor.ClusterWorkerContext; +import org.skywalking.apm.collector.actor.LocalAsyncWorkerRef; +import org.skywalking.apm.collector.actor.LocalWorkerContext; +import org.skywalking.apm.collector.actor.Role; +import org.skywalking.apm.collector.worker.segment.entity.Segment; +import org.skywalking.apm.collector.worker.segment.entity.SegmentAndJson; /** * @author pengys5 @@ -23,8 +27,7 @@ public abstract class AbstractPost extends AbstractLocalAsyncWorker { super(role, clusterContext, selfContext); } - @Override - final public void onWork(Object message) throws Exception { + @Override final public void onWork(Object message) throws Exception { onReceive(message); } @@ -34,15 +37,16 @@ public abstract class AbstractPost extends AbstractLocalAsyncWorker { private Logger logger = LogManager.getFormatterLogger(PostWithHttpServlet.class); + private final Gson gson = new Gson(); + private final LocalAsyncWorkerRef ownerWorkerRef; PostWithHttpServlet(LocalAsyncWorkerRef ownerWorkerRef) { this.ownerWorkerRef = ownerWorkerRef; } - @Override - final protected void doPost(HttpServletRequest request, - HttpServletResponse response) throws ServletException, IOException { + @Override final protected void doPost(HttpServletRequest request, + HttpServletResponse response) throws ServletException, IOException { JsonObject resJson = new JsonObject(); try { BufferedReader bufferedReader = request.getReader(); @@ -56,19 +60,34 @@ public abstract class AbstractPost extends AbstractLocalAsyncWorker { } private void streamReader(BufferedReader bufferedReader) throws Exception { - try (JsonReader reader = new JsonReader(bufferedReader)) { - readSegmentArray(reader); - } - } + Segment segment; + do { + int character; + StringBuilder builder = new StringBuilder(); + while ((character = bufferedReader.read()) != ' ') { + if (character == -1) { + return; + } + builder.append((char)character); + } - private void readSegmentArray(JsonReader reader) throws Exception { - reader.beginArray(); - while (reader.hasNext()) { - Segment segment = new Segment(); - segment.deserialize(reader); - ownerWorkerRef.tell(segment); + int length = Integer.valueOf(builder.toString()); + builder = new StringBuilder(); + + char[] buffer = new char[length]; + int readLength = bufferedReader.read(buffer, 0, length); + if (readLength != length) { + logger.error("The actual data length was different from the length in data head! "); + return; + } + builder.append(buffer); + + String segmentJsonStr = builder.toString(); + segment = gson.fromJson(segmentJsonStr, Segment.class); + + ownerWorkerRef.tell(new SegmentAndJson(segment, segmentJsonStr)); } - reader.endArray(); + while (segment != null); } } } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentPost.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentPost.java index 70e6fdda8..a066d9122 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentPost.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentPost.java @@ -2,6 +2,7 @@ package org.skywalking.apm.collector.worker.segment; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; +import org.skywalking.apm.collector.worker.segment.entity.SegmentAndJson; import org.skywalking.apm.util.StringUtil; import org.skywalking.apm.collector.actor.ClusterWorkerContext; import org.skywalking.apm.collector.actor.LocalWorkerContext; @@ -58,8 +59,9 @@ public class SegmentPost extends AbstractPost { @Override protected void onReceive(Object message) throws Exception { - if (message instanceof Segment) { - Segment segment = (Segment) message; + if (message instanceof SegmentAndJson) { + SegmentAndJson segmentAndJson = (SegmentAndJson) message; + Segment segment = segmentAndJson.getSegment(); try { validateData(segment); } catch (IllegalArgumentException e) { @@ -75,7 +77,7 @@ public class SegmentPost extends AbstractPost { logger.debug("minuteSlice: %s, hourSlice: %s, daySlice: %s, second:%s", minuteSlice, hourSlice, daySlice, second); SegmentWithTimeSlice segmentWithTimeSlice = new SegmentWithTimeSlice(segment, minuteSlice, hourSlice, daySlice, second); - getSelfContext().lookup(SegmentAnalysis.Role.INSTANCE).tell(segment); + getSelfContext().lookup(SegmentAnalysis.Role.INSTANCE).tell(segmentAndJson); getSelfContext().lookup(SegmentCostAnalysis.Role.INSTANCE).tell(segmentWithTimeSlice); getSelfContext().lookup(GlobalTraceAnalysis.Role.INSTANCE).tell(segmentWithTimeSlice); diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentAnalysis.java index f4c1af31d..1e41c21f2 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentAnalysis.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentAnalysis.java @@ -8,6 +8,7 @@ import org.skywalking.apm.collector.actor.selector.WorkerSelector; import org.skywalking.apm.collector.worker.RecordAnalysisMember; import org.skywalking.apm.collector.worker.config.WorkerConfig; import org.skywalking.apm.collector.worker.segment.entity.Segment; +import org.skywalking.apm.collector.worker.segment.entity.SegmentAndJson; import org.skywalking.apm.collector.worker.segment.persistence.SegmentSave; /** @@ -29,8 +30,8 @@ public class SegmentAnalysis extends RecordAnalysisMember { @Override public void analyse(Object message) throws Exception { if (message instanceof Segment) { - Segment segment = (Segment) message; - getSelfContext().lookup(SegmentSave.Role.INSTANCE).tell(segment); + SegmentAndJson segmentAndJson = (SegmentAndJson) message; + getSelfContext().lookup(SegmentSave.Role.INSTANCE).tell(segmentAndJson); } else { logger.error("unhandled message, message instance must Segment, but is %s", message.getClass().toString()); } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/DeserializeObject.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/DeserializeObject.java deleted file mode 100644 index d988d487f..000000000 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/DeserializeObject.java +++ /dev/null @@ -1,16 +0,0 @@ -package org.skywalking.apm.collector.worker.segment.entity; - -/** - * @author pengys5 - */ -public abstract class DeserializeObject { - private String jsonStr; - - public String getJsonStr() { - return jsonStr; - } - - public void setJsonStr(String jsonStr) { - this.jsonStr = jsonStr; - } -} diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/GlobalTraceId.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/GlobalTraceId.java index f2f15d71c..f3d248614 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/GlobalTraceId.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/GlobalTraceId.java @@ -1,22 +1,53 @@ package org.skywalking.apm.collector.worker.segment.entity; +import com.google.gson.TypeAdapter; +import com.google.gson.annotations.JsonAdapter; import com.google.gson.stream.JsonReader; - +import com.google.gson.stream.JsonWriter; import java.io.IOException; +import java.util.LinkedList; +import java.util.List; /** * @author pengys5 */ -public class GlobalTraceId extends DeserializeObject { - private String globalTraceId; +@JsonAdapter(GlobalTraceId.Serializer.class) +public class GlobalTraceId { - public String get() { - return globalTraceId; + public GlobalTraceId() { + globalTraceIds = new LinkedList<>(); } - public GlobalTraceId deserialize(JsonReader reader) throws IOException { - this.globalTraceId = reader.nextString(); - this.setJsonStr("\"" + globalTraceId + "\""); - return this; + private LinkedList globalTraceIds; + + public LinkedList get() { + return globalTraceIds; + } + + public static class Serializer extends TypeAdapter { + @Override public void write(JsonWriter out, GlobalTraceId value) throws IOException { + List globalTraceIds = value.globalTraceIds; + + if (globalTraceIds.size() > 0) { + out.beginArray(); + for (String globalTraceId : globalTraceIds) { + out.value(globalTraceId); + } + out.endArray(); + } + } + + @Override public GlobalTraceId read(JsonReader in) throws IOException { + GlobalTraceId globalTraceId = new GlobalTraceId(); + in.beginArray(); + try { + while (in.hasNext()) { + globalTraceId.get().add(in.nextString()); + } + } finally { + in.endArray(); + } + return globalTraceId; + } } } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/LogData.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/LogData.java index 8da6c7ca8..c97f31062 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/LogData.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/LogData.java @@ -1,16 +1,17 @@ package org.skywalking.apm.collector.worker.segment.entity; -import com.google.gson.stream.JsonReader; - -import java.io.IOException; -import java.util.HashMap; +import com.google.gson.annotations.SerializedName; import java.util.Map; /** * @author pengys5 */ -public class LogData extends DeserializeObject { +public class LogData { + + @SerializedName("tm") private long time; + + @SerializedName("fi") private Map fields; public long getTime() { @@ -20,40 +21,4 @@ public class LogData extends DeserializeObject { public Map getFields() { return fields; } - - public LogData deserialize(JsonReader reader) throws IOException { - StringBuilder stringBuilder = new StringBuilder(); - stringBuilder.append("{"); - - boolean first = true; - reader.beginObject(); - while (reader.hasNext()) { - switch (reader.nextName()) { - case "tm": - Long tm = reader.nextLong(); - this.time = tm; - JsonBuilder.INSTANCE.append(stringBuilder, "tm", tm, first); - break; - case "fi": - fields = new HashMap<>(); - reader.beginObject(); - - while (reader.hasNext()) { - String key = reader.nextName(); - String value = reader.nextString(); - fields.put(key, value); - } - reader.endObject(); - JsonBuilder.INSTANCE.append(stringBuilder, "fi", fields, first); - break; - default: - reader.skipValue(); - } - first = false; - } - reader.endObject(); - stringBuilder.append("}"); - this.setJsonStr(stringBuilder.toString()); - return this; - } } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Segment.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Segment.java index 7e8d2c3e3..13c5bb352 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Segment.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Segment.java @@ -1,22 +1,33 @@ package org.skywalking.apm.collector.worker.segment.entity; -import com.google.gson.stream.JsonReader; - -import java.io.IOException; -import java.util.ArrayList; +import com.google.gson.annotations.SerializedName; import java.util.List; /** * @author pengys5 */ -public class Segment extends DeserializeObject { +public class Segment { + + @SerializedName("ts") private String traceSegmentId; + + @SerializedName("st") private long startTime; + + @SerializedName("et") private long endTime; + + @SerializedName("rs") private List refs; + + @SerializedName("ss") private List spans; + + @SerializedName("ac") private String applicationCode; - private List relatedGlobalTraces; + + @SerializedName("gt") + private GlobalTraceId relatedGlobalTraces; public String getTraceSegmentId() { return traceSegmentId; @@ -42,86 +53,7 @@ public class Segment extends DeserializeObject { return spans; } - public List getRelatedGlobalTraces() { + public GlobalTraceId getRelatedGlobalTraces() { return relatedGlobalTraces; } - - public Segment deserialize(JsonReader reader) throws IOException { - StringBuilder stringBuilder = new StringBuilder(); - stringBuilder.append("{"); - - boolean first = true; - reader.beginObject(); - while (reader.hasNext()) { - switch (reader.nextName()) { - case "ts": - String ts = reader.nextString(); - this.traceSegmentId = ts; - JsonBuilder.INSTANCE.append(stringBuilder, "ts", ts, first); - break; - case "ac": - String ac = reader.nextString(); - this.applicationCode = ac; - JsonBuilder.INSTANCE.append(stringBuilder, "ac", ac, first); - break; - case "st": - long st = reader.nextLong(); - this.startTime = st; - JsonBuilder.INSTANCE.append(stringBuilder, "st", st, first); - break; - case "et": - long et = reader.nextLong(); - this.endTime = et; - JsonBuilder.INSTANCE.append(stringBuilder, "et", et, first); - break; - case "rs": - refs = new ArrayList<>(); - reader.beginArray(); - - while (reader.hasNext()) { - TraceSegmentRef ref = new TraceSegmentRef(); - ref.deserialize(reader); - refs.add(ref); - } - - reader.endArray(); - JsonBuilder.INSTANCE.append(stringBuilder, "rs", refs, first); - break; - case "ss": - spans = new ArrayList<>(); - reader.beginArray(); - - while (reader.hasNext()) { - Span span = new Span(); - span.deserialize(reader); - spans.add(span); - } - - reader.endArray(); - JsonBuilder.INSTANCE.append(stringBuilder, "ss", spans, first); - break; - case "gt": - relatedGlobalTraces = new ArrayList<>(); - reader.beginArray(); - - while (reader.hasNext()) { - GlobalTraceId globalTraceId = new GlobalTraceId(); - globalTraceId.deserialize(reader); - relatedGlobalTraces.add(globalTraceId); - } - JsonBuilder.INSTANCE.append(stringBuilder, "gt", relatedGlobalTraces, first); - - reader.endArray(); - break; - default: - reader.skipValue(); - } - first = false; - } - reader.endObject(); - - stringBuilder.append("}"); - this.setJsonStr(stringBuilder.toString()); - return this; - } } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndJson.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndJson.java new file mode 100644 index 000000000..0a8e08f56 --- /dev/null +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndJson.java @@ -0,0 +1,23 @@ +package org.skywalking.apm.collector.worker.segment.entity; + +/** + * @author pengys5 + */ +public class SegmentAndJson { + + private final Segment segment; + private final String jsonStr; + + public SegmentAndJson(Segment segment, String jsonStr) { + this.segment = segment; + this.jsonStr = jsonStr; + } + + public Segment getSegment() { + return segment; + } + + public String getJsonStr() { + return jsonStr; + } +} diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentDeserialize.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentDeserialize.java index 385456e66..981500f9c 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentDeserialize.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentDeserialize.java @@ -1,10 +1,9 @@ package org.skywalking.apm.collector.worker.segment.entity; +import com.google.gson.Gson; import com.google.gson.stream.JsonReader; - import java.io.FileReader; import java.io.IOException; -import java.io.StringReader; import java.util.ArrayList; import java.util.List; @@ -14,10 +13,10 @@ import java.util.List; public enum SegmentDeserialize { INSTANCE; + private final Gson gson = new Gson(); + public Segment deserializeSingle(String singleSegmentJsonStr) throws IOException { - JsonReader reader = new JsonReader(new StringReader(singleSegmentJsonStr)); - Segment segment = new Segment(); - segment.deserialize(reader); + Segment segment = gson.fromJson(singleSegmentJsonStr, Segment.class); return segment; } @@ -37,7 +36,7 @@ public enum SegmentDeserialize { reader.beginArray(); while (reader.hasNext()) { Segment segment = new Segment(); - segment.deserialize(reader); +// segment.deserialize(reader); segmentList.add(segment); } reader.endArray(); diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Span.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Span.java index edf184bed..716ee2fdf 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Span.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Span.java @@ -1,25 +1,39 @@ package org.skywalking.apm.collector.worker.segment.entity; -import com.google.gson.stream.JsonReader; - -import java.io.IOException; -import java.util.ArrayList; -import java.util.HashMap; +import com.google.gson.annotations.SerializedName; import java.util.List; import java.util.Map; /** * @author pengys5 */ -public class Span extends DeserializeObject { +public class Span { + + @SerializedName("si") private int spanId; + + @SerializedName("ps") private int parentSpanId; + + @SerializedName("st") private long startTime; + + @SerializedName("et") private long endTime; + + @SerializedName("on") private String operationName; + + @SerializedName("ts") private Map tagsWithStr; + + @SerializedName("tb") private Map tagsWithBool; + + @SerializedName("ti") private Map tagsWithInt; + + @SerializedName("lo") private List logs; public int getSpanId() { @@ -57,97 +71,4 @@ public class Span extends DeserializeObject { public List getLogs() { return logs; } - - public Span deserialize(JsonReader reader) throws IOException { - StringBuilder stringBuilder = new StringBuilder(); - stringBuilder.append("{"); - - boolean first = true; - reader.beginObject(); - while (reader.hasNext()) { - switch (reader.nextName()) { - case "si": - Integer si = reader.nextInt(); - this.spanId = si; - JsonBuilder.INSTANCE.append(stringBuilder, "si", si, first); - break; - case "ps": - Integer ps = reader.nextInt(); - this.parentSpanId = ps; - JsonBuilder.INSTANCE.append(stringBuilder, "ps", ps, first); - break; - case "st": - Long st = reader.nextLong(); - this.startTime = st; - JsonBuilder.INSTANCE.append(stringBuilder, "st", st, first); - break; - case "et": - Long et = reader.nextLong(); - this.endTime = et; - JsonBuilder.INSTANCE.append(stringBuilder, "et", et, first); - break; - case "on": - String on = reader.nextString(); - this.operationName = on; - JsonBuilder.INSTANCE.append(stringBuilder, "on", on, first); - break; - case "ts": - tagsWithStr = new HashMap<>(); - reader.beginObject(); - - while (reader.hasNext()) { - String key = reader.nextName(); - String value = reader.nextString(); - tagsWithStr.put(key, value); - } - reader.endObject(); - JsonBuilder.INSTANCE.append(stringBuilder, "ts", tagsWithStr, first); - break; - case "tb": - tagsWithBool = new HashMap<>(); - reader.beginObject(); - - while (reader.hasNext()) { - String key = reader.nextName(); - boolean value = reader.nextBoolean(); - tagsWithBool.put(key, value); - } - reader.endObject(); - JsonBuilder.INSTANCE.append(stringBuilder, "tb", tagsWithBool, first); - break; - case "ti": - tagsWithInt = new HashMap<>(); - reader.beginObject(); - - while (reader.hasNext()) { - String key = reader.nextName(); - Integer value = reader.nextInt(); - tagsWithInt.put(key, value); - } - reader.endObject(); - JsonBuilder.INSTANCE.append(stringBuilder, "ti", tagsWithInt, first); - break; - case "lo": - logs = new ArrayList<>(); - reader.beginArray(); - - while (reader.hasNext()) { - LogData logData = new LogData(); - logData.deserialize(reader); - logs.add(logData); - } - reader.endArray(); - JsonBuilder.INSTANCE.append(stringBuilder, "lo", logs, first); - break; - default: - reader.skipValue(); - } - first = false; - } - reader.endObject(); - - stringBuilder.append("}"); - this.setJsonStr(stringBuilder.toString()); - return this; - } } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRef.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRef.java index ee48dacfd..3629728e0 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRef.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRef.java @@ -1,20 +1,22 @@ package org.skywalking.apm.collector.worker.segment.entity; -import com.google.gson.stream.JsonReader; - -import java.io.IOException; +import com.google.gson.annotations.SerializedName; /** * @author pengys5 */ -public class TraceSegmentRef extends DeserializeObject { +public class TraceSegmentRef { + @SerializedName("ts") private String traceSegmentId; + @SerializedName("si") private int spanId = -1; + @SerializedName("ac") private String applicationCode; + @SerializedName("ph") private String peerHost; public String getTraceSegmentId() { @@ -32,44 +34,4 @@ public class TraceSegmentRef extends DeserializeObject { public String getPeerHost() { return peerHost; } - - public TraceSegmentRef deserialize(JsonReader reader) throws IOException { - StringBuilder stringBuilder = new StringBuilder(); - stringBuilder.append("{"); - - boolean first = true; - reader.beginObject(); - while (reader.hasNext()) { - switch (reader.nextName()) { - case "ts": - String ts = reader.nextString(); - this.traceSegmentId = ts; - JsonBuilder.INSTANCE.append(stringBuilder, "ts", ts, first); - break; - case "si": - Integer si = reader.nextInt(); - this.spanId = si; - JsonBuilder.INSTANCE.append(stringBuilder, "si", si, first); - break; - case "ac": - String ac = reader.nextString(); - this.applicationCode = ac; - JsonBuilder.INSTANCE.append(stringBuilder, "ac", ac, first); - break; - case "ph": - String ph = reader.nextString(); - this.peerHost = ph; - JsonBuilder.INSTANCE.append(stringBuilder, "ph", ph, first); - break; - default: - reader.skipValue(); - } - first = false; - } - reader.endObject(); - - stringBuilder.append("}"); - this.setJsonStr(stringBuilder.toString()); - return this; - } } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentSave.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentSave.java index 6c4cea9bb..6bd28cd1a 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentSave.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentSave.java @@ -1,5 +1,8 @@ package org.skywalking.apm.collector.worker.segment.persistence; +import java.util.LinkedList; +import java.util.List; +import java.util.Map; import org.elasticsearch.action.index.IndexRequestBuilder; import org.elasticsearch.client.Client; import org.skywalking.apm.collector.actor.AbstractLocalSyncWorkerProvider; @@ -10,12 +13,12 @@ import org.skywalking.apm.collector.actor.selector.WorkerSelector; import org.skywalking.apm.collector.worker.PersistenceMember; import org.skywalking.apm.collector.worker.config.CacheSizeConfig; import org.skywalking.apm.collector.worker.segment.SegmentIndex; -import org.skywalking.apm.collector.worker.segment.entity.Segment; -import org.skywalking.apm.collector.worker.storage.*; - -import java.util.LinkedList; -import java.util.List; -import java.util.Map; +import org.skywalking.apm.collector.worker.segment.entity.SegmentAndJson; +import org.skywalking.apm.collector.worker.storage.AbstractIndex; +import org.skywalking.apm.collector.worker.storage.EsClient; +import org.skywalking.apm.collector.worker.storage.PersistenceWorkerListener; +import org.skywalking.apm.collector.worker.storage.SegmentData; +import org.skywalking.apm.collector.worker.storage.SegmentPersistenceData; /** * @author pengys5 @@ -33,7 +36,7 @@ public class SegmentSave extends PersistenceMember= CacheSizeConfig.Cache.Persistence.SIZE) { persistence(data.asMap()); } @@ -69,8 +71,7 @@ public class SegmentSave extends PersistenceMember builderList) { + @Override final protected void prepareIndex(List builderList) { Map lastData = getPersistenceData().getLast().asMap(); Client client = EsClient.INSTANCE.getClient(); diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithGlobalTraceId.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithGlobalTraceId.java index 3cfc50d6a..5687a51eb 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithGlobalTraceId.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithGlobalTraceId.java @@ -82,12 +82,13 @@ public class SegmentTopSearchWithGlobalTraceId extends AbstractLocalSyncWorker { String segmentSource = client.prepareGet(SegmentIndex.INDEX, SegmentIndex.TYPE_RECORD, segId).get().getSourceAsString(); Segment segment = SegmentDeserialize.INSTANCE.deserializeSingle(segmentSource); - List distributedTraceIdList = segment.getRelatedGlobalTraces(); + List distributedTraceIdList = null; +// List distributedTraceIdList = segment.getRelatedGlobalTraces(); JsonArray distributedTraceIdArray = new JsonArray(); if (CollectionTools.isNotEmpty(distributedTraceIdList)) { for (GlobalTraceId distributedTraceId : distributedTraceIdList) { - distributedTraceIdArray.add(distributedTraceId.get()); +// distributedTraceIdArray.add(distributedTraceId.get()); } } topSegmentJson.add("traceIds", distributedTraceIdArray); diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithTimeSlice.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithTimeSlice.java index fc6ebfcc3..935825f48 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithTimeSlice.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearchWithTimeSlice.java @@ -90,12 +90,13 @@ public class SegmentTopSearchWithTimeSlice extends AbstractLocalSyncWorker { String segmentSource = EsClient.INSTANCE.getClient().prepareGet(SegmentIndex.INDEX, SegmentIndex.TYPE_RECORD, segId).get().getSourceAsString(); logger().debug("segmentSource:" + segmentSource); Segment segment = SegmentDeserialize.INSTANCE.deserializeSingle(segmentSource); - List distributedTraceIdList = segment.getRelatedGlobalTraces(); +// List distributedTraceIdList = segment.getRelatedGlobalTraces(); + List distributedTraceIdList = null; JsonArray distributedTraceIdArray = new JsonArray(); if (CollectionTools.isNotEmpty(distributedTraceIdList)) { for (GlobalTraceId distributedTraceId : distributedTraceIdList) { - distributedTraceIdArray.add(distributedTraceId.get()); +// distributedTraceIdArray.add(distributedTraceId.get()); } } topSegmentJson.add("traceIds", distributedTraceIdArray); diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/persistence/SpanSearchWithId.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/persistence/SpanSearchWithId.java index ebb96e8a4..4813e6ddd 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/persistence/SpanSearchWithId.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/persistence/SpanSearchWithId.java @@ -39,7 +39,6 @@ public class SpanSearchWithId extends AbstractLocalSyncWorker { for (Span span : spanList) { if (String.valueOf(span.getSpanId()).equals(search.spanId)) { - span.setJsonStr(""); String spanJsonStr = gson.toJson(span); dataJson = gson.fromJson(spanJsonStr, JsonObject.class); } diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/JsonFileReader.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/JsonFileReader.java index 8144819fc..70c2c50f6 100644 --- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/JsonFileReader.java +++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/JsonFileReader.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.worker.tools; +import com.google.gson.JsonArray; import com.google.gson.JsonElement; import com.google.gson.JsonParser; - import java.io.FileNotFoundException; import java.io.FileReader; @@ -15,6 +15,15 @@ public enum JsonFileReader { public String read(String path) throws FileNotFoundException { JsonParser jsonParser = new JsonParser(); JsonElement jsonElement = jsonParser.parse(new FileReader(path)); - return jsonElement.toString(); + + StringBuilder segmentBuilder = new StringBuilder(); + JsonArray segments = jsonElement.getAsJsonArray(); + for (int i = 0; i < segments.size(); i++) { + JsonElement segment = segments.get(i); + String segmentStr = segment.toString(); + segmentBuilder.append(segmentStr.length()).append(" ").append(segmentStr); + } + + return segmentBuilder.toString(); } } diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/AbstractPostTestCase.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/AbstractPostTestCase.java index 27b24b55e..c4306509d 100644 --- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/AbstractPostTestCase.java +++ b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/AbstractPostTestCase.java @@ -1,6 +1,13 @@ package org.skywalking.apm.collector.worker.httpserver; import com.google.gson.JsonObject; +import java.io.BufferedReader; +import java.io.FileReader; +import java.io.PrintWriter; +import java.io.StringReader; +import java.io.Writer; +import javax.servlet.http.HttpServletRequest; +import javax.servlet.http.HttpServletResponse; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; @@ -10,16 +17,21 @@ import org.powermock.core.classloader.annotations.PrepareForTest; import org.powermock.modules.junit4.PowerMockRunner; import org.skywalking.apm.collector.actor.ClusterWorkerContext; import org.skywalking.apm.collector.actor.LocalWorkerContext; +import org.skywalking.apm.collector.worker.segment.mock.SegmentMock; import static org.mockito.Matchers.anyString; -import static org.mockito.Mockito.*; +import static org.mockito.Mockito.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.spy; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; /** * @author pengys5 */ @RunWith(PowerMockRunner.class) -@PrepareForTest( {TestAbstractPost.class}) -@PowerMockIgnore( {"javax.management.*"}) +@PrepareForTest({TestAbstractPost.class}) +@PowerMockIgnore({"javax.management.*"}) public class AbstractPostTestCase { private TestAbstractPost post; @@ -43,4 +55,23 @@ public class AbstractPostTestCase { post.onWork(new JsonObject()); PowerMockito.verifyPrivate(post).invoke("saveException", any(IllegalArgumentException.class)); } + + @Test + public void testPostWithHttpServlet() throws Exception { + SegmentMock segmentMock = new SegmentMock(); + +// BufferedReader reader = new BufferedReader(new StringReader(segmentMock.mockCacheServiceExceptionSegmentAsString())); + BufferedReader reader = new BufferedReader(new StringReader(segmentMock.mockCacheServiceSegmentAsString())); + + HttpServletRequest request = mock(HttpServletRequest.class); + when(request.getReader()).thenReturn(reader); + + Writer writer = mock(Writer.class); + PrintWriter printWriter = new PrintWriter(writer); + HttpServletResponse response = mock(HttpServletResponse.class); + when(response.getWriter()).thenReturn(printWriter); + + AbstractPost.PostWithHttpServlet servlet = new AbstractPost.PostWithHttpServlet(null); + servlet.doPost(request, response); + } } diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/LogDataTestCase.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/LogDataTestCase.java deleted file mode 100644 index f6202ba9b..000000000 --- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/LogDataTestCase.java +++ /dev/null @@ -1,30 +0,0 @@ -package org.skywalking.apm.collector.worker.segment.entity; - -import com.google.gson.stream.JsonReader; -import org.junit.Assert; -import org.junit.Test; - -import java.io.IOException; -import java.io.StringReader; -import java.util.Map; - -/** - * @author pengys5 - */ -public class LogDataTestCase { - - @Test - public void deserialize() throws IOException { - LogData logData = new LogData(); - - JsonReader reader = new JsonReader(new StringReader("{\"tm\":1, \"fi\": {\"test1\":\"test1\",\"test2\":\"test2\"}, \"skip\":\"skip\"}")); - logData.deserialize(reader); - - Assert.assertEquals(1L, logData.getTime()); - - Map fields = logData.getFields(); - Assert.assertEquals("test1", fields.get("test1")); - Assert.assertEquals("test2", fields.get("test2")); - Assert.assertEquals(false, fields.containsKey("skip")); - } -} diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRefTestCase.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRefTestCase.java deleted file mode 100644 index e7a15962d..000000000 --- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRefTestCase.java +++ /dev/null @@ -1,26 +0,0 @@ -package org.skywalking.apm.collector.worker.segment.entity; - -import com.google.gson.stream.JsonReader; -import org.junit.Assert; -import org.junit.Test; - -import java.io.IOException; -import java.io.StringReader; - -/** - * @author pengys5 - */ -public class TraceSegmentRefTestCase { - - @Test - public void deserialize() throws IOException { - TraceSegmentRef traceSegmentRef = new TraceSegmentRef(); - JsonReader reader = new JsonReader(new StringReader("{\"ts\" :\"ts\",\"si\":0,\"ac\":\"ac\",\"ph\":\"ph\", \"skip\":\"skip\"}")); - traceSegmentRef.deserialize(reader); - - Assert.assertEquals("ts", traceSegmentRef.getTraceSegmentId()); - Assert.assertEquals("ac", traceSegmentRef.getApplicationCode()); - Assert.assertEquals("ph", traceSegmentRef.getPeerHost()); - Assert.assertEquals(0, traceSegmentRef.getSpanId()); - } -} diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/mock/SegmentMock.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/mock/SegmentMock.java index 131f40fbc..f14ca82ea 100644 --- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/mock/SegmentMock.java +++ b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/mock/SegmentMock.java @@ -35,6 +35,10 @@ public class SegmentMock { return JsonFileReader.INSTANCE.read(CacheServiceJsonFile); } + public String mockCacheServiceExceptionSegmentAsString() throws FileNotFoundException { + return JsonFileReader.INSTANCE.read(CacheServiceExceptionJsonFile); + } + public String mockPersistenceServiceSegmentAsString() throws FileNotFoundException { return JsonFileReader.INSTANCE.read(PersistenceServiceJsonFile); } diff --git a/apm-collector/apm-collector-worker/src/test/resources/json/segment/post/normal/cache-service.json b/apm-collector/apm-collector-worker/src/test/resources/json/segment/post/normal/cache-service.json index c1eb1768e..96d711ca6 100644 --- a/apm-collector/apm-collector-worker/src/test/resources/json/segment/post/normal/cache-service.json +++ b/apm-collector/apm-collector-worker/src/test/resources/json/segment/post/normal/cache-service.json @@ -366,7 +366,7 @@ ], "ac": "cache-service", "gt": [ - "Trace.1490922929254.1797892356.6003.69.2" + "Trace.1490922929254.1797892356.6003.69.2,Trace.1490922929254.1797892356.6003.69.3" ] } ] \ No newline at end of file