add timeslice abstract class

This commit is contained in:
pengys5 2017-03-10 11:52:47 +08:00
parent 1061dd6418
commit 813aef07c0
35 changed files with 489 additions and 372 deletions

View File

@ -27,7 +27,7 @@ public abstract class AbstractMember implements EventHandler<MessageHolder> {
this.actorRef = actorRef;
}
public abstract void beTold(Object message) throws Exception;
protected abstract void beTold(Object message) throws Exception;
/**
* Receive the message to analyse.

View File

@ -82,6 +82,14 @@ public abstract class AbstractWorker extends UntypedActor {
ClusterEvent.MemberUp memberUp = (ClusterEvent.MemberUp) message;
logger.info("receive ClusterEvent.MemberUp message, address: %s", memberUp.member().address().toString());
register(memberUp.member());
} else if (message instanceof ClusterEvent.MemberEvent) {
System.out.println("other event: " + message.getClass().getSimpleName());
} else if (message instanceof ClusterEvent.UnreachableMember) {
System.out.println("other event: " + message.getClass().getSimpleName());
} else if (message instanceof ClusterEvent.MemberJoined) {
System.out.println("other event: " + message.getClass().getSimpleName());
} else if (message instanceof ClusterEvent.ReachableMember) {
System.out.println("other event: " + message.getClass().getSimpleName());
} else {
logger.debug("worker class: %s, message class: %s", this.getClass().getName(), message.getClass().getName());
receive(message);
@ -101,6 +109,10 @@ public abstract class AbstractWorker extends UntypedActor {
selector.select(availableWorks, message).tell(message, getSelf());
}
public void tell(AbstractMember targetMember, Object message) throws Exception {
targetMember.beTold(message);
}
/**
* When member role is {@link WorkersListener#WorkName} then Select actor from context
* and send register message to {@link WorkersListener}

View File

@ -1,4 +1,4 @@
package com.a.eye.skywalking.collector.actor;
package com.a.eye.skywalking.collector.actor.selector;
/**
* @author pengys5
@ -6,11 +6,11 @@ package com.a.eye.skywalking.collector.actor;
public abstract class AbstractHashMessage {
private int hashCode;
public void setHashCode(String key) {
public AbstractHashMessage(String key) {
this.hashCode = key.hashCode();
}
public int getHashCode() {
protected int getHashCode() {
return hashCode;
}
}

View File

@ -1,6 +1,5 @@
package com.a.eye.skywalking.collector.actor.selector;
import com.a.eye.skywalking.collector.actor.AbstractHashMessage;
import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.WorkerRef;

View File

@ -24,6 +24,7 @@ public enum WorkersRefCenter {
private Map<ActorRef, WorkerRef> actorRefToWorkerRef = new ConcurrentHashMap<>();
public void register(ActorRef newActorRef, String workerRole) {
System.out.println("register: " + workerRole);
if (!roleToWorkerRef.containsKey(workerRole)) {
List<WorkerRef> actorList = Collections.synchronizedList(new ArrayList<WorkerRef>());
roleToWorkerRef.putIfAbsent(workerRole, actorList);

View File

@ -5,23 +5,29 @@ import com.a.eye.skywalking.collector.actor.WorkersCreator;
import com.a.eye.skywalking.collector.cluster.ClusterConfig;
import com.a.eye.skywalking.collector.cluster.ClusterConfigInitializer;
import com.a.eye.skywalking.collector.cluster.NoAvailableWorkerException;
import com.a.eye.skywalking.collector.worker.storage.EsClient;
import com.typesafe.config.Config;
import com.typesafe.config.ConfigFactory;
import java.net.UnknownHostException;
/**
* @author pengys5
*/
public class CollectorBootStartUp {
public static void main(String[] args) throws NoAvailableWorkerException, InterruptedException {
public static void main(String[] args) throws NoAvailableWorkerException, InterruptedException, UnknownHostException {
ClusterConfigInitializer.initialize("collector.config");
final Config config = ConfigFactory.parseString("akka.remote.netty.tcp.port=" + ClusterConfig.Cluster.Current.port).
withFallback(ConfigFactory.parseString("akka.cluster.roles = [" + ClusterConfig.Cluster.Current.roles + "]")).
withFallback(ConfigFactory.load());
final Config config = ConfigFactory.parseString("akka.remote.netty.tcp.hostname=" + ClusterConfig.Cluster.Current.hostname).
withFallback(ConfigFactory.parseString("akka.remote.netty.tcp.port=" + ClusterConfig.Cluster.Current.port)).
withFallback(ConfigFactory.parseString("akka.cluster.roles=" + ClusterConfig.Cluster.Current.roles)).
withFallback(ConfigFactory.parseString("akka.actor.provider=" + ClusterConfig.Cluster.provider)).
withFallback(ConfigFactory.parseString("akka.cluster.seed-nodes=" + ClusterConfig.Cluster.nodes)).
withFallback(ConfigFactory.load("application.conf"));
ActorSystem system = ActorSystem.create(ClusterConfig.Cluster.appname, config);
WorkersCreator.INSTANCE.boot(system);
EsClient.boot();
}
}

View File

@ -2,13 +2,12 @@ package com.a.eye.skywalking.collector.worker;
import akka.actor.ActorRef;
import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.storage.MetricData;
import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.util.Map;
/**
* @author pengys5
*/
@ -23,23 +22,15 @@ public abstract class MetricAnalysisMember extends AnalysisMember {
}
public void setMetric(String id, int second, Long value) throws Exception {
persistenceData.setMetric(id, second, value);
persistenceData.getElseCreate(id).setMetric(second, value);
if (persistenceData.size() >= WorkerConfig.Persistence.Data.size) {
aggregation();
}
}
public MetricPersistenceData pushOneMetric() {
if (persistenceData.getData().entrySet().iterator().hasNext()) {
Map.Entry<String, Map<String, Long>> entry = persistenceData.getData().entrySet().iterator().next();
MetricPersistenceData oneRecord = new MetricPersistenceData();
for (Map.Entry<String, Long> entry1 : entry.getValue().entrySet()) {
oneRecord.setMetric(entry.getKey(), entry1.getKey(), entry1.getValue());
}
oneRecord.setHashCode(entry.getKey());
persistenceData.getData().remove(entry.getKey());
return oneRecord;
public MetricData pushOne() {
if (persistenceData.iterator().hasNext()) {
return persistenceData.pushOne();
}
return null;
}

View File

@ -2,12 +2,21 @@ package com.a.eye.skywalking.collector.worker;
import akka.actor.ActorRef;
import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.storage.EsClient;
import com.a.eye.skywalking.collector.worker.storage.MetricData;
import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData;
import com.a.eye.skywalking.collector.worker.tools.PersistenceDataTools;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.elasticsearch.action.bulk.BulkRequestBuilder;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.get.MultiGetItemResponse;
import org.elasticsearch.action.get.MultiGetRequestBuilder;
import org.elasticsearch.action.get.MultiGetResponse;
import org.elasticsearch.client.Client;
import java.util.Iterator;
import java.util.Map;
/**
@ -25,35 +34,57 @@ public abstract class MetricPersistenceMember extends PersistenceMember {
@Override
public void analyse(Object message) throws Exception {
if (message instanceof MetricPersistenceData) {
MetricPersistenceData persistenceData = (MetricPersistenceData) message;
merge(persistenceData);
if (message instanceof MetricData) {
MetricData metricData = (MetricData) message;
persistenceData.getElseCreate(metricData.getId()).merge(metricData);
if (persistenceData.size() >= WorkerConfig.Persistence.Data.size) {
persistence();
}
} else {
logger.error("message unhandled");
}
}
public void merge(MetricPersistenceData receiveData) {
for (Map.Entry<String, Map<String, Long>> lineDate : receiveData.getData().entrySet()) {
for (Map.Entry<String, Long> columnDate : lineDate.getValue().entrySet()) {
persistenceData.setMetric(lineDate.getKey(), columnDate.getKey(), columnDate.getValue());
if (persistenceData.size() >= WorkerConfig.Persistence.Data.size) {
persistence();
}
protected void persistence() {
MultiGetResponse multiGetResponse = searchFromEs();
for (MultiGetItemResponse itemResponse : multiGetResponse) {
GetResponse response = itemResponse.getResponse();
if (response != null && response.isExists()) {
persistenceData.getElseCreate(response.getId()).merge(response.getSource());
}
}
boolean success = saveToEs();
if (success) {
persistenceData.clear();
}
}
protected void persistence() {
if (persistenceData.size() > 0) {
Map<String, Map<String, Object>> dataInDB = PersistenceDataTools.searchEs(esIndex(), esType(), persistenceData);
MetricPersistenceData dbData = PersistenceDataTools.dbData2PersistenceData(dataInDB);
PersistenceDataTools.mergeData(dbData, persistenceData);
public MultiGetResponse searchFromEs() {
Client client = EsClient.getClient();
MultiGetRequestBuilder multiGetRequestBuilder = client.prepareMultiGet();
boolean success = PersistenceDataTools.saveToEs(esIndex(), esType(), persistenceData);
if (success) {
persistenceData.clear();
}
Iterator<Map.Entry<String, MetricData>> iterator = persistenceData.iterator();
while (iterator.hasNext()) {
multiGetRequestBuilder.add(esIndex(), esType(), iterator.next().getKey());
}
MultiGetResponse multiGetResponse = multiGetRequestBuilder.get();
return multiGetResponse;
}
public boolean saveToEs() {
Client client = EsClient.getClient();
BulkRequestBuilder bulkRequest = client.prepareBulk();
logger.debug("persistenceData size: %s", persistenceData.size());
Iterator<Map.Entry<String, MetricData>> iterator = persistenceData.iterator();
while (iterator.hasNext()) {
MetricData metricData = iterator.next().getValue();
bulkRequest.add(client.prepareIndex(esIndex(), esType(), metricData.getId()).setSource(metricData.toMap()));
}
BulkResponse bulkResponse = bulkRequest.execute().actionGet();
return !bulkResponse.hasFailures();
}
}

View File

@ -2,14 +2,13 @@ package com.a.eye.skywalking.collector.worker;
import akka.actor.ActorRef;
import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.google.gson.JsonObject;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.util.Map;
/**
* @author pengys5
*/
@ -24,21 +23,15 @@ public abstract class RecordAnalysisMember extends AnalysisMember {
}
public void setRecord(String id, JsonObject record) throws Exception {
persistenceData.setMetric(id, record);
persistenceData.getElseCreate(id).setRecord(record);
if (persistenceData.size() >= WorkerConfig.Analysis.Data.size) {
aggregation();
}
}
public RecordPersistenceData pushOneRecord() {
if (persistenceData.getData().entrySet().iterator().hasNext()) {
Map.Entry<String, JsonObject> entry = persistenceData.getData().entrySet().iterator().next();
RecordPersistenceData oneRecord = new RecordPersistenceData();
oneRecord.setMetric(entry.getKey(), entry.getValue());
oneRecord.setHashCode(entry.getKey());
persistenceData.getData().remove(entry.getKey());
return oneRecord;
public RecordData pushOne() {
if (persistenceData.hasNext()) {
return persistenceData.pushOne();
}
return null;
}

View File

@ -2,13 +2,17 @@ package com.a.eye.skywalking.collector.worker;
import akka.actor.ActorRef;
import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.storage.EsClient;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.a.eye.skywalking.collector.worker.tools.PersistenceDataTools;
import com.google.gson.JsonObject;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.elasticsearch.action.bulk.BulkRequestBuilder;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.client.Client;
import java.util.Iterator;
import java.util.Map;
/**
@ -24,38 +28,39 @@ public abstract class RecordPersistenceMember extends PersistenceMember {
super(ringBuffer, actorRef);
}
public void setRecord(String id, JsonObject record) {
persistenceData.setMetric(id, record);
if (persistenceData.size() >= WorkerConfig.Persistence.Data.size) {
persistence();
}
}
@Override
public void analyse(Object message) throws Exception {
if (message instanceof RecordPersistenceData) {
RecordPersistenceData persistenceData = (RecordPersistenceData) message;
merge(persistenceData);
if (message instanceof RecordData) {
RecordData recordData = (RecordData) message;
persistenceData.getElseCreate(recordData.getId()).setRecord(recordData.getRecord());
if (persistenceData.size() >= WorkerConfig.Persistence.Data.size) {
persistence();
}
} else {
logger.error("message unhandled");
}
}
public void merge(RecordPersistenceData receiveData) {
for (Map.Entry<String, JsonObject> lineDate : receiveData.getData().entrySet()) {
persistenceData.setMetric(lineDate.getKey(), lineDate.getValue());
if (persistenceData.size() >= WorkerConfig.Persistence.Data.size) {
persistence();
}
protected void persistence() {
boolean success = saveToEs();
if (success) {
persistenceData.clear();
}
}
protected void persistence() {
if (persistenceData.size() > 0) {
boolean success = PersistenceDataTools.saveToEs(esIndex(), esType(), persistenceData);
if (success) {
persistenceData.clear();
}
public boolean saveToEs() {
Client client = EsClient.getClient();
BulkRequestBuilder bulkRequest = client.prepareBulk();
logger.debug("persistenceData size: %s", persistenceData.size());
Iterator<Map.Entry<String, RecordData>> iterator = persistenceData.iterator();
while (iterator.hasNext()) {
Map.Entry<String, RecordData> recordData = iterator.next();
bulkRequest.add(client.prepareIndex(esIndex(), esType(), recordData.getKey()).setSource(recordData.getValue().getRecord().toString()));
}
BulkResponse bulkResponse = bulkRequest.execute().actionGet();
return !bulkResponse.hasFailures();
}
}

View File

@ -21,27 +21,27 @@ public class WorkerConfig extends ClusterConfig {
public static class Worker {
public static class TraceSegmentReceiver {
public static int Num = 5;
public static int Num = 10;
}
public static class DAGNodeReceiver {
public static int Num = 5;
public static int Num = 10;
}
public static class NodeInstanceReceiver {
public static int Num = 5;
public static int Num = 10;
}
public static class ResponseCostReceiver {
public static int Num = 5;
public static int Num = 10;
}
public static class ResponseSummaryReceiver {
public static int Num = 5;
public static int Num = 10;
}
public static class DAGNodeRefReceiver {
public static int Num = 5;
public static int Num = 10;
}
}

View File

@ -1,6 +1,7 @@
package com.a.eye.skywalking.collector.worker.application;
import akka.actor.ActorRef;
import com.a.eye.skywalking.api.util.StringUtil;
import com.a.eye.skywalking.collector.actor.AbstractSyncMember;
import com.a.eye.skywalking.collector.actor.AbstractSyncMemberProvider;
import com.a.eye.skywalking.collector.worker.application.analysis.DAGNodeAnalysis;
@ -8,9 +9,9 @@ import com.a.eye.skywalking.collector.worker.application.analysis.NodeInstanceAn
import com.a.eye.skywalking.collector.worker.application.analysis.ResponseCostAnalysis;
import com.a.eye.skywalking.collector.worker.application.analysis.ResponseSummaryAnalysis;
import com.a.eye.skywalking.collector.worker.application.persistence.TraceSegmentRecordPersistence;
import com.a.eye.skywalking.collector.worker.tools.DateTools;
import com.a.eye.skywalking.collector.worker.receiver.TraceSegmentReceiver;
import com.a.eye.skywalking.trace.Span;
import com.a.eye.skywalking.trace.TraceSegment;
import com.a.eye.skywalking.trace.TraceSegmentRef;
import com.a.eye.skywalking.trace.tag.Tags;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@ -39,17 +40,16 @@ public class ApplicationMain extends AbstractSyncMember {
@Override
public void receive(Object message) throws Exception {
if (message instanceof TraceSegment) {
if (message instanceof TraceSegmentReceiver.TraceSegmentTimeSlice) {
logger.debug("begin translate TraceSegment Object to JsonObject");
TraceSegment traceSegment = (TraceSegment) message;
int second = DateTools.timeStampToSecond(traceSegment.getStartTime());
TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment = (TraceSegmentReceiver.TraceSegmentTimeSlice) message;
recordPersistence.beTold(traceSegment);
sendToDAGNodePersistence(traceSegment);
sendToNodeInstanceAnalysis(traceSegment);
sendToResponseCostPersistence(traceSegment, second);
sendToResponseSummaryPersistence(traceSegment, second);
sendToResponseCostPersistence(traceSegment);
sendToResponseSummaryPersistence(traceSegment);
}
}
@ -62,40 +62,42 @@ public class ApplicationMain extends AbstractSyncMember {
}
}
private void sendToDAGNodePersistence(TraceSegment traceSegment) throws Exception {
String code = traceSegment.getApplicationCode();
private void sendToDAGNodePersistence(TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment) throws Exception {
String code = traceSegment.getTraceSegment().getApplicationCode();
String component = null;
String layer = null;
for (Span span : traceSegment.getSpans()) {
for (Span span : traceSegment.getTraceSegment().getSpans()) {
if (span.getParentSpanId() == -1) {
component = Tags.COMPONENT.get(span);
layer = Tags.SPAN_LAYER.get(span);
}
}
DAGNodeAnalysis.Metric node = new DAGNodeAnalysis.Metric(code, component, layer);
DAGNodeAnalysis.Metric node = new DAGNodeAnalysis.Metric(traceSegment.getMinute(), traceSegment.getSecond(), code, component, layer);
dagNodeAnalysis.beTold(node);
}
private void sendToNodeInstanceAnalysis(TraceSegment traceSegment) throws Exception {
if (traceSegment.getPrimaryRef() != null) {
String code = traceSegment.getPrimaryRef().getApplicationCode();
String address = traceSegment.getPrimaryRef().getPeerHost();
private void sendToNodeInstanceAnalysis(TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment) throws Exception {
TraceSegmentRef traceSegmentRef = traceSegment.getTraceSegment().getPrimaryRef();
NodeInstanceAnalysis.Metric property = new NodeInstanceAnalysis.Metric(code, address);
if (traceSegmentRef != null && !StringUtil.isEmpty(traceSegmentRef.getApplicationCode())) {
String code = traceSegmentRef.getApplicationCode();
String address = traceSegmentRef.getPeerHost();
NodeInstanceAnalysis.Metric property = new NodeInstanceAnalysis.Metric(traceSegment.getMinute(), traceSegment.getSecond(), code, address);
nodeInstanceAnalysis.beTold(property);
}
}
private void sendToResponseCostPersistence(TraceSegment traceSegment, int second) throws Exception {
String code = traceSegment.getApplicationCode();
private void sendToResponseCostPersistence(TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment) throws Exception {
String code = traceSegment.getTraceSegment().getApplicationCode();
long startTime = -1;
long endTime = -1;
Boolean isError = false;
for (Span span : traceSegment.getSpans()) {
for (Span span : traceSegment.getTraceSegment().getSpans()) {
if (span.getParentSpanId() == -1) {
startTime = span.getStartTime();
endTime = span.getEndTime();
@ -103,21 +105,21 @@ public class ApplicationMain extends AbstractSyncMember {
}
}
ResponseCostAnalysis.Metric cost = new ResponseCostAnalysis.Metric(code, second, isError, startTime, endTime);
ResponseCostAnalysis.Metric cost = new ResponseCostAnalysis.Metric(traceSegment.getMinute(), traceSegment.getSecond(), code, isError, startTime, endTime);
responseCostAnalysis.beTold(cost);
}
private void sendToResponseSummaryPersistence(TraceSegment traceSegment, int second) throws Exception {
String code = traceSegment.getApplicationCode();
private void sendToResponseSummaryPersistence(TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment) throws Exception {
String code = traceSegment.getTraceSegment().getApplicationCode();
boolean isError = false;
for (Span span : traceSegment.getSpans()) {
for (Span span : traceSegment.getTraceSegment().getSpans()) {
if (span.getParentSpanId() == -1) {
isError = Tags.ERROR.get(span);
}
}
ResponseSummaryAnalysis.Metric summary = new ResponseSummaryAnalysis.Metric(code, second, isError);
ResponseSummaryAnalysis.Metric summary = new ResponseSummaryAnalysis.Metric(traceSegment.getMinute(), traceSegment.getSecond(), code, isError);
responseSummaryAnalysis.beTold(summary);
}
}

View File

@ -6,14 +6,14 @@ import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector;
import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.RecordAnalysisMember;
import com.a.eye.skywalking.collector.worker.application.receiver.DAGNodeReceiver;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.AbstractTimeSlice;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import com.a.eye.skywalking.collector.worker.tools.DateTools;
import com.google.gson.JsonObject;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.io.Serializable;
/**
* @author pengys5
*/
@ -31,11 +31,13 @@ public class DAGNodeAnalysis extends RecordAnalysisMember {
Metric metric = (Metric) message;
JsonObject propertyJsonObj = new JsonObject();
propertyJsonObj.addProperty("code", metric.code);
propertyJsonObj.addProperty(DateTools.Time_Slice_Column_Name, metric.getMinute());
propertyJsonObj.addProperty("component", metric.component);
propertyJsonObj.addProperty("layer", metric.layer);
String id = metric.getMinute() + "-" + metric.code;
logger.debug("dag node: %s", propertyJsonObj.toString());
setRecord(metric.code, propertyJsonObj);
setRecord(id, propertyJsonObj);
} else {
logger.error("message unhandled");
}
@ -43,8 +45,8 @@ public class DAGNodeAnalysis extends RecordAnalysisMember {
@Override
protected void aggregation() throws Exception {
RecordPersistenceData oneRecord;
while ((oneRecord = pushOneRecord()) != null) {
RecordData oneRecord;
while ((oneRecord = pushOne()) != null) {
tell(DAGNodeReceiver.Factory.INSTANCE, HashCodeSelector.INSTANCE, oneRecord);
}
}
@ -63,12 +65,13 @@ public class DAGNodeAnalysis extends RecordAnalysisMember {
}
}
public static class Metric implements Serializable {
public static class Metric extends AbstractTimeSlice {
private final String code;
private final String component;
private final String layer;
public Metric(String code, String component, String layer) {
public Metric(long minute, int second, String code, String component, String layer) {
super(minute, second);
this.code = code;
this.component = component;
this.layer = layer;

View File

@ -6,9 +6,10 @@ import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector;
import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.RecordAnalysisMember;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.receiver.DAGNodeReceiver;
import com.a.eye.skywalking.collector.worker.application.receiver.NodeInstanceReceiver;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.AbstractTimeSlice;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import com.a.eye.skywalking.collector.worker.tools.DateTools;
import com.google.gson.JsonObject;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
@ -31,9 +32,11 @@ public class NodeInstanceAnalysis extends RecordAnalysisMember {
Metric metric = (Metric) message;
JsonObject propertyJsonObj = new JsonObject();
propertyJsonObj.addProperty("code", metric.code);
propertyJsonObj.addProperty(DateTools.Time_Slice_Column_Name, metric.getMinute());
propertyJsonObj.addProperty("address", metric.address);
setRecord(metric.address, propertyJsonObj);
String id = metric.getMinute() + "-" + metric.address;
setRecord(id, propertyJsonObj);
logger.debug("node instance: %s", propertyJsonObj.toString());
} else {
logger.error("message unhandled");
@ -42,8 +45,8 @@ public class NodeInstanceAnalysis extends RecordAnalysisMember {
@Override
protected void aggregation() throws Exception {
RecordPersistenceData oneRecord;
while ((oneRecord = pushOneRecord()) != null) {
RecordData oneRecord;
while ((oneRecord = pushOne()) != null) {
tell(NodeInstanceReceiver.Factory.INSTANCE, HashCodeSelector.INSTANCE, oneRecord);
}
}
@ -62,11 +65,12 @@ public class NodeInstanceAnalysis extends RecordAnalysisMember {
}
}
public static class Metric {
public static class Metric extends AbstractTimeSlice{
private final String code;
private final String address;
public Metric(String code, String address) {
public Metric(long minute, int second, String code, String address) {
super(minute, second);
this.code = code;
this.address = address;
}

View File

@ -7,13 +7,12 @@ import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.MetricAnalysisMember;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.receiver.ResponseCostReceiver;
import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.AbstractTimeSlice;
import com.a.eye.skywalking.collector.worker.storage.MetricData;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.io.Serializable;
/**
* @author pengys5
*/
@ -29,9 +28,10 @@ public class ResponseCostAnalysis extends MetricAnalysisMember {
public void analyse(Object message) throws Exception {
if (message instanceof Metric) {
Metric metric = (Metric) message;
long cost = metric.startTime - metric.endTime;
long cost = metric.endTime - metric.startTime;
if (cost <= 1000 && !metric.isError) {
setMetric(metric.code, metric.second, 1L);
String id = metric.getMinute() + "-" + metric.code;
setMetric(id, metric.getSecond(), cost);
}
// logger.debug("response cost metric: %s", data.toString());
}
@ -39,8 +39,8 @@ public class ResponseCostAnalysis extends MetricAnalysisMember {
@Override
protected void aggregation() throws Exception {
MetricPersistenceData oneMetric;
while ((oneMetric = pushOneMetric()) != null) {
MetricData oneMetric;
while ((oneMetric = pushOne()) != null) {
tell(ResponseCostReceiver.Factory.INSTANCE, HashCodeSelector.INSTANCE, oneMetric);
}
}
@ -59,16 +59,15 @@ public class ResponseCostAnalysis extends MetricAnalysisMember {
}
}
public static class Metric implements Serializable {
public static class Metric extends AbstractTimeSlice {
private final String code;
private final int second;
private final Boolean isError;
private final Long startTime;
private final Long endTime;
public Metric(String code, int second, Boolean isError, Long startTime, Long endTime) {
public Metric(long minute, int second, String code, Boolean isError, Long startTime, Long endTime) {
super(minute, second);
this.code = code;
this.second = second;
this.isError = isError;
this.startTime = startTime;
this.endTime = endTime;

View File

@ -7,13 +7,12 @@ import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.MetricAnalysisMember;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.receiver.ResponseSummaryReceiver;
import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.AbstractTimeSlice;
import com.a.eye.skywalking.collector.worker.storage.MetricData;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.io.Serializable;
/**
* @author pengys5
*/
@ -29,16 +28,16 @@ public class ResponseSummaryAnalysis extends MetricAnalysisMember {
public void analyse(Object message) throws Exception {
if (message instanceof Metric) {
Metric metric = (Metric) message;
setMetric(metric.code, metric.second, 1L);
String id = metric.getMinute() + "-" + metric.code;
setMetric(id, metric.getSecond(), 1L);
// logger.debug("response summary metric: %s", data.toString());
}
}
@Override
protected void aggregation() throws Exception {
MetricPersistenceData oneMetric;
while ((oneMetric = pushOneMetric()) != null) {
MetricData oneMetric;
while ((oneMetric = pushOne()) != null) {
tell(ResponseSummaryReceiver.Factory.INSTANCE, HashCodeSelector.INSTANCE, oneMetric);
}
}
@ -57,14 +56,13 @@ public class ResponseSummaryAnalysis extends MetricAnalysisMember {
}
}
public static class Metric implements Serializable {
public static class Metric extends AbstractTimeSlice {
private final String code;
private final int second;
private final Boolean isError;
public Metric(String code, int second, Boolean isError) {
public Metric(long minute, int second, String code, Boolean isError) {
super(minute, second);
this.code = code;
this.second = second;
this.isError = isError;
}
}

View File

@ -6,6 +6,9 @@ import com.a.eye.skywalking.collector.actor.AbstractAsyncMemberProvider;
import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.RecordPersistenceMember;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.receiver.TraceSegmentReceiver;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import com.a.eye.skywalking.collector.worker.tools.DateTools;
import com.a.eye.skywalking.trace.Span;
import com.a.eye.skywalking.trace.TraceSegment;
import com.a.eye.skywalking.trace.TraceSegmentRef;
@ -41,12 +44,13 @@ public class TraceSegmentRecordPersistence extends RecordPersistenceMember {
@Override
public void analyse(Object message) throws Exception {
if (message instanceof TraceSegment) {
TraceSegment traceSegment = (TraceSegment) message;
JsonObject traceSegmentJsonObj = parseTraceSegment(traceSegment);
if (message instanceof TraceSegmentReceiver.TraceSegmentTimeSlice) {
TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment = (TraceSegmentReceiver.TraceSegmentTimeSlice) message;
JsonObject jsonObject = parseTraceSegment(traceSegment.getTraceSegment(), traceSegment.getMinute());
setRecord(traceSegmentJsonObj.get("segmentId").getAsString(), traceSegmentJsonObj);
logger.debug("segment record: %s", traceSegmentJsonObj.toString());
RecordData recordData = new RecordData(traceSegment.getTraceSegment().getTraceSegmentId());
recordData.setRecord(jsonObject);
super.analyse(recordData);
}
}
@ -64,9 +68,10 @@ public class TraceSegmentRecordPersistence extends RecordPersistenceMember {
}
}
private JsonObject parseTraceSegment(TraceSegment traceSegment) {
private JsonObject parseTraceSegment(TraceSegment traceSegment, long minute) {
JsonObject traceJsonObj = new JsonObject();
traceJsonObj.addProperty("segmentId", traceSegment.getTraceSegmentId());
traceJsonObj.addProperty(DateTools.Time_Slice_Column_Name, minute);
traceJsonObj.addProperty("startTime", traceSegment.getStartTime());
traceJsonObj.addProperty("endTime", traceSegment.getEndTime());
traceJsonObj.addProperty("appCode", traceSegment.getApplicationCode());

View File

@ -4,7 +4,7 @@ import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.persistence.DAGNodePersistence;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@ -25,7 +25,7 @@ public class DAGNodeReceiver extends AbstractWorker {
@Override
public void receive(Object message) throws Throwable {
if (message instanceof RecordPersistenceData) {
if (message instanceof RecordData) {
persistence.beTold(message);
} else {
logger.error("message unhandled");

View File

@ -4,7 +4,7 @@ import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.persistence.NodeInstancePersistence;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@ -25,7 +25,7 @@ public class NodeInstanceReceiver extends AbstractWorker {
@Override
public void receive(Object message) throws Throwable {
if (message instanceof RecordPersistenceData) {
if (message instanceof RecordData) {
persistence.beTold(message);
} else {
logger.error("message unhandled");

View File

@ -4,7 +4,7 @@ import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.persistence.ResponseCostPersistence;
import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.MetricData;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@ -25,7 +25,7 @@ public class ResponseCostReceiver extends AbstractWorker {
@Override
public void receive(Object message) throws Throwable {
if (message instanceof MetricPersistenceData) {
if (message instanceof MetricData) {
persistence.beTold(message);
} else {
logger.error("message unhandled");

View File

@ -4,7 +4,7 @@ import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.persistence.ResponseSummaryPersistence;
import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.MetricData;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@ -25,7 +25,7 @@ public class ResponseSummaryReceiver extends AbstractWorker {
@Override
public void receive(Object message) throws Throwable {
if (message instanceof MetricPersistenceData) {
if (message instanceof MetricData) {
persistence.beTold(message);
} else {
logger.error("message unhandled");

View File

@ -5,7 +5,8 @@ import com.a.eye.skywalking.api.util.StringUtil;
import com.a.eye.skywalking.collector.actor.AbstractSyncMember;
import com.a.eye.skywalking.collector.actor.AbstractSyncMemberProvider;
import com.a.eye.skywalking.collector.worker.applicationref.analysis.DAGNodeRefAnalysis;
import com.a.eye.skywalking.trace.TraceSegment;
import com.a.eye.skywalking.collector.worker.receiver.TraceSegmentReceiver;
import com.a.eye.skywalking.trace.TraceSegmentRef;
/**
* @author pengys5
@ -21,13 +22,14 @@ public class ApplicationRefMain extends AbstractSyncMember {
@Override
public void receive(Object message) throws Exception {
TraceSegment traceSegment = (TraceSegment) message;
TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment = (TraceSegmentReceiver.TraceSegmentTimeSlice) message;
if (traceSegment.getPrimaryRef() != null && !StringUtil.isEmpty(traceSegment.getPrimaryRef().getApplicationCode())) {
String front = traceSegment.getPrimaryRef().getApplicationCode();
String behind = traceSegment.getApplicationCode();
TraceSegmentRef traceSegmentRef = traceSegment.getTraceSegment().getPrimaryRef();
if (traceSegmentRef != null && !StringUtil.isEmpty(traceSegmentRef.getApplicationCode())) {
String front = traceSegmentRef.getApplicationCode();
String behind = traceSegment.getTraceSegment().getApplicationCode();
DAGNodeRefAnalysis.Metric nodeRef = new DAGNodeRefAnalysis.Metric(front, behind);
DAGNodeRefAnalysis.Metric nodeRef = new DAGNodeRefAnalysis.Metric(traceSegment.getMinute(), traceSegment.getSecond(), front, behind);
dagNodeRefAnalysis.beTold(nodeRef);
}
}

View File

@ -7,14 +7,14 @@ import com.a.eye.skywalking.collector.queue.MessageHolder;
import com.a.eye.skywalking.collector.worker.RecordAnalysisMember;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.applicationref.receiver.DAGNodeRefReceiver;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.AbstractTimeSlice;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import com.a.eye.skywalking.collector.worker.tools.DateTools;
import com.google.gson.JsonObject;
import com.lmax.disruptor.RingBuffer;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import java.io.Serializable;
/**
* @author pengys5
*/
@ -28,21 +28,23 @@ public class DAGNodeRefAnalysis extends RecordAnalysisMember {
@Override
public void analyse(Object message) throws Exception {
if (message instanceof RecordPersistenceData) {
if (message instanceof Metric) {
Metric metric = (Metric) message;
JsonObject propertyJsonObj = new JsonObject();
propertyJsonObj.addProperty("frontCode", metric.frontCode);
propertyJsonObj.addProperty("behindCode", metric.behindCode);
propertyJsonObj.addProperty(DateTools.Time_Slice_Column_Name, metric.getMinute());
setRecord(metric.frontCode + "-" + metric.behindCode, propertyJsonObj);
String id = metric.getMinute() + "-" + metric.frontCode + "-" + metric.behindCode;
setRecord(id, propertyJsonObj);
logger.debug("dag node ref: %s", propertyJsonObj.toString());
}
}
@Override
protected void aggregation() throws Exception {
RecordPersistenceData oneRecord;
while ((oneRecord = pushOneRecord()) != null) {
RecordData oneRecord;
while ((oneRecord = pushOne()) != null) {
tell(DAGNodeRefReceiver.Factory.INSTANCE, HashCodeSelector.INSTANCE, oneRecord);
}
}
@ -62,11 +64,12 @@ public class DAGNodeRefAnalysis extends RecordAnalysisMember {
}
}
public static class Metric implements Serializable {
public static class Metric extends AbstractTimeSlice {
private final String frontCode;
private final String behindCode;
public Metric(String frontCode, String behindCode) {
public Metric(long minute, int second, String frontCode, String behindCode) {
super(minute, second);
this.frontCode = frontCode;
this.behindCode = behindCode;
}

View File

@ -4,7 +4,7 @@ import com.a.eye.skywalking.collector.actor.AbstractWorker;
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.applicationref.persistence.DAGNodeRefPersistence;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.RecordData;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@ -25,7 +25,7 @@ public class DAGNodeRefReceiver extends AbstractWorker {
@Override
public void receive(Object message) throws Throwable {
if (message instanceof RecordPersistenceData) {
if (message instanceof RecordData) {
persistence.beTold(message);
} else {
logger.error("message unhandled");

View File

@ -5,6 +5,8 @@ import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
import com.a.eye.skywalking.collector.worker.WorkerConfig;
import com.a.eye.skywalking.collector.worker.application.ApplicationMain;
import com.a.eye.skywalking.collector.worker.applicationref.ApplicationRefMain;
import com.a.eye.skywalking.collector.worker.storage.AbstractTimeSlice;
import com.a.eye.skywalking.collector.worker.tools.DateTools;
import com.a.eye.skywalking.trace.TraceSegment;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
@ -31,9 +33,12 @@ public class TraceSegmentReceiver extends AbstractWorker {
if (message instanceof TraceSegment) {
TraceSegment traceSegment = (TraceSegment) message;
logger.debug("receive message instanceof TraceSegment, traceSegmentId is %s", traceSegment.getTraceSegmentId());
long timeSlice = DateTools.timeStampToTimeSlice(traceSegment.getStartTime());
int second = DateTools.timeStampToSecond(traceSegment.getStartTime());
applicationMain.beTold(traceSegment);
applicationRefMain.beTold(traceSegment);
TraceSegmentTimeSlice segmentTimeSlice = new TraceSegmentTimeSlice(timeSlice, second, traceSegment);
tell(applicationMain, segmentTimeSlice);
tell(applicationRefMain, segmentTimeSlice);
}
}
@ -50,4 +55,17 @@ public class TraceSegmentReceiver extends AbstractWorker {
return WorkerConfig.Worker.TraceSegmentReceiver.Num;
}
}
public static class TraceSegmentTimeSlice extends AbstractTimeSlice {
private final TraceSegment traceSegment;
public TraceSegmentTimeSlice(long timeSliceMinute, int second, TraceSegment traceSegment) {
super(timeSliceMinute, second);
this.traceSegment = traceSegment;
}
public TraceSegment getTraceSegment() {
return traceSegment;
}
}
}

View File

@ -1,14 +0,0 @@
package com.a.eye.skywalking.collector.worker.storage;
/**
* @author pengys5
*/
public abstract class AbstractMetricData {
private final String timeMinute;
private final int timeSecond;
public AbstractMetricData(String timeMinute, int timeSecond) {
this.timeMinute = timeMinute;
this.timeSecond = timeSecond;
}
}

View File

@ -1,9 +0,0 @@
package com.a.eye.skywalking.collector.worker.storage;
import java.io.Serializable;
/**
* @author pengys5
*/
public abstract class AbstractMetricStorage {
}

View File

@ -0,0 +1,22 @@
package com.a.eye.skywalking.collector.worker.storage;
/**
* @author pengys5
*/
public abstract class AbstractTimeSlice {
private final long minute;
private final int second;
public AbstractTimeSlice(long minute, int second) {
this.minute = minute;
this.second = second;
}
public long getMinute() {
return minute;
}
public int getSecond() {
return second;
}
}

View File

@ -0,0 +1,106 @@
package com.a.eye.skywalking.collector.worker.storage;
import com.a.eye.skywalking.collector.actor.selector.AbstractHashMessage;
import java.util.HashMap;
import java.util.Map;
/**
* @author pengys5
*/
public class MetricData extends AbstractHashMessage {
public MetricData(String key) {
super(key);
this.id = key;
}
private String id;
private static final String s10 = "s10";
private static final String s20 = "s20";
private static final String s30 = "s30";
private static final String s40 = "s40";
private static final String s50 = "s50";
private static final String s60 = "s60";
private Long s10Value = 0L;
private Long s20Value = 0L;
private Long s30Value = 0L;
private Long s40Value = 0L;
private Long s50Value = 0L;
private Long s60Value = 0L;
public void setMetric(int second, Long value) {
if (second <= 10) {
s10Value += value;
} else if (second > 10 && second <= 20) {
s20Value += value;
} else if (second > 20 && second <= 30) {
s30Value += value;
} else if (second > 30 && second <= 40) {
s40Value += value;
} else if (second > 40 && second <= 50) {
s50Value += value;
} else {
s60Value += value;
}
}
public void merge(MetricData metricData) {
s10Value += metricData.s10Value;
s20Value += metricData.s20Value;
s30Value += metricData.s30Value;
s40Value += metricData.s40Value;
s50Value += metricData.s50Value;
s60Value += metricData.s60Value;
}
public void merge(Map<String, Object> dbData) {
s10Value += Long.valueOf(dbData.get(s10).toString());
s20Value += Long.valueOf(dbData.get(s20).toString());
s30Value += Long.valueOf(dbData.get(s30).toString());
s40Value += Long.valueOf(dbData.get(s40).toString());
s50Value += Long.valueOf(dbData.get(s50).toString());
s60Value += Long.valueOf(dbData.get(s60).toString());
}
public Map<String, Long> toMap() {
Map<String, Long> map = new HashMap<>();
map.put(s10, s10Value);
map.put(s20, s20Value);
map.put(s30, s30Value);
map.put(s40, s40Value);
map.put(s50, s50Value);
map.put(s60, s60Value);
return map;
}
public String getId() {
return id;
}
protected Long getS10Value() {
return s10Value;
}
protected Long getS20Value() {
return s20Value;
}
protected Long getS30Value() {
return s30Value;
}
protected Long getS40Value() {
return s40Value;
}
protected Long getS50Value() {
return s50Value;
}
protected Long getS60Value() {
return s60Value;
}
}

View File

@ -1,43 +1,23 @@
package com.a.eye.skywalking.collector.worker.storage;
import com.a.eye.skywalking.collector.actor.AbstractHashMessage;
import com.a.eye.skywalking.collector.worker.tools.PersistenceDataTools;
import java.util.HashMap;
import java.util.Iterator;
import java.util.Map;
import java.util.Spliterator;
import java.util.function.Consumer;
/**
* @author pengys5
*/
public class MetricPersistenceData extends AbstractHashMessage {
public class MetricPersistenceData implements Iterable {
private Map<String, Map<String, Long>> persistenceData = new HashMap();
private Map<String, MetricData> persistenceData = new HashMap();
public void setMetric(String id, int second, Long value) {
if (persistenceData.containsKey(id)) {
String columnName = PersistenceDataTools.second2ColumnName(second);
Long metric = persistenceData.get(id).get(columnName);
persistenceData.get(id).put(columnName, metric + value);
} else {
Map<String, Long> metrics = PersistenceDataTools.getFilledPersistenceData();
metrics.put(PersistenceDataTools.second2ColumnName(second), value);
persistenceData.put(id, metrics);
public MetricData getElseCreate(String id) {
if (!persistenceData.containsKey(id)) {
persistenceData.put(id, new MetricData(id));
}
}
public void setMetric(String id, String column, Long value) {
if (persistenceData.containsKey(id)) {
Long metric = persistenceData.get(id).get(column);
persistenceData.get(id).put(column, metric + value);
} else {
Map<String, Long> metrics = PersistenceDataTools.getFilledPersistenceData();
metrics.put(column, value);
persistenceData.put(id, metrics);
}
}
public Map<String, Map<String, Long>> getData() {
return persistenceData;
return persistenceData.get(id);
}
public int size() {
@ -47,4 +27,25 @@ public class MetricPersistenceData extends AbstractHashMessage {
public void clear() {
persistenceData.clear();
}
public MetricData pushOne() {
MetricData one = persistenceData.entrySet().iterator().next().getValue();
persistenceData.remove(one.getId());
return one;
}
@Override
public void forEach(Consumer action) {
throw new UnsupportedOperationException("forEach");
}
@Override
public Spliterator spliterator() {
throw new UnsupportedOperationException("spliterator");
}
@Override
public Iterator<Map.Entry<String, MetricData>> iterator() {
return persistenceData.entrySet().iterator();
}
}

View File

@ -0,0 +1,30 @@
package com.a.eye.skywalking.collector.worker.storage;
import com.a.eye.skywalking.collector.actor.selector.AbstractHashMessage;
import com.google.gson.JsonObject;
/**
* @author pengys5
*/
public class RecordData extends AbstractHashMessage {
private String id;
private JsonObject record;
public RecordData(String key) {
super(key);
this.id = key;
}
public String getId() {
return id;
}
public JsonObject getRecord() {
return record;
}
public void setRecord(JsonObject record) {
this.record = record;
}
}

View File

@ -1,23 +1,23 @@
package com.a.eye.skywalking.collector.worker.storage;
import com.a.eye.skywalking.collector.actor.AbstractHashMessage;
import com.google.gson.JsonObject;
import java.util.HashMap;
import java.util.Iterator;
import java.util.Map;
import java.util.Spliterator;
import java.util.function.Consumer;
/**
* @author pengys5
*/
public class RecordPersistenceData extends AbstractHashMessage {
private Map<String, JsonObject> persistenceData = new HashMap();
public class RecordPersistenceData implements Iterable {
public void setMetric(String id, JsonObject record) {
persistenceData.put(id, record);
}
private Map<String, RecordData> persistenceData = new HashMap();
public Map<String, JsonObject> getData() {
return persistenceData;
public RecordData getElseCreate(String id) {
if (!persistenceData.containsKey(id)) {
persistenceData.put(id, new RecordData(id));
}
return persistenceData.get(id);
}
public int size() {
@ -27,4 +27,29 @@ public class RecordPersistenceData extends AbstractHashMessage {
public void clear() {
persistenceData.clear();
}
public boolean hasNext() {
return persistenceData.entrySet().iterator().hasNext();
}
public RecordData pushOne() {
RecordData one = persistenceData.entrySet().iterator().next().getValue();
persistenceData.remove(one.getId());
return one;
}
@Override
public void forEach(Consumer action) {
throw new UnsupportedOperationException("forEach");
}
@Override
public Spliterator spliterator() {
throw new UnsupportedOperationException("spliterator");
}
@Override
public Iterator<Map.Entry<String, RecordData>> iterator() {
return persistenceData.entrySet().iterator();
}
}

View File

@ -1,14 +1,27 @@
package com.a.eye.skywalking.collector.worker.tools;
import java.text.SimpleDateFormat;
import java.util.Calendar;
/**
* @author pengys5
*/
public class DateTools {
private static final SimpleDateFormat sdf = new SimpleDateFormat("yyyyMMddHHmm");
public static final String Time_Slice_Column_Name = "timeSlice";
public static int timeStampToSecond(long time) {
Calendar calendar = Calendar.getInstance();
calendar.setTimeInMillis(time);
return calendar.get(Calendar.SECOND);
}
public static long timeStampToTimeSlice(long time) {
Calendar calendar = Calendar.getInstance();
calendar.setTimeInMillis(time);
String timeStr = sdf.format(calendar.getTime());
return Long.valueOf(timeStr);
}
}

View File

@ -1,129 +0,0 @@
package com.a.eye.skywalking.collector.worker.tools;
import com.a.eye.skywalking.collector.worker.storage.EsClient;
import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData;
import com.a.eye.skywalking.collector.worker.storage.RecordPersistenceData;
import com.google.gson.JsonObject;
import org.apache.logging.log4j.LogManager;
import org.apache.logging.log4j.Logger;
import org.elasticsearch.action.bulk.BulkRequestBuilder;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.get.GetResponse;
import org.elasticsearch.action.get.MultiGetItemResponse;
import org.elasticsearch.action.get.MultiGetRequestBuilder;
import org.elasticsearch.action.get.MultiGetResponse;
import org.elasticsearch.client.Client;
import java.util.HashMap;
import java.util.Map;
/**
* @author pengys5
*/
public class PersistenceDataTools {
private static Logger logger = LogManager.getFormatterLogger(PersistenceDataTools.class);
private static final String s10 = "s10";
private static final String s20 = "s20";
private static final String s30 = "s30";
private static final String s40 = "s40";
private static final String s50 = "s50";
private static final String s60 = "s60";
public static Map<String, Long> getFilledPersistenceData() {
Map<String, Long> columns = new HashMap();
columns.put(s10, 0L);
columns.put(s20, 0L);
columns.put(s30, 0L);
columns.put(s40, 0L);
columns.put(s50, 0L);
columns.put(s60, 0L);
return columns;
}
public static String second2ColumnName(int second) {
if (second <= 10) {
return s10;
} else if (second > 10 && second <= 20) {
return s20;
} else if (second > 20 && second <= 30) {
return s30;
} else if (second > 30 && second <= 40) {
return s40;
} else if (second > 40 && second <= 50) {
return s50;
} else {
return s60;
}
}
public static boolean saveToEs(String esIndex, String esType, MetricPersistenceData persistenceData) {
Client client = EsClient.getClient();
BulkRequestBuilder bulkRequest = client.prepareBulk();
logger.debug("persistenceData size: %s", persistenceData.size());
for (Map.Entry<String, Map<String, Long>> entry : persistenceData.getData().entrySet()) {
bulkRequest.add(client.prepareIndex(esIndex, esType, entry.getKey()).setSource(entry.getValue()));
}
BulkResponse bulkResponse = bulkRequest.execute().actionGet();
return !bulkResponse.hasFailures();
}
public static boolean saveToEs(String esIndex, String esType, RecordPersistenceData persistenceData) {
Client client = EsClient.getClient();
BulkRequestBuilder bulkRequest = client.prepareBulk();
logger.debug("persistenceData size: %s", persistenceData.size());
for (Map.Entry<String, JsonObject> entry : persistenceData.getData().entrySet()) {
logger.debug("record: %s", entry.getValue().toString());
bulkRequest.add(client.prepareIndex(esIndex, esType, entry.getKey()).setSource(entry.getValue().toString()));
}
BulkResponse bulkResponse = bulkRequest.execute().actionGet();
return !bulkResponse.hasFailures();
}
public static Map<String, Map<String, Object>> searchEs(String esIndex, String esType, MetricPersistenceData persistenceData) {
Client client = EsClient.getClient();
Map<String, Map<String, Object>> dataInEs = new HashMap();
MultiGetRequestBuilder multiGetRequestBuilder = client.prepareMultiGet();
for (Map.Entry<String, Map<String, Long>> entry : persistenceData.getData().entrySet()) {
multiGetRequestBuilder.add(esIndex, esType, entry.getKey());
}
MultiGetResponse multiGetResponse = multiGetRequestBuilder.get();
for (MultiGetItemResponse itemResponse : multiGetResponse) {
GetResponse response = itemResponse.getResponse();
if (response != null && response.isExists()) {
dataInEs.put(response.getId(), response.getSource());
}
}
return dataInEs;
}
public static MetricPersistenceData dbData2PersistenceData(Map<String, Map<String, Object>> dbData) {
MetricPersistenceData persistenceData = new MetricPersistenceData();
for (Map.Entry<String, Map<String, Object>> entryLines : dbData.entrySet()) {
for (Map.Entry<String, Object> entryColumns : entryLines.getValue().entrySet()) {
persistenceData.setMetric(entryLines.getKey(), entryColumns.getKey(), Long.valueOf(entryColumns.getValue().toString()));
}
}
return persistenceData;
}
public static void mergeData(MetricPersistenceData dbData, MetricPersistenceData memoryData) {
for (Map.Entry<String, Map<String, Long>> memoryEntry : memoryData.getData().entrySet()) {
String id = memoryEntry.getKey();
if (dbData.getData().containsKey(id)) {
for (Map.Entry<String, Long> memoryMetricEntry : memoryEntry.getValue().entrySet()) {
memoryMetricEntry.setValue(dbData.getData().get(id).get(memoryMetricEntry.getKey()) + memoryMetricEntry.getValue());
}
}
}
}
}

View File

@ -23,8 +23,8 @@ import org.junit.Test;
*/
public class StartUpTestCase {
@Test
public void test() throws Exception {
System.out.println(TraceSegmentReceiver.class.getSimpleName());
ClusterConfigInitializer.initialize("collector.config");
System.out.println(ClusterConfig.Cluster.Current.roles);