From f377cfb82488fdbaf39ea9aa3f1d02900ac330bf Mon Sep 17 00:00:00 2001
From: pengys5 <8082209@qq.com>
Date: Fri, 7 Apr 2017 11:22:26 +0800
Subject: [PATCH] agg save mock finish
---
.../skywalking-collector-worker/pom.xml | 5 +
.../collector/worker/WorkerConfig.java | 122 ++++++------------
.../worker/globaltrace/entity/TreeNode.java | 23 ----
.../persistence/GlobalTraceAgg.java | 4 +-
.../persistence/GlobalTraceSave.java | 8 +-
.../worker/httpserver/HttpServer.java | 5 -
.../worker/node/persistence/NodeCompAgg.java | 8 +-
.../noderef/NodeRefGetWithTimeSlice.java | 94 --------------
.../NodeRefResSumGetWithTimeSlice.java | 94 --------------
.../noderef/persistence/NodeRefDayAgg.java | 4 +-
.../noderef/persistence/NodeRefHourAgg.java | 4 +-
.../noderef/persistence/NodeRefMinuteAgg.java | 4 +-
.../persistence/NodeRefResSumDayAgg.java | 4 +-
.../persistence/NodeRefResSumHourAgg.java | 4 +-
.../persistence/NodeRefResSumMinuteAgg.java | 4 +-
.../collector/worker/segment/SegmentPost.java | 3 +-
.../collector/worker/TimeSliceTestCase.java | 24 ++++
.../persistence/GlobalTraceAggTestCase.java | 85 ++++++++++++
.../persistence/GlobalTraceSaveTestCase.java | 51 ++++++++
.../worker/mock/MergeDataAnswer.java | 27 ++++
.../worker/mock/MetricDataAnswer.java | 25 ++++
.../node/persistence/NodeCompAggTestCase.java | 88 +++++++++++++
.../NodeMappingDayAggTestCase.java | 40 +++---
.../NodeMappingHourAggTestCase.java | 35 ++---
.../NodeMappingMinuteAggTestCase.java | 35 ++---
.../persistence/NodeRefDayAggTestCase.java | 86 ++++++++++++
.../persistence/NodeRefHourAggTestCase.java | 86 ++++++++++++
.../persistence/NodeRefMinuteAggTestCase.java | 86 ++++++++++++
.../NodeRefResSumDayAggTestCase.java | 86 ++++++++++++
.../NodeRefResSumHourAggTestCase.java | 86 ++++++++++++
.../NodeRefResSumMinuteAggTestCase.java | 86 ++++++++++++
.../worker/segment/SegmentPostTestCase.java | 105 +++++++++++++--
.../worker/tools/MergeDataAggTools.java | 22 ++++
.../worker/tools/MetricDataAggTools.java | 23 ++++
.../worker/tools/RecordDataAggTools.java | 22 ++++
35 files changed, 1085 insertions(+), 403 deletions(-)
delete mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/entity/TreeNode.java
delete mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefGetWithTimeSlice.java
delete mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefResSumGetWithTimeSlice.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/TimeSliceTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSaveTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/mock/MergeDataAnswer.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/mock/MetricDataAnswer.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/node/persistence/NodeCompAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefDayAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefHourAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefMinuteAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumDayAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumHourAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumMinuteAggTestCase.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/tools/MergeDataAggTools.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/tools/MetricDataAggTools.java
create mode 100644 skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/tools/RecordDataAggTools.java
diff --git a/skywalking-collector/skywalking-collector-worker/pom.xml b/skywalking-collector/skywalking-collector-worker/pom.xml
index 65c0e4d5b..b39688634 100644
--- a/skywalking-collector/skywalking-collector-worker/pom.xml
+++ b/skywalking-collector/skywalking-collector-worker/pom.xml
@@ -48,5 +48,10 @@
${project.version}
test
+
+ org.jetbrains
+ annotations
+ RELEASE
+
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java
index 6e6b93b96..1e9e69e6e 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java
@@ -19,43 +19,9 @@ public class WorkerConfig extends ClusterConfig {
}
}
- public static class Worker {
- public static class TraceSegmentReceiver {
- public static int Num = 10;
- }
-
- public static class DAGNodeReceiver {
- public static int Num = 10;
- }
-
- public static class NodeInstanceReceiver {
- public static int Num = 10;
- }
-
- public static class ResponseCostReceiver {
- public static int Num = 10;
- }
-
- public static class ResponseSummaryReceiver {
- public static int Num = 10;
- }
-
- public static class DAGNodeRefReceiver {
- public static int Num = 10;
- }
- }
-
public static class WorkerNum {
public static class Node {
- public static class NodeDayAgg {
- public static int Value = 10;
- }
-
- public static class NodeHourAgg {
- public static int Value = 10;
- }
-
- public static class NodeMinuteAgg {
+ public static class NodeCompAgg {
public static int Value = 10;
}
@@ -71,10 +37,52 @@ public class WorkerConfig extends ClusterConfig {
public static int Value = 10;
}
}
+
+ public static class NodeRef {
+ public static class NodeRefDayAgg {
+ public static int Value = 10;
+ }
+
+ public static class NodeRefHourAgg {
+ public static int Value = 10;
+ }
+
+ public static class NodeRefMinuteAgg {
+ public static int Value = 10;
+ }
+
+ public static class NodeRefResSumDayAgg {
+ public static int Value = 10;
+ }
+
+ public static class NodeRefResSumHourAgg {
+ public static int Value = 10;
+ }
+
+ public static class NodeRefResSumMinuteAgg {
+ public static int Value = 10;
+ }
+ }
+
+ public static class GlobalTrace {
+ public static class GlobalTraceAgg {
+ public static int Value = 10;
+ }
+ }
}
public static class Queue {
+ public static class GlobalTrace {
+ public static class GlobalTraceSave {
+ public static int Size = 1024;
+ }
+ }
+
public static class Segment {
+ public static class SegmentPost {
+ public static int Size = 1024;
+ }
+
public static class SegmentCostSave {
public static int Size = 1024;
}
@@ -160,52 +168,8 @@ public class WorkerConfig extends ClusterConfig {
}
}
-
- public static class Persistence {
- public static class DAGNodePersistence {
- public static int Size = 1024;
- }
-
- public static class NodeInstancePersistence {
- public static int Size = 1024;
- }
-
- public static class ResponseCostPersistence {
- public static int Size = 1024;
- }
-
- public static class ResponseSummaryPersistence {
- public static int Size = 1024;
- }
-
- public static class DAGNodeRefPersistence {
- public static int Size = 1024;
- }
- }
-
-
- public static class TraceSegmentRecordAnalysis {
- public static int Size = 1024;
- }
-
- public static class NodeInstanceAnalysis {
- public static int Size = 1024;
- }
-
- public static class DAGNodeAnalysis {
- public static int Size = 1024;
- }
-
- public static class ResponseCostAnalysis {
- public static int Size = 1024;
- }
-
public static class ResponseSummaryAnalysis {
public static int Size = 1024;
}
-
- public static class DAGNodeRefAnalysis {
- public static int Size = 1024;
- }
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/entity/TreeNode.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/entity/TreeNode.java
deleted file mode 100644
index 7f7f22cbc..000000000
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/entity/TreeNode.java
+++ /dev/null
@@ -1,23 +0,0 @@
-package com.a.eye.skywalking.collector.worker.globaltrace.entity;
-
-import java.util.ArrayList;
-import java.util.List;
-
-/**
- * @author pengys5
- */
-public class TreeNode {
-
- private String spanId;
-
- private List childNodes;
-
- public TreeNode(String spanId) {
- this.spanId = spanId;
- childNodes = new ArrayList<>();
- }
-
- public void addChild(TreeNode childNode) {
- childNodes.add(childNode);
- }
-}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAgg.java
index 0a6cf9079..5e246aa7a 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAgg.java
@@ -15,7 +15,7 @@ public class GlobalTraceAgg extends AbstractClusterWorker {
private Logger logger = LogManager.getFormatterLogger(GlobalTraceAgg.class);
- private GlobalTraceAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ GlobalTraceAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -48,7 +48,7 @@ public class GlobalTraceAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.Worker.DAGNodeReceiver.Num;
+ return WorkerConfig.WorkerNum.GlobalTrace.GlobalTraceAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSave.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSave.java
index 6954e9b4d..3c1f12199 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSave.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSave.java
@@ -4,7 +4,7 @@ package com.a.eye.skywalking.collector.worker.globaltrace.persistence;
import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider;
import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
-import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
+import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector;
import com.a.eye.skywalking.collector.actor.selector.WorkerSelector;
import com.a.eye.skywalking.collector.worker.MergePersistenceMember;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
@@ -15,7 +15,7 @@ import com.a.eye.skywalking.collector.worker.globaltrace.GlobalTraceIndex;
*/
public class GlobalTraceSave extends MergePersistenceMember {
- private GlobalTraceSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ GlobalTraceSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -39,7 +39,7 @@ public class GlobalTraceSave extends MergePersistenceMember {
@Override
public int queueSize() {
- return WorkerConfig.Queue.TraceSegmentRecordAnalysis.Size;
+ return WorkerConfig.Queue.GlobalTrace.GlobalTraceSave.Size;
}
@Override
@@ -58,7 +58,7 @@ public class GlobalTraceSave extends MergePersistenceMember {
@Override
public WorkerSelector workerSelector() {
- return new RollingSelector();
+ return new HashCodeSelector();
}
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/httpserver/HttpServer.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/httpserver/HttpServer.java
index 481ac94ee..144a9ab69 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/httpserver/HttpServer.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/httpserver/HttpServer.java
@@ -25,11 +25,6 @@ public enum HttpServer {
ServletsCreator.INSTANCE.boot(servletContextHandler, clusterContext);
-// ServerConnector serverConnector = new ServerConnector(server);
-// serverConnector.setHost("127.0.0.1");
-// serverConnector.setPort(7001);
-// serverConnector.setIdleTimeout(5000);
-
server.setHandler(servletContextHandler);
server.start();
server.join();
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/node/persistence/NodeCompAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/node/persistence/NodeCompAgg.java
index fd2fd46a0..2149a2931 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/node/persistence/NodeCompAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/node/persistence/NodeCompAgg.java
@@ -11,7 +11,7 @@ import com.a.eye.skywalking.collector.worker.storage.RecordData;
*/
public class NodeCompAgg extends AbstractClusterWorker {
- public NodeCompAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ NodeCompAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -24,9 +24,7 @@ public class NodeCompAgg extends AbstractClusterWorker {
protected void onWork(Object message) throws Exception {
if (message instanceof RecordData) {
getSelfContext().lookup(NodeCompSave.Role.INSTANCE).tell(message);
- } else {
- throw new IllegalArgumentException("message instance must RecordData");
- }
+ } else throw new IllegalArgumentException("message instance must RecordData");
}
public static class Factory extends AbstractClusterWorkerProvider {
@@ -44,7 +42,7 @@ public class NodeCompAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.WorkerNum.Node.NodeDayAgg.Value;
+ return WorkerConfig.WorkerNum.Node.NodeCompAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefGetWithTimeSlice.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefGetWithTimeSlice.java
deleted file mode 100644
index d02326b7a..000000000
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefGetWithTimeSlice.java
+++ /dev/null
@@ -1,94 +0,0 @@
-package com.a.eye.skywalking.collector.worker.noderef;
-
-import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
-import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
-import com.a.eye.skywalking.collector.actor.ProviderNotFoundException;
-import com.a.eye.skywalking.collector.actor.Role;
-import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
-import com.a.eye.skywalking.collector.actor.selector.WorkerSelector;
-import com.a.eye.skywalking.collector.worker.httpserver.AbstractGet;
-import com.a.eye.skywalking.collector.worker.httpserver.AbstractGetProvider;
-import com.a.eye.skywalking.collector.worker.noderef.persistence.NodeRefSearchWithTimeSlice;
-import com.a.eye.skywalking.collector.worker.tools.ParameterTools;
-import com.google.gson.JsonObject;
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-
-import java.util.Arrays;
-import java.util.Map;
-
-/**
- * @author pengys5
- */
-public class NodeRefGetWithTimeSlice extends AbstractGet {
-
- private Logger logger = LogManager.getFormatterLogger(NodeRefGetWithTimeSlice.class);
-
- private NodeRefGetWithTimeSlice(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- super(role, clusterContext, selfContext);
- }
-
- @Override
- public void preStart() throws ProviderNotFoundException {
- getClusterContext().findProvider(NodeRefSearchWithTimeSlice.WorkerRole.INSTANCE).create(this);
- }
-
- @Override
- protected void onSearch(Map request, JsonObject response) throws Exception {
- if (!request.containsKey("startTime") || !request.containsKey("endTime") || !request.containsKey("timeSliceType")) {
- throw new IllegalArgumentException("the request parameter must contains startTime,endTime,timeSliceType");
- }
- logger.debug("startTime: %s, endTime: %s, timeSliceType: %s", Arrays.toString(request.get("startTime")),
- Arrays.toString(request.get("endTime")), Arrays.toString(request.get("timeSliceType")));
-
- long startTime;
- try {
- startTime = Long.valueOf(ParameterTools.INSTANCE.toString(request, "startTime"));
- } catch (NumberFormatException e) {
- throw new IllegalArgumentException("the request parameter startTime must numeric with long type");
- }
-
- long endTime;
- try {
- endTime = Long.valueOf(ParameterTools.INSTANCE.toString(request, "endTime"));
- } catch (NumberFormatException e) {
- throw new IllegalArgumentException("the request parameter endTime must numeric with long type");
- }
-
- NodeRefSearchWithTimeSlice.RequestEntity requestEntity;
- requestEntity = new NodeRefSearchWithTimeSlice.RequestEntity(ParameterTools.INSTANCE.toString(request, "timeSliceType"), startTime, endTime);
- getSelfContext().lookup(NodeRefSearchWithTimeSlice.WorkerRole.INSTANCE).ask(requestEntity, response);
- }
-
- public static class Factory extends AbstractGetProvider {
-
- @Override
- public Role role() {
- return WorkerRole.INSTANCE;
- }
-
- @Override
- public NodeRefGetWithTimeSlice workerInstance(ClusterWorkerContext clusterContext) {
- return new NodeRefGetWithTimeSlice(role(), clusterContext, new LocalWorkerContext());
- }
-
- @Override
- public String servletPath() {
- return "/nodeRef/timeSlice";
- }
- }
-
- public enum WorkerRole implements Role {
- INSTANCE;
-
- @Override
- public String roleName() {
- return NodeRefGetWithTimeSlice.class.getSimpleName();
- }
-
- @Override
- public WorkerSelector workerSelector() {
- return new RollingSelector();
- }
- }
-}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefResSumGetWithTimeSlice.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefResSumGetWithTimeSlice.java
deleted file mode 100644
index 758ddda51..000000000
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/NodeRefResSumGetWithTimeSlice.java
+++ /dev/null
@@ -1,94 +0,0 @@
-package com.a.eye.skywalking.collector.worker.noderef;
-
-import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
-import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
-import com.a.eye.skywalking.collector.actor.ProviderNotFoundException;
-import com.a.eye.skywalking.collector.actor.Role;
-import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
-import com.a.eye.skywalking.collector.actor.selector.WorkerSelector;
-import com.a.eye.skywalking.collector.worker.httpserver.AbstractGet;
-import com.a.eye.skywalking.collector.worker.httpserver.AbstractGetProvider;
-import com.a.eye.skywalking.collector.worker.noderef.persistence.NodeRefResSumSearchWithTimeSlice;
-import com.a.eye.skywalking.collector.worker.tools.ParameterTools;
-import com.google.gson.JsonObject;
-import org.apache.logging.log4j.LogManager;
-import org.apache.logging.log4j.Logger;
-
-import java.util.Arrays;
-import java.util.Map;
-
-/**
- * @author pengys5
- */
-public class NodeRefResSumGetWithTimeSlice extends AbstractGet {
-
- private Logger logger = LogManager.getFormatterLogger(NodeRefResSumGetWithTimeSlice.class);
-
- private NodeRefResSumGetWithTimeSlice(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- super(role, clusterContext, selfContext);
- }
-
- @Override
- public void preStart() throws ProviderNotFoundException {
- getClusterContext().findProvider(NodeRefResSumSearchWithTimeSlice.WorkerRole.INSTANCE).create(this);
- }
-
- @Override
- protected void onSearch(Map request, JsonObject response) throws Exception {
- if (!request.containsKey("startTime") || !request.containsKey("endTime") || !request.containsKey("timeSliceType")) {
- throw new IllegalArgumentException("the request parameter must contains startTime,endTime,timeSliceType");
- }
- logger.debug("startTime: %s, endTime: %s, timeSliceType: %s", Arrays.toString(request.get("startTime")),
- Arrays.toString(request.get("endTime")), Arrays.toString(request.get("timeSliceType")));
-
- long startTime;
- try {
- startTime = Long.valueOf(ParameterTools.INSTANCE.toString(request, "startTime"));
- } catch (NumberFormatException e) {
- throw new IllegalArgumentException("the request parameter startTime must numeric with long type");
- }
-
- long endTime;
- try {
- endTime = Long.valueOf(ParameterTools.INSTANCE.toString(request, "endTime"));
- } catch (NumberFormatException e) {
- throw new IllegalArgumentException("the request parameter endTime must numeric with long type");
- }
-
- NodeRefResSumSearchWithTimeSlice.RequestEntity requestEntity;
- requestEntity = new NodeRefResSumSearchWithTimeSlice.RequestEntity(ParameterTools.INSTANCE.toString(request, "timeSliceType"), startTime, endTime);
- getSelfContext().lookup(NodeRefResSumSearchWithTimeSlice.WorkerRole.INSTANCE).ask(requestEntity, response);
- }
-
- public static class Factory extends AbstractGetProvider {
-
- @Override
- public Role role() {
- return WorkerRole.INSTANCE;
- }
-
- @Override
- public NodeRefResSumGetWithTimeSlice workerInstance(ClusterWorkerContext clusterContext) {
- return new NodeRefResSumGetWithTimeSlice(role(), clusterContext, new LocalWorkerContext());
- }
-
- @Override
- public String servletPath() {
- return "/nodeRef/resSum/timeSlice";
- }
- }
-
- public enum WorkerRole implements Role {
- INSTANCE;
-
- @Override
- public String roleName() {
- return NodeRefResSumGetWithTimeSlice.class.getSimpleName();
- }
-
- @Override
- public WorkerSelector workerSelector() {
- return new RollingSelector();
- }
- }
-}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefDayAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefDayAgg.java
index ce44bfcec..e045565ac 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefDayAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefDayAgg.java
@@ -15,7 +15,7 @@ public class NodeRefDayAgg extends AbstractClusterWorker {
private Logger logger = LogManager.getFormatterLogger(NodeRefDayAgg.class);
- public NodeRefDayAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ NodeRefDayAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -48,7 +48,7 @@ public class NodeRefDayAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.Worker.DAGNodeRefReceiver.Num;
+ return WorkerConfig.WorkerNum.NodeRef.NodeRefDayAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefHourAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefHourAgg.java
index 14c6890f7..78974e82b 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefHourAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefHourAgg.java
@@ -15,7 +15,7 @@ public class NodeRefHourAgg extends AbstractClusterWorker {
private Logger logger = LogManager.getFormatterLogger(NodeRefHourAgg.class);
- public NodeRefHourAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ NodeRefHourAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -48,7 +48,7 @@ public class NodeRefHourAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.Worker.DAGNodeRefReceiver.Num;
+ return WorkerConfig.WorkerNum.NodeRef.NodeRefHourAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefMinuteAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefMinuteAgg.java
index 5d1b4f124..b961e2ce7 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefMinuteAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefMinuteAgg.java
@@ -15,7 +15,7 @@ public class NodeRefMinuteAgg extends AbstractClusterWorker {
private Logger logger = LogManager.getFormatterLogger(NodeRefMinuteAgg.class);
- public NodeRefMinuteAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ NodeRefMinuteAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -48,7 +48,7 @@ public class NodeRefMinuteAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.Worker.DAGNodeRefReceiver.Num;
+ return WorkerConfig.WorkerNum.NodeRef.NodeRefMinuteAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumDayAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumDayAgg.java
index a03ea1caf..61fa1ad5c 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumDayAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumDayAgg.java
@@ -15,7 +15,7 @@ public class NodeRefResSumDayAgg extends AbstractClusterWorker {
private Logger logger = LogManager.getFormatterLogger(NodeRefResSumDayAgg.class);
- private NodeRefResSumDayAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ NodeRefResSumDayAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -48,7 +48,7 @@ public class NodeRefResSumDayAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.Worker.ResponseSummaryReceiver.Num;
+ return WorkerConfig.WorkerNum.NodeRef.NodeRefResSumDayAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumHourAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumHourAgg.java
index 5e4ef0c3e..39435bbd5 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumHourAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumHourAgg.java
@@ -15,7 +15,7 @@ public class NodeRefResSumHourAgg extends AbstractClusterWorker {
private Logger logger = LogManager.getFormatterLogger(NodeRefResSumHourAgg.class);
- private NodeRefResSumHourAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ NodeRefResSumHourAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -48,7 +48,7 @@ public class NodeRefResSumHourAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.Worker.ResponseSummaryReceiver.Num;
+ return WorkerConfig.WorkerNum.NodeRef.NodeRefResSumHourAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumMinuteAgg.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumMinuteAgg.java
index 1981f643a..7afb59d8b 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumMinuteAgg.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/noderef/persistence/NodeRefResSumMinuteAgg.java
@@ -15,7 +15,7 @@ public class NodeRefResSumMinuteAgg extends AbstractClusterWorker {
private Logger logger = LogManager.getFormatterLogger(NodeRefResSumMinuteAgg.class);
- private NodeRefResSumMinuteAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ NodeRefResSumMinuteAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
@@ -48,7 +48,7 @@ public class NodeRefResSumMinuteAgg extends AbstractClusterWorker {
@Override
public int workerNum() {
- return WorkerConfig.Worker.ResponseSummaryReceiver.Num;
+ return WorkerConfig.WorkerNum.NodeRef.NodeRefResSumMinuteAgg.Value;
}
}
diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/segment/SegmentPost.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/segment/SegmentPost.java
index da1317929..af96de4b2 100644
--- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/segment/SegmentPost.java
+++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/segment/SegmentPost.java
@@ -7,6 +7,7 @@ import com.a.eye.skywalking.collector.actor.ProviderNotFoundException;
import com.a.eye.skywalking.collector.actor.Role;
import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
import com.a.eye.skywalking.collector.actor.selector.WorkerSelector;
+import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.globaltrace.analysis.GlobalTraceAnalysis;
import com.a.eye.skywalking.collector.worker.httpserver.AbstractPost;
import com.a.eye.skywalking.collector.worker.httpserver.AbstractPostProvider;
@@ -134,7 +135,7 @@ public class SegmentPost extends AbstractPost {
@Override
public int queueSize() {
- return 128;
+ return WorkerConfig.Queue.Segment.SegmentPost.Size;
}
@Override
diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/TimeSliceTestCase.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/TimeSliceTestCase.java
new file mode 100644
index 000000000..289040b92
--- /dev/null
+++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/TimeSliceTestCase.java
@@ -0,0 +1,24 @@
+package com.a.eye.skywalking.collector.worker;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+/**
+ * @author pengys5
+ */
+public class TimeSliceTestCase {
+
+ @Test
+ public void test() {
+ TestTimeSlice timeSlice = new TestTimeSlice("A", 10L, 20L);
+ Assert.assertEquals("A", timeSlice.getSliceType());
+ Assert.assertEquals(10L, timeSlice.getStartTime());
+ Assert.assertEquals(20L, timeSlice.getEndTime());
+ }
+
+ class TestTimeSlice extends TimeSlice {
+ public TestTimeSlice(String sliceType, long startTime, long endTime) {
+ super(sliceType, startTime, endTime);
+ }
+ }
+}
diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAggTestCase.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAggTestCase.java
new file mode 100644
index 000000000..b8f87262f
--- /dev/null
+++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceAggTestCase.java
@@ -0,0 +1,85 @@
+package com.a.eye.skywalking.collector.worker.globaltrace.persistence;
+
+import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
+import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
+import com.a.eye.skywalking.collector.actor.ProviderNotFoundException;
+import com.a.eye.skywalking.collector.actor.WorkerRefs;
+import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector;
+import com.a.eye.skywalking.collector.worker.WorkerConfig;
+import com.a.eye.skywalking.collector.worker.mock.MergeDataAnswer;
+import com.a.eye.skywalking.collector.worker.storage.RecordData;
+import com.a.eye.skywalking.collector.worker.tools.MergeDataAggTools;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+import org.junit.runner.RunWith;
+import org.mockito.ArgumentCaptor;
+import org.mockito.Mockito;
+import org.powermock.api.mockito.PowerMockito;
+import org.powermock.core.classloader.annotations.PowerMockIgnore;
+import org.powermock.core.classloader.annotations.PrepareForTest;
+import org.powermock.modules.junit4.PowerMockRunner;
+
+import static org.mockito.Mockito.*;
+
+/**
+ * @author pengys5
+ */
+@RunWith(PowerMockRunner.class)
+@PrepareForTest({LocalWorkerContext.class})
+@PowerMockIgnore({"javax.management.*"})
+public class GlobalTraceAggTestCase {
+
+ private GlobalTraceAgg agg;
+ private MergeDataAnswer mergeDataAnswer;
+ private ClusterWorkerContext clusterWorkerContext;
+
+ @Before
+ public void init() throws Exception {
+ clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
+
+ LocalWorkerContext localWorkerContext = PowerMockito.mock(LocalWorkerContext.class);
+ WorkerRefs workerRefs = mock(WorkerRefs.class);
+
+ mergeDataAnswer = new MergeDataAnswer();
+ doAnswer(mergeDataAnswer).when(workerRefs).tell(Mockito.any(RecordData.class));
+
+ when(localWorkerContext.lookup(GlobalTraceSave.Role.INSTANCE)).thenReturn(workerRefs);
+ agg = new GlobalTraceAgg(GlobalTraceAgg.Role.INSTANCE, clusterWorkerContext, localWorkerContext);
+ }
+
+ @Test
+ public void testRole() {
+ Assert.assertEquals(GlobalTraceAgg.class.getSimpleName(), GlobalTraceAgg.Role.INSTANCE.roleName());
+ Assert.assertEquals(HashCodeSelector.class.getSimpleName(), GlobalTraceAgg.Role.INSTANCE.workerSelector().getClass().getSimpleName());
+ }
+
+ @Test
+ public void testFactory() {
+ Assert.assertEquals(GlobalTraceAgg.class.getSimpleName(), GlobalTraceAgg.Factory.INSTANCE.role().roleName());
+ Assert.assertEquals(GlobalTraceAgg.class.getSimpleName(), GlobalTraceAgg.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
+
+ int testSize = 10;
+ WorkerConfig.WorkerNum.GlobalTrace.GlobalTraceAgg.Value = testSize;
+ Assert.assertEquals(testSize, GlobalTraceAgg.Factory.INSTANCE.workerNum());
+ }
+
+ @Test
+ public void testPreStart() throws ProviderNotFoundException {
+ when(clusterWorkerContext.findProvider(GlobalTraceSave.Role.INSTANCE)).thenReturn(GlobalTraceSave.Factory.INSTANCE);
+
+ ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(GlobalTraceSave.Role.class);
+ agg.preStart();
+ verify(clusterWorkerContext).findProvider(argumentCaptor.capture());
+ }
+
+ @Test
+ public void testOnWork() throws Exception {
+ MergeDataAggTools.INSTANCE.testOnWork(agg, mergeDataAnswer);
+ }
+
+ @Test
+ public void testOnWorkError() throws Exception {
+ agg.onWork(new Object());
+ }
+}
diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSaveTestCase.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSaveTestCase.java
new file mode 100644
index 000000000..2a56ae1fa
--- /dev/null
+++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/globaltrace/persistence/GlobalTraceSaveTestCase.java
@@ -0,0 +1,51 @@
+package com.a.eye.skywalking.collector.worker.globaltrace.persistence;
+
+import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
+import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
+import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector;
+import com.a.eye.skywalking.collector.worker.WorkerConfig;
+import com.a.eye.skywalking.collector.worker.globaltrace.GlobalTraceIndex;
+import org.junit.Assert;
+import org.junit.Before;
+import org.junit.Test;
+
+/**
+ * @author pengys5
+ */
+public class GlobalTraceSaveTestCase {
+
+ private GlobalTraceSave save;
+
+ @Before
+ public void init() {
+ ClusterWorkerContext cluster = new ClusterWorkerContext(null);
+ LocalWorkerContext local = new LocalWorkerContext();
+ save = new GlobalTraceSave(GlobalTraceSave.Role.INSTANCE, cluster, local);
+ }
+
+ @Test
+ public void testEsIndex() {
+ Assert.assertEquals(GlobalTraceIndex.Index, save.esIndex());
+ }
+
+ @Test
+ public void testEsType() {
+ Assert.assertEquals(GlobalTraceIndex.Type_Record, save.esType());
+ }
+
+ @Test
+ public void testRole() {
+ Assert.assertEquals(GlobalTraceSave.class.getSimpleName(), GlobalTraceSave.Role.INSTANCE.roleName());
+ Assert.assertEquals(HashCodeSelector.class.getSimpleName(), GlobalTraceSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
+ }
+
+ @Test
+ public void testFactory() {
+ Assert.assertEquals(GlobalTraceSave.class.getSimpleName(), GlobalTraceSave.Factory.INSTANCE.role().roleName());
+ Assert.assertEquals(GlobalTraceSave.class.getSimpleName(), GlobalTraceSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
+
+ int testSize = 10;
+ WorkerConfig.Queue.GlobalTrace.GlobalTraceSave.Size = testSize;
+ Assert.assertEquals(testSize, GlobalTraceSave.Factory.INSTANCE.queueSize());
+ }
+}
diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/mock/MergeDataAnswer.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/mock/MergeDataAnswer.java
new file mode 100644
index 000000000..20f72e643
--- /dev/null
+++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/mock/MergeDataAnswer.java
@@ -0,0 +1,27 @@
+package com.a.eye.skywalking.collector.worker.mock;
+
+import com.a.eye.skywalking.collector.worker.storage.MergeData;
+import org.mockito.invocation.InvocationOnMock;
+import org.mockito.stubbing.Answer;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * @author pengys5
+ */
+public class MergeDataAnswer implements Answer