verride the {@link #onReceive(Map, JsonObject)} method to support a search service.
+ *
verride the {@link #onReceive(Map, JsonElement)} method to support a search service.
*
* @author pengys5
* @since v3.0-2017
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractServlet.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractServlet.java
index 17d0994e0..c51036cac 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractServlet.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractServlet.java
@@ -1,6 +1,6 @@
package org.skywalking.apm.collector.worker.httpserver;
-import com.google.gson.JsonObject;
+import com.google.gson.JsonElement;
import java.io.IOException;
import java.io.PrintWriter;
import java.util.Map;
@@ -31,6 +31,8 @@ public abstract class AbstractServlet extends AbstractLocalSyncWorker {
super(role, clusterContext, selfContext);
}
+ protected abstract Class extends JsonElement> responseClass();
+
/**
* Override this method to implementing business logic.
*
@@ -42,12 +44,16 @@ public abstract class AbstractServlet extends AbstractLocalSyncWorker {
*/
@Override protected void onWork(Object parameter,
Object response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
- JsonObject resJson = new JsonObject();
try {
+ JsonElement resJson = responseClass().newInstance();
onReceive((Map)parameter, resJson);
onSuccessResponse((HttpServletResponse)response, resJson);
} catch (IOException e) {
logger.error(e.getMessage(), e);
+ } catch (InstantiationException e) {
+ logger.error(e.getMessage(), e);
+ } catch (IllegalAccessException e) {
+ logger.error(e.getMessage(), e);
}
}
@@ -55,23 +61,22 @@ public abstract class AbstractServlet extends AbstractLocalSyncWorker {
* Override this method to implementing business logic.
*
* @param parameter {@link Map}, get the request parameter by key.
- * @param response {@link JsonObject}, set the response data as json object.
+ * @param response {@link JsonElement}, set the response data as json object.
* @throws ArgumentsParseException if the key could not contains in parameter
* @throws WorkerInvokeException if any error is detected when call(or ask) worker
* @throws WorkerNotFoundException if the worker reference could not found in context.
*/
protected abstract void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException;
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException;
/**
* Set the worker response and the success status into the servlet response object
*
* @param response {@link HttpServletResponse} object that contains the response the servlet reply to the client
- * @param resJson {@link JsonObject} object that contains the response from worker
+ * @param resJson {@link JsonElement} that contains the response from worker
* @throws IOException if any error is detected when the servlet handles the response.
*/
- protected void onSuccessResponse(HttpServletResponse response, JsonObject resJson) throws IOException {
- resJson.addProperty("isSuccess", true);
+ protected void onSuccessResponse(HttpServletResponse response, JsonElement resJson) throws IOException {
reply(response, resJson, HttpServletResponse.SC_OK);
}
@@ -82,14 +87,15 @@ public abstract class AbstractServlet extends AbstractLocalSyncWorker {
* @param response {@link HttpServletResponse} object that contains the response the servlet reply to the client
*/
protected void onErrorResponse(Exception exception, HttpServletResponse response) {
- JsonObject resJson = new JsonObject();
- resJson.addProperty("isSuccess", false);
- resJson.addProperty("reason", exception.getMessage());
-
+ response.setHeader("reason", exception.getMessage());
try {
- reply(response, resJson, HttpServletResponse.SC_INTERNAL_SERVER_ERROR);
+ reply(response, responseClass().newInstance(), HttpServletResponse.SC_INTERNAL_SERVER_ERROR);
} catch (IOException e) {
logger.error(e.getMessage(), e);
+ } catch (IllegalAccessException e) {
+ logger.error(e.getMessage(), e);
+ } catch (InstantiationException e) {
+ logger.error(e.getMessage(), e);
}
}
@@ -97,11 +103,11 @@ public abstract class AbstractServlet extends AbstractLocalSyncWorker {
* Build the response head and body
*
* @param response {@link HttpServletResponse} object that contains the response the servlet reply to the client
- * @param resJson {@link JsonObject} object that contains the response from worker
+ * @param resJson {@link JsonElement} that contains the response from worker
* @param status http status code
* @throws IOException if an input or output error is detected when the servlet handles the response
*/
- private void reply(HttpServletResponse response, JsonObject resJson, int status) throws IOException {
+ private void reply(HttpServletResponse response, JsonElement resJson, int status) throws IOException {
response.setContentType("text/json");
response.setCharacterEncoding("utf-8");
response.setStatus(status);
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractStreamPost.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractStreamPost.java
index 229c96745..c762d7803 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractStreamPost.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/AbstractStreamPost.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.httpserver;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.io.BufferedReader;
import java.io.IOException;
@@ -58,13 +59,13 @@ public abstract class AbstractStreamPost extends AbstractServlet {
* Override the default implementation, forbidden to call this method.
*
* @param parameter {@link Map}, get the request parameter by key.
- * @param response {@link JsonObject}, set the response data as json object.
+ * @param response {@link JsonElement}, set the response data as json object.
* @throws ArgumentsParseException if the key could not contains in parameter
* @throws WorkerInvokeException if any error is detected when call(or ask) worker
* @throws WorkerNotFoundException if the worker reference could not found in context.
*/
@Override final protected void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
throw new WorkerInvokeException("Use the other method with buffer reader parameter");
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/HttpServer.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/HttpServer.java
index cac22a086..00212e117 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/HttpServer.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/httpserver/HttpServer.java
@@ -29,6 +29,5 @@ public enum HttpServer {
server.setHandler(servletContextHandler);
server.start();
- server.join();
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeCompAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeCompAnalysis.java
index 4e8df2b08..60065a023 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeCompAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeCompAnalysis.java
@@ -9,12 +9,11 @@ import org.skywalking.apm.collector.actor.LocalWorkerContext;
import org.skywalking.apm.collector.actor.Role;
import org.skywalking.apm.collector.worker.RecordAnalysisMember;
import org.skywalking.apm.collector.worker.node.NodeCompIndex;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-import org.skywalking.apm.collector.worker.segment.entity.tag.Tags;
-import org.skywalking.apm.collector.worker.tools.ClientSpanIsLeafTools;
import org.skywalking.apm.collector.worker.tools.CollectionTools;
import org.skywalking.apm.collector.worker.tools.SpanPeersTools;
+import org.skywalking.apm.network.proto.SpanObject;
+import org.skywalking.apm.network.proto.SpanType;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -28,30 +27,29 @@ abstract class AbstractNodeCompAnalysis extends RecordAnalysisMember {
super(role, clusterContext, selfContext);
}
- final void analyseSpans(Segment segment) {
- List spanList = segment.getSpans();
+ final void analyseSpans(TraceSegmentObject segment) {
+ List spanList = segment.getSpansList();
logger.debug("node analysis span isNotEmpty %s", CollectionTools.isNotEmpty(spanList));
if (CollectionTools.isNotEmpty(spanList)) {
logger.debug("node analysis span list SIZE: %s", spanList.size());
- for (Span span : spanList) {
- String kind = Tags.SPAN_KIND.get(span);
- if (Tags.SPAN_KIND_CLIENT.equals(kind) && ClientSpanIsLeafTools.isLeaf(span.getSpanId(), spanList)) {
- String peers = SpanPeersTools.INSTANCE.getPeers(span);
+ for (SpanObject span : spanList) {
+ if (SpanType.Exit.equals(span.getSpanType())) {
+ int peers = SpanPeersTools.INSTANCE.getPeers(span);
JsonObject compJsonObj = new JsonObject();
compJsonObj.addProperty(NodeCompIndex.PEERS, peers);
- compJsonObj.addProperty(NodeCompIndex.NAME, Tags.COMPONENT.get(span));
+ compJsonObj.addProperty(NodeCompIndex.NAME, span.getComponent());
- set(peers, compJsonObj);
- } else if (Tags.SPAN_KIND_SERVER.equals(kind) && span.getParentSpanId() == -1) {
- String peers = segment.getApplicationCode();
+ set(String.valueOf(peers), compJsonObj);
+ } else if (SpanType.Entry.equals(span.getSpanType()) && span.getParentSpanId() == -1) {
+ int peers = segment.getApplicationId();
JsonObject compJsonObj = new JsonObject();
compJsonObj.addProperty(NodeCompIndex.PEERS, peers);
- compJsonObj.addProperty(NodeCompIndex.NAME, Tags.COMPONENT.get(span));
+ compJsonObj.addProperty(NodeCompIndex.NAME, span.getComponent());
- set(peers, compJsonObj);
+ set(String.valueOf(peers), compJsonObj);
}
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeMappingAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeMappingAnalysis.java
index 973ebe8d3..23521ffc8 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeMappingAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/AbstractNodeMappingAnalysis.java
@@ -10,9 +10,9 @@ import org.skywalking.apm.collector.actor.Role;
import org.skywalking.apm.collector.worker.Const;
import org.skywalking.apm.collector.worker.RecordAnalysisMember;
import org.skywalking.apm.collector.worker.node.NodeMappingIndex;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
-import org.skywalking.apm.collector.worker.segment.entity.TraceSegmentRef;
import org.skywalking.apm.collector.worker.tools.CollectionTools;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
+import org.skywalking.apm.network.proto.TraceSegmentReference;
/**
* @author pengys5
@@ -26,15 +26,15 @@ abstract class AbstractNodeMappingAnalysis extends RecordAnalysisMember {
super(role, clusterContext, selfContext);
}
- final void analyseRefs(Segment segment, long timeSlice) {
- List segmentRefList = segment.getRefs();
+ final void analyseRefs(TraceSegmentObject segment, long timeSlice) {
+ List segmentRefList = segment.getRefsList();
logger.debug("node mapping analysis refs isNotEmpty %s", CollectionTools.isNotEmpty(segmentRefList));
if (CollectionTools.isNotEmpty(segmentRefList)) {
logger.debug("node mapping analysis refs list SIZE: %s", segmentRefList.size());
- for (TraceSegmentRef segmentRef : segmentRefList) {
- String peers = Const.PEERS_FRONT_SPLIT + segmentRef.getPeerHost() + Const.PEERS_BEHIND_SPLIT;
- String code = segment.getApplicationCode();
+ for (TraceSegmentReference segmentRef : segmentRefList) {
+ String peers = Const.PEERS_FRONT_SPLIT + segmentRef.getNetworkAddress() + Const.PEERS_BEHIND_SPLIT;
+ int code = segment.getApplicationId();
JsonObject nodeMappingJsonObj = new JsonObject();
nodeMappingJsonObj.addProperty(NodeMappingIndex.CODE, code);
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeCompAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeCompAnalysis.java
index f0864716b..356303cff 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeCompAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeCompAnalysis.java
@@ -11,8 +11,8 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.config.WorkerConfig;
import org.skywalking.apm.collector.worker.node.persistence.NodeCompAgg;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -28,12 +28,12 @@ public class NodeCompAnalysis extends AbstractNodeCompAnalysis {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
analyseSpans(segment);
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingDayAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingDayAnalysis.java
index e3527289e..750f470a4 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingDayAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingDayAnalysis.java
@@ -11,8 +11,8 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.config.WorkerConfig;
import org.skywalking.apm.collector.worker.node.persistence.NodeMappingDayAgg;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -28,12 +28,12 @@ public class NodeMappingDayAnalysis extends AbstractNodeMappingAnalysis {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
analyseRefs(segment, segmentWithTimeSlice.getDay());
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingHourAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingHourAnalysis.java
index 3ef01d58a..2274a8388 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingHourAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingHourAnalysis.java
@@ -11,8 +11,8 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.config.WorkerConfig;
import org.skywalking.apm.collector.worker.node.persistence.NodeMappingHourAgg;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -28,12 +28,12 @@ public class NodeMappingHourAnalysis extends AbstractNodeMappingAnalysis {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
analyseRefs(segment, segmentWithTimeSlice.getHour());
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingMinuteAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingMinuteAnalysis.java
index f01e995c0..bf0351116 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingMinuteAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/node/analysis/NodeMappingMinuteAnalysis.java
@@ -11,8 +11,8 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.config.WorkerConfig;
import org.skywalking.apm.collector.worker.node.persistence.NodeMappingMinuteAgg;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -28,12 +28,12 @@ public class NodeMappingMinuteAnalysis extends AbstractNodeMappingAnalysis {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
analyseRefs(segment, segmentWithTimeSlice.getMinute());
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/NodeRefResSumGetGroupWithTimeSlice.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/NodeRefResSumGetGroupWithTimeSlice.java
index 323dc435e..30b5e797a 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/NodeRefResSumGetGroupWithTimeSlice.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/NodeRefResSumGetGroupWithTimeSlice.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.noderef;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.util.Arrays;
import java.util.Map;
@@ -30,13 +31,17 @@ public class NodeRefResSumGetGroupWithTimeSlice extends AbstractGet {
super(role, clusterContext, selfContext);
}
+ @Override protected Class extends JsonElement> responseClass() {
+ return JsonObject.class;
+ }
+
@Override
public void preStart() throws ProviderNotFoundException {
getClusterContext().findProvider(NodeRefResSumGroupWithTimeSlice.WorkerRole.INSTANCE).create(this);
}
@Override protected void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
if (!parameter.containsKey("startTime") || !parameter.containsKey("endTime") || !parameter.containsKey("timeSliceType")) {
throw new ArgumentsParseException("the request parameter must contains startTime,endTime,timeSliceType");
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/AbstractNodeRefAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/AbstractNodeRefAnalysis.java
index e98ded1ec..3251ef251 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/AbstractNodeRefAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/AbstractNodeRefAnalysis.java
@@ -10,12 +10,11 @@ import org.skywalking.apm.collector.actor.Role;
import org.skywalking.apm.collector.worker.Const;
import org.skywalking.apm.collector.worker.RecordAnalysisMember;
import org.skywalking.apm.collector.worker.noderef.NodeRefIndex;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-import org.skywalking.apm.collector.worker.segment.entity.tag.Tags;
-import org.skywalking.apm.collector.worker.tools.ClientSpanIsLeafTools;
import org.skywalking.apm.collector.worker.tools.CollectionTools;
import org.skywalking.apm.collector.worker.tools.SpanPeersTools;
+import org.skywalking.apm.network.proto.SpanObject;
+import org.skywalking.apm.network.proto.SpanType;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -29,21 +28,21 @@ abstract class AbstractNodeRefAnalysis extends RecordAnalysisMember {
super(role, clusterContext, selfContext);
}
- final void analyseNodeRef(Segment segment, long timeSlice, long minute, long hour, long day,
+ final void analyseNodeRef(TraceSegmentObject segment, long timeSlice, long minute, long hour, long day,
int second) {
- List spanList = segment.getSpans();
+ List spanList = segment.getSpansList();
if (CollectionTools.isNotEmpty(spanList)) {
- for (Span span : spanList) {
+ for (SpanObject span : spanList) {
JsonObject dataJsonObj = new JsonObject();
dataJsonObj.addProperty(NodeRefIndex.TIME_SLICE, timeSlice);
dataJsonObj.addProperty(NodeRefIndex.FRONT_IS_REAL_CODE, true);
dataJsonObj.addProperty(NodeRefIndex.BEHIND_IS_REAL_CODE, true);
- if (Tags.SPAN_KIND_CLIENT.equals(Tags.SPAN_KIND.get(span)) && ClientSpanIsLeafTools.isLeaf(span.getSpanId(), spanList)) {
- String front = segment.getApplicationCode();
+ if (SpanType.Exit.equals(span.getSpanType())) {
+ int front = segment.getApplicationId();
dataJsonObj.addProperty(NodeRefIndex.FRONT, front);
- String behind = SpanPeersTools.INSTANCE.getPeers(span);
+ int behind = SpanPeersTools.INSTANCE.getPeers(span);
dataJsonObj.addProperty(NodeRefIndex.BEHIND, behind);
dataJsonObj.addProperty(NodeRefIndex.BEHIND_IS_REAL_CODE, false);
@@ -51,9 +50,9 @@ abstract class AbstractNodeRefAnalysis extends RecordAnalysisMember {
logger.debug("dag node ref: %s", dataJsonObj.toString());
set(id, dataJsonObj);
buildNodeRefResRecordData(id, span, minute, hour, day, second);
- } else if (Tags.SPAN_KIND_SERVER.equals(Tags.SPAN_KIND.get(span))) {
- if (span.getParentSpanId() == -1 && CollectionTools.isEmpty(segment.getRefs())) {
- String behind = segment.getApplicationCode();
+ } else if (SpanType.Entry.equals(span.getSpanType())) {
+ if (span.getParentSpanId() == -1 && segment.getRefsCount() == 0) {
+ int behind = segment.getApplicationId();
dataJsonObj.addProperty(NodeRefIndex.BEHIND, behind);
String front = Const.USER_CODE;
@@ -68,13 +67,13 @@ abstract class AbstractNodeRefAnalysis extends RecordAnalysisMember {
}
}
- private void buildNodeRefResRecordData(String nodeRefId, Span span, long minute, long hour, long day,
+ private void buildNodeRefResRecordData(String nodeRefId, SpanObject span, long minute, long hour, long day,
int second) {
AbstractNodeRefResSumAnalysis.NodeRefResRecord refResRecord = new AbstractNodeRefResSumAnalysis.NodeRefResRecord(minute, hour, day, second);
refResRecord.setStartTime(span.getStartTime());
refResRecord.setEndTime(span.getEndTime());
refResRecord.setNodeRefId(nodeRefId);
- refResRecord.setError(Tags.ERROR.get(span));
+ refResRecord.setError(span.getIsError());
sendToResSumAnalysis(refResRecord);
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefDayAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefDayAnalysis.java
index 65fb7b379..cfe84bcf0 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefDayAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefDayAnalysis.java
@@ -13,8 +13,8 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.config.WorkerConfig;
import org.skywalking.apm.collector.worker.noderef.persistence.NodeRefDayAgg;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -36,9 +36,9 @@ public class NodeRefDayAnalysis extends AbstractNodeRefAnalysis {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
long minute = segmentWithTimeSlice.getMinute();
long hour = segmentWithTimeSlice.getHour();
@@ -46,7 +46,7 @@ public class NodeRefDayAnalysis extends AbstractNodeRefAnalysis {
int second = segmentWithTimeSlice.getSecond();
analyseNodeRef(segment, segmentWithTimeSlice.getDay(), minute, hour, day, second);
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefHourAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefHourAnalysis.java
index b07d7a09e..336081d60 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefHourAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefHourAnalysis.java
@@ -13,8 +13,8 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.config.WorkerConfig;
import org.skywalking.apm.collector.worker.noderef.persistence.NodeRefHourAgg;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -36,9 +36,9 @@ public class NodeRefHourAnalysis extends AbstractNodeRefAnalysis {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
long minute = segmentWithTimeSlice.getMinute();
long hour = segmentWithTimeSlice.getHour();
@@ -46,7 +46,7 @@ public class NodeRefHourAnalysis extends AbstractNodeRefAnalysis {
int second = segmentWithTimeSlice.getSecond();
analyseNodeRef(segment, segmentWithTimeSlice.getHour(), minute, hour, day, second);
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefMinuteAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefMinuteAnalysis.java
index f8bcd2247..eb9822a0b 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefMinuteAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/noderef/analysis/NodeRefMinuteAnalysis.java
@@ -13,8 +13,8 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.config.WorkerConfig;
import org.skywalking.apm.collector.worker.noderef.persistence.NodeRefMinuteAgg;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -36,16 +36,16 @@ public class NodeRefMinuteAnalysis extends AbstractNodeRefAnalysis {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
long minute = segmentWithTimeSlice.getMinute();
long hour = segmentWithTimeSlice.getHour();
long day = segmentWithTimeSlice.getDay();
int second = segmentWithTimeSlice.getSecond();
analyseNodeRef(segment, segmentWithTimeSlice.getMinute(), minute, hour, day, second);
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentIndex.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentIndex.java
index 01589147b..f2917da81 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentIndex.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentIndex.java
@@ -1,18 +1,19 @@
package org.skywalking.apm.collector.worker.segment;
+import java.io.IOException;
import org.elasticsearch.common.xcontent.XContentBuilder;
import org.elasticsearch.common.xcontent.XContentFactory;
import org.skywalking.apm.collector.worker.config.EsConfig;
import org.skywalking.apm.collector.worker.storage.AbstractIndex;
-import java.io.IOException;
-
/**
* @author pengys5
*/
public class SegmentIndex extends AbstractIndex {
public static final String INDEX = "segment_idx";
+ public static final String TRACE_SEGMENT_ID = "traceSegmentId";
+ public static final String SEGMENT_OBJ_BLOB = "segmentObjBlob";
@Override
public String index() {
@@ -34,31 +35,20 @@ public class SegmentIndex extends AbstractIndex {
return XContentFactory.jsonBuilder()
.startObject()
.startObject("properties")
- .startObject("traceSegmentId")
+ .startObject(TRACE_SEGMENT_ID)
.field("type", "keyword")
.endObject()
- .startObject("startTime")
- .field("type", "date")
- .field("index", "not_analyzed")
- .endObject()
- .startObject("endTime")
- .field("type", "date")
- .field("index", "not_analyzed")
- .endObject()
- .startObject("applicationCode")
- .field("type", "keyword")
- .endObject()
- .startObject("minute")
+ .startObject(TYPE_MINUTE)
.field("type", "long")
- .field("index", "not_analyzed")
.endObject()
- .startObject("hour")
+ .startObject(TYPE_HOUR)
.field("type", "long")
- .field("index", "not_analyzed")
.endObject()
- .startObject("day")
+ .startObject(TYPE_DAY)
.field("type", "long")
- .field("index", "not_analyzed")
+ .endObject()
+ .startObject(SEGMENT_OBJ_BLOB)
+ .field("type", "binary")
.endObject()
.endObject()
.endObject();
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/SegmentReceiver.java
similarity index 60%
rename from apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentPost.java
rename to apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentReceiver.java
index eab460f74..d1199de08 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/SegmentReceiver.java
@@ -1,10 +1,12 @@
package org.skywalking.apm.collector.worker.segment;
-import com.google.gson.JsonObject;
-import java.io.BufferedReader;
-import java.io.IOException;
+import com.google.protobuf.ByteString;
+import com.google.protobuf.InvalidProtocolBufferException;
+import java.util.List;
+import org.apache.commons.codec.binary.Base64;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
+import org.skywalking.apm.collector.actor.AbstractLocalSyncWorkerProvider;
import org.skywalking.apm.collector.actor.ClusterWorkerContext;
import org.skywalking.apm.collector.actor.LocalWorkerContext;
import org.skywalking.apm.collector.actor.ProviderNotFoundException;
@@ -14,8 +16,7 @@ import org.skywalking.apm.collector.actor.WorkerNotFoundException;
import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.globaltrace.analysis.GlobalTraceAnalysis;
-import org.skywalking.apm.collector.worker.httpserver.AbstractStreamPost;
-import org.skywalking.apm.collector.worker.httpserver.AbstractStreamPostProvider;
+import org.skywalking.apm.collector.worker.grpcserver.AbstractReceiver;
import org.skywalking.apm.collector.worker.httpserver.ArgumentsParseException;
import org.skywalking.apm.collector.worker.node.analysis.NodeCompAnalysis;
import org.skywalking.apm.collector.worker.node.analysis.NodeMappingDayAnalysis;
@@ -27,20 +28,21 @@ import org.skywalking.apm.collector.worker.noderef.analysis.NodeRefMinuteAnalysi
import org.skywalking.apm.collector.worker.segment.analysis.SegmentAnalysis;
import org.skywalking.apm.collector.worker.segment.analysis.SegmentCostAnalysis;
import org.skywalking.apm.collector.worker.segment.analysis.SegmentExceptionAnalysis;
-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.entity.SegmentDeserialize;
+import org.skywalking.apm.collector.worker.segment.entity.SegmentAndBase64;
import org.skywalking.apm.collector.worker.storage.AbstractTimeSlice;
import org.skywalking.apm.collector.worker.tools.DateTools;
+import org.skywalking.apm.network.proto.SpanObject;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
+import org.skywalking.apm.network.proto.UpstreamSegment;
import org.skywalking.apm.util.StringUtil;
/**
* @author pengys5
*/
-public class SegmentPost extends AbstractStreamPost {
- private static final Logger logger = LogManager.getFormatterLogger(SegmentPost.class);
+public class SegmentReceiver extends AbstractReceiver {
+ private static final Logger logger = LogManager.getFormatterLogger(SegmentReceiver.class);
- public SegmentPost(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ public SegmentReceiver(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -67,59 +69,47 @@ public class SegmentPost extends AbstractStreamPost {
* Read segment's buffer from buffer reader by stream mode. when finish read one segment then send to analysis.
* This method in there, so post servlet just can receive segments data.
*/
- @Override protected void onReceive(BufferedReader bufferedReader,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
- Segment segment;
- try {
- do {
- int character;
- StringBuilder builder = new StringBuilder();
- while ((character = bufferedReader.read()) != ' ') {
- if (character == -1) {
- return;
- }
- builder.append((char)character);
- }
+ @Override protected void onReceive(
+ Object request) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ if (request instanceof UpstreamSegment) {
+ UpstreamSegment upstreamSegment = (UpstreamSegment)request;
+ ByteString segmentByte = upstreamSegment.getSegment();
+ List globalTraceIds = upstreamSegment.getGlobalTraceIdsList();
- int length = Integer.valueOf(builder.toString());
- builder = new StringBuilder();
+ String segmentBase64 = new String(Base64.encodeBase64(segmentByte.toByteArray()));
- 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 = SegmentDeserialize.INSTANCE.deserializeSingle(segmentJsonStr);
- tellWorkers(new SegmentAndJson(segment, segmentJsonStr));
+ TraceSegmentObject segment;
+ try {
+ segment = TraceSegmentObject.parseFrom(segmentByte);
+ } catch (InvalidProtocolBufferException e) {
+ throw new ArgumentsParseException(e.getMessage(), e);
}
- while (segment != null);
- } catch (IOException e) {
- throw new ArgumentsParseException(e.getMessage(), e);
+ tellWorkers(new SegmentAndBase64(segment, segmentBase64), globalTraceIds);
}
}
- private void tellWorkers(SegmentAndJson segmentAndJson) throws WorkerNotFoundException, WorkerInvokeException {
- Segment segment = segmentAndJson.getSegment();
+ private void tellWorkers(
+ SegmentAndBase64 segmentAndBase64,
+ List globalTraceIds) throws WorkerNotFoundException, WorkerInvokeException {
+ TraceSegmentObject segment = segmentAndBase64.getObject();
try {
validateData(segment);
- } catch (IllegalArgumentException e) {
+ } catch (ArgumentsParseException e) {
+ logger.error(e.getMessage(), e);
return;
}
logger.debug("receive message instanceof TraceSegment, traceSegmentId is %s", segment.getTraceSegmentId());
+ SpanObject firstSpan = segment.getSpans(segment.getSpansCount() - 1);
- long minuteSlice = DateTools.getMinuteSlice(segment.getStartTime());
- long hourSlice = DateTools.getHourSlice(segment.getStartTime());
- long daySlice = DateTools.getDaySlice(segment.getStartTime());
- int second = DateTools.getSecond(segment.getStartTime());
+ long minuteSlice = DateTools.getMinuteSlice(firstSpan.getStartTime());
+ long hourSlice = DateTools.getHourSlice(firstSpan.getStartTime());
+ long daySlice = DateTools.getDaySlice(firstSpan.getStartTime());
+ int second = DateTools.getSecond(firstSpan.getStartTime());
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(segmentAndJson);
+ SegmentWithTimeSlice segmentWithTimeSlice = new SegmentWithTimeSlice(segment, globalTraceIds, minuteSlice, hourSlice, daySlice, second);
+ getSelfContext().lookup(SegmentAnalysis.Role.INSTANCE).tell(segmentAndBase64);
getSelfContext().lookup(SegmentCostAnalysis.Role.INSTANCE).tell(segmentWithTimeSlice);
getSelfContext().lookup(GlobalTraceAnalysis.Role.INSTANCE).tell(segmentWithTimeSlice);
@@ -145,29 +135,31 @@ public class SegmentPost extends AbstractStreamPost {
getSelfContext().lookup(NodeMappingDayAnalysis.Role.INSTANCE).tell(segmentWithTimeSlice);
}
- private void validateData(Segment segment) {
+ private void validateData(TraceSegmentObject segment) throws ArgumentsParseException {
if (StringUtil.isEmpty(segment.getTraceSegmentId())) {
- throw new IllegalArgumentException("traceSegmentId required");
+ throw new ArgumentsParseException("traceSegmentId required");
}
- if (0 == segment.getStartTime()) {
- throw new IllegalArgumentException("startTime required");
+ if (segment.getSpansCount() < 1) {
+ throw new ArgumentsParseException("must contain at least one span");
+ }
+ SpanObject firstSpan = segment.getSpans(segment.getSpansCount() - 1);
+ if (firstSpan.getSpanId() != 0 && firstSpan.getParentSpanId() != -1) {
+ throw new ArgumentsParseException("first span id must equals 0 and parent span id must equals -1");
+ }
+ if (0 == firstSpan.getStartTime()) {
+ throw new ArgumentsParseException("startTime required");
}
}
- public static class Factory extends AbstractStreamPostProvider {
- @Override
- public String servletPath() {
- return "/segments";
- }
-
+ public static class Factory extends AbstractLocalSyncWorkerProvider {
@Override
public Role role() {
return WorkerRole.INSTANCE;
}
@Override
- public SegmentPost workerInstance(ClusterWorkerContext clusterContext) {
- return new SegmentPost(role(), clusterContext, new LocalWorkerContext());
+ public SegmentReceiver workerInstance(ClusterWorkerContext clusterContext) {
+ return new SegmentReceiver(role(), clusterContext, new LocalWorkerContext());
}
}
@@ -176,7 +168,7 @@ public class SegmentPost extends AbstractStreamPost {
@Override
public String roleName() {
- return SegmentPost.class.getSimpleName();
+ return SegmentReceiver.class.getSimpleName();
}
@Override
@@ -186,15 +178,23 @@ public class SegmentPost extends AbstractStreamPost {
}
public static class SegmentWithTimeSlice extends AbstractTimeSlice {
- private final Segment segment;
+ private final TraceSegmentObject segment;
- public SegmentWithTimeSlice(Segment segment, long minute, long hour, long day, int second) {
+ private final List globalTraceIds;
+
+ public SegmentWithTimeSlice(TraceSegmentObject segment, List globalTraceIds, long minute, long hour,
+ long day, int second) {
super(minute, hour, day, second);
this.segment = segment;
+ this.globalTraceIds = globalTraceIds;
}
- public Segment getSegment() {
+ public TraceSegmentObject getSegment() {
return segment;
}
+
+ public List getGlobalTraceIds() {
+ return globalTraceIds;
+ }
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentTopGet.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentTopGet.java
index da67dfe48..9896b59c6 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentTopGet.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/SegmentTopGet.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.segment;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.util.Arrays;
import java.util.Map;
@@ -30,13 +31,17 @@ public class SegmentTopGet extends AbstractGet {
super(role, clusterContext, selfContext);
}
+ @Override protected Class extends JsonElement> responseClass() {
+ return JsonObject.class;
+ }
+
@Override
public void preStart() throws ProviderNotFoundException {
getClusterContext().findProvider(SegmentTopSearch.WorkerRole.INSTANCE).create(this);
}
@Override protected void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
if (!parameter.containsKey("startTime") || !parameter.containsKey("endTime") || !parameter.containsKey("from") || !parameter.containsKey("limit")) {
throw new ArgumentsParseException("the request parameter must contains startTime, endTime, from, limit");
}
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 d2a5eb938..ffc568d04 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
@@ -13,7 +13,7 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
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.SegmentAndJson;
+import org.skywalking.apm.collector.worker.segment.entity.SegmentAndBase64;
import org.skywalking.apm.collector.worker.segment.persistence.SegmentSave;
/**
@@ -34,11 +34,11 @@ public class SegmentAnalysis extends RecordAnalysisMember {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentAndJson) {
- SegmentAndJson segmentAndJson = (SegmentAndJson)message;
+ if (message instanceof SegmentAndBase64) {
+ SegmentAndBase64 segmentAndBase64 = (SegmentAndBase64)message;
try {
- getSelfContext().lookup(SegmentSave.Role.INSTANCE).tell(segmentAndJson);
+ getSelfContext().lookup(SegmentSave.Role.INSTANCE).tell(segmentAndBase64);
} catch (WorkerInvokeException | WorkerNotFoundException e) {
e.printStackTrace();
logger.error(e.getMessage(), e);
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentCostAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentCostAnalysis.java
index 2f4f7c79d..72161f8a2 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentCostAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentCostAnalysis.java
@@ -14,11 +14,11 @@ 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.SegmentCostIndex;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
-import org.skywalking.apm.collector.worker.segment.entity.Span;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
import org.skywalking.apm.collector.worker.segment.persistence.SegmentCostSave;
import org.skywalking.apm.collector.worker.tools.CollectionTools;
+import org.skywalking.apm.network.proto.SpanObject;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -38,36 +38,37 @@ public class SegmentCostAnalysis extends RecordAnalysisMember {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
- if (CollectionTools.isNotEmpty(segment.getSpans())) {
- for (Span span : segment.getSpans()) {
+ if (CollectionTools.isNotEmpty(segment.getSpansList())) {
+ for (SpanObject span : segment.getSpansList()) {
if (span.getParentSpanId() == -1) {
- JsonObject dataJsonObj = new JsonObject();
- dataJsonObj.addProperty(SegmentCostIndex.SEG_ID, segment.getTraceSegmentId());
- dataJsonObj.addProperty(SegmentCostIndex.START_TIME, span.getStartTime());
- dataJsonObj.addProperty(SegmentCostIndex.END_TIME, span.getEndTime());
- if (segment.getRelatedGlobalTraces().get() != null && segment.getRelatedGlobalTraces().get().size() > 0) {
- dataJsonObj.addProperty(SegmentCostIndex.GLOBAL_TRACE_ID, segment.getRelatedGlobalTraces().get().get(0));
- }
- dataJsonObj.addProperty(SegmentCostIndex.OPERATION_NAME, span.getOperationName());
- dataJsonObj.addProperty(SegmentCostIndex.TIME_SLICE, segmentWithTimeSlice.getMinute());
+ for (String globalTraceId : segmentWithTimeSlice.getGlobalTraceIds()) {
+ segment.getGlobalTraceIdsList();
+ JsonObject dataJsonObj = new JsonObject();
+ dataJsonObj.addProperty(SegmentCostIndex.SEG_ID, segment.getTraceSegmentId());
+ dataJsonObj.addProperty(SegmentCostIndex.START_TIME, span.getStartTime());
+ dataJsonObj.addProperty(SegmentCostIndex.END_TIME, span.getEndTime());
+ dataJsonObj.addProperty(SegmentCostIndex.GLOBAL_TRACE_ID, globalTraceId);
+ dataJsonObj.addProperty(SegmentCostIndex.OPERATION_NAME, span.getOperationName());
+ dataJsonObj.addProperty(SegmentCostIndex.TIME_SLICE, segmentWithTimeSlice.getMinute());
- long startTime = span.getStartTime();
- long endTime = span.getEndTime();
- long cost = endTime - startTime;
- if (cost == 0) {
- cost = 1;
+ long startTime = span.getStartTime();
+ long endTime = span.getEndTime();
+ long cost = endTime - startTime;
+ if (cost == 0) {
+ cost = 1;
+ }
+ dataJsonObj.addProperty(SegmentCostIndex.COST, cost);
+ set(segment.getTraceSegmentId(), dataJsonObj);
}
- dataJsonObj.addProperty(SegmentCostIndex.COST, cost);
- set(segment.getTraceSegmentId(), dataJsonObj);
}
}
}
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentExceptionAnalysis.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentExceptionAnalysis.java
index d990a105a..09df42fdd 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentExceptionAnalysis.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/analysis/SegmentExceptionAnalysis.java
@@ -16,13 +16,12 @@ 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.SegmentExceptionIndex;
-import org.skywalking.apm.collector.worker.segment.SegmentPost;
-import org.skywalking.apm.collector.worker.segment.entity.LogData;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-import org.skywalking.apm.collector.worker.segment.entity.tag.Tags;
+import org.skywalking.apm.collector.worker.segment.SegmentReceiver;
import org.skywalking.apm.collector.worker.segment.persistence.SegmentExceptionSave;
import org.skywalking.apm.collector.worker.tools.CollectionTools;
+import org.skywalking.apm.network.proto.LogMessage;
+import org.skywalking.apm.network.proto.SpanObject;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -42,13 +41,13 @@ public class SegmentExceptionAnalysis extends RecordAnalysisMember {
@Override
public void analyse(Object message) {
- if (message instanceof SegmentPost.SegmentWithTimeSlice) {
- SegmentPost.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentPost.SegmentWithTimeSlice)message;
- Segment segment = segmentWithTimeSlice.getSegment();
+ if (message instanceof SegmentReceiver.SegmentWithTimeSlice) {
+ SegmentReceiver.SegmentWithTimeSlice segmentWithTimeSlice = (SegmentReceiver.SegmentWithTimeSlice)message;
+ TraceSegmentObject segment = segmentWithTimeSlice.getSegment();
- if (CollectionTools.isNotEmpty(segment.getSpans())) {
- for (Span span : segment.getSpans()) {
- boolean isError = Tags.ERROR.get(span);
+ if (CollectionTools.isNotEmpty(segment.getSpansList())) {
+ for (SpanObject span : segment.getSpansList()) {
+ boolean isError = span.getIsError();
JsonObject dataJsonObj = new JsonObject();
dataJsonObj.addProperty(SegmentExceptionIndex.IS_ERROR, isError);
@@ -56,11 +55,11 @@ public class SegmentExceptionAnalysis extends RecordAnalysisMember {
JsonArray errorKind = new JsonArray();
if (isError) {
- List logDataList = span.getLogs();
- for (LogData logData : logDataList) {
- if (logData.getFields().containsKey("error.kind")) {
- errorKind.add(String.valueOf(logData.getFields().get("error.kind")));
- }
+ List logMessages = span.getLogsList();
+ for (LogMessage logMessage : logMessages) {
+// if (logMessage.getFields().containsKey("error.kind")) {
+// errorKind.add(String.valueOf(logData.getFields().get("error.kind")));
+// }
}
}
dataJsonObj.add(SegmentExceptionIndex.ERROR_KIND, errorKind);
@@ -68,7 +67,7 @@ public class SegmentExceptionAnalysis extends RecordAnalysisMember {
}
}
} else {
- logger.error("unhandled message, message instance must SegmentPost.SegmentWithTimeSlice, but is %s", message.getClass().toString());
+ logger.error("unhandled message, message instance must SegmentReceiver.SegmentWithTimeSlice, but is %s", message.getClass().toString());
}
}
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
deleted file mode 100644
index f3d248614..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/GlobalTraceId.java
+++ /dev/null
@@ -1,53 +0,0 @@
-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
- */
-@JsonAdapter(GlobalTraceId.Serializer.class)
-public class GlobalTraceId {
-
- public GlobalTraceId() {
- globalTraceIds = new LinkedList<>();
- }
-
- 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
deleted file mode 100644
index c97f31062..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/LogData.java
+++ /dev/null
@@ -1,24 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity;
-
-import com.google.gson.annotations.SerializedName;
-import java.util.Map;
-
-/**
- * @author pengys5
- */
-public class LogData {
-
- @SerializedName("tm")
- private long time;
-
- @SerializedName("fi")
- private Map fields;
-
- public long getTime() {
- return time;
- }
-
- public Map getFields() {
- return fields;
- }
-}
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
deleted file mode 100644
index 13c5bb352..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Segment.java
+++ /dev/null
@@ -1,59 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity;
-
-import com.google.gson.annotations.SerializedName;
-import java.util.List;
-
-/**
- * @author pengys5
- */
-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;
-
- @SerializedName("gt")
- private GlobalTraceId relatedGlobalTraces;
-
- public String getTraceSegmentId() {
- return traceSegmentId;
- }
-
- public long getStartTime() {
- return startTime;
- }
-
- public long getEndTime() {
- return endTime;
- }
-
- public String getApplicationCode() {
- return applicationCode;
- }
-
- public List getRefs() {
- return refs;
- }
-
- public List getSpans() {
- return spans;
- }
-
- public GlobalTraceId getRelatedGlobalTraces() {
- return relatedGlobalTraces;
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndBase64.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndBase64.java
new file mode 100644
index 000000000..38e4ac7b4
--- /dev/null
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndBase64.java
@@ -0,0 +1,34 @@
+package org.skywalking.apm.collector.worker.segment.entity;
+
+import com.google.gson.JsonObject;
+import org.skywalking.apm.collector.worker.segment.SegmentIndex;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
+
+/**
+ * @author pengys5
+ */
+public class SegmentAndBase64 {
+
+ private final TraceSegmentObject object;
+ private final String base64;
+
+ public SegmentAndBase64(TraceSegmentObject object, String base64) {
+ this.object = object;
+ this.base64 = base64;
+ }
+
+ public TraceSegmentObject getObject() {
+ return object;
+ }
+
+ public String getBase64() {
+ return base64;
+ }
+
+ public String getSegmentJsonStr() {
+ JsonObject segmentJson = new JsonObject();
+ segmentJson.addProperty(SegmentIndex.TRACE_SEGMENT_ID, object.getTraceSegmentId());
+ segmentJson.addProperty(SegmentIndex.SEGMENT_OBJ_BLOB, base64);
+ return segmentJson.toString();
+ }
+}
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
deleted file mode 100644
index 0a8e08f56..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/SegmentAndJson.java
+++ /dev/null
@@ -1,23 +0,0 @@
-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 af338afaa..b9c64f05b 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,7 +1,10 @@
package org.skywalking.apm.collector.worker.segment.entity;
-import com.google.gson.Gson;
+import com.sun.org.apache.xml.internal.security.utils.Base64;
import java.io.IOException;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* The SegmentDeserialize provides single segment json string deserialize and segment array file
@@ -13,17 +16,21 @@ import java.io.IOException;
public enum SegmentDeserialize {
INSTANCE;
- private final Gson gson = new Gson();
+ private final Logger logger = LogManager.getFormatterLogger(SegmentDeserialize.class);
/**
- * Single segment json string deserialize.
+ * Segment object binary value as a base64 encoded string deserialize.
*
- * @param singleSegmentJsonStr a segment json string
- * @return an {@link Segment}
- * @throws IOException if json string illegal or file broken.
+ * @param segmentObjBlob , to be a binary value as a base64 encoded string
+ * @return an {@link TraceSegmentObject}
*/
- public Segment deserializeSingle(String singleSegmentJsonStr) throws IOException {
- Segment segment = gson.fromJson(singleSegmentJsonStr, Segment.class);
- return segment;
+ public TraceSegmentObject deserializeSingle(String segmentObjBlob) {
+ try {
+ byte[] decode = Base64.decode(segmentObjBlob);
+ return TraceSegmentObject.parseFrom(decode);
+ } catch (Exception e) {
+ logger.error(e.getMessage(), e);
+ }
+ return null;
}
}
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
deleted file mode 100644
index 716ee2fdf..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/Span.java
+++ /dev/null
@@ -1,74 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity;
-
-import com.google.gson.annotations.SerializedName;
-import java.util.List;
-import java.util.Map;
-
-/**
- * @author pengys5
- */
-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() {
- return spanId;
- }
-
- public int getParentSpanId() {
- return parentSpanId;
- }
-
- public long getStartTime() {
- return startTime;
- }
-
- public long getEndTime() {
- return endTime;
- }
-
- public String getOperationName() {
- return operationName;
- }
-
- public String getStrTag(String key) {
- return tagsWithStr.get(key);
- }
-
- public Boolean getBoolTag(String key) {
- return tagsWithBool.get(key);
- }
-
- public Integer getIntTag(String key) {
- return tagsWithInt.get(key);
- }
-
- public List getLogs() {
- return logs;
- }
-}
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
deleted file mode 100644
index 3629728e0..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/TraceSegmentRef.java
+++ /dev/null
@@ -1,37 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity;
-
-import com.google.gson.annotations.SerializedName;
-
-/**
- * @author pengys5
- */
-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() {
- return traceSegmentId;
- }
-
- public int getSpanId() {
- return spanId;
- }
-
- public String getApplicationCode() {
- return applicationCode;
- }
-
- public String getPeerHost() {
- return peerHost;
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/AbstractTag.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/AbstractTag.java
deleted file mode 100644
index f3d6197f1..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/AbstractTag.java
+++ /dev/null
@@ -1,16 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity.tag;
-
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-
-public abstract class AbstractTag {
- /**
- * The key of this Tag.
- */
- protected final String key;
-
- public AbstractTag(String tagKey) {
- this.key = tagKey;
- }
-
- public abstract T get(Span span);
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/BooleanTag.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/BooleanTag.java
deleted file mode 100644
index 8c3fb0e6d..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/BooleanTag.java
+++ /dev/null
@@ -1,35 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity.tag;
-
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-
-/**
- * Do the same thing as {@link StringTag}, just with a {@link Boolean} value.
- *
- * Created by wusheng on 2017/2/17.
- */
-public class BooleanTag extends AbstractTag {
-
- private boolean defaultValue;
-
- public BooleanTag(String key, boolean defaultValue) {
- super(key);
- this.defaultValue = defaultValue;
- }
-
- /**
- * Get a tag value, type of {@link Boolean}. After akka-message/serialize, all tags values are type of {@link
- * String}, convert to {@link Boolean}, if necessary.
- *
- * @param span
- * @return tag value
- */
- @Override
- public Boolean get(Span span) {
- Boolean tagValue = span.getBoolTag(super.key);
- if (tagValue == null) {
- return defaultValue;
- } else {
- return tagValue;
- }
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/IntTag.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/IntTag.java
deleted file mode 100644
index 169fe58ca..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/IntTag.java
+++ /dev/null
@@ -1,31 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity.tag;
-
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-
-/**
- * Do the same thing as {@link StringTag}, just with a {@link Integer} value.
- *
- * Created by wusheng on 2017/2/18.
- */
-public class IntTag extends AbstractTag {
- public IntTag(String key) {
- super(key);
- }
-
- /**
- * Get a tag value, type of {@link Integer}.
- * After akka-message/serialize, all tags values are type of {@link String}, convert to {@link Integer}, if necessary.
- *
- * @param span
- * @return tag value
- */
- @Override
- public Integer get(Span span) {
- Integer tagValue = span.getIntTag(super.key);
- if (tagValue == null) {
- return null;
- } else {
- return tagValue;
- }
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/ShortTag.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/ShortTag.java
deleted file mode 100644
index de1ec580f..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/ShortTag.java
+++ /dev/null
@@ -1,31 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity.tag;
-
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-
-/**
- * Do the same thing as {@link StringTag}, just with a {@link Short} value.
- *
- * Created by wusheng on 2017/2/17.
- */
-public class ShortTag extends AbstractTag {
- public ShortTag(String key) {
- super(key);
- }
-
- /**
- * Get a tag value, type of {@link Short}.
- * After akka-message/serialize, all tags values are type of {@link String}, convert to {@link Short}, if necessary.
- *
- * @param span
- * @return tag value
- */
- @Override
- public Short get(Span span) {
- Integer tagValue = span.getIntTag(super.key);
- if (tagValue == null) {
- return null;
- } else {
- return Short.valueOf(tagValue.toString());
- }
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/StringTag.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/StringTag.java
deleted file mode 100644
index 370c7617a..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/StringTag.java
+++ /dev/null
@@ -1,21 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity.tag;
-
-
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-
-/**
- * A subclass of {@link AbstractTag},
- * represent a tag with a {@link String} value.
- *
- * Created by wusheng on 2017/2/17.
- */
-public class StringTag extends AbstractTag {
- public StringTag(String tagKey) {
- super(tagKey);
- }
-
- @Override
- public String get(Span span) {
- return span.getStrTag(super.key);
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/Tags.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/Tags.java
deleted file mode 100644
index 007c63869..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/entity/tag/Tags.java
+++ /dev/null
@@ -1,97 +0,0 @@
-package org.skywalking.apm.collector.worker.segment.entity.tag;
-
-
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-
-/**
- * The span tags are supported by sky-walking engine.
- * As default, all tags will be stored, but these ones have particular meanings.
- *
- * Created by wusheng on 2017/2/17.
- */
-public final class Tags {
- private Tags() {
- }
-
- /**
- * URL records the url of the incoming request.
- */
- public static final StringTag URL = new StringTag("url");
-
- /**
- * STATUS_CODE records the http status code of the response.
- */
- public static final IntTag STATUS_CODE = new IntTag("status_code");
-
- /**
- * SPAN_KIND hints at the relationship between spans, e.g. client/server.
- */
- public static final StringTag SPAN_KIND = new StringTag("span.kind");
-
- /**
- * A constant for setting the span kind to indicate that it represents a server span.
- */
- public static final String SPAN_KIND_SERVER = "server";
-
- /**
- * A constant for setting the span kind to indicate that it represents a client span.
- */
- public static final String SPAN_KIND_CLIENT = "client";
-
- /**
- * SPAN_LAYER represents the kind of span.
- *
- * e.g.
- * db=database;
- * rpc=Remote Procedure Call Framework, like motan, thift;
- * nosql=something like redis/memcache
- */
- public static final class SPAN_LAYER {
- private static StringTag SPAN_LAYER_TAG = new StringTag("span.layer");
-
- public static String get(Span span) {
- return SPAN_LAYER_TAG.get(span);
- }
- }
-
- /**
- * COMPONENT is a low-cardinality identifier of the module, library, or package that is instrumented.
- * Like dubbo/dubbox/motan
- */
- public static final StringTag COMPONENT = new StringTag("component");
-
- /**
- * ERROR indicates whether a Span ended in an error state.
- */
- public static final BooleanTag ERROR = new BooleanTag("error", false);
-
- /**
- * PEER_HOST records host address (ip:port, or ip1:port1,ip2:port2) of the peer, maybe IPV4, IPV6 or hostname.
- */
- public static final StringTag PEER_HOST = new StringTag("peer.host");
-
- /**
- * PEER_PORT records remote port of the peer
- */
- public static final IntTag PEER_PORT = new IntTag("peer.port");
-
- /**
- * PEERS records multiple host address and port of remote
- */
- public static final StringTag PEERS = new StringTag("peers");
-
- /**
- * DB_TYPE records database type, such as sql, redis, cassandra and so on.
- */
- public static final StringTag DB_TYPE = new StringTag("db.type");
-
- /**
- * DB_INSTANCE records database instance name.
- */
- public static final StringTag DB_INSTANCE = new StringTag("db.instance");
-
- /**
- * DB_STATEMENT records the sql statement of the database access.
- */
- public static final StringTag DB_STATEMENT = new StringTag("db.statement");
-}
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 ba94bab86..a74072fed 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,6 @@
package org.skywalking.apm.collector.worker.segment.persistence;
+import com.google.gson.JsonObject;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
@@ -13,7 +14,7 @@ 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.SegmentAndJson;
+import org.skywalking.apm.collector.worker.segment.entity.SegmentAndBase64;
import org.skywalking.apm.collector.worker.storage.AbstractIndex;
import org.skywalking.apm.collector.worker.storage.EsClient;
import org.skywalking.apm.collector.worker.storage.PersistenceWorkerListener;
@@ -46,11 +47,11 @@ public class SegmentSave extends PersistenceMember= CacheSizeConfig.Cache.Persistence.SIZE) {
persistence(data.asMap());
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearch.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearch.java
index 15ff20174..00d1631a3 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearch.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/segment/persistence/SegmentTopSearch.java
@@ -2,8 +2,8 @@ package org.skywalking.apm.collector.worker.segment.persistence;
import com.google.gson.JsonArray;
import com.google.gson.JsonObject;
-import java.io.IOException;
import java.util.List;
+import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.search.SearchRequestBuilder;
import org.elasticsearch.action.search.SearchResponse;
import org.elasticsearch.action.search.SearchType;
@@ -25,10 +25,10 @@ import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.segment.SegmentCostIndex;
import org.skywalking.apm.collector.worker.segment.SegmentExceptionIndex;
import org.skywalking.apm.collector.worker.segment.SegmentIndex;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
import org.skywalking.apm.collector.worker.segment.entity.SegmentDeserialize;
import org.skywalking.apm.collector.worker.storage.EsClient;
import org.skywalking.apm.collector.worker.tools.CollectionTools;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
import org.skywalking.apm.util.StringUtil;
/**
@@ -102,15 +102,11 @@ public class SegmentTopSearch extends AbstractLocalSyncWorker {
topSegmentJson.addProperty(SegmentCostIndex.OPERATION_NAME, (String)searchHit.getSource().get(SegmentCostIndex.OPERATION_NAME));
topSegmentJson.addProperty(SegmentCostIndex.COST, (Number)searchHit.getSource().get(SegmentCostIndex.COST));
- String segmentSource = EsClient.INSTANCE.getClient().prepareGet(SegmentIndex.INDEX, SegmentIndex.TYPE_RECORD, segId).get().getSourceAsString();
- logger().debug("segmentSource:" + segmentSource);
- Segment segment;
- try {
- segment = SegmentDeserialize.INSTANCE.deserializeSingle(segmentSource);
- } catch (IOException e) {
- throw new WorkerException(e.getMessage(), e);
- }
- List distributedTraceIdList = segment.getRelatedGlobalTraces().get();
+ GetResponse getResponse = EsClient.INSTANCE.getClient().prepareGet(SegmentIndex.INDEX, SegmentIndex.TYPE_RECORD, segId).get();
+ String segmentObjBlob = (String)getResponse.getSource().get(SegmentIndex.SEGMENT_OBJ_BLOB);
+
+ TraceSegmentObject segment = SegmentDeserialize.INSTANCE.deserializeSingle(segmentObjBlob);
+ List distributedTraceIdList = segment.getGlobalTraceIdsList();
JsonArray distributedTraceIdArray = new JsonArray();
if (CollectionTools.isNotEmpty(distributedTraceIdList)) {
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/SpanGetWithId.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/SpanGetWithId.java
index bdcb5f592..0e41c43b4 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/SpanGetWithId.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/span/SpanGetWithId.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.span;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.util.Arrays;
import java.util.Map;
@@ -30,13 +31,17 @@ public class SpanGetWithId extends AbstractGet {
super(role, clusterContext, selfContext);
}
+ @Override protected Class extends JsonElement> responseClass() {
+ return JsonObject.class;
+ }
+
@Override
public void preStart() throws ProviderNotFoundException {
getClusterContext().findProvider(SpanSearchWithId.WorkerRole.INSTANCE).create(this);
}
@Override protected void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
if (!parameter.containsKey("segId") || !parameter.containsKey("spanId")) {
throw new ArgumentsParseException("the request parameter must contains segId, spanId");
}
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 74da73e47..d62ab7907 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
@@ -2,7 +2,6 @@ package org.skywalking.apm.collector.worker.span.persistence;
import com.google.gson.Gson;
import com.google.gson.JsonObject;
-import java.io.IOException;
import java.util.List;
import org.elasticsearch.action.get.GetResponse;
import org.skywalking.apm.collector.actor.AbstractLocalSyncWorker;
@@ -15,10 +14,10 @@ import org.skywalking.apm.collector.actor.selector.RollingSelector;
import org.skywalking.apm.collector.actor.selector.WorkerSelector;
import org.skywalking.apm.collector.worker.Const;
import org.skywalking.apm.collector.worker.segment.SegmentIndex;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
import org.skywalking.apm.collector.worker.segment.entity.SegmentDeserialize;
-import org.skywalking.apm.collector.worker.segment.entity.Span;
import org.skywalking.apm.collector.worker.storage.GetResponseFromEs;
+import org.skywalking.apm.network.proto.SpanObject;
+import org.skywalking.apm.network.proto.TraceSegmentObject;
/**
* @author pengys5
@@ -36,18 +35,15 @@ public class SpanSearchWithId extends AbstractLocalSyncWorker {
if (request instanceof RequestEntity) {
RequestEntity search = (RequestEntity)request;
GetResponse getResponse = GetResponseFromEs.INSTANCE.get(SegmentIndex.INDEX, SegmentIndex.TYPE_RECORD, search.segId);
- Segment segment;
- try {
- segment = SegmentDeserialize.INSTANCE.deserializeSingle(getResponse.getSourceAsString());
- } catch (IOException e) {
- throw new WorkerException(e.getMessage(), e);
- }
- List spanList = segment.getSpans();
+ String segmentObjBlob = (String)getResponse.getSource().get(SegmentIndex.SEGMENT_OBJ_BLOB);
+
+ TraceSegmentObject segment = SegmentDeserialize.INSTANCE.deserializeSingle(segmentObjBlob);
+ List spanList = segment.getSpansList();
getResponse.getSource();
JsonObject dataJson = new JsonObject();
- for (Span span : spanList) {
+ for (SpanObject span : spanList) {
if (String.valueOf(span.getSpanId()).equals(search.spanId)) {
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/ClientSpanIsLeafTools.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/ClientSpanIsLeafTools.java
deleted file mode 100644
index 24a344f7e..000000000
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/ClientSpanIsLeafTools.java
+++ /dev/null
@@ -1,27 +0,0 @@
-package org.skywalking.apm.collector.worker.tools;
-
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-import org.skywalking.apm.collector.worker.segment.entity.tag.Tags;
-
-import java.util.List;
-
-/**
- * @author pengys5
- */
-public class ClientSpanIsLeafTools {
- private static final Logger logger = LogManager.getFormatterLogger(ClientSpanIsLeafTools.class);
-
- public static boolean isLeaf(int spanId, List spanList) {
- boolean isLeaf = true;
- for (Span span : spanList) {
- if (span.getParentSpanId() == spanId && Tags.SPAN_KIND_CLIENT.equals(Tags.SPAN_KIND.get(span))) {
- logger.debug("current spanId=%s, merge spanId=%s, span kind=%s", spanId, span.getSpanId(), Tags.SPAN_KIND.get(span));
- isLeaf = false;
- }
- }
-
- return isLeaf;
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/SpanPeersTools.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/SpanPeersTools.java
index 92fbcdda6..e48cb2106 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/SpanPeersTools.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tools/SpanPeersTools.java
@@ -1,9 +1,6 @@
package org.skywalking.apm.collector.worker.tools;
-import org.skywalking.apm.util.StringUtil;
-import org.skywalking.apm.collector.worker.Const;
-import org.skywalking.apm.collector.worker.segment.entity.Span;
-import org.skywalking.apm.collector.worker.segment.entity.tag.Tags;
+import org.skywalking.apm.network.proto.SpanObject;
/**
* @author pengys5
@@ -11,13 +8,11 @@ import org.skywalking.apm.collector.worker.segment.entity.tag.Tags;
public enum SpanPeersTools {
INSTANCE;
- public String getPeers(Span span) {
- if (StringUtil.isEmpty(Tags.PEERS.get(span))) {
- String host = Tags.PEER_HOST.get(span);
- int port = Tags.PEER_PORT.get(span);
- return Const.PEERS_FRONT_SPLIT + host + ":" + port + Const.PEERS_BEHIND_SPLIT;
+ public int getPeers(SpanObject span) {
+ if (span.getPeerId() == 0) {
+ return 0; //TODO exchange peer to peer id
} else {
- return Const.PEERS_FRONT_SPLIT + Tags.PEERS.get(span) + Const.PEERS_BEHIND_SPLIT;
+ return span.getPeerId();
}
}
}
diff --git a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tracedag/TraceDagGetWithTimeSlice.java b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tracedag/TraceDagGetWithTimeSlice.java
index 2f0c39b36..cecffbebe 100644
--- a/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tracedag/TraceDagGetWithTimeSlice.java
+++ b/apm-collector/apm-collector-worker/src/main/java/org/skywalking/apm/collector/worker/tracedag/TraceDagGetWithTimeSlice.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.tracedag;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.util.Arrays;
import java.util.Map;
@@ -34,6 +35,10 @@ public class TraceDagGetWithTimeSlice extends AbstractGet {
super(role, clusterContext, selfContext);
}
+ @Override protected Class extends JsonElement> responseClass() {
+ return JsonObject.class;
+ }
+
@Override
public void preStart() throws ProviderNotFoundException {
getClusterContext().findProvider(NodeCompLoad.WorkerRole.INSTANCE).create(this);
@@ -43,7 +48,7 @@ public class TraceDagGetWithTimeSlice extends AbstractGet {
}
@Override protected void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
if (!parameter.containsKey("startTime") || !parameter.containsKey("endTime") || !parameter.containsKey("timeSliceType")) {
throw new ArgumentsParseException("the request parameter must contains startTime,endTime,timeSliceType");
}
@@ -87,7 +92,7 @@ public class TraceDagGetWithTimeSlice extends AbstractGet {
JsonObject result = getBuilder().build(compResponse.get(Const.RESULT).getAsJsonArray(), nodeMappingResponse.get(Const.RESULT).getAsJsonArray(),
nodeRefResponse.get(Const.RESULT).getAsJsonArray(), resSumResponse.get(Const.RESULT).getAsJsonArray());
- response.add(Const.RESULT, result);
+ ((JsonObject)response).add(Const.RESULT, result);
}
private JsonObject getNewResponse() {
diff --git a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/io.grpc.BindableService b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/io.grpc.BindableService
new file mode 100644
index 000000000..094c4208f
--- /dev/null
+++ b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/io.grpc.BindableService
@@ -0,0 +1 @@
+org.skywalking.apm.collector.worker.grpcserver.TraceSegmentServiceImpl
\ No newline at end of file
diff --git a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractClusterWorkerProvider b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractClusterWorkerProvider
index 073dc4de3..40be71f4d 100644
--- a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractClusterWorkerProvider
+++ b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractClusterWorkerProvider
@@ -1,3 +1,4 @@
+org.skywalking.apm.collector.worker.grpcserver.GRPCAddressRegister$Factory
org.skywalking.apm.collector.worker.globaltrace.persistence.GlobalTraceAgg$Factory
org.skywalking.apm.collector.worker.noderef.persistence.NodeRefDayAgg$Factory
diff --git a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractLocalWorkerProvider b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractLocalWorkerProvider
index ca946f74e..9847c2360 100644
--- a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractLocalWorkerProvider
+++ b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.actor.AbstractLocalWorkerProvider
@@ -2,6 +2,7 @@ org.skywalking.apm.collector.worker.segment.analysis.SegmentAnalysis$Factory
org.skywalking.apm.collector.worker.segment.analysis.SegmentCostAnalysis$Factory
org.skywalking.apm.collector.worker.segment.analysis.SegmentExceptionAnalysis$Factory
+org.skywalking.apm.collector.worker.segment.SegmentReceiver$Factory
org.skywalking.apm.collector.worker.segment.persistence.SegmentSave$Factory
org.skywalking.apm.collector.worker.segment.persistence.SegmentCostSave$Factory
org.skywalking.apm.collector.worker.segment.persistence.SegmentExceptionSave$Factory
diff --git a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.config.ConfigProvider b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.config.ConfigProvider
index b266ffacd..aec02c4db 100644
--- a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.config.ConfigProvider
+++ b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.config.ConfigProvider
@@ -1,5 +1,6 @@
org.skywalking.apm.collector.cluster.ClusterConfigProvider
org.skywalking.apm.collector.worker.config.EsConfigProvider
org.skywalking.apm.collector.worker.config.HttpConfigProvider
+org.skywalking.apm.collector.worker.config.GRPCConfigProvider
org.skywalking.apm.collector.worker.config.CacheSizeConfigProvider
org.skywalking.apm.collector.worker.config.WorkerConfigProvider
\ No newline at end of file
diff --git a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractGetProvider b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractGetProvider
index c24779b4a..3fb933e71 100644
--- a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractGetProvider
+++ b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractGetProvider
@@ -2,4 +2,6 @@ org.skywalking.apm.collector.worker.noderef.NodeRefResSumGetGroupWithTimeSlice$F
org.skywalking.apm.collector.worker.segment.SegmentTopGet$Factory
org.skywalking.apm.collector.worker.globaltrace.GlobalTraceGetWithGlobalId$Factory
org.skywalking.apm.collector.worker.span.SpanGetWithId$Factory
-org.skywalking.apm.collector.worker.tracedag.TraceDagGetWithTimeSlice$Factory
\ No newline at end of file
+org.skywalking.apm.collector.worker.tracedag.TraceDagGetWithTimeSlice$Factory
+
+org.skywalking.apm.collector.worker.grpcserver.GRPCAddressGet$Factory
\ No newline at end of file
diff --git a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractStreamPostProvider b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractStreamPostProvider
index 6d8279b06..e69de29bb 100644
--- a/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractStreamPostProvider
+++ b/apm-collector/apm-collector-worker/src/main/resources/META-INF/services/org.skywalking.apm.collector.worker.httpserver.AbstractStreamPostProvider
@@ -1 +0,0 @@
-org.skywalking.apm.collector.worker.segment.SegmentPost$Factory
\ No newline at end of file
diff --git a/apm-collector/apm-collector-worker/src/main/resources/collector.config b/apm-collector/apm-collector-worker/src/main/resources/collector.config
index c007f3ef7..abdc725fd 100644
--- a/apm-collector/apm-collector-worker/src/main/resources/collector.config
+++ b/apm-collector/apm-collector-worker/src/main/resources/collector.config
@@ -6,7 +6,7 @@ cluster.current.port = 11800
# In this version, all members have same roles, and everyone of them is listening the status of others.
#The routers do not send message to nodes, which is unreachable, caused by network trouble, jvm crash or any other reasons.
-cluster.current.roles=WorkersListener
+cluster.current.roles=WorkersListener,RPCAddressListener
#Initial contact points of the cluster, e.g. seed_nodes = 127.0.0.1:11800, 127.0.0.1:11801.
#The nodes to join automatically at startup.
@@ -27,7 +27,7 @@ es.cluster.nodes=127.0.0.1:9300
#auto: create index when it doesn't exist.
# forced: delete and create.
# manual: do nothing.
-es.index.initialize.mode=auto
+es.index.initialize.mode=forced
# Config of shards or replicas in Elasticsearch.
es.index.shards.number=2
es.index.replicas.number=0
@@ -39,6 +39,10 @@ http.port=12800
# Web context path
http.contextPath=/
+# GRPC services
+grpc.hostname=127.0.0.1
+grpc.port=22800
+
# Cache size of analysis worker. The value determines whether sending to next worker and clear, or not.
cache.analysis.size=1024
# Cache size of persistence worker. The value determines whether save data and clear, or not.
diff --git a/apm-collector/apm-collector-worker/src/main/resources/log4j2.xml b/apm-collector/apm-collector-worker/src/main/resources/log4j2.xml
index 78e4e5a4a..afee48333 100644
--- a/apm-collector/apm-collector-worker/src/main/resources/log4j2.xml
+++ b/apm-collector/apm-collector-worker/src/main/resources/log4j2.xml
@@ -14,16 +14,19 @@
+
+
+
-
-
+
+
-
-
+
+
-
+
diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/PostWithHttpServletTestCase.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/PostWithHttpServletTestCase.java
deleted file mode 100644
index ac28e1925..000000000
--- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/PostWithHttpServletTestCase.java
+++ /dev/null
@@ -1,89 +0,0 @@
-package org.skywalking.apm.collector.worker.httpserver;
-
-import java.io.BufferedReader;
-import java.io.PrintWriter;
-import java.io.StringReader;
-import javax.servlet.http.HttpServletRequest;
-import javax.servlet.http.HttpServletResponse;
-import org.junit.Assert;
-import org.junit.Before;
-import org.junit.Test;
-import org.mockito.invocation.InvocationOnMock;
-import org.mockito.stubbing.Answer;
-import org.skywalking.apm.collector.actor.LocalSyncWorkerRef;
-import org.skywalking.apm.collector.actor.WorkerInvokeException;
-import org.skywalking.apm.collector.worker.segment.entity.Segment;
-
-import static org.mockito.Matchers.anyInt;
-import static org.mockito.Mockito.any;
-import static org.mockito.Mockito.anyString;
-import static org.mockito.Mockito.doAnswer;
-import static org.mockito.Mockito.doThrow;
-import static org.mockito.Mockito.mock;
-import static org.mockito.Mockito.when;
-
-/**
- * @author pengys5
- */
-public class PostWithHttpServletTestCase {
-
- private LocalSyncWorkerRef workerRef;
- private AbstractPost.PostWithHttpServlet servlet;
- private HttpServletRequest request;
- private HttpServletResponse response;
- private PrintWriter writer;
-
- @Before
- public void init() throws Exception {
- workerRef = mock(LocalSyncWorkerRef.class);
- servlet = new AbstractPost.PostWithHttpServlet(workerRef);
-
- request = mock(HttpServletRequest.class);
- response = mock(HttpServletResponse.class);
-
- writer = mock(PrintWriter.class);
- when(response.getWriter()).thenReturn(writer);
- }
-
- @Test
- public void testDoPost() throws Exception {
-
- doAnswer(new Answer() {
- @Override
- public Object answer(InvocationOnMock invocation) throws Throwable {
- Integer status = (Integer)invocation.getArguments()[0];
- Assert.assertEquals(new Integer(200), status);
- return null;
- }
- }).when(response).setStatus(anyInt());
-
- doAnswer(new Answer() {
- @Override
- public Object answer(InvocationOnMock invocation) throws Throwable {
- Segment segment = (Segment)invocation.getArguments()[0];
- Assert.assertEquals("TestTest2", segment.getTraceSegmentId());
- return null;
- }
- }).when(workerRef).tell(any(Segment.class));
-
- BufferedReader bufferedReader = new BufferedReader(new StringReader("[{\"ts\":\"TestTest2\"}]"));
-
- when(request.getReader()).thenReturn(bufferedReader);
-
- servlet.doPost(request, response);
- }
-
- @Test
- public void testDoPostError() throws Exception {
- doAnswer(new Answer() {
- @Override
- public Object answer(InvocationOnMock invocation) throws Throwable {
- Integer status = (Integer)invocation.getArguments()[0];
- Assert.assertEquals(new Integer(500), status);
- return null;
- }
- }).when(response).setStatus(anyInt());
- doThrow(new WorkerInvokeException("")).when(workerRef).tell(anyString());
- servlet.doPost(request, response);
- }
-}
diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractGet.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractGet.java
index c9ae71754..d7104a1f2 100644
--- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractGet.java
+++ b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractGet.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.httpserver;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.util.Map;
import org.skywalking.apm.collector.actor.ClusterWorkerContext;
@@ -19,13 +20,17 @@ public class TestAbstractGet extends AbstractGet {
super(role, clusterContext, selfContext);
}
+ @Override protected Class extends JsonElement> responseClass() {
+ return JsonObject.class;
+ }
+
@Override
public void preStart() throws ProviderNotFoundException {
super.preStart();
}
@Override protected void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
}
diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractPost.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractPost.java
index 17b178ea7..9d1d36da8 100644
--- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractPost.java
+++ b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractPost.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.httpserver;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.util.Map;
import org.skywalking.apm.collector.actor.ClusterWorkerContext;
@@ -19,13 +20,17 @@ public class TestAbstractPost extends AbstractPost {
super(role, clusterContext, selfContext);
}
+ @Override protected Class extends JsonElement> responseClass() {
+ return JsonObject.class;
+ }
+
@Override
public void preStart() throws ProviderNotFoundException {
super.preStart();
}
@Override protected void onReceive(Map parameter,
- JsonObject response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
+ JsonElement response) throws ArgumentsParseException, WorkerInvokeException, WorkerNotFoundException {
}
diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractStreamPost.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractStreamPost.java
index 58352d2a3..3311a3ed6 100644
--- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractStreamPost.java
+++ b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/httpserver/TestAbstractStreamPost.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.worker.httpserver;
+import com.google.gson.JsonElement;
import com.google.gson.JsonObject;
import java.io.BufferedReader;
import org.skywalking.apm.collector.actor.ClusterWorkerContext;
@@ -19,6 +20,10 @@ public class TestAbstractStreamPost extends AbstractStreamPost {
super(role, clusterContext, selfContext);
}
+ @Override protected Class extends JsonElement> responseClass() {
+ return JsonObject.class;
+ }
+
@Override
public void preStart() throws ProviderNotFoundException {
super.preStart();
diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentRealPost.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentRealPost.java
index 0abb2727b..071264abe 100644
--- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentRealPost.java
+++ b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentRealPost.java
@@ -1,32 +1,50 @@
package org.skywalking.apm.collector.worker.segment;
+import io.grpc.ManagedChannel;
+import io.grpc.ManagedChannelBuilder;
+import io.grpc.stub.StreamObserver;
+import java.util.List;
+import org.apache.logging.log4j.LogManager;
+import org.apache.logging.log4j.Logger;
import org.skywalking.apm.collector.worker.segment.mock.SegmentMock;
+import org.skywalking.apm.network.proto.Downstream;
+import org.skywalking.apm.network.proto.TraceSegmentServiceGrpc;
+import org.skywalking.apm.network.proto.UpstreamSegment;
/**
* @author pengys5
*/
public class SegmentRealPost {
+ private static Logger logger = LogManager.getFormatterLogger(SegmentRealPost.class);
+
public static void main(String[] args) throws Exception {
- SegmentMock mock = new SegmentMock();
-// String cacheServiceExceptionSegmentAsString = mock.mockCacheServiceExceptionSegmentAsString();
-// HttpClientTools.INSTANCE.post("http://localhost:7001/segments", cacheServiceExceptionSegmentAsString);
-//
-// String portalServiceExceptionSegmentAsString = mock.mockPortalServiceExceptionSegmentAsString();
-// HttpClientTools.INSTANCE.post("http://localhost:7001/segments", portalServiceExceptionSegmentAsString);
+ ManagedChannel channel = ManagedChannelBuilder.forAddress("localhost", 22800)
+ .usePlaintext(true)
+ .build();
- String cacheServiceSegmentAsString = mock.mockCacheServiceSegmentAsString();
- System.out.println(cacheServiceSegmentAsString);
- HttpClientTools.INSTANCE.post("http://localhost:12800/segments", cacheServiceSegmentAsString);
+ TraceSegmentServiceGrpc.TraceSegmentServiceStub stub = TraceSegmentServiceGrpc.newStub(channel);
+ StreamObserver observer = stub.collect(new StreamObserver() {
+ @Override public void onNext(Downstream downstream) {
- String persistenceServiceSegmentAsString = mock.mockPersistenceServiceSegmentAsString();
- HttpClientTools.INSTANCE.post("http://localhost:12800/segments", persistenceServiceSegmentAsString);
+ }
- String portalServiceSegmentAsString = mock.mockPortalServiceSegmentAsString();
- HttpClientTools.INSTANCE.post("http://localhost:12800/segments", portalServiceSegmentAsString);
+ @Override public void onError(Throwable throwable) {
-// String specialSegmentAsString = mock.mockSpecialSegmentAsString();
-// HttpClientTools.INSTANCE.post("http://localhost:7001/segments", specialSegmentAsString);
+ }
+ @Override public void onCompleted() {
+
+ }
+ });
+
+ List upstreamSegmentList = SegmentMock.mockPortalServiceSegment();
+ logger.debug("upstreamSegmentList size: %s", upstreamSegmentList.size());
+ upstreamSegmentList.forEach(upstreamSegment -> {
+ observer.onNext(upstreamSegment);
+ });
+ observer.onCompleted();
+
+ Thread.sleep(2000);
}
}
diff --git a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentPostTestCase.java b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentReceiverTestCase.java
similarity index 94%
rename from apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentPostTestCase.java
rename to apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentReceiverTestCase.java
index 644f4b976..3f82c26fc 100644
--- a/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentPostTestCase.java
+++ b/apm-collector/apm-collector-worker/src/test/java/org/skywalking/apm/collector/worker/segment/SegmentReceiverTestCase.java
@@ -58,12 +58,12 @@ import static org.powermock.api.mockito.PowerMockito.mock;
@RunWith(PowerMockRunner.class)
@PrepareForTest({LocalWorkerContext.class, WorkerRef.class})
@PowerMockIgnore({"javax.management.*"})
-public class SegmentPostTestCase {
+public class SegmentReceiverTestCase {
- private Logger logger = LogManager.getFormatterLogger(SegmentPostTestCase.class);
+ private Logger logger = LogManager.getFormatterLogger(SegmentReceiverTestCase.class);
private SegmentMock segmentMock;
- private SegmentPost segmentPost;
+ private SegmentReceiver segmentReceiver;
private LocalWorkerContext localWorkerContext;
private ClusterWorkerContext clusterWorkerContext;
@@ -76,7 +76,7 @@ public class SegmentPostTestCase {
clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
localWorkerContext = new LocalWorkerContext();
- segmentPost = new SegmentPost(SegmentPost.WorkerRole.INSTANCE, clusterWorkerContext, localWorkerContext);
+ segmentReceiver = new SegmentReceiver(SegmentReceiver.WorkerRole.INSTANCE, clusterWorkerContext, localWorkerContext);
initNodeNodeMappingAnalysis();
initNodeCompAnalysis();
@@ -89,16 +89,15 @@ public class SegmentPostTestCase {
@Test
public void testRole() {
- Assert.assertEquals(SegmentPost.class.getSimpleName(), SegmentPost.WorkerRole.INSTANCE.roleName());
- Assert.assertEquals(RollingSelector.class.getSimpleName(), SegmentPost.WorkerRole.INSTANCE.workerSelector().getClass().getSimpleName());
+ Assert.assertEquals(SegmentReceiver.class.getSimpleName(), SegmentReceiver.WorkerRole.INSTANCE.roleName());
+ Assert.assertEquals(RollingSelector.class.getSimpleName(), SegmentReceiver.WorkerRole.INSTANCE.workerSelector().getClass().getSimpleName());
}
@Test
public void testFactory() {
- SegmentPost.Factory factory = new SegmentPost.Factory();
- Assert.assertEquals(SegmentPost.class.getSimpleName(), factory.role().roleName());
- Assert.assertEquals(SegmentPost.class.getSimpleName(), factory.workerInstance(null).getClass().getSimpleName());
- Assert.assertEquals("/segments", factory.servletPath());
+ SegmentReceiver.Factory factory = new SegmentReceiver.Factory();
+ Assert.assertEquals(SegmentReceiver.class.getSimpleName(), factory.role().roleName());
+ Assert.assertEquals(SegmentReceiver.class.getSimpleName(), factory.workerInstance(null).getClass().getSimpleName());
}
@Test
@@ -141,7 +140,7 @@ public class SegmentPostTestCase {
ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(Role.class);
- segmentPost.preStart();
+ segmentReceiver.preStart();
verify(clusterWorkerContext, times(17)).findProvider(argumentCaptor.capture());
Assert.assertEquals(GlobalTraceAnalysis.Role.INSTANCE.roleName(), argumentCaptor.getAllValues().get(0).roleName());
@@ -178,7 +177,7 @@ public class SegmentPostTestCase {
JsonObject response = new JsonObject();
BufferedReader reader = new BufferedReader(new StringReader(jsonStr.length() + " " + jsonStr));
- segmentPost.onReceive(reader, response);
+// segmentReceiver.onReceive(reader, response);
}
private SegmentSaveAnswer segmentSaveAnswer_1;
@@ -304,7 +303,7 @@ public class SegmentPostTestCase {
public void testOnReceive() throws Exception {
String cacheServiceSegmentAsString = segmentMock.mockCacheServiceSegmentAsString();
- segmentPost.onReceive(new BufferedReader(new StringReader(cacheServiceSegmentAsString)), new JsonObject());
+// segmentReceiver.onReceive(new BufferedReader(new StringReader(cacheServiceSegmentAsString)), new JsonObject());
Assert.assertEquals(DateTools.changeToUTCSlice(201703310915L), segmentSaveAnswer_1.minute);
Assert.assertEquals(DateTools.changeToUTCSlice(201703310900L), segmentSaveAnswer_1.hour);
@@ -341,11 +340,11 @@ public class SegmentPostTestCase {
public class SegmentOtherAnswer implements Answer