diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java deleted file mode 100644 index 865534c11..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java +++ /dev/null @@ -1,25 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -import akka.actor.ActorSystem; -import akka.actor.Props; - -/** - * @author pengys5 - */ -public abstract class AbstractClusterWorkerProvider extends AbstractWorkerProvider { - - @Override - public void createWorker(ActorSystem system) { - if (workerClass() == null) { - throw new IllegalArgumentException("cannot createInstance() with nothing obtained from workerClass()"); - } - if (workerNum() <= 0) { - throw new IllegalArgumentException("cannot createInstance() with obtained from workerNum() must greater than 0"); - } - - for (int i = 1; i <= workerNum(); i++) { - system.actorOf(Props.create(workerClass()), roleName() + "_" + i); - } - } - -} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalWorker.java deleted file mode 100644 index d7e9c4bc6..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalWorker.java +++ /dev/null @@ -1,28 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; - -/** - * @author pengys5 - */ -public abstract class AbstractLocalWorker implements Worker { - - /** - * Receive the message to analyse. - * - * @param message is the data send from the forward worker - * @throws Throwable is the exception thrown by that worker implementation processing - */ - public abstract void receive(Object message) throws Throwable; - - /** - * Send analysed data to next Worker. - * - * @param targetWorkerProvider is the worker provider to create worker instance. - * @param message is the data used to send to next worker. - * @throws Throwable - */ - public void tell(AbstractLocalWorkerProvider targetWorkerProvider, T message) throws Throwable { - LocalSystem.actorFor(targetWorkerProvider.getClass(), targetWorkerProvider.roleName()); - } -} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalWorkerProvider.java deleted file mode 100644 index a45bffb48..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalWorkerProvider.java +++ /dev/null @@ -1,28 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -import akka.actor.ActorSystem; - -/** - * @author pengys5 - */ -public abstract class AbstractLocalWorkerProvider extends AbstractWorkerProvider { - - /** - * Use {@link ActorSystem} to Create worker instance with the {@link #workerClass()} method returned class. - * - * @param system is a akka {@link ActorSystem} instance. - */ - @Override - public void createWorker(LocalSystem system) { - if (workerClass() == null) { - throw new IllegalArgumentException("cannot createInstance() with nothing obtained from workerClass()"); - } - if (workerNum() <= 0) { - throw new IllegalArgumentException("cannot createInstance() with obtained from workerNum() must greater than 0"); - } - - for (int i = 1; i <= workerNum(); i++) { - LocalSystem.actorOf(getClass(), roleName()); - } - } -} 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 new file mode 100644 index 000000000..9b27882da --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMember.java @@ -0,0 +1,44 @@ +package com.a.eye.skywalking.collector.actor; + +import akka.actor.ActorRef; +import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; +import com.a.eye.skywalking.collector.cluster.WorkersRefCenter; + +import java.util.List; + +/** + * @author pengys5 + */ +public abstract class AbstractMember { + + private ActorRef actorRef; + + public ActorRef getSelf() { + return actorRef; + } + + public void creatorRef(ActorRef actorRef) { + this.actorRef = actorRef; + } + + /** + * Receive the message to analyse. + * + * @param message is the data send from the forward worker + * @throws Throwable is the exception thrown by that worker implementation processing + */ + public abstract void receive(Object message) throws Throwable; + + /** + * Send analysed data to next Worker. + * + * @param targetWorkerProvider is the worker provider to create worker instance. + * @param selector is the selector to select a same role worker instance form cluster. + * @param message is the data used to send to next worker. + * @throws Throwable + */ + public void tell(AbstractWorkerProvider targetWorkerProvider, WorkerSelector selector, T message) throws Throwable { + List availableWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName()); + selector.select(availableWorks, message).tell(message, getSelf()); + } +} 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 new file mode 100644 index 000000000..51336d31c --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractMemberProvider.java @@ -0,0 +1,28 @@ +package com.a.eye.skywalking.collector.actor; + +import akka.actor.ActorRef; + +/** + * @author pengys5 + */ +public abstract class AbstractMemberProvider { + public abstract Class memberClass(); + + public void createWorker(MemberSystem system, ActorRef actorRef) { + if (memberClass() == null) { + throw new IllegalArgumentException("cannot createInstance() with nothing obtained from memberClass()"); + } + + AbstractMember member = system.memberOf(memberClass(), roleName()); + member.creatorRef(actorRef); + } + + /** + * Use {@link #memberClass()} method returned class's simple name as a role name. + * + * @return is role of Worker + */ + protected String roleName() { + return memberClass().getSimpleName(); + } +} 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 37be6e090..0d1a15b32 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 @@ -35,7 +35,9 @@ import java.util.List; * } * }}} */ -public abstract class AbstractWorker extends UntypedActor implements Worker{ +public abstract class AbstractWorker extends UntypedActor { + + private MemberSystem memberSystem = new MemberSystem(); /** * Receive the message to analyse. @@ -76,13 +78,8 @@ public abstract class AbstractWorker extends UntypedActor implements Worker{ * @throws Throwable */ public void tell(AbstractWorkerProvider targetWorkerProvider, WorkerSelector selector, T message) throws Throwable { - if (targetWorkerProvider instanceof AbstractLocalWorkerProvider) { - Worker worker = LocalSystem.actorFor(targetWorkerProvider.getClass(), targetWorkerProvider.roleName()); - worker.receive(message); - } else if (targetWorkerProvider instanceof AbstractClusterWorkerProvider) { - List availableWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName()); - selector.select(availableWorks, message).tell(message, getSelf()); - } + List availableWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName()); + selector.select(availableWorks, message).tell(message, getSelf()); } /** @@ -97,4 +94,8 @@ public abstract class AbstractWorker extends UntypedActor implements Worker{ getContext().actorSelection(member.address() + "/user/" + WorkersListener.WorkName).tell(registerMessage, getSelf()); } } + + public MemberSystem getMemberContext() { + return memberSystem; + } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java index 079cf5c6b..17e58f012 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java @@ -31,7 +31,18 @@ public abstract class AbstractWorkerProvider { public abstract int workerNum(); - public abstract void createWorker(T system); + public void createWorker(ActorSystem system) { + if (workerClass() == null) { + throw new IllegalArgumentException("cannot createInstance() with nothing obtained from workerClass()"); + } + if (workerNum() <= 0) { + throw new IllegalArgumentException("cannot createInstance() with obtained from workerNum() must greater than 0"); + } + + for (int i = 1; i <= workerNum(); i++) { + system.actorOf(Props.create(workerClass()), roleName() + "_" + i); + } + } /** * Use {@link #workerClass()} method returned class's simple name as a role name. diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LocalSystem.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LocalSystem.java deleted file mode 100644 index ef5b5dd15..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LocalSystem.java +++ /dev/null @@ -1,27 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -import java.util.HashMap; -import java.util.Map; - -/** - * @author pengys5 - */ -public class LocalSystem { - - private static Map context = new HashMap(); - - public static void actorOf(Class clazz, String role) { - try { - Worker classInstance = (Worker) clazz.newInstance(); - context.put(clazz.getName() + "_" + role, classInstance); - } catch (InstantiationException e) { - e.printStackTrace(); - } catch (IllegalAccessException e) { - e.printStackTrace(); - } - } - - public static Worker actorFor(Class clazz, String role) { - return context.get(clazz.getName() + "_" + role); - } -} 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 new file mode 100644 index 000000000..c745c46be --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/MemberSystem.java @@ -0,0 +1,29 @@ +package com.a.eye.skywalking.collector.actor; + +import java.util.HashMap; +import java.util.Map; + +/** + * @author pengys5 + */ +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 AbstractMember memberFor(String role) { + return memberMap.get(role); + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Worker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Worker.java deleted file mode 100644 index da6e3fc1c..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Worker.java +++ /dev/null @@ -1,9 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -/** - * @author pengys5 - */ -public interface Worker { - - public void receive(Object message) throws Throwable; -} 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 d80008b3d..97295eb28 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 @@ -23,11 +23,5 @@ public enum WorkersCreator { for (AbstractClusterWorkerProvider provider : clusterServiceLoader) { provider.createWorker(system); } - - LocalSystem localSystem = new LocalSystem(); - ServiceLoader localServiceLoader = ServiceLoader.load(AbstractLocalWorkerProvider.class); - for (AbstractLocalWorkerProvider provider : localServiceLoader) { - provider.createWorker(localSystem); - } } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractLocalWorkerProvider b/skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractMemberProvider similarity index 100% rename from skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractLocalWorkerProvider rename to skywalking-collector/skywalking-collector-cluster/src/test/resources/META-INF/services/com.a.eye.skywalking.collector.actor.AbstractMemberProvider 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 new file mode 100644 index 000000000..8b36383a2 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorker.java @@ -0,0 +1,30 @@ +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/metric/ApplicationDiscoverFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorkerFactory.java similarity index 55% rename from skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/metric/ApplicationDiscoverFactory.java rename to skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorkerFactory.java index 102526cee..f001b2423 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/metric/ApplicationDiscoverFactory.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationWorkerFactory.java @@ -1,15 +1,14 @@ -package com.a.eye.skywalking.collector.worker.metric; +package com.a.eye.skywalking.collector.worker.application; import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; /** * @author pengys5 */ -public class ApplicationDiscoverFactory extends AbstractWorkerProvider { - +public class ApplicationWorkerFactory extends AbstractWorkerProvider { @Override public Class workerClass() { - return ApplicationDiscoverMetric.class; + return ApplicationWorker.class; } @Override 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 new file mode 100644 index 000000000..f4ed569cc --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverFactory.java @@ -0,0 +1,14 @@ +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/metric/ApplicationDiscoverMetric.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverMember.java similarity index 72% rename from skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/metric/ApplicationDiscoverMetric.java rename to skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverMember.java index b18230817..e1d9d0674 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/metric/ApplicationDiscoverMetric.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/member/ApplicationDiscoverMember.java @@ -1,17 +1,17 @@ -package com.a.eye.skywalking.collector.worker.metric; +package com.a.eye.skywalking.collector.worker.application.member; -import com.a.eye.skywalking.collector.actor.AbstractWorker; +import com.a.eye.skywalking.collector.actor.AbstractMember; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; -import com.a.eye.skywalking.collector.worker.persistence.ApplicationMessage; -import com.a.eye.skywalking.collector.worker.persistence.ApplicationPersistenceFactory; +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 ApplicationDiscoverMetric extends AbstractWorker { +public class ApplicationDiscoverMember extends AbstractMember { @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/persistence/ApplicationMessage.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationMessage.java similarity index 90% rename from skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationMessage.java rename to skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationMessage.java index 13f0ce5e9..2ce38e064 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationMessage.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationMessage.java @@ -1,4 +1,4 @@ -package com.a.eye.skywalking.collector.worker.persistence; +package com.a.eye.skywalking.collector.worker.application.persistence; /** * @author pengys5 diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistence.java similarity index 63% rename from skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationPersistence.java rename to skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistence.java index a17fef4cc..5ebc5e3f9 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationPersistence.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistence.java @@ -1,6 +1,7 @@ -package com.a.eye.skywalking.collector.worker.persistence; +package com.a.eye.skywalking.collector.worker.application.persistence; -import com.a.eye.skywalking.collector.actor.AbstractWorker; +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; @@ -8,7 +9,7 @@ import java.util.Map; /** * @author pengys5 */ -public class ApplicationPersistence extends AbstractWorker { +public class ApplicationPersistence extends PersistenceWorker { private Map appData = new HashMap(); diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationPersistenceFactory.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistenceFactory.java similarity index 82% rename from skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationPersistenceFactory.java rename to skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistenceFactory.java index 99d1cfd64..5b208783b 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationPersistenceFactory.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/ApplicationPersistenceFactory.java @@ -1,4 +1,4 @@ -package com.a.eye.skywalking.collector.worker.persistence; +package com.a.eye.skywalking.collector.worker.application.persistence; import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider; 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 new file mode 100644 index 000000000..339790a1b --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorker.java @@ -0,0 +1,13 @@ +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 new file mode 100644 index 000000000..7c19beb07 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefWorkerFactory.java @@ -0,0 +1,18 @@ +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/persistence/ApplicationRefRecord.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/ApplicationRefRecord.java index a53958f55..d72f7f7f9 100644 --- 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 @@ -2,11 +2,9 @@ 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; -import java.util.HashMap; -import java.util.Map; - /** * @author pengys5 */ 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 new file mode 100644 index 000000000..6c678d013 --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/persistence/PersistenceWorker.java @@ -0,0 +1,9 @@ +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 new file mode 100644 index 000000000..db0a10a7e --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiver.java @@ -0,0 +1,17 @@ +package com.a.eye.skywalking.collector.worker.receiver; + +import com.a.eye.skywalking.collector.actor.AbstractWorker; +import com.a.eye.skywalking.trace.TraceSegment; + +/** + * @author pengys5 + */ +public class TraceSegmentReceiver extends AbstractWorker { + + @Override + public void receive(Object message) throws Throwable { + if (message instanceof TraceSegment) { + + } + } +} 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 new file mode 100644 index 000000000..174ea544b --- /dev/null +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/receiver/TraceSegmentReceiverFactory.java @@ -0,0 +1,19 @@ +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; + } +}