From 813aef07c07f9a489adc33df1e7107da5a8e3ca2 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Fri, 10 Mar 2017 11:52:47 +0800 Subject: [PATCH] add timeslice abstract class --- .../collector/actor/AbstractMember.java | 2 +- .../collector/actor/AbstractWorker.java | 12 ++ .../{ => selector}/AbstractHashMessage.java | 6 +- .../actor/selector/HashCodeSelector.java | 1 - .../collector/cluster/WorkersRefCenter.java | 1 + .../worker/CollectorBootStartUp.java | 16 ++- .../worker/MetricAnalysisMember.java | 19 +-- .../worker/MetricPersistenceMember.java | 71 +++++++--- .../worker/RecordAnalysisMember.java | 17 +-- .../worker/RecordPersistenceMember.java | 53 +++---- .../collector/worker/WorkerConfig.java | 12 +- .../worker/application/ApplicationMain.java | 50 +++---- .../application/analysis/DAGNodeAnalysis.java | 19 +-- .../analysis/NodeInstanceAnalysis.java | 18 ++- .../analysis/ResponseCostAnalysis.java | 21 ++- .../analysis/ResponseSummaryAnalysis.java | 20 ++- .../TraceSegmentRecordPersistence.java | 17 ++- .../application/receiver/DAGNodeReceiver.java | 4 +- .../receiver/NodeInstanceReceiver.java | 4 +- .../receiver/ResponseCostReceiver.java | 4 +- .../receiver/ResponseSummaryReceiver.java | 4 +- .../applicationref/ApplicationRefMain.java | 14 +- .../analysis/DAGNodeRefAnalysis.java | 21 +-- .../receiver/DAGNodeRefReceiver.java | 4 +- .../worker/receiver/TraceSegmentReceiver.java | 22 ++- .../worker/storage/AbstractMetricData.java | 14 -- .../worker/storage/AbstractMetricStorage.java | 9 -- .../worker/storage/AbstractTimeSlice.java | 22 +++ .../collector/worker/storage/MetricData.java | 106 ++++++++++++++ .../worker/storage/MetricPersistenceData.java | 59 ++++---- .../collector/worker/storage/RecordData.java | 30 ++++ .../worker/storage/RecordPersistenceData.java | 45 ++++-- .../collector/worker/tools/DateTools.java | 13 ++ .../worker/tools/PersistenceDataTools.java | 129 ------------------ .../collector/worker/StartUpTestCase.java | 2 +- 35 files changed, 489 insertions(+), 372 deletions(-) rename skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/{ => selector}/AbstractHashMessage.java (57%) delete mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricData.java delete mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricStorage.java create mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractTimeSlice.java create mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricData.java create mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordData.java delete mode 100644 skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/PersistenceDataTools.java diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMember.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMember.java index eb39ecf5d..6f7c72a6d 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMember.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMember.java @@ -27,7 +27,7 @@ public abstract class AbstractMember implements EventHandler { this.actorRef = actorRef; } - public abstract void beTold(Object message) throws Exception; + protected abstract void beTold(Object message) throws Exception; /** * Receive the message to analyse. diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java index 670d0a1ba..a483b0f74 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java @@ -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} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractHashMessage.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java similarity index 57% rename from skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractHashMessage.java rename to skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java index bd0ec3292..0a5f64a14 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractHashMessage.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java @@ -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; } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java index 33d6b81f3..61877163d 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java @@ -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; diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java index d9d0007c6..99a209855 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java @@ -24,6 +24,7 @@ public enum WorkersRefCenter { private Map actorRefToWorkerRef = new ConcurrentHashMap<>(); public void register(ActorRef newActorRef, String workerRole) { + System.out.println("register: " + workerRole); if (!roleToWorkerRef.containsKey(workerRole)) { List actorList = Collections.synchronizedList(new ArrayList()); roleToWorkerRef.putIfAbsent(workerRole, actorList); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/CollectorBootStartUp.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/CollectorBootStartUp.java index be13321e6..ceab05e73 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/CollectorBootStartUp.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/CollectorBootStartUp.java @@ -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(); } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java index 24c35f1f8..8dd0a9eda 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java @@ -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> entry = persistenceData.getData().entrySet().iterator().next(); - MetricPersistenceData oneRecord = new MetricPersistenceData(); - for (Map.Entry 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; } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java index be981ed44..d1e88f573 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java @@ -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> lineDate : receiveData.getData().entrySet()) { - for (Map.Entry 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> 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> 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> 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(); } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java index 83a8693b3..272cc9656 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java @@ -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 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; } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java index 47e49a72b..36db75d24 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java @@ -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 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> iterator = persistenceData.iterator(); + + while (iterator.hasNext()) { + Map.Entry recordData = iterator.next(); + bulkRequest.add(client.prepareIndex(esIndex(), esType(), recordData.getKey()).setSource(recordData.getValue().getRecord().toString())); } + + BulkResponse bulkResponse = bulkRequest.execute().actionGet(); + return !bulkResponse.hasFailures(); } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java index 6c60941b0..d92ba1340 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/WorkerConfig.java @@ -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; } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java index 6a55494f4..d488088a7 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java @@ -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); } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java index 705dd2224..4f409a542 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java @@ -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; diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java index ad24fa7d4..fe7dd8d92 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java @@ -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; } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java index 0b703d39a..8eba11256 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java @@ -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; diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java index 86cc1bb7f..08a360177 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java @@ -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; } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java index 0fbcde70b..d6a5dfcf1 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java @@ -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()); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java index 56f6d410a..e63bd6edf 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java @@ -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"); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java index ee2a773ef..977601519 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java @@ -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"); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java index 2037da3e5..f086d393e 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java @@ -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"); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java index a9a918955..e264e6e0d 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java @@ -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"); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java index 9f3c66afc..5fd19b558 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java @@ -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); } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java index bbce14a0b..4fb35491e 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java @@ -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; } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java index c7d35795f..1e4b53850 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java @@ -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"); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiver.java index 02a5952d0..0e16f6b1e 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiver.java @@ -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; + } + } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricData.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricData.java deleted file mode 100644 index d58bde0ec..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricData.java +++ /dev/null @@ -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; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricStorage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricStorage.java deleted file mode 100644 index 99939dc3b..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractMetricStorage.java +++ /dev/null @@ -1,9 +0,0 @@ -package com.a.eye.skywalking.collector.worker.storage; - -import java.io.Serializable; - -/** - * @author pengys5 - */ -public abstract class AbstractMetricStorage { -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractTimeSlice.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractTimeSlice.java new file mode 100644 index 000000000..bb0ade6e1 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/AbstractTimeSlice.java @@ -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; + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricData.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricData.java new file mode 100644 index 000000000..89341c910 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricData.java @@ -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 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 toMap() { + Map 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; + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricPersistenceData.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricPersistenceData.java index feb19d2a9..0b2467f68 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricPersistenceData.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/MetricPersistenceData.java @@ -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> persistenceData = new HashMap(); + private Map 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 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 metrics = PersistenceDataTools.getFilledPersistenceData(); - metrics.put(column, value); - persistenceData.put(id, metrics); - } - } - - public Map> 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> iterator() { + return persistenceData.entrySet().iterator(); + } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordData.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordData.java new file mode 100644 index 000000000..8eba5110b --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordData.java @@ -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; + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordPersistenceData.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordPersistenceData.java index 37e81c941..cfcf54505 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordPersistenceData.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/storage/RecordPersistenceData.java @@ -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 persistenceData = new HashMap(); +public class RecordPersistenceData implements Iterable { - public void setMetric(String id, JsonObject record) { - persistenceData.put(id, record); - } + private Map persistenceData = new HashMap(); - public Map 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> iterator() { + return persistenceData.entrySet().iterator(); + } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/DateTools.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/DateTools.java index 8f9e34a42..bf61c67cf 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/DateTools.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/DateTools.java @@ -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); + } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/PersistenceDataTools.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/PersistenceDataTools.java deleted file mode 100644 index 21fc595ee..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/PersistenceDataTools.java +++ /dev/null @@ -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 getFilledPersistenceData() { - Map 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> 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 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> searchEs(String esIndex, String esType, MetricPersistenceData persistenceData) { - Client client = EsClient.getClient(); - Map> dataInEs = new HashMap(); - - MultiGetRequestBuilder multiGetRequestBuilder = client.prepareMultiGet(); - for (Map.Entry> 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> dbData) { - MetricPersistenceData persistenceData = new MetricPersistenceData(); - - for (Map.Entry> entryLines : dbData.entrySet()) { - for (Map.Entry 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> memoryEntry : memoryData.getData().entrySet()) { - String id = memoryEntry.getKey(); - if (dbData.getData().containsKey(id)) { - for (Map.Entry memoryMetricEntry : memoryEntry.getValue().entrySet()) { - memoryMetricEntry.setValue(dbData.getData().get(id).get(memoryMetricEntry.getKey()) + memoryMetricEntry.getValue()); - } - } - } - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java index 6024079bf..3e24b08bd 100644 --- a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java +++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java @@ -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);