From 7970dce9fb36112d65cb89bd95c4558618a1c4ba Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Fri, 17 Mar 2017 13:04:45 +0800 Subject: [PATCH] refactor the way how to create worker instance --- .../skywalking/collector/CollectorSystem.java | 17 +++++--------- .../actor/AbstractClusterWorker.java | 4 ++-- .../actor/AbstractClusterWorkerProvider.java | 8 ++----- .../actor/AbstractLocalAsyncWorker.java | 22 +++++++++++-------- .../AbstractLocalAsyncWorkerProvider.java | 19 ++++++---------- .../actor/AbstractLocalSyncWorker.java | 4 ++-- .../AbstractLocalSyncWorkerProvider.java | 10 ++------- .../collector/actor/AbstractLocalWorker.java | 4 ++-- .../collector/actor/AbstractWorker.java | 7 +++--- .../actor/AbstractWorkerProvider.java | 13 +++++------ .../collector/actor/ClusterWorkerContext.java | 18 +++++++-------- .../actor/ClusterWorkerRefCounter.java | 6 ++--- .../skywalking/collector/actor/Context.java | 2 +- .../actor/DuplicateProviderException.java | 7 ------ .../collector/actor/LocalWorkerContext.java | 2 +- .../eye/skywalking/collector/actor/Role.java | 6 ++--- .../actor/UsedRoleNameException.java | 7 ++++++ .../collector/actor/WorkerContext.java | 16 ++++++-------- .../actor/selector/HashCodeSelector.java | 12 ++++++---- .../collector/actor/TestClusterWorker.java | 19 ++++++++-------- .../actor/TestClusterWorkerTestCase.java | 12 +++++----- .../collector/actor/TestLocalAsyncWorker.java | 16 +++++++------- .../collector/actor/TestLocalSyncWorker.java | 16 +++++++------- .../actor/TestLocalSyncWorkerTestCase.java | 3 +-- .../role/TraceSegmentReceiverRole.java | 6 ++--- .../collector/worker/AnalysisMember.java | 11 +++++----- .../worker/MetricAnalysisMember.java | 5 +++-- .../worker/MetricPersistenceMember.java | 8 +++---- .../collector/worker/PersistenceMember.java | 11 +++++----- .../worker/RecordAnalysisMember.java | 5 +++-- .../worker/RecordPersistenceMember.java | 5 +++-- .../application/analysis/DAGNodeAnalysis.java | 15 +++++++------ .../analysis/NodeInstanceAnalysis.java | 15 +++++++------ .../analysis/ResponseCostAnalysis.java | 15 +++++++------ .../analysis/ResponseSummaryAnalysis.java | 15 +++++++------ .../persistence/DAGNodePersistence.java | 15 +++++++------ .../persistence/NodeInstancePersistence.java | 17 +++++++------- .../persistence/ResponseCostPersistence.java | 15 +++++++------ .../ResponseSummaryPersistence.java | 15 +++++++------ .../TraceSegmentRecordPersistence.java | 15 +++++++------ .../application/receiver/DAGNodeReceiver.java | 20 ++++++++--------- .../receiver/NodeInstanceReceiver.java | 20 ++++++++--------- .../receiver/ResponseCostReceiver.java | 20 ++++++++--------- .../receiver/ResponseSummaryReceiver.java | 20 ++++++++--------- .../applicationref/ApplicationRefMain.java | 20 ++++++++--------- .../analysis/DAGNodeRefAnalysis.java | 15 +++++++------ .../persistence/DAGNodeRefPersistence.java | 15 +++++++------ .../receiver/DAGNodeRefReceiver.java | 20 ++++++++--------- .../worker/receiver/TraceSegmentReceiver.java | 15 +++++-------- .../collector/worker/StartUpTestCase.java | 5 +---- 50 files changed, 293 insertions(+), 315 deletions(-) delete mode 100644 skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/DuplicateProviderException.java create mode 100644 skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/UsedRoleNameException.java diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/CollectorSystem.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/CollectorSystem.java index cd0b988a7..1c599d5cc 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/CollectorSystem.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/CollectorSystem.java @@ -2,10 +2,7 @@ package com.a.eye.skywalking.collector; import akka.actor.ActorSystem; import akka.actor.Props; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider; -import com.a.eye.skywalking.collector.actor.AbstractLocalWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; -import com.a.eye.skywalking.collector.actor.DuplicateProviderException; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.cluster.ClusterConfig; import com.a.eye.skywalking.collector.cluster.ClusterConfigInitializer; import com.a.eye.skywalking.collector.cluster.WorkersListener; @@ -19,9 +16,7 @@ import java.util.ServiceLoader; /** * @author pengys5 */ -public enum CollectorSystem { - INSTANCE; - +public class CollectorSystem { private Logger logger = LogManager.getFormatterLogger(CollectorSystem.class); private ClusterWorkerContext clusterContext; @@ -33,7 +28,7 @@ public enum CollectorSystem { public void boot() throws Exception { createAkkaSystem(); createListener(); - createLocalProvider(); + loadLocalProviders(); createClusterWorker(); } @@ -63,14 +58,14 @@ public enum CollectorSystem { private void createClusterWorker() throws Exception { ServiceLoader clusterServiceLoader = ServiceLoader.load(AbstractClusterWorkerProvider.class); for (AbstractClusterWorkerProvider provider : clusterServiceLoader) { - logger.info("create {%s} worker {%s} using java service loader", provider.workerNum(), provider.workerClass().getName()); + logger.info("create {%s} worker using java service loader", provider.workerNum()); for (int i = 1; i <= provider.workerNum(); i++) { - provider.create(clusterContext, null); + provider.create(clusterContext, new LocalWorkerContext()); } } } - private void createLocalProvider() throws DuplicateProviderException { + private void loadLocalProviders() throws UsedRoleNameException { ServiceLoader clusterServiceLoader = ServiceLoader.load(AbstractLocalWorkerProvider.class); for (AbstractLocalWorkerProvider provider : clusterServiceLoader) { clusterContext.putProvider(provider); diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java index fda0ddf62..b72190d94 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java @@ -15,8 +15,8 @@ import org.apache.logging.log4j.Logger; */ public abstract class AbstractClusterWorker extends AbstractWorker { - public AbstractClusterWorker(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public AbstractClusterWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } static class WorkerWithAkka extends UntypedActor { 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 index c6e4e0c3e..b37e87584 100644 --- 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 @@ -3,8 +3,6 @@ package com.a.eye.skywalking.collector.actor; import akka.actor.ActorRef; import akka.actor.Props; -import java.lang.reflect.Constructor; - /** * @author pengys5 */ @@ -13,12 +11,10 @@ public abstract class AbstractClusterWorkerProvider[]{Role.class, ClusterWorkerContext.class}); - workerConstructor.setAccessible(true); - T clusterWorker = (T) workerConstructor.newInstance(role(), clusterContext); + T clusterWorker = (T) workerInstance(clusterContext); clusterWorker.preStart(); ActorRef actorRef = clusterContext.getAkkaSystem().actorOf(Props.create(AbstractClusterWorker.WorkerWithAkka.class, clusterWorker), role() + "_" + num); diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java index 227768c62..b9b167aca 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java @@ -10,8 +10,8 @@ import com.lmax.disruptor.RingBuffer; */ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker { - public AbstractLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public AbstractLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } static class WorkerWithDisruptor implements EventHandler { @@ -19,17 +19,21 @@ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker { private RingBuffer ringBuffer; private AbstractLocalAsyncWorker asyncWorker; - private WorkerWithDisruptor(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) { + public WorkerWithDisruptor(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) { this.ringBuffer = ringBuffer; this.asyncWorker = asyncWorker; } - public void onEvent(MessageHolder event, long sequence, boolean endOfBatch) throws Exception { - Object message = event.getMessage(); - event.reset(); - asyncWorker.work(message); - if (endOfBatch) { - asyncWorker.work(new EndOfBatchCommand()); + public void onEvent(MessageHolder event, long sequence, boolean endOfBatch) { + try { + Object message = event.getMessage(); + event.reset(); + asyncWorker.work(message); + if (endOfBatch) { + asyncWorker.work(new EndOfBatchCommand()); + } + } catch (Exception e) { + e.printStackTrace(); } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorkerProvider.java index a70f74355..9633e8d31 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorkerProvider.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorkerProvider.java @@ -6,8 +6,6 @@ import com.a.eye.skywalking.collector.queue.MessageHolderFactory; import com.lmax.disruptor.RingBuffer; import com.lmax.disruptor.dsl.Disruptor; -import java.lang.reflect.Constructor; - /** * @author pengys5 */ @@ -16,24 +14,21 @@ public abstract class AbstractLocalAsyncWorkerProvider[]{Role.class, ClusterWorkerContext.class}); - workerConstructor.setAccessible(true); - T localAsyncWorker = (T) workerConstructor.newInstance(role(), clusterContext); + final public WorkerRef onCreate(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException { + T localAsyncWorker = (T) workerInstance(clusterContext); localAsyncWorker.preStart(); - Constructor memberConstructor = AbstractLocalAsyncWorker.WorkerWithDisruptor.class.getDeclaredConstructor(new Class[]{RingBuffer.class, AbstractLocalAsyncWorker.class}); - memberConstructor.setAccessible(true); - // Specify the size of the ring buffer, must be power of 2. int bufferSize = queueSize(); + if (!((((bufferSize - 1) & bufferSize) == 0) && bufferSize != 0)) { + throw new IllegalArgumentException("queue size must be power of 2"); + } + // Construct the Disruptor Disruptor disruptor = new Disruptor(MessageHolderFactory.INSTANCE, bufferSize, DaemonThreadFactory.INSTANCE); RingBuffer ringBuffer = disruptor.getRingBuffer(); - T.WorkerWithDisruptor disruptorWorker = (T.WorkerWithDisruptor) memberConstructor.newInstance(ringBuffer, localAsyncWorker); + T.WorkerWithDisruptor disruptorWorker = new T.WorkerWithDisruptor(ringBuffer, localAsyncWorker); // Connect the handler disruptor.handleEventsWith(disruptorWorker); diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorker.java index c0baf3f01..09817d04f 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorker.java @@ -4,7 +4,7 @@ package com.a.eye.skywalking.collector.actor; * @author pengys5 */ public abstract class AbstractLocalSyncWorker extends AbstractLocalWorker { - public AbstractLocalSyncWorker(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public AbstractLocalSyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorkerProvider.java index 912d5e345..20e50854c 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorkerProvider.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalSyncWorkerProvider.java @@ -1,19 +1,13 @@ package com.a.eye.skywalking.collector.actor; -import java.lang.reflect.Constructor; - /** * @author pengys5 */ public abstract class AbstractLocalSyncWorkerProvider extends AbstractLocalWorkerProvider { @Override - final public WorkerRef create(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws Exception { - validate(); - - Constructor workerConstructor = workerClass().getDeclaredConstructor(new Class[]{Role.class, ClusterWorkerContext.class}); - workerConstructor.setAccessible(true); - T localSyncWorker = (T) workerConstructor.newInstance(role(), clusterContext); + final public WorkerRef onCreate(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException { + T localSyncWorker = (T) workerInstance(clusterContext); localSyncWorker.preStart(); LocalSyncWorkerRef workerRef = new LocalSyncWorkerRef(role(), localSyncWorker); 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 index 8feea4ce2..f939e9fa8 100644 --- 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 @@ -4,7 +4,7 @@ package com.a.eye.skywalking.collector.actor; * @author pengys5 */ public abstract class AbstractLocalWorker extends AbstractWorker { - public AbstractLocalWorker(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public AbstractLocalWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } } 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 4c9d78c1d..a830708db 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 @@ -5,18 +5,19 @@ package com.a.eye.skywalking.collector.actor; */ public abstract class AbstractWorker { - private final LocalWorkerContext selfContext = new LocalWorkerContext(); + private final LocalWorkerContext selfContext; private final Role role; private final ClusterWorkerContext clusterContext; - public AbstractWorker(Role role, ClusterWorkerContext clusterContext) { + public AbstractWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { this.role = role; this.clusterContext = clusterContext; + this.selfContext = selfContext; } - public abstract void preStart() throws Exception; + public abstract void preStart() throws ProviderNotFountException; public abstract void work(Object message) throws Exception; 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 9d2125f5e..dd0b4aee8 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 @@ -1,7 +1,5 @@ package com.a.eye.skywalking.collector.actor; -import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; - /** * @author pengys5 */ @@ -9,13 +7,14 @@ public abstract class AbstractWorkerProvider implement public abstract Role role(); - public abstract Class workerClass(); + public abstract T workerInstance(ClusterWorkerContext clusterContext); -// public abstract WorkerSelector selector(); + public abstract WorkerRef onCreate(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException; - final void validate() throws Exception { - if (workerClass() == null) { - throw new IllegalArgumentException("cannot createInstance() with nothing obtained from workerClass()"); + final public WorkerRef create(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException { + if (workerInstance(clusterContext) == null) { + throw new IllegalArgumentException("cannot get worker instance with nothing obtained from workerInstance()"); } + return onCreate(clusterContext, localContext); } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerContext.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerContext.java index edbf71e0a..abd107d04 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerContext.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerContext.java @@ -28,21 +28,21 @@ public class ClusterWorkerContext extends WorkerContext { @Override public AbstractWorkerProvider findProvider(Role role) throws ProviderNotFountException { - logger.debug("find role of %s provider from ClusterWorkerContext", role.name()); - if (providers.containsKey(role.name())) { - return providers.get(role.name()); + logger.debug("find role of %s provider from ClusterWorkerContext", role.roleName()); + if (providers.containsKey(role.roleName())) { + return providers.get(role.roleName()); } else { - throw new ProviderNotFountException("role=" + role.name() + ", no available provider."); + throw new ProviderNotFountException("role=" + role.roleName() + ", no available provider."); } } @Override - public void putProvider(AbstractWorkerProvider provider) throws DuplicateProviderException { - logger.debug("put role of %s provider into ClusterWorkerContext", provider.role().name()); - if (providers.containsKey(provider.role().name())) { - throw new DuplicateProviderException("provider with role=" + provider.role().name() + " duplicate each other."); + public void putProvider(AbstractWorkerProvider provider) throws UsedRoleNameException { + logger.debug("put role of %s provider into ClusterWorkerContext", provider.role().roleName()); + if (providers.containsKey(provider.role().roleName())) { + throw new UsedRoleNameException("provider with role=" + provider.role().roleName() + " duplicate each other."); } else { - providers.put(provider.role().name(), provider); + providers.put(provider.role().roleName(), provider); } } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerRefCounter.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerRefCounter.java index 7067d923c..18cfb6c9e 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerRefCounter.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/ClusterWorkerRefCounter.java @@ -13,10 +13,10 @@ public enum ClusterWorkerRefCounter { private Map counter = new ConcurrentHashMap<>(); public int incrementAndGet(Role role) { - if (!counter.containsKey(role.name())) { + if (!counter.containsKey(role.roleName())) { AtomicInteger atomic = new AtomicInteger(0); - counter.putIfAbsent(role.name(), atomic); + counter.putIfAbsent(role.roleName(), atomic); } - return counter.get(role.name()).incrementAndGet(); + return counter.get(role.roleName()).incrementAndGet(); } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Context.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Context.java index f214c2129..9cd13576e 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Context.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Context.java @@ -7,7 +7,7 @@ public interface Context { AbstractWorkerProvider findProvider(Role role) throws ProviderNotFountException; - void putProvider(AbstractWorkerProvider provider) throws DuplicateProviderException; + void putProvider(AbstractWorkerProvider provider) throws UsedRoleNameException; WorkerRefs lookup(Role role) throws WorkerNotFountException; diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/DuplicateProviderException.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/DuplicateProviderException.java deleted file mode 100644 index bb3d258a4..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/DuplicateProviderException.java +++ /dev/null @@ -1,7 +0,0 @@ -package com.a.eye.skywalking.collector.actor; - -public class DuplicateProviderException extends Exception { - public DuplicateProviderException(String message){ - super(message); - } -} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LocalWorkerContext.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LocalWorkerContext.java index c058784f1..09b3db05a 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LocalWorkerContext.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LocalWorkerContext.java @@ -11,7 +11,7 @@ public class LocalWorkerContext extends WorkerContext { } @Override - final public void putProvider(AbstractWorkerProvider provider) throws DuplicateProviderException { + final public void putProvider(AbstractWorkerProvider provider) throws UsedRoleNameException { } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Role.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Role.java index ff8c47245..7255d98ea 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Role.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Role.java @@ -5,9 +5,9 @@ import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; /** * @author pengys5 */ -public abstract class Role { +public interface Role { - public abstract String name(); + String roleName(); - public abstract WorkerSelector workerSelector(); + WorkerSelector workerSelector(); } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/UsedRoleNameException.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/UsedRoleNameException.java new file mode 100644 index 000000000..a597b7f6f --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/UsedRoleNameException.java @@ -0,0 +1,7 @@ +package com.a.eye.skywalking.collector.actor; + +public class UsedRoleNameException extends Exception { + public UsedRoleNameException(String message){ + super(message); + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerContext.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerContext.java index e292b911f..3573067ef 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerContext.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerContext.java @@ -1,7 +1,6 @@ package com.a.eye.skywalking.collector.actor; import java.util.ArrayList; -import java.util.Collections; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -15,25 +14,24 @@ public abstract class WorkerContext implements Context { @Override final public WorkerRefs lookup(Role role) throws WorkerNotFountException { - if (roleWorkers.containsKey(role.name())) { - WorkerRefs refs = new WorkerRefs(roleWorkers.get(role.name()), role.workerSelector()); + if (roleWorkers.containsKey(role.roleName())) { + WorkerRefs refs = new WorkerRefs(roleWorkers.get(role.roleName()), role.workerSelector()); return refs; } else { - throw new WorkerNotFountException("role=" + role.name() + ", no available worker."); + throw new WorkerNotFountException("role=" + role.roleName() + ", no available worker."); } } @Override final public void put(WorkerRef workerRef) { - if (!roleWorkers.containsKey(workerRef.getRole().name())) { - List actorList = Collections.synchronizedList(new ArrayList()); - roleWorkers.putIfAbsent(workerRef.getRole().name(), actorList); + if (!roleWorkers.containsKey(workerRef.getRole().roleName())) { + roleWorkers.putIfAbsent(workerRef.getRole().roleName(), new ArrayList()); } - roleWorkers.get(workerRef.getRole().name()).add(workerRef); + roleWorkers.get(workerRef.getRole().roleName()).add(workerRef); } @Override final public void remove(WorkerRef workerRef) { - roleWorkers.remove(workerRef); + roleWorkers.remove(workerRef.getRole().roleName()); } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java index e74cc3b62..a3d46df68 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java @@ -11,9 +11,13 @@ public class HashCodeSelector implements WorkerSelector { @Override public WorkerRef select(List members, Object message) { - AbstractHashMessage hashMessage = (AbstractHashMessage) message; - int size = members.size(); - int selectIndex = Math.abs(hashMessage.getHashCode()) % size; - return members.get(selectIndex); + if (message instanceof AbstractHashMessage) { + AbstractHashMessage hashMessage = (AbstractHashMessage) message; + int size = members.size(); + int selectIndex = Math.abs(hashMessage.getHashCode()) % size; + return members.get(selectIndex); + } else { + throw new IllegalArgumentException("the message send into HashCodeSelector must implementation of AbstractHashMessage"); + } } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorker.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorker.java index ae80eefd9..aeee348fe 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorker.java @@ -1,5 +1,6 @@ package com.a.eye.skywalking.collector.actor; + import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; @@ -8,12 +9,12 @@ import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; */ public class TestClusterWorker extends AbstractClusterWorker { - public TestClusterWorker(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public TestClusterWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(TestLocalSyncWorker.TestLocalSyncWorkerRole.INSTANCE).create(getClusterContext(), getSelfContext()); getClusterContext().findProvider(TestLocalAsyncWorker.TestLocalASyncWorkerRole.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -42,20 +43,20 @@ public class TestClusterWorker extends AbstractClusterWorker { @Override public Role role() { - return new TestClusterWorkerRole(); + return TestClusterWorkerRole.INSTANCE; } @Override - public Class workerClass() { - return TestClusterWorker.class; + public TestClusterWorker workerInstance(ClusterWorkerContext clusterContext) { + return new TestClusterWorker(role(), clusterContext, new LocalWorkerContext()); } } - public static class TestClusterWorkerRole extends Role { - public static TestClusterWorkerRole INSTANCE = new TestClusterWorkerRole(); + public enum TestClusterWorkerRole implements Role { + INSTANCE; @Override - public String name() { + public String roleName() { return TestClusterWorker.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorkerTestCase.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorkerTestCase.java index 81fd13d70..3b056621d 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorkerTestCase.java +++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestClusterWorkerTestCase.java @@ -10,19 +10,19 @@ import org.junit.Test; */ public class TestClusterWorkerTestCase { - @Before + private CollectorSystem collectorSystem; + public void createSystem() throws Exception { - CollectorSystem.INSTANCE.boot(); + collectorSystem = new CollectorSystem(); + collectorSystem.boot(); } - @After public void terminateSystem() { - CollectorSystem.INSTANCE.terminate(); + collectorSystem.terminate(); } - @Test public void testTellWorker() throws Exception { - WorkerRefs workerRefs = CollectorSystem.INSTANCE.getClusterContext().lookup(TestClusterWorker.TestClusterWorkerRole.INSTANCE); + WorkerRefs workerRefs = collectorSystem.getClusterContext().lookup(TestClusterWorker.TestClusterWorkerRole.INSTANCE); workerRefs.tell("Print"); workerRefs.tell("TellLocalWorker"); workerRefs.tell("TellLocalAsyncWorker"); diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalAsyncWorker.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalAsyncWorker.java index 59bff85dc..e2e07a13d 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalAsyncWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalAsyncWorker.java @@ -8,12 +8,12 @@ import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; */ public class TestLocalAsyncWorker extends AbstractLocalAsyncWorker { - public TestLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public TestLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { } @@ -37,16 +37,16 @@ public class TestLocalAsyncWorker extends AbstractLocalAsyncWorker { } @Override - public Class workerClass() { - return TestLocalAsyncWorker.class; + public TestLocalAsyncWorker workerInstance(ClusterWorkerContext clusterContext) { + return new TestLocalAsyncWorker(role(), clusterContext, new LocalWorkerContext()); } } - public static class TestLocalASyncWorkerRole extends Role { - public static TestLocalASyncWorkerRole INSTANCE = new TestLocalASyncWorkerRole(); + public enum TestLocalASyncWorkerRole implements Role { + INSTANCE; @Override - public String name() { + public String roleName() { return TestLocalAsyncWorker.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorker.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorker.java index 64340004e..51128712f 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorker.java @@ -8,12 +8,12 @@ import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; */ public class TestLocalSyncWorker extends AbstractLocalSyncWorker { - public TestLocalSyncWorker(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public TestLocalSyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { } @@ -33,16 +33,16 @@ public class TestLocalSyncWorker extends AbstractLocalSyncWorker { } @Override - public Class workerClass() { - return TestLocalSyncWorker.class; + public TestLocalSyncWorker workerInstance(ClusterWorkerContext clusterContext) { + return new TestLocalSyncWorker(role(), clusterContext, new LocalWorkerContext()); } } - public static class TestLocalSyncWorkerRole extends Role { - public static TestLocalSyncWorkerRole INSTANCE = new TestLocalSyncWorkerRole(); + public enum TestLocalSyncWorkerRole implements Role { + INSTANCE; @Override - public String name() { + public String roleName() { return TestLocalSyncWorker.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorkerTestCase.java b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorkerTestCase.java index e80a94586..9f61ccadc 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorkerTestCase.java +++ b/skywalking-collector/skywalking-collector-cluster/src/test/java/com/a/eye/skywalking/collector/actor/TestLocalSyncWorkerTestCase.java @@ -12,12 +12,11 @@ public class TestLocalSyncWorkerTestCase { @Before public void createSystem() throws Exception { - CollectorSystem.INSTANCE.boot(); } @After public void terminateSystem() { - CollectorSystem.INSTANCE.terminate(); + } @Test diff --git a/skywalking-collector/skywalking-collector-role/src/main/java/com/a/eye/skywalking/collector/role/TraceSegmentReceiverRole.java b/skywalking-collector/skywalking-collector-role/src/main/java/com/a/eye/skywalking/collector/role/TraceSegmentReceiverRole.java index d1de89ee5..1e77636ff 100644 --- a/skywalking-collector/skywalking-collector-role/src/main/java/com/a/eye/skywalking/collector/role/TraceSegmentReceiverRole.java +++ b/skywalking-collector/skywalking-collector-role/src/main/java/com/a/eye/skywalking/collector/role/TraceSegmentReceiverRole.java @@ -7,11 +7,11 @@ import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; /** * @author pengys5 */ -public class TraceSegmentReceiverRole extends Role { - public static TraceSegmentReceiverRole INSTANCE = new TraceSegmentReceiverRole(); +public enum TraceSegmentReceiverRole implements Role { + INSTANCE; @Override - public String name() { + public String roleName() { return "TraceSegmentReceiver"; } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/AnalysisMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/AnalysisMember.java index f0b93f8cc..c823ada00 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/AnalysisMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/AnalysisMember.java @@ -1,8 +1,6 @@ package com.a.eye.skywalking.collector.worker; -import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorker; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; -import com.a.eye.skywalking.collector.actor.Role; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.queue.EndOfBatchCommand; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -14,14 +12,15 @@ public abstract class AnalysisMember extends AbstractLocalAsyncWorker { private Logger logger = LogManager.getFormatterLogger(AnalysisMember.class); - public AnalysisMember(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public AnalysisMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } public abstract void analyse(Object message) throws Exception; @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { + } @Override diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java index ad7e432e5..113627fbd 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricAnalysisMember.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker; import akka.actor.ActorRef; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.Role; import com.a.eye.skywalking.collector.queue.MessageHolder; import com.a.eye.skywalking.collector.worker.storage.MetricData; @@ -19,8 +20,8 @@ public abstract class MetricAnalysisMember extends AnalysisMember { protected MetricPersistenceData persistenceData = new MetricPersistenceData(); - public MetricAnalysisMember(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public MetricAnalysisMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } public void setMetric(String id, int second, Long value) throws Exception { diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java index e24161d60..acc9c0bca 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/MetricPersistenceMember.java @@ -1,13 +1,11 @@ package com.a.eye.skywalking.collector.worker; -import akka.actor.ActorRef; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.Role; -import com.a.eye.skywalking.collector.queue.MessageHolder; import com.a.eye.skywalking.collector.worker.storage.EsClient; import com.a.eye.skywalking.collector.worker.storage.MetricData; import com.a.eye.skywalking.collector.worker.storage.MetricPersistenceData; -import com.lmax.disruptor.RingBuffer; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; import org.elasticsearch.action.bulk.BulkRequestBuilder; @@ -30,8 +28,8 @@ public abstract class MetricPersistenceMember extends PersistenceMember { protected MetricPersistenceData persistenceData = new MetricPersistenceData(); - public MetricPersistenceMember(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public MetricPersistenceMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceMember.java index 52dd2b2f3..5bcb3fb54 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/PersistenceMember.java @@ -1,8 +1,6 @@ package com.a.eye.skywalking.collector.worker; -import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorker; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; -import com.a.eye.skywalking.collector.actor.Role; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.queue.EndOfBatchCommand; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -14,8 +12,8 @@ public abstract class PersistenceMember extends AbstractLocalAsyncWorker { private Logger logger = LogManager.getFormatterLogger(PersistenceMember.class); - public PersistenceMember(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public PersistenceMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } public abstract String esIndex(); @@ -25,7 +23,8 @@ public abstract class PersistenceMember extends AbstractLocalAsyncWorker { public abstract void analyse(Object message) throws Exception; @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { + } @Override diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java index e7bf3435d..69a4fde43 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordAnalysisMember.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker; import akka.actor.ActorRef; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.Role; import com.a.eye.skywalking.collector.queue.MessageHolder; import com.a.eye.skywalking.collector.worker.storage.RecordData; @@ -20,8 +21,8 @@ public abstract class RecordAnalysisMember extends AnalysisMember { private RecordPersistenceData persistenceData = new RecordPersistenceData(); - public RecordAnalysisMember(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public RecordAnalysisMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } public void setRecord(String id, JsonObject record) throws Exception { diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java index ab72f19b4..108b6a896 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/RecordPersistenceMember.java @@ -1,6 +1,7 @@ package com.a.eye.skywalking.collector.worker; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.Role; import com.a.eye.skywalking.collector.worker.storage.EsClient; import com.a.eye.skywalking.collector.worker.storage.RecordData; @@ -23,8 +24,8 @@ public abstract class RecordPersistenceMember extends PersistenceMember { protected RecordPersistenceData persistenceData = new RecordPersistenceData(); - public RecordPersistenceMember(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public RecordPersistenceMember(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java index 40c3eaeda..f91a5785c 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/DAGNodeAnalysis.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.analysis; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.RecordAnalysisMember; @@ -20,8 +21,8 @@ public class DAGNodeAnalysis extends RecordAnalysisMember { private Logger logger = LogManager.getFormatterLogger(DAGNodeAnalysis.class); - public DAGNodeAnalysis(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public DAGNodeAnalysis(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -59,8 +60,8 @@ public class DAGNodeAnalysis extends RecordAnalysisMember { } @Override - public Class workerClass() { - return DAGNodeAnalysis.class; + public DAGNodeAnalysis workerInstance(ClusterWorkerContext clusterContext) { + return new DAGNodeAnalysis(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -69,11 +70,11 @@ public class DAGNodeAnalysis extends RecordAnalysisMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return DAGNodeAnalysis.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java index 2eaa79e30..56aeeb0aa 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/NodeInstanceAnalysis.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.analysis; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.RecordAnalysisMember; @@ -21,8 +22,8 @@ public class NodeInstanceAnalysis extends RecordAnalysisMember { private Logger logger = LogManager.getFormatterLogger(NodeInstanceAnalysis.class); - public NodeInstanceAnalysis(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public NodeInstanceAnalysis(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -59,8 +60,8 @@ public class NodeInstanceAnalysis extends RecordAnalysisMember { } @Override - public Class workerClass() { - return NodeInstanceAnalysis.class; + public NodeInstanceAnalysis workerInstance(ClusterWorkerContext clusterContext) { + return new NodeInstanceAnalysis(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -69,11 +70,11 @@ public class NodeInstanceAnalysis extends RecordAnalysisMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return NodeInstanceAnalysis.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java index b41324e4c..fd31469d5 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseCostAnalysis.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.analysis; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.MetricAnalysisMember; @@ -19,8 +20,8 @@ public class ResponseCostAnalysis extends MetricAnalysisMember { private Logger logger = LogManager.getFormatterLogger(ResponseCostAnalysis.class); - public ResponseCostAnalysis(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public ResponseCostAnalysis(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -53,8 +54,8 @@ public class ResponseCostAnalysis extends MetricAnalysisMember { } @Override - public Class workerClass() { - return ResponseCostAnalysis.class; + public ResponseCostAnalysis workerInstance(ClusterWorkerContext clusterContext) { + return new ResponseCostAnalysis(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -63,11 +64,11 @@ public class ResponseCostAnalysis extends MetricAnalysisMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return ResponseCostAnalysis.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java index a26cf64e6..9dae5cccb 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/analysis/ResponseSummaryAnalysis.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.analysis; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.MetricAnalysisMember; @@ -19,8 +20,8 @@ public class ResponseSummaryAnalysis extends MetricAnalysisMember { private Logger logger = LogManager.getFormatterLogger(ResponseSummaryAnalysis.class); - public ResponseSummaryAnalysis(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public ResponseSummaryAnalysis(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -50,8 +51,8 @@ public class ResponseSummaryAnalysis extends MetricAnalysisMember { } @Override - public Class workerClass() { - return ResponseSummaryAnalysis.class; + public ResponseSummaryAnalysis workerInstance(ClusterWorkerContext clusterContext) { + return new ResponseSummaryAnalysis(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -60,11 +61,11 @@ public class ResponseSummaryAnalysis extends MetricAnalysisMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return ResponseSummaryAnalysis.class.getSimpleName(); } 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 index 9a4f6358f..941f03809 100644 --- 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 @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.persistence; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.RecordPersistenceMember; @@ -16,8 +17,8 @@ public class DAGNodePersistence extends RecordPersistenceMember { private Logger logger = LogManager.getFormatterLogger(DAGNodePersistence.class); - public DAGNodePersistence(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public DAGNodePersistence(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -39,8 +40,8 @@ public class DAGNodePersistence extends RecordPersistenceMember { } @Override - public Class workerClass() { - return DAGNodePersistence.class; + public DAGNodePersistence workerInstance(ClusterWorkerContext clusterContext) { + return new DAGNodePersistence(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -49,11 +50,11 @@ public class DAGNodePersistence extends RecordPersistenceMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return DAGNodePersistence.class.getSimpleName(); } 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 index bf2d3d298..0b70b33b1 100644 --- 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 @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.persistence; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.RecordPersistenceMember; @@ -16,8 +17,8 @@ public class NodeInstancePersistence extends RecordPersistenceMember { private Logger logger = LogManager.getFormatterLogger(NodeInstancePersistence.class); - public NodeInstancePersistence(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public NodeInstancePersistence(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -35,12 +36,12 @@ public class NodeInstancePersistence extends RecordPersistenceMember { @Override public Role role() { - return NodeInstancePersistence.Role.INSTANCE; + return Role.INSTANCE; } @Override - public Class workerClass() { - return NodeInstancePersistence.class; + public NodeInstancePersistence workerInstance(ClusterWorkerContext clusterContext) { + return new NodeInstancePersistence(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -49,11 +50,11 @@ public class NodeInstancePersistence extends RecordPersistenceMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return NodeInstancePersistence.class.getSimpleName(); } 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 index 54edbcc29..e950a13b0 100644 --- 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 @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.persistence; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.MetricPersistenceMember; @@ -16,8 +17,8 @@ public class ResponseCostPersistence extends MetricPersistenceMember { private Logger logger = LogManager.getFormatterLogger(ResponseCostPersistence.class); - public ResponseCostPersistence(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public ResponseCostPersistence(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -39,8 +40,8 @@ public class ResponseCostPersistence extends MetricPersistenceMember { } @Override - public Class workerClass() { - return ResponseCostPersistence.class; + public ResponseCostPersistence workerInstance(ClusterWorkerContext clusterContext) { + return new ResponseCostPersistence(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -49,11 +50,11 @@ public class ResponseCostPersistence extends MetricPersistenceMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return ResponseCostPersistence.class.getSimpleName(); } 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 index 8bae69017..6517b2ace 100644 --- 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 @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.application.persistence; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.MetricPersistenceMember; @@ -16,8 +17,8 @@ public class ResponseSummaryPersistence extends MetricPersistenceMember { private Logger logger = LogManager.getFormatterLogger(ResponseSummaryPersistence.class); - public ResponseSummaryPersistence(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public ResponseSummaryPersistence(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -39,8 +40,8 @@ public class ResponseSummaryPersistence extends MetricPersistenceMember { } @Override - public Class workerClass() { - return ResponseSummaryPersistence.class; + public ResponseSummaryPersistence workerInstance(ClusterWorkerContext clusterContext) { + return new ResponseSummaryPersistence(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -49,11 +50,11 @@ public class ResponseSummaryPersistence extends MetricPersistenceMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return ResponseSummaryPersistence.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java index a365fb634..70c664fbf 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/persistence/TraceSegmentRecordPersistence.java @@ -3,6 +3,7 @@ package com.a.eye.skywalking.collector.worker.application.persistence; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.RecordPersistenceMember; @@ -38,8 +39,8 @@ public class TraceSegmentRecordPersistence extends RecordPersistenceMember { return "trace_segment"; } - public TraceSegmentRecordPersistence(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public TraceSegmentRecordPersistence(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -68,16 +69,16 @@ public class TraceSegmentRecordPersistence extends RecordPersistenceMember { } @Override - public Class workerClass() { - return TraceSegmentRecordPersistence.class; + public TraceSegmentRecordPersistence workerInstance(ClusterWorkerContext clusterContext) { + return new TraceSegmentRecordPersistence(role(), clusterContext, new LocalWorkerContext()); } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return TraceSegmentRecordPersistence.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java index 9cd89968f..fb9cf3fab 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/DAGNodeReceiver.java @@ -1,8 +1,6 @@ package com.a.eye.skywalking.collector.worker.application.receiver; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorker; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.WorkerConfig; @@ -18,12 +16,12 @@ public class DAGNodeReceiver extends AbstractClusterWorker { private Logger logger = LogManager.getFormatterLogger(DAGNodeReceiver.class); - public DAGNodeReceiver(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public DAGNodeReceiver(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(DAGNodePersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -45,8 +43,8 @@ public class DAGNodeReceiver extends AbstractClusterWorker { } @Override - public Class workerClass() { - return DAGNodeReceiver.class; + public DAGNodeReceiver workerInstance(ClusterWorkerContext clusterContext) { + return new DAGNodeReceiver(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -55,11 +53,11 @@ public class DAGNodeReceiver extends AbstractClusterWorker { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return DAGNodeReceiver.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java index aed8aea97..739b51461 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/NodeInstanceReceiver.java @@ -1,8 +1,6 @@ package com.a.eye.skywalking.collector.worker.application.receiver; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorker; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.WorkerConfig; @@ -18,12 +16,12 @@ public class NodeInstanceReceiver extends AbstractClusterWorker { private Logger logger = LogManager.getFormatterLogger(NodeInstanceReceiver.class); - public NodeInstanceReceiver(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public NodeInstanceReceiver(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(NodeInstancePersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -45,8 +43,8 @@ public class NodeInstanceReceiver extends AbstractClusterWorker { } @Override - public Class workerClass() { - return NodeInstanceReceiver.class; + public NodeInstanceReceiver workerInstance(ClusterWorkerContext clusterContext) { + return new NodeInstanceReceiver(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -55,11 +53,11 @@ public class NodeInstanceReceiver extends AbstractClusterWorker { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return NodeInstanceReceiver.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java index 2f62d2de7..788b19a03 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseCostReceiver.java @@ -1,8 +1,6 @@ package com.a.eye.skywalking.collector.worker.application.receiver; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorker; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.WorkerConfig; @@ -18,12 +16,12 @@ public class ResponseCostReceiver extends AbstractClusterWorker { private Logger logger = LogManager.getFormatterLogger(ResponseCostReceiver.class); - public ResponseCostReceiver(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public ResponseCostReceiver(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(ResponseCostPersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -45,8 +43,8 @@ public class ResponseCostReceiver extends AbstractClusterWorker { } @Override - public Class workerClass() { - return ResponseCostReceiver.class; + public ResponseCostReceiver workerInstance(ClusterWorkerContext clusterContext) { + return new ResponseCostReceiver(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -55,11 +53,11 @@ public class ResponseCostReceiver extends AbstractClusterWorker { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return ResponseCostReceiver.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java index 58cdc2b33..e533305b5 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/receiver/ResponseSummaryReceiver.java @@ -1,8 +1,6 @@ package com.a.eye.skywalking.collector.worker.application.receiver; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorker; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.WorkerConfig; @@ -18,12 +16,12 @@ public class ResponseSummaryReceiver extends AbstractClusterWorker { private Logger logger = LogManager.getFormatterLogger(ResponseSummaryReceiver.class); - public ResponseSummaryReceiver(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public ResponseSummaryReceiver(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(ResponseSummaryPersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -45,8 +43,8 @@ public class ResponseSummaryReceiver extends AbstractClusterWorker { } @Override - public Class workerClass() { - return ResponseSummaryReceiver.class; + public ResponseSummaryReceiver workerInstance(ClusterWorkerContext clusterContext) { + return new ResponseSummaryReceiver(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -55,11 +53,11 @@ public class ResponseSummaryReceiver extends AbstractClusterWorker { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return ResponseSummaryReceiver.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java index b133d94f9..db484a475 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/ApplicationRefMain.java @@ -1,9 +1,7 @@ package com.a.eye.skywalking.collector.worker.applicationref; import com.a.eye.skywalking.api.util.StringUtil; -import com.a.eye.skywalking.collector.actor.AbstractLocalSyncWorker; -import com.a.eye.skywalking.collector.actor.AbstractLocalSyncWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.applicationref.analysis.DAGNodeRefAnalysis; @@ -17,12 +15,12 @@ public class ApplicationRefMain extends AbstractLocalSyncWorker { private DAGNodeRefAnalysis dagNodeRefAnalysis; - public ApplicationRefMain(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public ApplicationRefMain(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(DAGNodeRefAnalysis.Role.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -49,16 +47,16 @@ public class ApplicationRefMain extends AbstractLocalSyncWorker { } @Override - public Class workerClass() { - return ApplicationRefMain.class; + public ApplicationRefMain workerInstance(ClusterWorkerContext clusterContext) { + return new ApplicationRefMain(role(), clusterContext, new LocalWorkerContext()); } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return ApplicationRefMain.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java index 8e2fe5bc1..3b026dbd2 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/analysis/DAGNodeRefAnalysis.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.applicationref.analysis; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.HashCodeSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.RecordAnalysisMember; @@ -21,8 +22,8 @@ public class DAGNodeRefAnalysis extends RecordAnalysisMember { private Logger logger = LogManager.getFormatterLogger(DAGNodeRefAnalysis.class); - public DAGNodeRefAnalysis(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public DAGNodeRefAnalysis(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -58,8 +59,8 @@ public class DAGNodeRefAnalysis extends RecordAnalysisMember { } @Override - public Class workerClass() { - return DAGNodeRefAnalysis.class; + public DAGNodeRefAnalysis workerInstance(ClusterWorkerContext clusterContext) { + return new DAGNodeRefAnalysis(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -68,11 +69,11 @@ public class DAGNodeRefAnalysis extends RecordAnalysisMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return DAGNodeRefAnalysis.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/persistence/DAGNodeRefPersistence.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/persistence/DAGNodeRefPersistence.java index 22274f3eb..901d31c1d 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/persistence/DAGNodeRefPersistence.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/persistence/DAGNodeRefPersistence.java @@ -2,6 +2,7 @@ package com.a.eye.skywalking.collector.worker.applicationref.persistence; import com.a.eye.skywalking.collector.actor.AbstractLocalAsyncWorkerProvider; import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.LocalWorkerContext; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.RecordPersistenceMember; @@ -16,8 +17,8 @@ public class DAGNodeRefPersistence extends RecordPersistenceMember { private Logger logger = LogManager.getFormatterLogger(DAGNodeRefPersistence.class); - public DAGNodeRefPersistence(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public DAGNodeRefPersistence(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override @@ -40,8 +41,8 @@ public class DAGNodeRefPersistence extends RecordPersistenceMember { } @Override - public Class workerClass() { - return DAGNodeRefPersistence.class; + public DAGNodeRefPersistence workerInstance(ClusterWorkerContext clusterContext) { + return new DAGNodeRefPersistence(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -50,11 +51,11 @@ public class DAGNodeRefPersistence extends RecordPersistenceMember { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return DAGNodeRefPersistence.class.getSimpleName(); } diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java index 1cfa80b19..4c10298f1 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/applicationref/receiver/DAGNodeRefReceiver.java @@ -1,8 +1,6 @@ package com.a.eye.skywalking.collector.worker.applicationref.receiver; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorker; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.actor.selector.RollingSelector; import com.a.eye.skywalking.collector.actor.selector.WorkerSelector; import com.a.eye.skywalking.collector.worker.WorkerConfig; @@ -20,12 +18,12 @@ public class DAGNodeRefReceiver extends AbstractClusterWorker { private DAGNodeRefPersistence persistence; - public DAGNodeRefReceiver(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public DAGNodeRefReceiver(com.a.eye.skywalking.collector.actor.Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(DAGNodeRefPersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -47,8 +45,8 @@ public class DAGNodeRefReceiver extends AbstractClusterWorker { } @Override - public Class workerClass() { - return DAGNodeRefReceiver.class; + public DAGNodeRefReceiver workerInstance(ClusterWorkerContext clusterContext) { + return new DAGNodeRefReceiver(role(), clusterContext, new LocalWorkerContext()); } @Override @@ -57,11 +55,11 @@ public class DAGNodeRefReceiver extends AbstractClusterWorker { } } - public static class Role extends com.a.eye.skywalking.collector.actor.Role { - public static Role INSTANCE = new Role(); + public enum Role implements com.a.eye.skywalking.collector.actor.Role { + INSTANCE; @Override - public String name() { + public String roleName() { return DAGNodeRefReceiver.class.getSimpleName(); } 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 ae85b9bc7..679dc55b4 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,9 +1,6 @@ package com.a.eye.skywalking.collector.worker.receiver; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorker; -import com.a.eye.skywalking.collector.actor.AbstractClusterWorkerProvider; -import com.a.eye.skywalking.collector.actor.ClusterWorkerContext; -import com.a.eye.skywalking.collector.actor.Role; +import com.a.eye.skywalking.collector.actor.*; import com.a.eye.skywalking.collector.role.TraceSegmentReceiverRole; import com.a.eye.skywalking.collector.worker.WorkerConfig; import com.a.eye.skywalking.collector.worker.application.ApplicationMain; @@ -21,12 +18,12 @@ public class TraceSegmentReceiver extends AbstractClusterWorker { private Logger logger = LogManager.getFormatterLogger(TraceSegmentReceiver.class); - public TraceSegmentReceiver(Role role, ClusterWorkerContext clusterContext) throws Exception { - super(role, clusterContext); + public TraceSegmentReceiver(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { + super(role, clusterContext, selfContext); } @Override - public void preStart() throws Exception { + public void preStart() throws ProviderNotFountException { getClusterContext().findProvider(ApplicationMain.Role.INSTANCE).create(getClusterContext(), getSelfContext()); getClusterContext().findProvider(ApplicationRefMain.Role.INSTANCE).create(getClusterContext(), getSelfContext()); } @@ -59,8 +56,8 @@ public class TraceSegmentReceiver extends AbstractClusterWorker { } @Override - public Class workerClass() { - return TraceSegmentReceiver.class; + public TraceSegmentReceiver workerInstance(ClusterWorkerContext clusterContext) { + return new TraceSegmentReceiver(role(), clusterContext, new LocalWorkerContext()); } } diff --git a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java index b22d68b3a..bd0bbd5a9 100644 --- a/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java +++ b/skywalking-collector/skywalking-collector-worker/src/test/java/com/a/eye/skywalking/collector/worker/StartUpTestCase.java @@ -3,10 +3,8 @@ package com.a.eye.skywalking.collector.worker; import akka.actor.ActorRef; import akka.actor.ActorSelection; import akka.actor.ActorSystem; -import com.a.eye.skywalking.collector.actor.WorkersCreator; import com.a.eye.skywalking.collector.cluster.ClusterConfig; import com.a.eye.skywalking.collector.cluster.ClusterConfigInitializer; -import com.a.eye.skywalking.collector.cluster.WorkersListener; import com.a.eye.skywalking.collector.worker.receiver.TraceSegmentReceiver; import com.a.eye.skywalking.collector.worker.storage.EsClient; import com.a.eye.skywalking.sniffer.mock.trace.TraceSegmentBuilderFactory; @@ -16,7 +14,6 @@ import com.a.eye.skywalking.trace.proto.SegmentRefMessage; import com.a.eye.skywalking.trace.tag.Tags; import com.typesafe.config.Config; import com.typesafe.config.ConfigFactory; -import org.junit.Test; /** * @author pengys5 @@ -35,7 +32,7 @@ public class StartUpTestCase { withFallback(ConfigFactory.parseString("akka.cluster.seed-nodes=" + ClusterConfig.Cluster.nodes)). withFallback(ConfigFactory.load("application.conf")); ActorSystem system = ActorSystem.create(ClusterConfig.Cluster.appname, config); - WorkersCreator.INSTANCE.boot(system); +// WorkersCreator.INSTANCE.boot(system); EsClient.boot();