mock persistence
This commit is contained in:
parent
75be151485
commit
7f2295b0a8
|
|
@ -44,7 +44,7 @@ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker {
|
|||
asyncWorker.allocateJob(new EndOfBatchCommand());
|
||||
}
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
asyncWorker.saveException(e);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -33,4 +33,8 @@ public abstract class AbstractWorker {
|
|||
final public static AbstractWorker noOwner() {
|
||||
return null;
|
||||
}
|
||||
|
||||
final protected void saveException(Exception e) {
|
||||
// e.printStackTrace();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -2,17 +2,13 @@ package com.a.eye.skywalking.collector.worker;
|
|||
|
||||
import com.a.eye.skywalking.collector.actor.*;
|
||||
import com.a.eye.skywalking.collector.queue.EndOfBatchCommand;
|
||||
import org.apache.logging.log4j.LogManager;
|
||||
import org.apache.logging.log4j.Logger;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AnalysisMember extends AbstractLocalAsyncWorker {
|
||||
|
||||
private Logger logger = LogManager.getFormatterLogger(AnalysisMember.class);
|
||||
|
||||
public AnalysisMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
AnalysisMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -20,7 +16,7 @@ public abstract class AnalysisMember extends AbstractLocalAsyncWorker {
|
|||
|
||||
@Override
|
||||
public void preStart() throws ProviderNotFoundException {
|
||||
|
||||
super.preStart();
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -31,7 +27,7 @@ public abstract class AnalysisMember extends AbstractLocalAsyncWorker {
|
|||
try {
|
||||
analyse(message);
|
||||
} catch (Exception e) {
|
||||
e.printStackTrace();
|
||||
saveException(e);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -26,18 +26,23 @@ public abstract class MergePersistenceMember extends PersistenceMember {
|
|||
|
||||
private Logger logger = LogManager.getFormatterLogger(MergePersistenceMember.class);
|
||||
|
||||
private MergePersistenceData persistenceData = new MergePersistenceData();
|
||||
private MergePersistenceData persistenceData;
|
||||
|
||||
protected MergePersistenceMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
persistenceData = new MergePersistenceData();
|
||||
}
|
||||
|
||||
private MergePersistenceData getPersistenceData() {
|
||||
return persistenceData;
|
||||
}
|
||||
|
||||
@Override
|
||||
final public void analyse(Object message) throws Exception {
|
||||
if (message instanceof MergeData) {
|
||||
MergeData mergeData = (MergeData) message;
|
||||
persistenceData.getElseCreate(mergeData.getId()).merge(mergeData);
|
||||
if (persistenceData.size() >= WorkerConfig.Persistence.Data.size) {
|
||||
getPersistenceData().getElseCreate(mergeData.getId()).merge(mergeData);
|
||||
if (getPersistenceData().size() >= WorkerConfig.Persistence.Data.size) {
|
||||
persistence();
|
||||
}
|
||||
} else {
|
||||
|
|
@ -50,13 +55,13 @@ public abstract class MergePersistenceMember extends PersistenceMember {
|
|||
for (MultiGetItemResponse itemResponse : multiGetResponse) {
|
||||
GetResponse response = itemResponse.getResponse();
|
||||
if (response != null && response.isExists()) {
|
||||
persistenceData.getElseCreate(response.getId()).merge(response.getSource());
|
||||
getPersistenceData().getElseCreate(response.getId()).merge(response.getSource());
|
||||
}
|
||||
}
|
||||
|
||||
boolean success = saveToEs();
|
||||
if (success) {
|
||||
persistenceData.clear();
|
||||
getPersistenceData().clear();
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -64,7 +69,7 @@ public abstract class MergePersistenceMember extends PersistenceMember {
|
|||
Client client = EsClient.INSTANCE.getClient();
|
||||
MultiGetRequestBuilder multiGetRequestBuilder = client.prepareMultiGet();
|
||||
|
||||
Iterator<Map.Entry<String, MergeData>> iterator = persistenceData.iterator();
|
||||
Iterator<Map.Entry<String, MergeData>> iterator = getPersistenceData().iterator();
|
||||
|
||||
while (iterator.hasNext()) {
|
||||
multiGetRequestBuilder.add(esIndex(), esType(), iterator.next().getKey());
|
||||
|
|
@ -76,9 +81,9 @@ public abstract class MergePersistenceMember extends PersistenceMember {
|
|||
private boolean saveToEs() {
|
||||
Client client = EsClient.INSTANCE.getClient();
|
||||
BulkRequestBuilder bulkRequest = client.prepareBulk();
|
||||
logger.debug("persistenceData size: %s", persistenceData.size());
|
||||
logger.debug("persistenceData size: %s", getPersistenceData().size());
|
||||
|
||||
Iterator<Map.Entry<String, MergeData>> iterator = persistenceData.iterator();
|
||||
Iterator<Map.Entry<String, MergeData>> iterator = getPersistenceData().iterator();
|
||||
while (iterator.hasNext()) {
|
||||
MergeData mergeData = iterator.next().getValue();
|
||||
bulkRequest.add(client.prepareIndex(esIndex(), esType(), mergeData.getId()).setSource(mergeData.toMap()));
|
||||
|
|
|
|||
|
|
@ -89,15 +89,7 @@ public class WorkerConfig extends ClusterConfig {
|
|||
}
|
||||
|
||||
public static class Node {
|
||||
public static class NodeDayAnalysis {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeHourAnalysis {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeMinuteAnalysis {
|
||||
public static class NodeCompAnalysis {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
|
|
@ -124,6 +116,48 @@ public class WorkerConfig extends ClusterConfig {
|
|||
public static class NodeMappingMinuteAnalysis {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeCompSave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeMappingDaySave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeMappingHourSave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeMappingMinuteSave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
}
|
||||
|
||||
public static class NodeRef {
|
||||
public static class NodeRefDaySave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeRefHourSave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeRefMinuteSave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeRefResSumDaySave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeRefResSumHourSave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
|
||||
public static class NodeRefResSumMinuteSave {
|
||||
public static int Size = 1024;
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
|
|
|
|||
|
|
@ -21,7 +21,7 @@ import java.util.List;
|
|||
*/
|
||||
public class GlobalTraceAnalysis extends MergeAnalysisMember {
|
||||
|
||||
private GlobalTraceAnalysis(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
GlobalTraceAnalysis(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -9,7 +9,7 @@ import com.a.eye.skywalking.collector.worker.segment.SegmentIndex;
|
|||
import com.a.eye.skywalking.collector.worker.segment.logic.Segment;
|
||||
import com.a.eye.skywalking.collector.worker.segment.logic.SegmentDeserialize;
|
||||
import com.a.eye.skywalking.collector.worker.segment.logic.SpanView;
|
||||
import com.a.eye.skywalking.collector.worker.storage.EsClient;
|
||||
import com.a.eye.skywalking.collector.worker.storage.GetResponseFromEs;
|
||||
import com.a.eye.skywalking.collector.worker.storage.MergeData;
|
||||
import com.a.eye.skywalking.collector.worker.tools.CollectionTools;
|
||||
import com.a.eye.skywalking.trace.Span;
|
||||
|
|
@ -18,7 +18,6 @@ import com.google.gson.Gson;
|
|||
import com.google.gson.JsonObject;
|
||||
import org.apache.logging.log4j.LogManager;
|
||||
import org.apache.logging.log4j.Logger;
|
||||
import org.elasticsearch.client.Client;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
|
|
@ -33,17 +32,17 @@ public class GlobalTraceSearchWithGlobalId extends AbstractLocalSyncWorker {
|
|||
|
||||
private Gson gson = new Gson();
|
||||
|
||||
public GlobalTraceSearchWithGlobalId(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
GlobalTraceSearchWithGlobalId(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void onWork(Object request, Object response) throws Exception {
|
||||
if (request instanceof String) {
|
||||
Client client = EsClient.INSTANCE.getClient();
|
||||
String globalId = (String) request;
|
||||
String globalTraceData = client.prepareGet(GlobalTraceIndex.Index, GlobalTraceIndex.Type_Record, globalId).get().getSourceAsString();
|
||||
String globalTraceData = GetResponseFromEs.INSTANCE.get(GlobalTraceIndex.Index, GlobalTraceIndex.Type_Record, globalId).getSourceAsString();
|
||||
JsonObject globalTraceObj = gson.fromJson(globalTraceData, JsonObject.class);
|
||||
logger.debug("globalTraceObj: %s", globalTraceObj);
|
||||
|
||||
String subSegIdsStr = globalTraceObj.get(GlobalTraceIndex.SubSegIds).getAsString();
|
||||
String[] subSegIds = subSegIdsStr.split(MergeData.Split);
|
||||
|
|
@ -51,7 +50,8 @@ public class GlobalTraceSearchWithGlobalId extends AbstractLocalSyncWorker {
|
|||
List<SpanView> spanViewList = new ArrayList<>();
|
||||
for (String subSegId : subSegIds) {
|
||||
logger.debug("subSegId: %s", subSegId);
|
||||
String segmentSource = client.prepareGet(SegmentIndex.Index, SegmentIndex.Type_Record, subSegId).get().getSourceAsString();
|
||||
String segmentSource = GetResponseFromEs.INSTANCE.get(SegmentIndex.Index, SegmentIndex.Type_Record, subSegId).getSourceAsString();
|
||||
logger.debug("segmentSource: %s", segmentSource);
|
||||
Segment segment = SegmentDeserialize.INSTANCE.deserializeFromES(segmentSource);
|
||||
String segmentId = segment.getTraceSegmentId();
|
||||
List<TraceSegmentRef> refsList = segment.getRefs();
|
||||
|
|
@ -62,19 +62,23 @@ public class GlobalTraceSearchWithGlobalId extends AbstractLocalSyncWorker {
|
|||
}
|
||||
}
|
||||
|
||||
SpanView rootSpan = findRoot(spanViewList);
|
||||
findChild(rootSpan, spanViewList, rootSpan.getStartTime());
|
||||
|
||||
List<SpanView> viewList = new ArrayList<>();
|
||||
viewList.add(rootSpan);
|
||||
|
||||
Gson gson = new Gson();
|
||||
String globalTraceStr = gson.toJson(viewList);
|
||||
JsonObject responseObj = (JsonObject) response;
|
||||
responseObj.addProperty("result", globalTraceStr);
|
||||
responseObj.addProperty("result", buildTree(spanViewList));
|
||||
}
|
||||
}
|
||||
|
||||
private String buildTree(List<SpanView> spanViewList) {
|
||||
SpanView rootSpan = findRoot(spanViewList);
|
||||
assert rootSpan != null;
|
||||
findChild(rootSpan, spanViewList, rootSpan.getStartTime());
|
||||
|
||||
List<SpanView> viewList = new ArrayList<>();
|
||||
viewList.add(rootSpan);
|
||||
|
||||
Gson gson = new Gson();
|
||||
return gson.toJson(viewList);
|
||||
}
|
||||
|
||||
private SpanView findRoot(List<SpanView> spanViewList) {
|
||||
for (SpanView spanView : spanViewList) {
|
||||
if (StringUtil.isEmpty(spanView.getParentSpanSegId())) {
|
||||
|
|
|
|||
|
|
@ -50,7 +50,7 @@ public class NodeCompAnalysis extends AbstractNodeCompAnalysis {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Node.NodeDayAnalysis.Size;
|
||||
return WorkerConfig.Queue.Node.NodeCompAnalysis.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -43,7 +43,7 @@ public class NodeCompSave extends RecordPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.DAGNodePersistence.Size;
|
||||
return WorkerConfig.Queue.Node.NodeCompSave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -11,7 +11,7 @@ import com.a.eye.skywalking.collector.worker.storage.RecordData;
|
|||
*/
|
||||
public class NodeMappingDayAgg extends AbstractClusterWorker {
|
||||
|
||||
public NodeMappingDayAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeMappingDayAgg(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.node.NodeMappingIndex;
|
|||
*/
|
||||
public class NodeMappingDaySave extends RecordPersistenceMember {
|
||||
|
||||
public NodeMappingDaySave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeMappingDaySave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -43,7 +43,7 @@ public class NodeMappingDaySave extends RecordPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.DAGNodePersistence.Size;
|
||||
return WorkerConfig.Queue.Node.NodeMappingDaySave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.node.NodeMappingIndex;
|
|||
*/
|
||||
public class NodeMappingHourSave extends RecordPersistenceMember {
|
||||
|
||||
public NodeMappingHourSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeMappingHourSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -43,7 +43,7 @@ public class NodeMappingHourSave extends RecordPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.DAGNodePersistence.Size;
|
||||
return WorkerConfig.Queue.Node.NodeMappingHourSave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.node.NodeMappingIndex;
|
|||
*/
|
||||
public class NodeMappingMinuteSave extends RecordPersistenceMember {
|
||||
|
||||
public NodeMappingMinuteSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeMappingMinuteSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -43,7 +43,7 @@ public class NodeMappingMinuteSave extends RecordPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.DAGNodePersistence.Size;
|
||||
return WorkerConfig.Queue.Node.NodeMappingMinuteSave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.noderef.NodeRefIndex;
|
|||
*/
|
||||
public class NodeRefDaySave extends RecordPersistenceMember {
|
||||
|
||||
public NodeRefDaySave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeRefDaySave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -44,7 +44,7 @@ public class NodeRefDaySave extends RecordPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.DAGNodeRefPersistence.Size;
|
||||
return WorkerConfig.Queue.NodeRef.NodeRefDaySave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.noderef.NodeRefIndex;
|
|||
*/
|
||||
public class NodeRefHourSave extends RecordPersistenceMember {
|
||||
|
||||
public NodeRefHourSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeRefHourSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -44,7 +44,7 @@ public class NodeRefHourSave extends RecordPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.DAGNodeRefPersistence.Size;
|
||||
return WorkerConfig.Queue.NodeRef.NodeRefHourSave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.noderef.NodeRefIndex;
|
|||
*/
|
||||
public class NodeRefMinuteSave extends RecordPersistenceMember {
|
||||
|
||||
public NodeRefMinuteSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeRefMinuteSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -44,7 +44,7 @@ public class NodeRefMinuteSave extends RecordPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.DAGNodeRefPersistence.Size;
|
||||
return WorkerConfig.Queue.NodeRef.NodeRefMinuteSave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.noderef.NodeRefResSumIndex;
|
|||
*/
|
||||
public class NodeRefResSumDaySave extends MetricPersistenceMember {
|
||||
|
||||
private NodeRefResSumDaySave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeRefResSumDaySave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -43,7 +43,7 @@ public class NodeRefResSumDaySave extends MetricPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.ResponseSummaryPersistence.Size;
|
||||
return WorkerConfig.Queue.NodeRef.NodeRefResSumDaySave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.noderef.NodeRefResSumIndex;
|
|||
*/
|
||||
public class NodeRefResSumHourSave extends MetricPersistenceMember {
|
||||
|
||||
private NodeRefResSumHourSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeRefResSumHourSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -43,7 +43,7 @@ public class NodeRefResSumHourSave extends MetricPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.ResponseSummaryPersistence.Size;
|
||||
return WorkerConfig.Queue.NodeRef.NodeRefResSumHourSave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@ import com.a.eye.skywalking.collector.worker.noderef.NodeRefResSumIndex;
|
|||
*/
|
||||
public class NodeRefResSumMinuteSave extends MetricPersistenceMember {
|
||||
|
||||
private NodeRefResSumMinuteSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
NodeRefResSumMinuteSave(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -43,7 +43,7 @@ public class NodeRefResSumMinuteSave extends MetricPersistenceMember {
|
|||
|
||||
@Override
|
||||
public int queueSize() {
|
||||
return WorkerConfig.Queue.Persistence.ResponseSummaryPersistence.Size;
|
||||
return WorkerConfig.Queue.NodeRef.NodeRefResSumMinuteSave.Size;
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -31,7 +31,7 @@ public class SegmentCostIndex extends AbstractIndex {
|
|||
|
||||
@Override
|
||||
public XContentBuilder createMappingBuilder() throws IOException {
|
||||
XContentBuilder mappingBuilder = XContentFactory.jsonBuilder()
|
||||
return XContentFactory.jsonBuilder()
|
||||
.startObject()
|
||||
.startObject("properties")
|
||||
.startObject(SegId)
|
||||
|
|
@ -50,12 +50,11 @@ public class SegmentCostIndex extends AbstractIndex {
|
|||
.field("type", "string")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject(Cost)
|
||||
.startObject(Cost)
|
||||
.field("type", "long")
|
||||
.field("index", "not_analyzed")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.endObject()
|
||||
.endObject();
|
||||
return mappingBuilder;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,8 +1,6 @@
|
|||
package com.a.eye.skywalking.collector.worker.segment;
|
||||
|
||||
import com.a.eye.skywalking.collector.worker.storage.AbstractIndex;
|
||||
import org.apache.logging.log4j.LogManager;
|
||||
import org.apache.logging.log4j.Logger;
|
||||
import org.elasticsearch.common.xcontent.XContentBuilder;
|
||||
import org.elasticsearch.common.xcontent.XContentFactory;
|
||||
|
||||
|
|
@ -13,8 +11,6 @@ import java.io.IOException;
|
|||
*/
|
||||
public class SegmentIndex extends AbstractIndex {
|
||||
|
||||
private Logger logger = LogManager.getFormatterLogger(SegmentIndex.class);
|
||||
|
||||
public static final String Index = "segment_idx";
|
||||
|
||||
@Override
|
||||
|
|
@ -29,7 +25,7 @@ public class SegmentIndex extends AbstractIndex {
|
|||
|
||||
@Override
|
||||
public XContentBuilder createMappingBuilder() throws IOException {
|
||||
XContentBuilder mappingBuilder = XContentFactory.jsonBuilder()
|
||||
return XContentFactory.jsonBuilder()
|
||||
.startObject()
|
||||
.startObject("properties")
|
||||
.startObject("traceSegmentId")
|
||||
|
|
@ -60,54 +56,7 @@ public class SegmentIndex extends AbstractIndex {
|
|||
.field("type", "long")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startArray("refs")
|
||||
.startObject("traceSegmentId")
|
||||
.field("type", "String")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject("spanId")
|
||||
.field("type", "integer")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject("applicationCode")
|
||||
.field("type", "String")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject("peerHost")
|
||||
.field("type", "String")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.endArray()
|
||||
.startArray("refs")
|
||||
.startObject("spanId")
|
||||
.field("type", "integer")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject("parentSpanId")
|
||||
.field("type", "integer")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject("startTime")
|
||||
.field("type", "date")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject("endTime")
|
||||
.field("type", "date")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.startObject("operationName")
|
||||
.field("type", "String")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.endArray()
|
||||
.startArray("relatedGlobalTraces")
|
||||
.startObject("id")
|
||||
.field("type", "String")
|
||||
.field("index", "not_analyzed")
|
||||
.endObject()
|
||||
.endArray()
|
||||
.endObject()
|
||||
.endObject();
|
||||
return mappingBuilder;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -40,11 +40,6 @@ public class SpanGetWithId extends AbstractGet {
|
|||
}
|
||||
logger.debug("segId: %s, spanId: %s", Arrays.toString(request.get("segId")), Arrays.toString(request.get("spanId")));
|
||||
|
||||
int maxCost = -1;
|
||||
if (request.containsKey("maxCost")) {
|
||||
maxCost = Integer.valueOf(ParameterTools.INSTANCE.toString(request, "maxCost"));
|
||||
}
|
||||
|
||||
String segId = ParameterTools.INSTANCE.toString(request, "segId");
|
||||
String spanId = ParameterTools.INSTANCE.toString(request, "spanId");
|
||||
|
||||
|
|
|
|||
|
|
@ -3,10 +3,11 @@ package com.a.eye.skywalking.collector.worker.span.persistence;
|
|||
import com.a.eye.skywalking.collector.actor.*;
|
||||
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.Const;
|
||||
import com.a.eye.skywalking.collector.worker.segment.SegmentIndex;
|
||||
import com.a.eye.skywalking.collector.worker.segment.logic.Segment;
|
||||
import com.a.eye.skywalking.collector.worker.segment.logic.SegmentDeserialize;
|
||||
import com.a.eye.skywalking.collector.worker.storage.EsClient;
|
||||
import com.a.eye.skywalking.collector.worker.storage.GetResponseFromEs;
|
||||
import com.a.eye.skywalking.trace.Span;
|
||||
import com.google.gson.Gson;
|
||||
import com.google.gson.JsonObject;
|
||||
|
|
@ -21,7 +22,7 @@ public class SpanSearchWithId extends AbstractLocalSyncWorker {
|
|||
|
||||
private Gson gson = new Gson();
|
||||
|
||||
private SpanSearchWithId(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
SpanSearchWithId(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
|
|
@ -29,7 +30,7 @@ public class SpanSearchWithId extends AbstractLocalSyncWorker {
|
|||
protected void onWork(Object request, Object response) throws Exception {
|
||||
if (request instanceof RequestEntity) {
|
||||
RequestEntity search = (RequestEntity) request;
|
||||
GetResponse getResponse = EsClient.INSTANCE.getClient().prepareGet(SegmentIndex.Index, SegmentIndex.Type_Record, search.segId).get();
|
||||
GetResponse getResponse = GetResponseFromEs.INSTANCE.get(SegmentIndex.Index, SegmentIndex.Type_Record, search.segId);
|
||||
Segment segment = SegmentDeserialize.INSTANCE.deserializeFromES(getResponse.getSourceAsString());
|
||||
List<Span> spanList = segment.getSpans();
|
||||
|
||||
|
|
@ -44,7 +45,7 @@ public class SpanSearchWithId extends AbstractLocalSyncWorker {
|
|||
}
|
||||
|
||||
JsonObject resJsonObj = (JsonObject) response;
|
||||
resJsonObj.add("result", dataJson);
|
||||
resJsonObj.add(Const.RESULT, dataJson);
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -0,0 +1,14 @@
|
|||
package com.a.eye.skywalking.collector.worker.storage;
|
||||
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public enum GetResponseFromEs {
|
||||
INSTANCE;
|
||||
|
||||
public GetResponse get(String index, String type, String id) {
|
||||
return EsClient.INSTANCE.getClient().prepareGet(index, type, id).get();
|
||||
}
|
||||
}
|
||||
|
|
@ -1,18 +1,35 @@
|
|||
package com.a.eye.skywalking.collector.worker;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
|
||||
import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
|
||||
import com.a.eye.skywalking.collector.queue.EndOfBatchCommand;
|
||||
import org.junit.After;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.Mockito;
|
||||
import org.mockito.invocation.InvocationOnMock;
|
||||
import org.mockito.stubbing.Answer;
|
||||
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(TestAnalysisMember.class)
|
||||
@PowerMockIgnore({"javax.management.*"})
|
||||
public class AnalysisMemberTestCase {
|
||||
|
||||
@Test
|
||||
public void testCommandOnWork() throws Exception {
|
||||
AnalysisMember member = mock(AnalysisMember.class);
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
TestAnalysisMember member = PowerMockito.spy(new TestAnalysisMember(TestAnalysisMember.Role.INSTANCE, clusterWorkerContext, localWorkerContext));
|
||||
|
||||
EndOfBatchCommand command = new EndOfBatchCommand();
|
||||
member.onWork(command);
|
||||
|
|
@ -22,11 +39,58 @@ public class AnalysisMemberTestCase {
|
|||
|
||||
@Test
|
||||
public void testAnalyse() throws Exception {
|
||||
AnalysisMember member = mock(AnalysisMember.class);
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
TestAnalysisMember member = PowerMockito.spy(new TestAnalysisMember(TestAnalysisMember.Role.INSTANCE, clusterWorkerContext, localWorkerContext));
|
||||
|
||||
Object message = new Object();
|
||||
member.onWork(message);
|
||||
verify(member, never()).aggregation();
|
||||
verify(member, times(1)).analyse(anyObject());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testPreStart() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
TestAnalysisMember member = PowerMockito.spy(new TestAnalysisMember(TestAnalysisMember.Role.INSTANCE, clusterWorkerContext, localWorkerContext));
|
||||
member.preStart();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOnWorkException() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
TestAnalysisMember member = PowerMockito.spy(new TestAnalysisMember(TestAnalysisMember.Role.INSTANCE, clusterWorkerContext, localWorkerContext));
|
||||
|
||||
doThrow(new TestException()).when(member).analyse(anyObject());
|
||||
|
||||
ExceptionAnswer answer = new ExceptionAnswer();
|
||||
PowerMockito.when(member, "saveException", any(TestException.class)).thenAnswer(answer);
|
||||
|
||||
member.onWork(new Object());
|
||||
|
||||
Assert.assertEquals(true, answer.isTestException);
|
||||
}
|
||||
|
||||
class TestException extends Exception {
|
||||
|
||||
}
|
||||
|
||||
|
||||
class ExceptionAnswer implements Answer {
|
||||
|
||||
boolean isTestException = false;
|
||||
|
||||
@Override
|
||||
public Object answer(InvocationOnMock invocation) throws Throwable {
|
||||
Object obj = invocation.getArguments()[0];
|
||||
if (obj instanceof TestException) {
|
||||
isTestException = true;
|
||||
} else {
|
||||
isTestException = false;
|
||||
}
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,5 +1,7 @@
|
|||
package com.a.eye.skywalking.collector.worker;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
|
||||
import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
|
||||
import com.a.eye.skywalking.collector.worker.storage.MergeData;
|
||||
import com.a.eye.skywalking.collector.worker.storage.MergePersistenceData;
|
||||
import org.junit.Assert;
|
||||
|
|
@ -7,6 +9,8 @@ import org.junit.Before;
|
|||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
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;
|
||||
|
||||
|
|
@ -16,38 +20,42 @@ import static org.powermock.api.mockito.PowerMockito.*;
|
|||
* @author pengys5
|
||||
*/
|
||||
@RunWith(PowerMockRunner.class)
|
||||
@PrepareForTest(MergeAnalysisMember.class)
|
||||
@PrepareForTest(TestMergeAnalysisMember.class)
|
||||
@PowerMockIgnore({"javax.management.*"})
|
||||
public class MergeAnalysisMemberTestCase {
|
||||
|
||||
private MergeAnalysisMember member;
|
||||
private TestMergeAnalysisMember mergeAnalysisMember;
|
||||
private MergePersistenceData persistenceData;
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
member = mock(MergeAnalysisMember.class);
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
mergeAnalysisMember = PowerMockito.spy(new TestMergeAnalysisMember(TestMergeAnalysisMember.Role.INSTANCE, clusterWorkerContext, localWorkerContext));
|
||||
|
||||
persistenceData = mock(MergePersistenceData.class);
|
||||
MergeData mergeData = mock(MergeData.class);
|
||||
|
||||
when(member, "getPersistenceData").thenReturn(persistenceData);
|
||||
when(mergeAnalysisMember, "getPersistenceData").thenReturn(persistenceData);
|
||||
when(persistenceData.getElseCreate(Mockito.anyString())).thenReturn(mergeData);
|
||||
|
||||
doCallRealMethod().when(member).setMergeData(Mockito.anyString(), Mockito.anyString(), Mockito.anyString());
|
||||
doCallRealMethod().when(mergeAnalysisMember).setMergeData(Mockito.anyString(), Mockito.anyString(), Mockito.anyString());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSetMergeDataNotFull() throws Exception {
|
||||
when(persistenceData.size()).thenReturn(WorkerConfig.Persistence.Data.size - 1);
|
||||
|
||||
member.setMergeData("segment_1", "column", "value");
|
||||
Mockito.verify(member, Mockito.never()).aggregation();
|
||||
mergeAnalysisMember.setMergeData("segment_1", "column", "value");
|
||||
Mockito.verify(mergeAnalysisMember, Mockito.never()).aggregation();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testSetMergeDataFull() throws Exception {
|
||||
when(persistenceData.size()).thenReturn(WorkerConfig.Persistence.Data.size);
|
||||
|
||||
member.setMergeData("segment_1", "column", "value");
|
||||
Mockito.verify(member, Mockito.times(1)).aggregation();
|
||||
mergeAnalysisMember.setMergeData("segment_1", "column", "value");
|
||||
Mockito.verify(mergeAnalysisMember, Mockito.times(1)).aggregation();
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -55,9 +63,10 @@ public class MergeAnalysisMemberTestCase {
|
|||
MergePersistenceData persistenceData = new MergePersistenceData();
|
||||
persistenceData.getElseCreate("segment_1").setMergeData("column", "value");
|
||||
|
||||
when(member, "getPersistenceData").thenReturn(persistenceData);
|
||||
doCallRealMethod().when(member).pushOne();
|
||||
when(mergeAnalysisMember, "getPersistenceData").thenReturn(persistenceData);
|
||||
doCallRealMethod().when(mergeAnalysisMember).pushOne();
|
||||
|
||||
Assert.assertEquals("segment_1", member.pushOne().getId());
|
||||
Assert.assertEquals("segment_1", mergeAnalysisMember.pushOne().getId());
|
||||
Assert.assertEquals(null, mergeAnalysisMember.pushOne());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,59 @@
|
|||
package com.a.eye.skywalking.collector.worker;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
|
||||
import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
|
||||
import com.a.eye.skywalking.collector.queue.EndOfBatchCommand;
|
||||
import com.a.eye.skywalking.collector.worker.mock.MockEsBulkClient;
|
||||
import com.a.eye.skywalking.collector.worker.storage.EsClient;
|
||||
import com.a.eye.skywalking.collector.worker.storage.MergeData;
|
||||
import com.a.eye.skywalking.collector.worker.storage.MergePersistenceData;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
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.powermock.api.mockito.PowerMockito.*;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
@RunWith(PowerMockRunner.class)
|
||||
@PrepareForTest({TestMergePersistenceMember.class, EsClient.class})
|
||||
@PowerMockIgnore({"javax.management.*"})
|
||||
public class MergePersistenceMemberTestCase {
|
||||
|
||||
private TestMergePersistenceMember mergePersistenceMember;
|
||||
private MergePersistenceData persistenceData;
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
MockEsBulkClient mockEsBulkClient = new MockEsBulkClient();
|
||||
mockEsBulkClient.createMock();
|
||||
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
mergePersistenceMember = PowerMockito.spy(new TestMergePersistenceMember(TestMergePersistenceMember.Role.INSTANCE, clusterWorkerContext, localWorkerContext));
|
||||
|
||||
persistenceData = mock(MergePersistenceData.class);
|
||||
MergeData mergeData = mock(MergeData.class);
|
||||
|
||||
when(mergePersistenceMember, "getPersistenceData").thenReturn(persistenceData);
|
||||
when(persistenceData.getElseCreate(Mockito.anyString())).thenReturn(mergeData);
|
||||
|
||||
doCallRealMethod().when(mergePersistenceMember).analyse(Mockito.any(MergeData.class));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAnalyse() throws Exception {
|
||||
String id = "2016" + Const.ID_SPLIT + "A" + Const.ID_SPLIT + "B";
|
||||
MergeData mergeData = new MergeData(id);
|
||||
mergeData.setMergeData("Column", "Value");
|
||||
|
||||
// mergePersistenceMember.analyse(mergeData);
|
||||
// mergePersistenceMember.onWork(new EndOfBatchCommand());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,44 @@
|
|||
package com.a.eye.skywalking.collector.worker;
|
||||
|
||||
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.selector.WorkerSelector;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class TestAnalysisMember extends AnalysisMember {
|
||||
TestAnalysisMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void analyse(Object message) throws Exception {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void preStart() throws ProviderNotFoundException {
|
||||
super.preStart();
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void aggregation() throws Exception {
|
||||
|
||||
}
|
||||
|
||||
public enum Role implements com.a.eye.skywalking.collector.actor.Role {
|
||||
INSTANCE;
|
||||
|
||||
@Override
|
||||
public String roleName() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public WorkerSelector workerSelector() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,45 @@
|
|||
package com.a.eye.skywalking.collector.worker;
|
||||
|
||||
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.selector.WorkerSelector;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class TestMergeAnalysisMember extends MergeAnalysisMember {
|
||||
|
||||
TestMergeAnalysisMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void analyse(Object message) throws Exception {
|
||||
|
||||
}
|
||||
|
||||
@Override
|
||||
public void preStart() throws ProviderNotFoundException {
|
||||
super.preStart();
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void aggregation() throws Exception {
|
||||
|
||||
}
|
||||
|
||||
public enum Role implements com.a.eye.skywalking.collector.actor.Role {
|
||||
INSTANCE;
|
||||
|
||||
@Override
|
||||
public String roleName() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public WorkerSelector workerSelector() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,38 @@
|
|||
package com.a.eye.skywalking.collector.worker;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
|
||||
import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
|
||||
import com.a.eye.skywalking.collector.actor.selector.WorkerSelector;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class TestMergePersistenceMember extends MergePersistenceMember {
|
||||
TestMergePersistenceMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
@Override
|
||||
public String esIndex() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String esType() {
|
||||
return null;
|
||||
}
|
||||
|
||||
public enum Role implements com.a.eye.skywalking.collector.actor.Role {
|
||||
INSTANCE;
|
||||
|
||||
@Override
|
||||
public String roleName() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public WorkerSelector workerSelector() {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,25 @@
|
|||
package com.a.eye.skywalking.collector.worker.globaltrace;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class GlobalTraceIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
GlobalTraceIndex index = new GlobalTraceIndex();
|
||||
Assert.assertEquals("global_trace_idx", index.index());
|
||||
Assert.assertEquals(true, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
GlobalTraceIndex index = new GlobalTraceIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"subSegIds\":{\"type\":\"text\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,72 @@
|
|||
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.RollingSelector;
|
||||
import com.a.eye.skywalking.collector.worker.globaltrace.GlobalTraceIndex;
|
||||
import com.a.eye.skywalking.collector.worker.segment.SegmentIndex;
|
||||
import com.a.eye.skywalking.collector.worker.storage.GetResponseFromEs;
|
||||
import com.google.gson.JsonObject;
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
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 org.powermock.reflect.Whitebox;
|
||||
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
@RunWith(PowerMockRunner.class)
|
||||
@PrepareForTest({GetResponseFromEs.class})
|
||||
@PowerMockIgnore({"javax.management.*"})
|
||||
public class GlobalTraceSearchWithGlobalIdTestCase {
|
||||
|
||||
private GetResponseFromEs getResponseFromEs;
|
||||
|
||||
private String global_Str = "{\"subSegIds\":\"Segment.1491277162066.18986177.70531.27.1\"}";
|
||||
private String seg_str = "{\"ts\":\"Segment.1491277162066.18986177.70531.27.1\",\"st\":1491277162066,\"et\":1491277165743,\"ss\":[{\"si\":0,\"ps\":-1,\"st\":1491277162141,\"et\":1491277162144,\"on\":\"Jedis/getClient\",\"ts\":{\"span.layer\":\"db\",\"component\":\"Redis\",\"db.type\":\"Redis\",\"peer.host\":\"127.0.0.1\",\"span.kind\":\"client\"},\"tb\":{},\"ti\":{\"peer.port\":6379},\"lo\":[]}],\"ac\":\"cache-service\",\"gt\":[\"Trace.1491277147443.-1562443425.70539.65.2\"],\"sampled\":true,\"minute\":201704041139,\"hour\":201704041100,\"day\":201704040000,\"aggId\":null}";
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
getResponseFromEs = PowerMockito.mock(GetResponseFromEs.class);
|
||||
Whitebox.setInternalState(GetResponseFromEs.class, "INSTANCE", getResponseFromEs);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(GlobalTraceSearchWithGlobalId.class.getSimpleName(), GlobalTraceSearchWithGlobalId.WorkerRole.INSTANCE.roleName());
|
||||
Assert.assertEquals(RollingSelector.class.getSimpleName(), GlobalTraceSearchWithGlobalId.WorkerRole.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(GlobalTraceSearchWithGlobalId.class.getSimpleName(), GlobalTraceSearchWithGlobalId.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(GlobalTraceSearchWithGlobalId.class.getSimpleName(), GlobalTraceSearchWithGlobalId.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOnWork() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
GlobalTraceSearchWithGlobalId globalTraceSearchWithGlobalId = new GlobalTraceSearchWithGlobalId(GlobalTraceSearchWithGlobalId.WorkerRole.INSTANCE, clusterWorkerContext, localWorkerContext);
|
||||
|
||||
GetResponse getResponse = mock(GetResponse.class);
|
||||
when(getResponseFromEs.get(GlobalTraceIndex.Index, GlobalTraceIndex.Type_Record, "Trace.1491277147443.-1562443425.70539.65.2")).thenReturn(getResponse);
|
||||
when(getResponse.getSourceAsString()).thenReturn(global_Str);
|
||||
|
||||
GetResponse segResponse = mock(GetResponse.class);
|
||||
when(getResponseFromEs.get(SegmentIndex.Index, SegmentIndex.Type_Record, "Segment.1491277162066.18986177.70531.27.1")).thenReturn(segResponse);
|
||||
when(segResponse.getSourceAsString()).thenReturn(seg_str);
|
||||
|
||||
JsonObject response = new JsonObject();
|
||||
globalTraceSearchWithGlobalId.onWork("Trace.1491277147443.-1562443425.70539.65.2", response);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,35 @@
|
|||
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.worker.node.NodeMappingIndex;
|
||||
import com.a.eye.skywalking.collector.worker.node.persistence.NodeMappingSearchWithTimeSlice;
|
||||
import com.a.eye.skywalking.collector.worker.storage.EsClient;
|
||||
import com.google.gson.JsonArray;
|
||||
import com.google.gson.JsonObject;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class GlobalTraceSearchWithGlobalIdUseDB {
|
||||
|
||||
public static void main(String[] args) throws Exception {
|
||||
EsClient.INSTANCE.boot();
|
||||
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
GlobalTraceSearchWithGlobalId globalTraceSearchWithGlobalId =
|
||||
new GlobalTraceSearchWithGlobalId(GlobalTraceSearchWithGlobalId.WorkerRole.INSTANCE, clusterWorkerContext, localWorkerContext);
|
||||
|
||||
JsonObject response = new JsonObject();
|
||||
globalTraceSearchWithGlobalId.onWork("Trace.1491277147443.-1562443425.70539.65.2", response);
|
||||
|
||||
JsonArray nodeArray = response.get("result").getAsJsonArray();
|
||||
System.out.println(nodeArray.size());
|
||||
System.out.println(nodeArray.toString());
|
||||
for (int i = 0; i < nodeArray.size(); i++) {
|
||||
JsonObject nodeJsonObj = nodeArray.get(i).getAsJsonObject();
|
||||
System.out.println(nodeJsonObj);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,22 @@
|
|||
package com.a.eye.skywalking.collector.worker.mock;
|
||||
|
||||
import com.a.eye.skywalking.collector.worker.storage.EsClient;
|
||||
import org.elasticsearch.client.Client;
|
||||
import org.powermock.api.mockito.PowerMockito;
|
||||
import org.powermock.reflect.Whitebox;
|
||||
|
||||
import static org.powermock.api.mockito.PowerMockito.when;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class MockEsClient {
|
||||
|
||||
public Client mock() {
|
||||
Client client = PowerMockito.mock(Client.class);
|
||||
EsClient esClient = PowerMockito.mock(EsClient.class);
|
||||
Whitebox.setInternalState(EsClient.class, "INSTANCE", esClient);
|
||||
when(esClient.getClient()).thenReturn(client);
|
||||
return client;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,27 @@
|
|||
package com.a.eye.skywalking.collector.worker.mock;
|
||||
|
||||
import org.elasticsearch.action.get.GetRequestBuilder;
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
import org.elasticsearch.client.Client;
|
||||
|
||||
import static org.mockito.Mockito.*;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class MockGetResponse {
|
||||
|
||||
public GetResponse mockito() {
|
||||
MockEsClient mockEsClient = new MockEsClient();
|
||||
Client client = mockEsClient.mock();
|
||||
|
||||
GetRequestBuilder builder = mock(GetRequestBuilder.class);
|
||||
GetResponse getResponse = mock(GetResponse.class);
|
||||
when(builder.get()).thenReturn(getResponse);
|
||||
|
||||
when(client.prepareGet(anyString(), anyString(), anyString())).thenReturn(builder);
|
||||
|
||||
|
||||
return getResponse;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,26 @@
|
|||
package com.a.eye.skywalking.collector.worker.node;
|
||||
|
||||
import com.a.eye.skywalking.collector.worker.globaltrace.GlobalTraceIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeCompIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
NodeCompIndex index = new NodeCompIndex();
|
||||
Assert.assertEquals("node_comp_idx", index.index());
|
||||
Assert.assertEquals(false, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
NodeCompIndex index = new NodeCompIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"name\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"peers\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"aggId\":{\"type\":\"string\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,25 @@
|
|||
package com.a.eye.skywalking.collector.worker.node;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
NodeMappingIndex index = new NodeMappingIndex();
|
||||
Assert.assertEquals("node_mapping_idx", index.index());
|
||||
Assert.assertEquals(false, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
NodeMappingIndex index = new NodeMappingIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"code\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"peers\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"aggId\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"timeSlice\":{\"type\":\"long\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -2,12 +2,14 @@ package com.a.eye.skywalking.collector.worker.node.analysis;
|
|||
|
||||
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.RollingSelector;
|
||||
import com.a.eye.skywalking.collector.queue.EndOfBatchCommand;
|
||||
import com.a.eye.skywalking.collector.worker.WorkerConfig;
|
||||
import com.a.eye.skywalking.collector.worker.mock.RecordDataAnswer;
|
||||
import com.a.eye.skywalking.collector.worker.node.persistence.NodeMappingDayAgg;
|
||||
import com.a.eye.skywalking.collector.worker.noderef.analysis.NodeRefResSumDayAnalysis;
|
||||
import com.a.eye.skywalking.collector.worker.segment.SegmentPost;
|
||||
import com.a.eye.skywalking.collector.worker.segment.mock.SegmentMock;
|
||||
import com.a.eye.skywalking.collector.worker.storage.RecordData;
|
||||
|
|
@ -25,6 +27,7 @@ import org.powermock.modules.junit4.PowerMockRunner;
|
|||
import java.util.List;
|
||||
|
||||
import static org.mockito.Mockito.doAnswer;
|
||||
import static org.mockito.Mockito.doThrow;
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.powermock.api.mockito.PowerMockito.when;
|
||||
|
||||
|
|
@ -39,10 +42,11 @@ public class NodeMappingDayAnalysisTestCase {
|
|||
private NodeMappingDayAnalysis nodeMappingDayAnalysis;
|
||||
private SegmentMock segmentMock = new SegmentMock();
|
||||
private RecordDataAnswer recordDataAnswer;
|
||||
private ClusterWorkerContext clusterWorkerContext;
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
WorkerRefs workerRefs = mock(WorkerRefs.class);
|
||||
recordDataAnswer = new RecordDataAnswer();
|
||||
doAnswer(recordDataAnswer).when(workerRefs).tell(Mockito.any(RecordData.class));
|
||||
|
|
@ -69,6 +73,12 @@ public class NodeMappingDayAnalysisTestCase {
|
|||
Assert.assertEquals(testSize, NodeMappingDayAnalysis.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
|
||||
@Test(expected = Exception.class)
|
||||
public void testPreStart() throws ProviderNotFoundException {
|
||||
when(clusterWorkerContext.findProvider(NodeRefResSumDayAnalysis.Role.INSTANCE)).thenThrow(new Exception());
|
||||
nodeMappingDayAnalysis.preStart();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testAnalyse() throws Exception {
|
||||
List<SegmentPost.SegmentWithTimeSlice> cacheServiceSegment = segmentMock.mockCacheServiceSegmentSegmentTimeSlice();
|
||||
|
|
|
|||
|
|
@ -0,0 +1,52 @@
|
|||
package com.a.eye.skywalking.collector.worker.node.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.node.NodeCompIndex;
|
||||
import com.a.eye.skywalking.collector.worker.node.NodeMappingIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeCompSaveTestCase {
|
||||
|
||||
private NodeCompSave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeCompSave(NodeCompSave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeCompIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeCompIndex.Type_Record, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeCompSave.class.getSimpleName(), NodeCompSave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeCompSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeCompSave.class.getSimpleName(), NodeCompSave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeCompSave.class.getSimpleName(), NodeCompSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.Node.NodeCompSave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeCompSave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -1,8 +1,6 @@
|
|||
package com.a.eye.skywalking.collector.worker.node.persistence;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.ClusterWorkerContext;
|
||||
import com.a.eye.skywalking.collector.actor.LocalWorkerContext;
|
||||
import com.a.eye.skywalking.collector.actor.WorkerRefs;
|
||||
import com.a.eye.skywalking.collector.actor.*;
|
||||
import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector;
|
||||
import com.a.eye.skywalking.collector.worker.Const;
|
||||
import com.a.eye.skywalking.collector.worker.WorkerConfig;
|
||||
|
|
@ -34,10 +32,11 @@ public class NodeMappingDayAggTestCase {
|
|||
|
||||
private NodeMappingDayAgg nodeMappingDayAgg;
|
||||
private RecordDataAnswer recordDataAnswer;
|
||||
private ClusterWorkerContext clusterWorkerContext;
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
|
||||
LocalWorkerContext localWorkerContext = PowerMockito.mock(LocalWorkerContext.class);
|
||||
WorkerRefs workerRefs = mock(WorkerRefs.class);
|
||||
|
|
@ -65,6 +64,12 @@ public class NodeMappingDayAggTestCase {
|
|||
Assert.assertEquals(testSize, NodeMappingDayAgg.Factory.INSTANCE.workerNum());
|
||||
}
|
||||
|
||||
@Test(expected = Exception.class)
|
||||
public void testPreStart() throws ProviderNotFoundException {
|
||||
when(clusterWorkerContext.findProvider(NodeMappingDaySave.Role.INSTANCE)).thenThrow(new Exception());
|
||||
nodeMappingDayAgg.preStart();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOnWork() throws Exception {
|
||||
String id = "2017" + Const.ID_SPLIT + "TestNodeMappingDayAgg";
|
||||
|
|
|
|||
|
|
@ -0,0 +1,51 @@
|
|||
package com.a.eye.skywalking.collector.worker.node.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.node.NodeMappingIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingDaySaveTestCase {
|
||||
|
||||
private NodeMappingDaySave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeMappingDaySave(NodeMappingDaySave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeMappingIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeMappingIndex.Type_Day, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeMappingDaySave.class.getSimpleName(), NodeMappingDaySave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeMappingDaySave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeMappingDaySave.class.getSimpleName(), NodeMappingDaySave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeMappingDaySave.class.getSimpleName(), NodeMappingDaySave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.Node.NodeMappingDaySave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeMappingDaySave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.node.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.Const;
|
||||
|
|
@ -34,10 +35,11 @@ public class NodeMappingHourAggTestCase {
|
|||
|
||||
private NodeMappingHourAgg nodeMappingHourAgg;
|
||||
private RecordDataAnswer recordDataAnswer;
|
||||
private ClusterWorkerContext clusterWorkerContext;
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
|
||||
LocalWorkerContext localWorkerContext = PowerMockito.mock(LocalWorkerContext.class);
|
||||
WorkerRefs workerRefs = mock(WorkerRefs.class);
|
||||
|
|
@ -65,6 +67,12 @@ public class NodeMappingHourAggTestCase {
|
|||
Assert.assertEquals(testSize, NodeMappingHourAgg.Factory.INSTANCE.workerNum());
|
||||
}
|
||||
|
||||
@Test(expected = Exception.class)
|
||||
public void testPreStart() throws ProviderNotFoundException {
|
||||
when(clusterWorkerContext.findProvider(NodeMappingHourSave.Role.INSTANCE)).thenThrow(new Exception());
|
||||
nodeMappingHourAgg.preStart();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOnWork() throws Exception {
|
||||
String id = "2017" + Const.ID_SPLIT + "TestNodeMappingHourAgg";
|
||||
|
|
|
|||
|
|
@ -0,0 +1,51 @@
|
|||
package com.a.eye.skywalking.collector.worker.node.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.node.NodeMappingIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingHourSaveTestCase {
|
||||
|
||||
private NodeMappingHourSave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeMappingHourSave(NodeMappingHourSave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeMappingIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeMappingIndex.Type_Hour, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeMappingHourSave.class.getSimpleName(), NodeMappingHourSave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeMappingHourSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeMappingHourSave.class.getSimpleName(), NodeMappingHourSave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeMappingHourSave.class.getSimpleName(), NodeMappingHourSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.Node.NodeMappingHourSave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeMappingHourSave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.node.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.Const;
|
||||
|
|
@ -34,10 +35,11 @@ public class NodeMappingMinuteAggTestCase {
|
|||
|
||||
private NodeMappingMinuteAgg nodeMappingMinuteAgg;
|
||||
private RecordDataAnswer recordDataAnswer;
|
||||
private ClusterWorkerContext clusterWorkerContext;
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
clusterWorkerContext = PowerMockito.mock(ClusterWorkerContext.class);
|
||||
|
||||
LocalWorkerContext localWorkerContext = PowerMockito.mock(LocalWorkerContext.class);
|
||||
WorkerRefs workerRefs = mock(WorkerRefs.class);
|
||||
|
|
@ -65,6 +67,12 @@ public class NodeMappingMinuteAggTestCase {
|
|||
Assert.assertEquals(testSize, NodeMappingMinuteAgg.Factory.INSTANCE.workerNum());
|
||||
}
|
||||
|
||||
@Test(expected = Exception.class)
|
||||
public void testPreStart() throws ProviderNotFoundException {
|
||||
when(clusterWorkerContext.findProvider(NodeMappingDaySave.Role.INSTANCE)).thenThrow(new Exception());
|
||||
nodeMappingMinuteAgg.preStart();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOnWork() throws Exception {
|
||||
String id = "2017" + Const.ID_SPLIT + "TestNodeMappingMinuteAgg";
|
||||
|
|
|
|||
|
|
@ -0,0 +1,51 @@
|
|||
package com.a.eye.skywalking.collector.worker.node.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.node.NodeMappingIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeMappingMinuteSaveTestCase {
|
||||
|
||||
private NodeMappingMinuteSave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeMappingMinuteSave(NodeMappingMinuteSave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeMappingIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeMappingIndex.Type_Minute, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeMappingMinuteSave.class.getSimpleName(), NodeMappingMinuteSave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeMappingMinuteSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeMappingMinuteSave.class.getSimpleName(), NodeMappingMinuteSave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeMappingMinuteSave.class.getSimpleName(), NodeMappingMinuteSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.Node.NodeMappingMinuteSave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeMappingMinuteSave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,26 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef;
|
||||
|
||||
import com.a.eye.skywalking.collector.worker.globaltrace.GlobalTraceIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
NodeRefIndex index = new NodeRefIndex();
|
||||
Assert.assertEquals("node_ref_idx", index.index());
|
||||
Assert.assertEquals(false, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
NodeRefIndex index = new NodeRefIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"front\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"frontIsRealCode\":{\"type\":\"boolean\",\"index\":\"not_analyzed\"},\"behind\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"behindIsRealCode\":{\"type\":\"boolean\",\"index\":\"not_analyzed\"},\"aggId\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"timeSlice\":{\"type\":\"long\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,25 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefResSumIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
NodeRefResSumIndex index = new NodeRefResSumIndex();
|
||||
Assert.assertEquals("node_ref_res_sum_idx", index.index());
|
||||
Assert.assertEquals(false, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
NodeRefResSumIndex index = new NodeRefResSumIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"oneSecondLess\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"threeSecondLess\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"fiveSecondLess\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"fiveSecondGreater\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"error\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"summary\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"aggId\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"timeSlice\":{\"type\":\"long\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,51 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef.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.noderef.NodeRefIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefDaySaveTestCase {
|
||||
|
||||
private NodeRefDaySave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeRefDaySave(NodeRefDaySave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeRefIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeRefIndex.Type_Day, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeRefDaySave.class.getSimpleName(), NodeRefDaySave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeRefDaySave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeRefDaySave.class.getSimpleName(), NodeRefDaySave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeRefDaySave.class.getSimpleName(), NodeRefDaySave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.NodeRef.NodeRefDaySave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeRefDaySave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,51 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef.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.noderef.NodeRefIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefHourSaveTestCase {
|
||||
|
||||
private NodeRefHourSave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeRefHourSave(NodeRefHourSave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeRefIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeRefIndex.Type_Hour, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeRefHourSave.class.getSimpleName(), NodeRefHourSave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeRefHourSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeRefHourSave.class.getSimpleName(), NodeRefHourSave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeRefHourSave.class.getSimpleName(), NodeRefHourSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.NodeRef.NodeRefHourSave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeRefHourSave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,51 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef.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.noderef.NodeRefIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefMinuteSaveTestCase {
|
||||
|
||||
private NodeRefMinuteSave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeRefMinuteSave(NodeRefMinuteSave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeRefIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeRefIndex.Type_Minute, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeRefMinuteSave.class.getSimpleName(), NodeRefMinuteSave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeRefMinuteSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeRefMinuteSave.class.getSimpleName(), NodeRefMinuteSave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeRefMinuteSave.class.getSimpleName(), NodeRefMinuteSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.NodeRef.NodeRefMinuteSave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeRefMinuteSave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,50 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef.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.noderef.NodeRefResSumIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefResSumDaySaveTestCase {
|
||||
private NodeRefResSumDaySave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeRefResSumDaySave(NodeRefResSumDaySave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeRefResSumIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeRefResSumIndex.Type_Day, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeRefResSumDaySave.class.getSimpleName(), NodeRefResSumDaySave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeRefResSumDaySave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeRefResSumDaySave.class.getSimpleName(), NodeRefResSumDaySave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeRefResSumDaySave.class.getSimpleName(), NodeRefResSumDaySave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.NodeRef.NodeRefResSumDaySave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeRefResSumDaySave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,51 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef.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.noderef.NodeRefResSumIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefResSumHourSaveTestCase {
|
||||
|
||||
private NodeRefResSumHourSave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeRefResSumHourSave(NodeRefResSumHourSave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeRefResSumIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeRefResSumIndex.Type_Hour, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeRefResSumHourSave.class.getSimpleName(), NodeRefResSumHourSave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeRefResSumHourSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeRefResSumHourSave.class.getSimpleName(), NodeRefResSumHourSave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeRefResSumHourSave.class.getSimpleName(), NodeRefResSumHourSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.NodeRef.NodeRefResSumHourSave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeRefResSumHourSave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,52 @@
|
|||
package com.a.eye.skywalking.collector.worker.noderef.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.noderef.NodeRefIndex;
|
||||
import com.a.eye.skywalking.collector.worker.noderef.NodeRefResSumIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class NodeRefResSumMinuteSaveTestCase {
|
||||
|
||||
private NodeRefResSumMinuteSave save;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
ClusterWorkerContext cluster = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext local = new LocalWorkerContext();
|
||||
save = new NodeRefResSumMinuteSave(NodeRefResSumMinuteSave.Role.INSTANCE, cluster, local);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsIndex() {
|
||||
Assert.assertEquals(NodeRefResSumIndex.Index, save.esIndex());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testEsType() {
|
||||
Assert.assertEquals(NodeRefResSumIndex.Type_Minute, save.esType());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(NodeRefResSumMinuteSave.class.getSimpleName(), NodeRefResSumMinuteSave.Role.INSTANCE.roleName());
|
||||
Assert.assertEquals(HashCodeSelector.class.getSimpleName(), NodeRefResSumMinuteSave.Role.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(NodeRefResSumMinuteSave.class.getSimpleName(), NodeRefResSumMinuteSave.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(NodeRefResSumMinuteSave.class.getSimpleName(), NodeRefResSumMinuteSave.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
|
||||
int testSize = 10;
|
||||
WorkerConfig.Queue.NodeRef.NodeRefResSumMinuteSave.Size = testSize;
|
||||
Assert.assertEquals(testSize, NodeRefResSumMinuteSave.Factory.INSTANCE.queueSize());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,26 @@
|
|||
package com.a.eye.skywalking.collector.worker.segment;
|
||||
|
||||
import com.a.eye.skywalking.collector.worker.globaltrace.GlobalTraceIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class SegmentCostIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
SegmentCostIndex index = new SegmentCostIndex();
|
||||
Assert.assertEquals("segment_cost_idx", index.index());
|
||||
Assert.assertEquals(true, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
SegmentCostIndex index = new SegmentCostIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"segId\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"startTime\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"EndTime\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"operationName\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"cost\":{\"type\":\"long\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,26 @@
|
|||
package com.a.eye.skywalking.collector.worker.segment;
|
||||
|
||||
import com.a.eye.skywalking.collector.worker.globaltrace.GlobalTraceIndex;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class SegmentExceptionIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
SegmentExceptionIndex index = new SegmentExceptionIndex();
|
||||
Assert.assertEquals("segment_exp_idx", index.index());
|
||||
Assert.assertEquals(true, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
SegmentExceptionIndex index = new SegmentExceptionIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"segId\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"isError\":{\"type\":\"boolean\",\"index\":\"not_analyzed\"},\"errorKind\":{\"type\":\"string\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,25 @@
|
|||
package com.a.eye.skywalking.collector.worker.segment;
|
||||
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
import java.io.IOException;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class SegmentIndexTestCase {
|
||||
|
||||
@Test
|
||||
public void test() {
|
||||
SegmentIndex index = new SegmentIndex();
|
||||
Assert.assertEquals("segment_idx", index.index());
|
||||
Assert.assertEquals(true, index.isRecord());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testBuilder() throws IOException {
|
||||
SegmentIndex index = new SegmentIndex();
|
||||
Assert.assertEquals("{\"properties\":{\"traceSegmentId\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"startTime\":{\"type\":\"date\",\"index\":\"not_analyzed\"},\"endTime\":{\"type\":\"date\",\"index\":\"not_analyzed\"},\"applicationCode\":{\"type\":\"string\",\"index\":\"not_analyzed\"},\"minute\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"hour\":{\"type\":\"long\",\"index\":\"not_analyzed\"},\"day\":{\"type\":\"long\",\"index\":\"not_analyzed\"}}}", index.createMappingBuilder().string());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,17 @@
|
|||
package com.a.eye.skywalking.collector.worker.span;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Test;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class SpanGetWithIdTestCase {
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(SpanGetWithId.class.getSimpleName(), SpanGetWithId.WorkerRole.INSTANCE.roleName());
|
||||
Assert.assertEquals(RollingSelector.class.getSimpleName(), SpanGetWithId.WorkerRole.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,88 @@
|
|||
package com.a.eye.skywalking.collector.worker.span.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.RollingSelector;
|
||||
import com.a.eye.skywalking.collector.worker.Const;
|
||||
import com.a.eye.skywalking.collector.worker.segment.SegmentIndex;
|
||||
import com.a.eye.skywalking.collector.worker.storage.GetResponseFromEs;
|
||||
import com.a.eye.skywalking.trace.Span;
|
||||
import com.a.eye.skywalking.trace.TraceSegment;
|
||||
import com.google.gson.Gson;
|
||||
import com.google.gson.JsonObject;
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
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 org.powermock.reflect.Whitebox;
|
||||
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.when;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
@RunWith(PowerMockRunner.class)
|
||||
@PrepareForTest({GetResponseFromEs.class})
|
||||
@PowerMockIgnore({"javax.management.*"})
|
||||
public class SpanSearchWithIdTestCase {
|
||||
|
||||
private GetResponseFromEs getResponseFromEs;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
getResponseFromEs = PowerMockito.mock(GetResponseFromEs.class);
|
||||
Whitebox.setInternalState(GetResponseFromEs.class, "INSTANCE", getResponseFromEs);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRole() {
|
||||
Assert.assertEquals(SpanSearchWithId.class.getSimpleName(), SpanSearchWithId.WorkerRole.INSTANCE.roleName());
|
||||
Assert.assertEquals(RollingSelector.class.getSimpleName(), SpanSearchWithId.WorkerRole.INSTANCE.workerSelector().getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testFactory() {
|
||||
Assert.assertEquals(SpanSearchWithId.class.getSimpleName(), SpanSearchWithId.Factory.INSTANCE.role().roleName());
|
||||
Assert.assertEquals(SpanSearchWithId.class.getSimpleName(), SpanSearchWithId.Factory.INSTANCE.workerInstance(null).getClass().getSimpleName());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOnWork() throws Exception {
|
||||
ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext(null);
|
||||
LocalWorkerContext localWorkerContext = new LocalWorkerContext();
|
||||
SpanSearchWithId spanSearchWithId = new SpanSearchWithId(SpanSearchWithId.WorkerRole.INSTANCE, clusterWorkerContext, localWorkerContext);
|
||||
|
||||
TraceSegment segment = create();
|
||||
Gson gson = new Gson();
|
||||
String sourceString = gson.toJson(segment);
|
||||
|
||||
GetResponse getResponse = mock(GetResponse.class);
|
||||
when(getResponseFromEs.get(SegmentIndex.Index, SegmentIndex.Type_Record, "1")).thenReturn(getResponse);
|
||||
when(getResponse.getSourceAsString()).thenReturn(sourceString);
|
||||
|
||||
SpanSearchWithId.RequestEntity request = new SpanSearchWithId.RequestEntity("1", "0");
|
||||
JsonObject response = new JsonObject();
|
||||
spanSearchWithId.onWork(request, response);
|
||||
|
||||
JsonObject segJsonObj = response.get(Const.RESULT).getAsJsonObject();
|
||||
String value = segJsonObj.get("ts").getAsJsonObject().get("Tag").getAsString();
|
||||
Assert.assertEquals("Value", value);
|
||||
}
|
||||
|
||||
private TraceSegment create() {
|
||||
TraceSegment segment = new TraceSegment();
|
||||
|
||||
Span span = new Span();
|
||||
span.setTag("Tag", "Value");
|
||||
span.finish(segment);
|
||||
segment.finish();
|
||||
|
||||
return segment;
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,34 @@
|
|||
package com.a.eye.skywalking.collector.worker.storage;
|
||||
|
||||
import com.a.eye.skywalking.collector.worker.mock.MockGetResponse;
|
||||
import org.elasticsearch.action.get.GetResponse;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.powermock.core.classloader.annotations.PowerMockIgnore;
|
||||
import org.powermock.core.classloader.annotations.PrepareForTest;
|
||||
import org.powermock.modules.junit4.PowerMockRunner;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
@RunWith(PowerMockRunner.class)
|
||||
@PrepareForTest({EsClient.class})
|
||||
@PowerMockIgnore({"javax.management.*"})
|
||||
public class GetResponseFromEsTestCase {
|
||||
|
||||
private GetResponse getResponse;
|
||||
|
||||
@Before
|
||||
public void init() {
|
||||
MockGetResponse mockGetResponse = new MockGetResponse();
|
||||
getResponse = mockGetResponse.mockito();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testGet() {
|
||||
GetResponse response = GetResponseFromEs.INSTANCE.get("INDEX", "TYPE", "1");
|
||||
Assert.assertEquals(getResponse, response);
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue