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 9b27882da..9c6f743ba 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 @@ -11,16 +11,26 @@ import java.util.List; */ public abstract class AbstractMember { + private MemberSystem memberSystem; + private ActorRef actorRef; + public MemberSystem memberContext() { + return memberSystem; + } + public ActorRef getSelf() { return actorRef; } - public void creatorRef(ActorRef actorRef) { + public AbstractMember(MemberSystem memberSystem, ActorRef actorRef) { + this.memberSystem = memberSystem; this.actorRef = actorRef; } + + public abstract void preStart() throws Throwable; + /** * Receive the message to analyse. * diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMemberProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMemberProvider.java index 51336d31c..9f47d8d35 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMemberProvider.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMemberProvider.java @@ -2,19 +2,23 @@ package com.a.eye.skywalking.collector.actor; import akka.actor.ActorRef; +import java.lang.reflect.Constructor; + /** * @author pengys5 */ public abstract class AbstractMemberProvider { public abstract Class memberClass(); - public void createWorker(MemberSystem system, ActorRef actorRef) { + public void createWorker(MemberSystem system, ActorRef actorRef) throws Exception { if (memberClass() == null) { throw new IllegalArgumentException("cannot createInstance() with nothing obtained from memberClass()"); } - AbstractMember member = system.memberOf(memberClass(), roleName()); - member.creatorRef(actorRef); + Constructor memberConstructor = memberClass().getDeclaredConstructor(new Class[]{MemberSystem.class, ActorRef.class}); + memberConstructor.setAccessible(true); + AbstractMember member = (AbstractMember) memberConstructor.newInstance(system, actorRef); + system.memberOf(member, roleName()); } /** 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 0d1a15b32..f98072a44 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 @@ -95,7 +95,7 @@ public abstract class AbstractWorker extends UntypedActor { } } - public MemberSystem getMemberContext() { + public MemberSystem memberContext() { return memberSystem; } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/MemberSystem.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/MemberSystem.java index c745c46be..3c407bbaa 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/MemberSystem.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/MemberSystem.java @@ -10,17 +10,8 @@ public class MemberSystem { private Map memberMap = new HashMap(); - public AbstractMember memberOf(Class clazz, String role) { - try { - AbstractMember member = (AbstractMember) clazz.newInstance(); - memberMap.put(role, member); - return member; - } catch (InstantiationException e) { - e.printStackTrace(); - } catch (IllegalAccessException e) { - e.printStackTrace(); - } - return null; + public void memberOf(AbstractMember member, String role) { + memberMap.put(role, member); } public AbstractMember memberFor(String role) { diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java index 97295eb28..9a063d95d 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkersCreator.java @@ -19,8 +19,8 @@ public enum WorkersCreator { * @param system is create by akka {@link ActorSystem} */ public void boot(ActorSystem system) { - ServiceLoader clusterServiceLoader = ServiceLoader.load(AbstractClusterWorkerProvider.class); - for (AbstractClusterWorkerProvider provider : clusterServiceLoader) { + ServiceLoader clusterServiceLoader = ServiceLoader.load(AbstractWorkerProvider.class); + for (AbstractWorkerProvider provider : clusterServiceLoader) { provider.createWorker(system); } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/LocalSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/LocalSelector.java new file mode 100644 index 000000000..03f57ee62 --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/LocalSelector.java @@ -0,0 +1,36 @@ +package com.a.eye.skywalking.collector.actor.selector; + +import com.a.eye.skywalking.collector.actor.AbstractWorker; +import com.a.eye.skywalking.collector.actor.WorkerRef; + +import java.util.List; + +/** + * The LocalSelector is a simple implementation of {@link WorkerSelector}. + * It choose {@link WorkerRef} nearly random, by round-robin. + * + * @author wusheng + */ +public enum LocalSelector implements WorkerSelector { + INSTANCE; + + /** + * A simple round variable. + */ + private int index = 0; + + /** + * Use round-robin to select {@link WorkerRef}. + * + * @param members given {@link WorkerRef} list, which size is greater than 0; + * @param message the {@link AbstractWorker} is going to send. + * @return the selected {@link WorkerRef} + */ + @Override + public WorkerRef select(List members, Object message) { + int size = members.size(); + index++; + int selectIndex = Math.abs(index) % size; + return members.get(selectIndex); + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractMemberProvider b/skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractMemberProvider deleted file mode 100644 index 8d06ff8e3..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractMemberProvider +++ /dev/null @@ -1 +0,0 @@ -com.a.eye.skywalking.collector.actor.SpiTestWorkerFactory \ No newline at end of file diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider b/skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractWorkerProvider similarity index 100% rename from skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider rename to skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractWorkerProvider diff --git a/skywalking-collector/skywalking-collector-worker/pom.xml b/skywalking-collector/skywalking-collector-worker/pom.xml index c158a0d58..cd53c5494 100644 --- a/skywalking-collector/skywalking-collector-worker/pom.xml +++ b/skywalking-collector/skywalking-collector-worker/pom.xml @@ -23,5 +23,10 @@ gson 2.8.0 + + org.elasticsearch.client + transport + 5.2.2 + \ No newline at end of file diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/Metric.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/Metric.java deleted file mode 100644 index 97539de42..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/Metric.java +++ /dev/null @@ -1,40 +0,0 @@ -package com.a.eye.skywalking.collector.worker; - -/** - * @author pengys5 - */ -public class Metric { - private String timeSlice; - private String metricName; - private Long metricValue; - - public Metric(String timeSlice, String metricName, Long metricValue) { - this.timeSlice = timeSlice; - this.metricName = metricName; - this.metricValue = metricValue; - } - - public String getTimeSlice() { - return timeSlice; - } - - public void setTimeSlice(String timeSlice) { - this.timeSlice = timeSlice; - } - - public String getMetricName() { - return metricName; - } - - public void setMetricName(String metricName) { - this.metricName = metricName; - } - - public Long getMetricValue() { - return metricValue; - } - - public void setMetricValue(Long metricValue) { - this.metricValue = metricValue; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricCollection.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricCollection.java deleted file mode 100644 index fb5628317..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricCollection.java +++ /dev/null @@ -1,22 +0,0 @@ -package com.a.eye.skywalking.collector.worker; - -import java.util.HashMap; -import java.util.Map; - -/** - * @author pengys5 - */ -public class MetricCollection { - private Map metricMap = new HashMap(); - - public void put(String timeSlice, String name, Long value) { - String timeSliceName = name + timeSlice; - - if (metricMap.containsKey(timeSliceName)) { - Long metric = metricMap.get(timeSliceName).getMetricValue(); - metricMap.get(timeSliceName).setMetricValue(metric + value); - } else { - metricMap.put(timeSliceName, new Metric(timeSlice, name, value)); - } - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceCommand.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceCommand.java new file mode 100644 index 000000000..8b5724457 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceCommand.java @@ -0,0 +1,7 @@ +package com.a.eye.skywalking.collector.worker; + +/** + * @author pengys5 + */ +public class PersistenceCommand { +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceWorker.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceWorker.java new file mode 100644 index 000000000..e3049af66 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceWorker.java @@ -0,0 +1,75 @@ +package com.a.eye.skywalking.collector.worker; + +import com.a.eye.skywalking.collector.actor.AbstractWorker; +import com.a.eye.skywalking.collector.worker.PersistenceCommand; +import com.a.eye.skywalking.collector.worker.tools.EsClient; +import com.google.gson.JsonObject; +import org.elasticsearch.action.bulk.BulkRequestBuilder; +import org.elasticsearch.action.bulk.BulkResponse; + +import java.util.HashMap; +import java.util.Map; + +/** + * @author pengys5 + */ +public abstract class PersistenceWorker extends AbstractWorker { + + private long lastPersistenceTimestamp = 0; + + private Map persistenceData = new HashMap(); + + public abstract String esIndex(); + + public abstract String esType(); + + public void putData(String id, JsonObject data) { + persistenceData.put(id, data); + if (persistenceData.size() >= 1000) { + persistence(true); + } + } + + public boolean containsId(String id) { + return persistenceData.containsKey(id); + } + + public JsonObject getData(String id) { + return persistenceData.get(id); + } + + public abstract void analyse(Object message) throws Throwable; + + @Override + public void receive(Object message) throws Throwable { + if (message instanceof PersistenceCommand) { + persistence(false); + } else { + analyse(message); + } + } + + private void persistence(boolean dataFull) { + long now = System.currentTimeMillis(); + if (now - lastPersistenceTimestamp > 5000 || dataFull) { + boolean success = saveToEs(); + if (success) { + persistenceData.clear(); + lastPersistenceTimestamp = now; + } + } + } + + private boolean saveToEs() { + BulkRequestBuilder bulkRequest = EsClient.client().prepareBulk(); + + for (Map.Entry entry : persistenceData.entrySet()) { + String id = entry.getKey(); + JsonObject data = entry.getValue(); + bulkRequest.add(EsClient.client().prepareIndex(esIndex(), esType(), id).setSource(data.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/RecordCollection.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordCollection.java deleted file mode 100644 index ec068d75f..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordCollection.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker; - -import com.google.gson.JsonObject; - -import java.util.HashMap; -import java.util.Map; - -/** - * @author pengys5 - */ -public class RecordCollection { - - private Map recordMap = new HashMap(); - - public void put(String timeSlice, String primaryKey, JsonObject valueObj) { - recordMap.put(timeSlice + primaryKey, valueObj); - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/TimeSliceMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/TimeSliceMessage.java deleted file mode 100644 index ca85d6e17..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/TimeSliceMessage.java +++ /dev/null @@ -1,16 +0,0 @@ -package com.a.eye.skywalking.collector.worker; - -/** - * @author pengys5 - */ -public abstract class TimeSliceMessage { - private final String timeSlice; - - public TimeSliceMessage(String timeSlice) { - this.timeSlice = timeSlice; - } - - public String getTimeSlice() { - return timeSlice; - } -} 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 551152217..41453510f 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 @@ -7,4 +7,15 @@ import com.a.eye.skywalking.collector.cluster.ClusterConfig; */ public class WorkerConfig extends ClusterConfig { + public static class WorkerNum { + public static int TraceSegmentReceiver_Num = 1; + + public static int DAGNodePersistence_Num = 5; + public static int DAGNodeRefPersistence_Num = 5; + public static int NodeInstancePersistence_Num = 5; + public static int ResponseCostPersistence_Num = 5; + public static int ResponseSummaryPersistence_Num = 5; + public static int TraceSegmentRecordPersistence_Num = 5; + } + } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMember.java new file mode 100644 index 000000000..5caae1fbd --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMember.java @@ -0,0 +1,108 @@ +package com.a.eye.skywalking.collector.worker.application; + +import akka.actor.ActorRef; +import com.a.eye.skywalking.collector.actor.AbstractMember; +import com.a.eye.skywalking.collector.actor.AbstractMemberProvider; +import com.a.eye.skywalking.collector.actor.MemberSystem; +import com.a.eye.skywalking.collector.actor.selector.RollingSelector; +import com.a.eye.skywalking.collector.worker.application.metric.TraceSegmentRecordMember; +import com.a.eye.skywalking.collector.worker.application.persistence.DAGNodePersistence; +import com.a.eye.skywalking.collector.worker.application.persistence.NodeInstancePersistence; +import com.a.eye.skywalking.collector.worker.application.persistence.ResponseCostPersistence; +import com.a.eye.skywalking.collector.worker.application.persistence.ResponseSummaryPersistence; +import com.a.eye.skywalking.trace.Span; +import com.a.eye.skywalking.trace.TraceSegment; +import com.a.eye.skywalking.trace.tag.Tags; + +/** + * @author pengys5 + */ +public class ApplicationMember extends AbstractMember { + + public ApplicationMember(MemberSystem memberSystem, ActorRef actorRef) { + super(memberSystem, actorRef); + } + + @Override + public void preStart() throws Throwable { + TraceSegmentRecordMember.Factory factory = new TraceSegmentRecordMember.Factory(); + factory.createWorker(memberContext(), getSelf()); + } + + @Override + public void receive(Object message) throws Throwable { + if (message instanceof TraceSegment) { + TraceSegment traceSegment = (TraceSegment) message; + AbstractMember discoverMember = memberContext().memberFor(TraceSegmentRecordMember.class.getSimpleName()); + discoverMember.receive(traceSegment); + + sendToDAGNodePersistence(traceSegment); + sendToNodeInstancePersistence(traceSegment); + sendToResponseCostPersistence(traceSegment); + sendToResponseSummaryPersistence(traceSegment); + } + } + + public static class Factory extends AbstractMemberProvider { + @Override + public Class memberClass() { + return ApplicationMember.class; + } + } + + private void sendToDAGNodePersistence(TraceSegment traceSegment) throws Throwable { + String code = traceSegment.getApplicationCode(); + + String component = null; + String layer = null; + for (Span span : traceSegment.getSpans()) { + if (span.getParentSpanId() == -1) { + component = Tags.COMPONENT.get(span); + layer = Tags.SPAN_LAYER.get(span); + } + } + + DAGNodePersistence.Metric node = new DAGNodePersistence.Metric(code, component, layer); + tell(new NodeInstancePersistence.Factory(), RollingSelector.INSTANCE, node); + } + + private void sendToNodeInstancePersistence(TraceSegment traceSegment) throws Throwable { + String code = traceSegment.getPrimaryRef().getApplicationCode(); + String address = traceSegment.getPrimaryRef().getPeerHost(); + + NodeInstancePersistence.Metric property = new NodeInstancePersistence.Metric(code, address); + tell(new NodeInstancePersistence.Factory(), RollingSelector.INSTANCE, property); + } + + private void sendToResponseCostPersistence(TraceSegment traceSegment) throws Throwable { + String code = traceSegment.getApplicationCode(); + long startTime = -1; + long endTime = -1; + boolean isError = false; + + for (Span span : traceSegment.getSpans()) { + if (span.getParentSpanId() == -1) { + startTime = span.getStartTime(); + endTime = span.getEndTime(); + isError = Tags.ERROR.get(span); + } + } + + ResponseCostPersistence.Metric cost = new ResponseCostPersistence.Metric(code, isError, startTime, endTime); + tell(new ResponseCostPersistence.Factory(), RollingSelector.INSTANCE, cost); + } + + private void sendToResponseSummaryPersistence(TraceSegment traceSegment) throws Throwable { + String code = traceSegment.getApplicationCode(); + boolean isError = false; + + for (Span span : traceSegment.getSpans()) { + if (span.getParentSpanId() == -1) { + isError = Tags.ERROR.get(span); + } + } + + ResponseSummaryPersistence.Metric summary = new ResponseSummaryPersistence.Metric(code, isError); + tell(new ResponseSummaryPersistence.Factory(), RollingSelector.INSTANCE, summary); + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorker.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorker.java deleted file mode 100644 index 8b36383a2..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorker.java +++ /dev/null @@ -1,30 +0,0 @@ -package com.a.eye.skywalking.collector.worker.application; - -import com.a.eye.skywalking.collector.actor.AbstractMember; -import com.a.eye.skywalking.collector.actor.AbstractWorker; -import com.a.eye.skywalking.collector.worker.application.member.ApplicationDiscoverFactory; -import com.a.eye.skywalking.collector.worker.application.member.ApplicationDiscoverMember; -import com.a.eye.skywalking.trace.TraceSegment; - -/** - * @author pengys5 - */ -public class ApplicationWorker extends AbstractWorker { - - @Override - public void preStart() throws Exception { - ApplicationDiscoverFactory factory = new ApplicationDiscoverFactory(); - factory.createWorker(getMemberContext(), getSelf()); - - super.preStart(); - } - - @Override - public void receive(Object message) throws Throwable { - if (message instanceof TraceSegment) { - TraceSegment traceSegment = (TraceSegment) message; - AbstractMember discoverMember = getMemberContext().memberFor(ApplicationDiscoverMember.class.getSimpleName()); - discoverMember.receive(traceSegment); - } - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorkerFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorkerFactory.java deleted file mode 100644 index f001b2423..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorkerFactory.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker.application; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; - -/** - * @author pengys5 - */ -public class ApplicationWorkerFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return ApplicationWorker.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverFactory.java deleted file mode 100644 index f4ed569cc..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverFactory.java +++ /dev/null @@ -1,14 +0,0 @@ -package com.a.eye.skywalking.collector.worker.application.member; - -import com.a.eye.skywalking.collector.actor.AbstractMemberProvider; - -/** - * @author pengys5 - */ -public class ApplicationDiscoverFactory extends AbstractMemberProvider { - - @Override - public Class memberClass() { - return ApplicationDiscoverMember.class; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverMember.java deleted file mode 100644 index e1d9d0674..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverMember.java +++ /dev/null @@ -1,31 +0,0 @@ -package com.a.eye.skywalking.collector.worker.application.member; - - -import com.a.eye.skywalking.collector.actor.AbstractMember; -import com.a.eye.skywalking.collector.actor.selector.RollingSelector; -import com.a.eye.skywalking.collector.worker.application.persistence.ApplicationMessage; -import com.a.eye.skywalking.collector.worker.application.persistence.ApplicationPersistenceFactory; -import com.a.eye.skywalking.trace.TraceSegment; -import com.a.eye.skywalking.trace.tag.Tags; - -/** - * @author pengys5 - */ -public class ApplicationDiscoverMember extends AbstractMember { - - @Override - public void receive(Object message) throws Throwable { - if (message instanceof TraceSegment) { - TraceSegment traceSegment = (TraceSegment) message; - String code = traceSegment.getApplicationCode(); - String component = Tags.COMPONENT.get(traceSegment.getSpans().get(0)); - String host = Tags.PEER_HOST.get(traceSegment.getSpans().get(0)); - int port = Tags.PEER_PORT.get(traceSegment.getSpans().get(0)); - String layer = Tags.SPAN_LAYER.get(traceSegment.getSpans().get(0)); - - ApplicationMessage applicationMessage = new ApplicationMessage(code, component, host, layer); - tell(new ApplicationPersistenceFactory(), RollingSelector.INSTANCE, applicationMessage); - } - } - -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecord.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/metric/TraceSegmentRecordMember.java similarity index 76% rename from skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecord.java rename to skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/metric/TraceSegmentRecordMember.java index c7ad1b55c..703946cb0 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecord.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/metric/TraceSegmentRecordMember.java @@ -1,7 +1,12 @@ -package com.a.eye.skywalking.collector.worker.persistence; +package com.a.eye.skywalking.collector.worker.application.metric; -import com.a.eye.skywalking.collector.actor.AbstractWorker; -import com.a.eye.skywalking.collector.worker.RecordCollection; + +import akka.actor.ActorRef; +import com.a.eye.skywalking.collector.actor.AbstractMember; +import com.a.eye.skywalking.collector.actor.AbstractMemberProvider; +import com.a.eye.skywalking.collector.actor.MemberSystem; +import com.a.eye.skywalking.collector.actor.selector.LocalSelector; +import com.a.eye.skywalking.collector.worker.application.persistence.TraceSegmentRecordPersistence; import com.a.eye.skywalking.trace.Span; import com.a.eye.skywalking.trace.TraceSegment; import com.a.eye.skywalking.trace.TraceSegmentRef; @@ -14,17 +19,30 @@ import java.util.Map; /** * @author pengys5 */ -public class AppTraceSegmentRecord extends AbstractWorker { +public class TraceSegmentRecordMember extends AbstractMember { - private RecordCollection recordCollection = new RecordCollection(); + public TraceSegmentRecordMember(MemberSystem memberSystem, ActorRef actorRef) { + super(memberSystem, actorRef); + } + + @Override + public void preStart() throws Throwable { + } @Override public void receive(Object message) throws Throwable { if (message instanceof TraceSegment) { TraceSegment traceSegment = (TraceSegment) message; + JsonObject traceSegmentJsonObj = parseTraceSegment(traceSegment); - JsonObject traceJsonObj = parseTraceSegment(traceSegment); - recordCollection.put("", traceSegment.getTraceSegmentId(), traceJsonObj); + tell(new TraceSegmentRecordPersistence.Factory(), LocalSelector.INSTANCE, traceSegmentJsonObj); + } + } + + public static class Factory extends AbstractMemberProvider { + @Override + public Class memberClass() { + return TraceSegmentRecordMember.class; } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationMessage.java deleted file mode 100644 index 2ce38e064..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationMessage.java +++ /dev/null @@ -1,34 +0,0 @@ -package com.a.eye.skywalking.collector.worker.application.persistence; - -/** - * @author pengys5 - */ -public class ApplicationMessage { - private final String code; - private final String component; - private final String host; - private final String layer; - - public ApplicationMessage(String code, String component, String host, String layer) { - this.code = code; - this.component = component; - this.host = host; - this.layer = layer; - } - - public String getCode() { - return code; - } - - public String getComponent() { - return component; - } - - public String getHost() { - return host; - } - - public String getLayer() { - return layer; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistence.java deleted file mode 100644 index 5ebc5e3f9..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistence.java +++ /dev/null @@ -1,25 +0,0 @@ -package com.a.eye.skywalking.collector.worker.application.persistence; - -import com.a.eye.skywalking.collector.worker.persistence.PersistenceMessage; -import com.a.eye.skywalking.collector.worker.persistence.PersistenceWorker; - -import java.util.HashMap; -import java.util.Map; - -/** - * @author pengys5 - */ -public class ApplicationPersistence extends PersistenceWorker { - - private Map appData = new HashMap(); - - @Override - public void receive(Object message) throws Throwable { - if (message instanceof ApplicationMessage) { - ApplicationMessage applicationMessage = (ApplicationMessage) message; - appData.put(applicationMessage.getCode(), applicationMessage); - } else if (message instanceof PersistenceMessage) { - - } - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistenceFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistenceFactory.java deleted file mode 100644 index 5b208783b..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistenceFactory.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker.application.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; - -/** - * @author pengys5 - */ -public class ApplicationPersistenceFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return ApplicationPersistence.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/DAGNodePersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/DAGNodePersistence.java new file mode 100644 index 000000000..ce6119e0e --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/DAGNodePersistence.java @@ -0,0 +1,59 @@ +package com.a.eye.skywalking.collector.worker.application.persistence; + +import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; +import com.a.eye.skywalking.collector.worker.WorkerConfig; +import com.a.eye.skywalking.collector.worker.PersistenceWorker; +import com.google.gson.JsonObject; + +/** + * @author pengys5 + */ +public class DAGNodePersistence extends PersistenceWorker { + + @Override + public String esIndex() { + return "application"; + } + + @Override + public String esType() { + return "dag_node"; + } + + @Override + public void analyse(Object message) throws Throwable { + if (message instanceof Metric) { + Metric metric = (Metric) message; + JsonObject propertyJsonObj = new JsonObject(); + propertyJsonObj.addProperty("code", metric.code); + propertyJsonObj.addProperty("component", metric.component); + propertyJsonObj.addProperty("layer", metric.layer); + + putData(metric.code, propertyJsonObj); + } + } + + public static class Factory extends AbstractWorkerProvider { + @Override + public Class workerClass() { + return DAGNodePersistence.class; + } + + @Override + public int workerNum() { + return WorkerConfig.WorkerNum.DAGNodePersistence_Num; + } + } + + public static class Metric { + private final String code; + private final String component; + private final String layer; + + public Metric(String code, String component, String layer) { + 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/persistence/NodeInstancePersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/NodeInstancePersistence.java new file mode 100644 index 000000000..c0895a79e --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/NodeInstancePersistence.java @@ -0,0 +1,56 @@ +package com.a.eye.skywalking.collector.worker.application.persistence; + +import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; +import com.a.eye.skywalking.collector.worker.WorkerConfig; +import com.a.eye.skywalking.collector.worker.PersistenceWorker; +import com.google.gson.JsonObject; + +/** + * @author pengys5 + */ +public class NodeInstancePersistence extends PersistenceWorker { + + @Override + public String esIndex() { + return "application"; + } + + @Override + public String esType() { + return "node_instance"; + } + + @Override + public void analyse(Object message) throws Throwable { + if (message instanceof Metric) { + Metric metric = (Metric) message; + JsonObject propertyJsonObj = new JsonObject(); + propertyJsonObj.addProperty("code", metric.code); + propertyJsonObj.addProperty("address", metric.address); + + putData(metric.address, propertyJsonObj); + } + } + + public static class Factory extends AbstractWorkerProvider { + @Override + public Class workerClass() { + return NodeInstancePersistence.class; + } + + @Override + public int workerNum() { + return WorkerConfig.WorkerNum.NodeInstancePersistence_Num; + } + } + + public static class Metric { + private final String code; + private final String address; + + public Metric(String code, String address) { + this.code = code; + this.address = address; + } + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ResponseCostPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ResponseCostPersistence.java new file mode 100644 index 000000000..56892ba9c --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ResponseCostPersistence.java @@ -0,0 +1,82 @@ +package com.a.eye.skywalking.collector.worker.application.persistence; + +import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; +import com.a.eye.skywalking.collector.worker.WorkerConfig; +import com.a.eye.skywalking.collector.worker.PersistenceWorker; +import com.google.gson.JsonObject; + +/** + * @author pengys5 + */ +public class ResponseCostPersistence extends PersistenceWorker { + + @Override + public String esIndex() { + return "application_metric"; + } + + @Override + public String esType() { + return "response_cost"; + } + + @Override + public void analyse(Object message) throws Throwable { + if (message instanceof Metric) { + Metric metric = (Metric) message; + long cost = metric.startTime - metric.endTime; + JsonObject data; + if (containsId(metric.code)) { + data = getData(metric.code); + } else { + data = new JsonObject(); + } + + String propertyKey = ""; + + if (cost <= 1000 && !metric.isError) { + propertyKey = "one_second_less"; + } else if (cost > 1000 && cost <= 3000 && !metric.isError) { + propertyKey = "three_second_less"; + } else if (cost > 3000 && cost <= 5000 && !metric.isError) { + propertyKey = "five_second_less"; + } else if (cost > 5000 && cost <= 5000 && !metric.isError) { + propertyKey = "slow"; + } else { + propertyKey = "error"; + } + + if (data.has(propertyKey)) { + data.addProperty(propertyKey, data.get(propertyKey).getAsLong() + 1); + } else { + data.addProperty(propertyKey, 1); + } + } + } + + public static class Factory extends AbstractWorkerProvider { + @Override + public Class workerClass() { + return ResponseCostPersistence.class; + } + + @Override + public int workerNum() { + return WorkerConfig.WorkerNum.ResponseCostPersistence_Num; + } + } + + public static class Metric { + private final String code; + private final Boolean isError; + private final Long startTime; + private final Long endTime; + + public Metric(String code, Boolean isError, Long startTime, Long endTime) { + this.code = code; + 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/persistence/ResponseSummaryPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ResponseSummaryPersistence.java new file mode 100644 index 000000000..fa146c31d --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ResponseSummaryPersistence.java @@ -0,0 +1,72 @@ +package com.a.eye.skywalking.collector.worker.application.persistence; + +import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; +import com.a.eye.skywalking.collector.worker.PersistenceWorker; +import com.google.gson.JsonObject; + +import static com.a.eye.skywalking.collector.worker.WorkerConfig.WorkerNum.ResponseSummaryPersistence_Num; + +/** + * @author pengys5 + */ +public class ResponseSummaryPersistence extends PersistenceWorker { + + @Override + public String esIndex() { + return "application_metric"; + } + + @Override + public String esType() { + return "response_summary"; + } + + @Override + public void analyse(Object message) throws Throwable { + if (message instanceof Metric) { + Metric metric = (Metric) message; + + JsonObject data; + if (containsId(metric.code)) { + data = getData(metric.code); + } else { + data = new JsonObject(); + } + + String propertyKey = ""; + if (metric.isError) { + propertyKey = "error"; + } else { + propertyKey = "success"; + } + + if (data.has(propertyKey)) { + data.addProperty(propertyKey, data.get(propertyKey).getAsLong() + 1); + } else { + data.addProperty(propertyKey, 1); + } + } + } + + public static class Factory extends AbstractWorkerProvider { + @Override + public Class workerClass() { + return ResponseSummaryPersistence.class; + } + + @Override + public int workerNum() { + return ResponseSummaryPersistence_Num; + } + } + + public static class Metric { + private final String code; + private final Boolean isError; + + public Metric(String code, Boolean isError) { + this.code = code; + 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 new file mode 100644 index 000000000..2f3697eac --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java @@ -0,0 +1,43 @@ +package com.a.eye.skywalking.collector.worker.application.persistence; + +import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; +import com.a.eye.skywalking.collector.worker.PersistenceWorker; +import com.google.gson.JsonObject; + +import static com.a.eye.skywalking.collector.worker.WorkerConfig.WorkerNum.TraceSegmentRecordPersistence_Num; + +/** + * @author pengys5 + */ +public class TraceSegmentRecordPersistence extends PersistenceWorker { + + @Override + public String esIndex() { + return "application_record"; + } + + @Override + public String esType() { + return "trace_segment"; + } + + @Override + public void analyse(Object message) throws Throwable { + if (message instanceof JsonObject) { + JsonObject traceSegmentJsonObj = (JsonObject) message; + putData(traceSegmentJsonObj.get("segmentId").getAsString(), traceSegmentJsonObj); + } + } + + public static class Factory extends AbstractWorkerProvider { + @Override + public Class workerClass() { + return TraceSegmentRecordPersistence.class; + } + + @Override + public int workerNum() { + return TraceSegmentRecordPersistence_Num; + } + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMember.java new file mode 100644 index 000000000..d6802a1fa --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMember.java @@ -0,0 +1,44 @@ +package com.a.eye.skywalking.collector.worker.applicationref; + +import akka.actor.ActorRef; +import com.a.eye.skywalking.collector.actor.AbstractMember; +import com.a.eye.skywalking.collector.actor.AbstractMemberProvider; +import com.a.eye.skywalking.collector.actor.MemberSystem; +import com.a.eye.skywalking.collector.actor.selector.RollingSelector; +import com.a.eye.skywalking.collector.worker.applicationref.presistence.DAGNodeRefPersistence; +import com.a.eye.skywalking.trace.TraceSegment; + +/** + * @author pengys5 + */ +public class ApplicationRefMember extends AbstractMember { + + public ApplicationRefMember(MemberSystem memberSystem, ActorRef actorRef) { + super(memberSystem, actorRef); + } + + @Override + public void preStart() throws Throwable { + + } + + @Override + public void receive(Object message) throws Throwable { + TraceSegment traceSegment = (TraceSegment) message; + + if (traceSegment.getPrimaryRef() != null) { + String front = traceSegment.getPrimaryRef().getApplicationCode(); + String behind = traceSegment.getApplicationCode(); + + DAGNodeRefPersistence.Metric nodeRef = new DAGNodeRefPersistence.Metric(front, behind); + tell(new DAGNodeRefPersistence.Factory(), RollingSelector.INSTANCE, nodeRef); + } + } + + public static class Factory extends AbstractMemberProvider { + @Override + public Class memberClass() { + return ApplicationRefMember.class; + } + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorker.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorker.java deleted file mode 100644 index 339790a1b..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorker.java +++ /dev/null @@ -1,13 +0,0 @@ -package com.a.eye.skywalking.collector.worker.applicationref; - -import com.a.eye.skywalking.collector.actor.AbstractWorker; - -/** - * @author pengys5 - */ -public class ApplicationRefWorker extends AbstractWorker { - @Override - public void receive(Object message) throws Throwable { - - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorkerFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorkerFactory.java deleted file mode 100644 index 7c19beb07..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorkerFactory.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker.applicationref; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; - -/** - * @author pengys5 - */ -public class ApplicationRefWorkerFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return ApplicationRefWorker.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/presistence/DAGNodeRefPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/presistence/DAGNodeRefPersistence.java new file mode 100644 index 000000000..36f135a76 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/presistence/DAGNodeRefPersistence.java @@ -0,0 +1,56 @@ +package com.a.eye.skywalking.collector.worker.applicationref.presistence; + +import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; +import com.a.eye.skywalking.collector.worker.PersistenceWorker; +import com.a.eye.skywalking.collector.worker.WorkerConfig; +import com.google.gson.JsonObject; + +/** + * @author pengys5 + */ +public class DAGNodeRefPersistence extends PersistenceWorker { + + @Override + public String esIndex() { + return "node_ref"; + } + + @Override + public String esType() { + return "node_ref"; + } + + @Override + public void analyse(Object message) throws Throwable { + if (message instanceof Metric) { + Metric metric = (Metric) message; + JsonObject propertyJsonObj = new JsonObject(); + propertyJsonObj.addProperty("frontCode", metric.frontCode); + propertyJsonObj.addProperty("behindCode", metric.behindCode); + + putData(metric.frontCode + "-" + metric.behindCode, propertyJsonObj); + } + } + + public static class Factory extends AbstractWorkerProvider { + @Override + public Class workerClass() { + return DAGNodeRefPersistence.class; + } + + @Override + public int workerNum() { + return WorkerConfig.WorkerNum.DAGNodeRefPersistence_Num; + } + } + + public static class Metric { + private final String frontCode; + private final String behindCode; + + public Metric(String frontCode, String behindCode) { + this.frontCode = frontCode; + this.behindCode = behindCode; + } + } +} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCost.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCost.java deleted file mode 100644 index f6aa70164..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCost.java +++ /dev/null @@ -1,39 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorker; -import com.a.eye.skywalking.collector.worker.MetricCollection; - -/** - * @author pengys5 - */ -public class AppResponseCost extends AbstractWorker { - - private MetricCollection oneSecondsLessMetric = new MetricCollection(); - - private MetricCollection threeSecondsLessMetric = new MetricCollection(); - - private MetricCollection fiveSecondsLessMetric = new MetricCollection(); - - private MetricCollection slowSecondsLessMetric = new MetricCollection(); - - private MetricCollection errorSecondsLessMetric = new MetricCollection(); - - @Override - public void receive(Object message) throws Throwable { - if (message instanceof AppResponseSummaryMessage) { - AppResponseCostMessage costMessage = (AppResponseCostMessage) message; - long cost = costMessage.getEndTime() - costMessage.getStartTime(); - if (cost <= 1000 && !costMessage.getError()) { - oneSecondsLessMetric.put(costMessage.getTimeSlice(), costMessage.getCode(), cost); - } else if (cost > 1000 && cost <= 3000 && !costMessage.getError()) { - threeSecondsLessMetric.put(costMessage.getTimeSlice(), costMessage.getCode(), cost); - } else if (cost > 3000 && cost <= 5000 && !costMessage.getError()) { - fiveSecondsLessMetric.put(costMessage.getTimeSlice(), costMessage.getCode(), cost); - } else if (cost > 5000 && cost <= 5000 && !costMessage.getError()) { - slowSecondsLessMetric.put(costMessage.getTimeSlice(), costMessage.getCode(), cost); - } else { - errorSecondsLessMetric.put(costMessage.getTimeSlice(), costMessage.getCode(), cost); - } - } - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCostFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCostFactory.java deleted file mode 100644 index 3396ce806..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCostFactory.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; - -/** - * @author pengys5 - */ -public class AppResponseCostFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return AppResponseCost.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCostMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCostMessage.java deleted file mode 100644 index e5a368122..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseCostMessage.java +++ /dev/null @@ -1,37 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.worker.TimeSliceMessage; - -/** - * @author pengys5 - */ -public class AppResponseCostMessage extends TimeSliceMessage { - private final String code; - private final Boolean isError; - private final Long startTime; - private final Long endTime; - - public AppResponseCostMessage(String timeSlice, String code, Boolean isError, Long startTime, Long endTime) { - super(timeSlice); - this.code = code; - this.isError = isError; - this.startTime = startTime; - this.endTime = endTime; - } - - public String getCode() { - return code; - } - - public Boolean getError() { - return isError; - } - - public Long getStartTime() { - return startTime; - } - - public Long getEndTime() { - return endTime; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummary.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummary.java deleted file mode 100644 index 03bebcbbc..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummary.java +++ /dev/null @@ -1,29 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorker; -import com.a.eye.skywalking.collector.worker.MetricCollection; - -/** - * @author pengys5 - */ -public class AppResponseSummary extends AbstractWorker { - - private MetricCollection summaryMetric = new MetricCollection(); - - private MetricCollection errorSummaryMetric = new MetricCollection(); - - private MetricCollection successSummaryMetric = new MetricCollection(); - - @Override - public void receive(Object message) throws Throwable { - if (message instanceof AppResponseSummaryMessage) { - AppResponseSummaryMessage summaryMessage = (AppResponseSummaryMessage) message; - summaryMetric.put(summaryMessage.getTimeSlice(), summaryMessage.getCode(), 1l); - if (summaryMessage.getError()) { - errorSummaryMetric.put(summaryMessage.getTimeSlice(), summaryMessage.getCode(), 1l); - } else { - successSummaryMetric.put(summaryMessage.getTimeSlice(), summaryMessage.getCode(), 1l); - } - } - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummaryFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummaryFactory.java deleted file mode 100644 index 2a7f95b14..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummaryFactory.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; - -/** - * @author pengys5 - */ -public class AppResponseSummaryFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return AppResponseSummary.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummaryMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummaryMessage.java deleted file mode 100644 index 13f30a1e2..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppResponseSummaryMessage.java +++ /dev/null @@ -1,25 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.worker.TimeSliceMessage; - -/** - * @author pengys5 - */ -public class AppResponseSummaryMessage extends TimeSliceMessage { - private final String code; - private final Boolean isError; - - public AppResponseSummaryMessage(String timeSlice, String code, Boolean isError) { - super(timeSlice); - this.code = code; - this.isError = isError; - } - - public String getCode() { - return code; - } - - public Boolean getError() { - return isError; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecordFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecordFactory.java deleted file mode 100644 index 56412daee..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecordFactory.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; - -/** - * @author pengys5 - */ -public class AppTraceSegmentRecordFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return AppTraceSegmentRecord.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecordMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecordMessage.java deleted file mode 100644 index e24661865..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/AppTraceSegmentRecordMessage.java +++ /dev/null @@ -1,12 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.worker.TimeSliceMessage; - -/** - * @author pengys5 - */ -public class AppTraceSegmentRecordMessage extends TimeSliceMessage { - public AppTraceSegmentRecordMessage(String timeSlice) { - super(timeSlice); - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecord.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecord.java deleted file mode 100644 index d72f7f7f9..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecord.java +++ /dev/null @@ -1,24 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorker; -import com.a.eye.skywalking.collector.worker.RecordCollection; -import com.a.eye.skywalking.collector.worker.application.persistence.ApplicationMessage; -import com.google.gson.JsonObject; - -/** - * @author pengys5 - */ -public class ApplicationRefRecord extends AbstractWorker { - - private RecordCollection refRecord = new RecordCollection(); - - @Override - public void receive(Object message) throws Throwable { - if (message instanceof ApplicationMessage) { - ApplicationRefRecordMessage applicationMessage = (ApplicationRefRecordMessage) message; - refRecord.put("", applicationMessage.getCode() + "-" + applicationMessage.getRefCode(), new JsonObject()); - } else if (message instanceof PersistenceMessage) { - - } - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecordFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecordFactory.java deleted file mode 100644 index fe482668c..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecordFactory.java +++ /dev/null @@ -1,18 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; - -/** - * @author pengys5 - */ -public class ApplicationRefRecordFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return ApplicationRefRecord.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecordMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecordMessage.java deleted file mode 100644 index 45761c26d..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecordMessage.java +++ /dev/null @@ -1,22 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -/** - * @author pengys5 - */ -public class ApplicationRefRecordMessage { - private final String code; - private final String refCode; - - public ApplicationRefRecordMessage(String code, String refCode) { - this.code = code; - this.refCode = refCode; - } - - public String getCode() { - return code; - } - - public String getRefCode() { - return refCode; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/PersistenceMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/PersistenceMessage.java deleted file mode 100644 index 75a587a34..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/PersistenceMessage.java +++ /dev/null @@ -1,7 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -/** - * @author pengys5 - */ -public class PersistenceMessage { -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/PersistenceWorker.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/PersistenceWorker.java deleted file mode 100644 index 6c678d013..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/PersistenceWorker.java +++ /dev/null @@ -1,9 +0,0 @@ -package com.a.eye.skywalking.collector.worker.persistence; - -import com.a.eye.skywalking.collector.actor.AbstractWorker; - -/** - * @author pengys5 - */ -public abstract class PersistenceWorker extends AbstractWorker { -} 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 db0a10a7e..a66bf5515 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 @@ -1,6 +1,11 @@ package com.a.eye.skywalking.collector.worker.receiver; +import com.a.eye.skywalking.collector.actor.AbstractMember; 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.ApplicationMember; +import com.a.eye.skywalking.collector.worker.applicationref.ApplicationRefMember; import com.a.eye.skywalking.trace.TraceSegment; /** @@ -8,10 +13,34 @@ import com.a.eye.skywalking.trace.TraceSegment; */ public class TraceSegmentReceiver extends AbstractWorker { + @Override + public void preStart() throws Exception { + ApplicationMember.Factory factory = new ApplicationMember.Factory(); + factory.createWorker(memberContext(), getSelf()); + } + @Override public void receive(Object message) throws Throwable { if (message instanceof TraceSegment) { - + TraceSegment traceSegment = (TraceSegment) message; + + AbstractMember applicationMember = memberContext().memberFor(ApplicationMember.class.getSimpleName()); + applicationMember.receive(traceSegment); + + AbstractMember applicationRefMember = memberContext().memberFor(ApplicationRefMember.class.getSimpleName()); + applicationRefMember.receive(traceSegment); + } + } + + public class Factory extends AbstractWorkerProvider { + @Override + public Class workerClass() { + return TraceSegmentReceiver.class; + } + + @Override + public int workerNum() { + return WorkerConfig.WorkerNum.TraceSegmentReceiver_Num; } } } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiverFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiverFactory.java deleted file mode 100644 index 174ea544b..000000000 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiverFactory.java +++ /dev/null @@ -1,19 +0,0 @@ -package com.a.eye.skywalking.collector.worker.receiver; - -import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; -import com.a.eye.skywalking.collector.worker.application.member.ApplicationDiscoverMember; - -/** - * @author pengys5 - */ -public class TraceSegmentReceiverFactory extends AbstractWorkerProvider { - @Override - public Class workerClass() { - return ApplicationDiscoverMember.class; - } - - @Override - public int workerNum() { - return 0; - } -} diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/EsClient.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/EsClient.java new file mode 100644 index 000000000..b5984ff33 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/tools/EsClient.java @@ -0,0 +1,30 @@ +package com.a.eye.skywalking.collector.worker.tools; + +import org.elasticsearch.client.transport.TransportClient; +import org.elasticsearch.common.settings.Settings; +import org.elasticsearch.common.transport.InetSocketTransportAddress; +import org.elasticsearch.transport.client.PreBuiltTransportClient; + +import java.net.InetAddress; +import java.net.UnknownHostException; + +/** + * @author pengys5 + */ +public class EsClient { + + private static TransportClient client; + + public void boot() throws UnknownHostException { + Settings settings = Settings.builder() + .put("cluster.name", "myClusterName").build(); + + client = new PreBuiltTransportClient(Settings.EMPTY) + .addTransportAddress(new InetSocketTransportAddress(InetAddress.getByName("host1"), 9300)) + .addTransportAddress(new InetSocketTransportAddress(InetAddress.getByName("host1"), 9300)); + } + + public static TransportClient client() { + return client; + } +}