From baf542a7a34b3bf797ada576eb0861f57337b6ca Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Fri, 17 Mar 2017 18:40:22 +0800 Subject: [PATCH] refactor ClusterContext and LocalContext implements Lookup interface that let the user of context just use lookup or findprovider method. --- .../skywalking/collector/CollectorSystem.java | 13 ++++++----- .../actor/AbstractClusterWorkerProvider.java | 8 +++---- .../AbstractLocalAsyncWorkerProvider.java | 10 +++++--- .../actor/AbstractLocalSyncWorker.java | 6 +++++ .../AbstractLocalSyncWorkerProvider.java | 9 +++++--- .../collector/actor/AbstractWorker.java | 8 +++++-- .../actor/AbstractWorkerProvider.java | 23 ++++++++++++++++--- .../skywalking/collector/actor/Context.java | 6 +---- .../skywalking/collector/actor/LookUp.java | 11 +++++++++ .../skywalking/collector/actor/Promise.java | 18 +++++++++++++++ .../skywalking/collector/actor/Provider.java | 2 +- .../collector/actor/TestClusterWorker.java | 4 ++-- .../actor/TestClusterWorkerTestCase.java | 3 +++ .../collector/actor/TestLocalSyncWorker.java | 3 ++- .../worker/application/ApplicationMain.java | 13 ++++++----- .../application/receiver/DAGNodeReceiver.java | 2 +- .../receiver/NodeInstanceReceiver.java | 2 +- .../receiver/ResponseCostReceiver.java | 2 +- .../receiver/ResponseSummaryReceiver.java | 2 +- .../applicationref/ApplicationRefMain.java | 5 ++-- .../receiver/DAGNodeRefReceiver.java | 2 +- .../worker/receiver/TraceSegmentReceiver.java | 4 ++-- 22 files changed, 111 insertions(+), 45 deletions(-) create mode 100644 skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LookUp.java create mode 100644 skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Promise.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 1c599d5cc..4140dea62 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 @@ -21,16 +21,15 @@ public class CollectorSystem { private ClusterWorkerContext clusterContext; - public ClusterWorkerContext getClusterContext() { + public LookUp getClusterContext() { return clusterContext; } - public void boot() throws Exception { + public void boot() throws UsedRoleNameException, ProviderNotFountException { createAkkaSystem(); createListener(); loadLocalProviders(); - - createClusterWorker(); + createClusterWorkers(); } public void terminate() { @@ -55,12 +54,13 @@ public class CollectorSystem { clusterContext.getAkkaSystem().actorOf(Props.create(WorkersListener.class, clusterContext), WorkersListener.WorkName); } - private void createClusterWorker() throws Exception { + private void createClusterWorkers() throws ProviderNotFountException { ServiceLoader clusterServiceLoader = ServiceLoader.load(AbstractClusterWorkerProvider.class); for (AbstractClusterWorkerProvider provider : clusterServiceLoader) { logger.info("create {%s} worker using java service loader", provider.workerNum()); + provider.setClusterContext(clusterContext); for (int i = 1; i <= provider.workerNum(); i++) { - provider.create(clusterContext, new LocalWorkerContext()); + provider.create(AbstractWorker.noOwner()); } } } @@ -68,6 +68,7 @@ public class CollectorSystem { private void loadLocalProviders() throws UsedRoleNameException { ServiceLoader clusterServiceLoader = ServiceLoader.load(AbstractLocalWorkerProvider.class); for (AbstractLocalWorkerProvider provider : clusterServiceLoader) { + provider.setClusterContext(clusterContext); clusterContext.putProvider(provider); } } 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 b37e87584..39a983843 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 @@ -11,16 +11,16 @@ public abstract class AbstractClusterWorkerProvider extends AbstractLocalWorkerProvider { @Override - final public WorkerRef onCreate(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException { - T localSyncWorker = (T) workerInstance(clusterContext); + final public WorkerRef onCreate(LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException { + T localSyncWorker = (T) workerInstance(getClusterContext()); localSyncWorker.preStart(); LocalSyncWorkerRef workerRef = new LocalSyncWorkerRef(role(), localSyncWorker); - localContext.put(workerRef); + + if (localContext != null) { + localContext.put(workerRef); + } return workerRef; } } 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 a830708db..6af326564 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 @@ -21,15 +21,19 @@ public abstract class AbstractWorker { public abstract void work(Object message) throws Exception; - final public LocalWorkerContext getSelfContext() { + final public LookUp getSelfContext() { return selfContext; } - final public ClusterWorkerContext getClusterContext() { + final public LookUp getClusterContext() { return clusterContext; } final public Role getRole() { return role; } + + final public static AbstractWorker noOwner() { + return null; + } } 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 dd0b4aee8..8cc670ab0 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 @@ -5,16 +5,33 @@ package com.a.eye.skywalking.collector.actor; */ public abstract class AbstractWorkerProvider implements Provider { + private ClusterWorkerContext clusterContext; + public abstract Role role(); public abstract T workerInstance(ClusterWorkerContext clusterContext); - public abstract WorkerRef onCreate(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException; + public abstract WorkerRef onCreate(LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException; - final public WorkerRef create(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException { + final public void setClusterContext(ClusterWorkerContext clusterContext) { + this.clusterContext = clusterContext; + } + + final protected ClusterWorkerContext getClusterContext() { + return clusterContext; + } + + final public WorkerRef create(AbstractWorker workerOwner) throws IllegalArgumentException, ProviderNotFountException { if (workerInstance(clusterContext) == null) { throw new IllegalArgumentException("cannot get worker instance with nothing obtained from workerInstance()"); } - return onCreate(clusterContext, localContext); + + if (workerOwner == null) { + return onCreate(null); + } else if (workerOwner.getSelfContext() instanceof LocalWorkerContext) { + return onCreate((LocalWorkerContext) workerOwner.getSelfContext()); + } else { + throw new IllegalArgumentException("the argument of workerOwner is Illegal"); + } } } 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 9cd13576e..2411812e7 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 @@ -3,14 +3,10 @@ package com.a.eye.skywalking.collector.actor; /** * @author pengys5 */ -public interface Context { - - AbstractWorkerProvider findProvider(Role role) throws ProviderNotFountException; +public interface Context extends LookUp { void putProvider(AbstractWorkerProvider provider) throws UsedRoleNameException; - WorkerRefs lookup(Role role) throws WorkerNotFountException; - void put(WorkerRef workerRef); void remove(WorkerRef workerRef); diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LookUp.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LookUp.java new file mode 100644 index 000000000..05bb5ee5a --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/LookUp.java @@ -0,0 +1,11 @@ +package com.a.eye.skywalking.collector.actor; + +/** + * @author pengys5 + */ +public interface LookUp { + + WorkerRefs lookup(Role role) throws WorkerNotFountException; + + Provider findProvider(Role role) throws ProviderNotFountException; +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Promise.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Promise.java new file mode 100644 index 000000000..7e5ff7adf --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Promise.java @@ -0,0 +1,18 @@ +package com.a.eye.skywalking.collector.actor; + +/** + * @author pengys5 + */ +public class Promise { + private boolean isTold = false; + private Object value; + + protected void completed(Object value) { + this.value = value; + isTold = true; + } + + public boolean isTold() { + return isTold; + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Provider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Provider.java index dbe8de929..0d767033b 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Provider.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/Provider.java @@ -5,5 +5,5 @@ package com.a.eye.skywalking.collector.actor; */ public interface Provider { - WorkerRef create(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws Exception; + WorkerRef create(AbstractWorker workerOwner) throws IllegalArgumentException, ProviderNotFountException; } 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 aeee348fe..82ed8057c 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 @@ -15,8 +15,8 @@ public class TestClusterWorker extends AbstractClusterWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(TestLocalSyncWorker.TestLocalSyncWorkerRole.INSTANCE).create(getClusterContext(), getSelfContext()); - getClusterContext().findProvider(TestLocalAsyncWorker.TestLocalASyncWorkerRole.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(TestLocalSyncWorker.TestLocalSyncWorkerRole.INSTANCE).create(this); + getClusterContext().findProvider(TestLocalAsyncWorker.TestLocalASyncWorkerRole.INSTANCE).create(this); } @Override 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 3b056621d..bdb75e8f5 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 @@ -12,15 +12,18 @@ public class TestClusterWorkerTestCase { private CollectorSystem collectorSystem; +// @Before public void createSystem() throws Exception { collectorSystem = new CollectorSystem(); collectorSystem.boot(); } +// @Before public void terminateSystem() { collectorSystem.terminate(); } +// @Test public void testTellWorker() throws Exception { WorkerRefs workerRefs = collectorSystem.getClusterContext().lookup(TestClusterWorker.TestClusterWorkerRole.INSTANCE); workerRefs.tell("Print"); 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 51128712f..185173852 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 @@ -18,12 +18,13 @@ public class TestLocalSyncWorker extends AbstractLocalSyncWorker { } @Override - public void work(Object message) throws Exception { + public Object onWork(Object message) throws Exception { if (message.equals("TellLocalWorker")) { System.out.println("hello! "); } else { System.out.println("unhandled"); } + return "Hello"; } public static class Factory extends AbstractLocalSyncWorkerProvider { diff --git a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java index 20da536ba..74b9669bf 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java +++ b/skywalking-collector/skywalking-collector-worker/src/main/java/com/a/eye/skywalking/collector/worker/application/ApplicationMain.java @@ -29,15 +29,15 @@ public class ApplicationMain extends AbstractLocalSyncWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(DAGNodeAnalysis.Role.INSTANCE).create(getClusterContext(), getSelfContext()); - getClusterContext().findProvider(NodeInstanceAnalysis.Role.INSTANCE).create(getClusterContext(), getSelfContext()); - getClusterContext().findProvider(ResponseCostAnalysis.Role.INSTANCE).create(getClusterContext(), getSelfContext()); - getClusterContext().findProvider(ResponseSummaryAnalysis.Role.INSTANCE).create(getClusterContext(), getSelfContext()); - getClusterContext().findProvider(TraceSegmentRecordPersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(DAGNodeAnalysis.Role.INSTANCE).create(this); + getClusterContext().findProvider(NodeInstanceAnalysis.Role.INSTANCE).create(this); + getClusterContext().findProvider(ResponseCostAnalysis.Role.INSTANCE).create(this); + getClusterContext().findProvider(ResponseSummaryAnalysis.Role.INSTANCE).create(this); + getClusterContext().findProvider(TraceSegmentRecordPersistence.Role.INSTANCE).create(this); } @Override - public void work(Object message) throws Exception { + public Object onWork(Object message) throws Exception { if (message instanceof TraceSegmentReceiver.TraceSegmentTimeSlice) { logger.debug("begin translate TraceSegment Object to JsonObject"); TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment = (TraceSegmentReceiver.TraceSegmentTimeSlice) message; @@ -49,6 +49,7 @@ public class ApplicationMain extends AbstractLocalSyncWorker { sendToResponseCostPersistence(traceSegment); sendToResponseSummaryPersistence(traceSegment); } + return null; } public static class Factory extends AbstractLocalSyncWorkerProvider { 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 fb9cf3fab..b44e18404 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 @@ -22,7 +22,7 @@ public class DAGNodeReceiver extends AbstractClusterWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(DAGNodePersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(DAGNodePersistence.Role.INSTANCE).create(this); } @Override 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 739b51461..baefbe616 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 @@ -22,7 +22,7 @@ public class NodeInstanceReceiver extends AbstractClusterWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(NodeInstancePersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(NodeInstancePersistence.Role.INSTANCE).create(this); } @Override 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 788b19a03..b80f9154e 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 @@ -22,7 +22,7 @@ public class ResponseCostReceiver extends AbstractClusterWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(ResponseCostPersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(ResponseCostPersistence.Role.INSTANCE).create(this); } @Override 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 e533305b5..f4cd961d4 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 @@ -22,7 +22,7 @@ public class ResponseSummaryReceiver extends AbstractClusterWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(ResponseSummaryPersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(ResponseSummaryPersistence.Role.INSTANCE).create(this); } @Override 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 db484a475..a02b3f59b 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 @@ -21,11 +21,11 @@ public class ApplicationRefMain extends AbstractLocalSyncWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(DAGNodeRefAnalysis.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(DAGNodeRefAnalysis.Role.INSTANCE).create(this); } @Override - public void work(Object message) throws Exception { + public Object onWork(Object message) throws Exception { TraceSegmentReceiver.TraceSegmentTimeSlice traceSegment = (TraceSegmentReceiver.TraceSegmentTimeSlice) message; TraceSegmentRef traceSegmentRef = traceSegment.getTraceSegment().getPrimaryRef(); @@ -36,6 +36,7 @@ public class ApplicationRefMain extends AbstractLocalSyncWorker { DAGNodeRefAnalysis.Metric nodeRef = new DAGNodeRefAnalysis.Metric(traceSegment.getMinute(), traceSegment.getSecond(), front, behind); getSelfContext().lookup(DAGNodeRefAnalysis.Role.INSTANCE).tell(nodeRef); } + return null; } public static class Factory extends AbstractLocalSyncWorkerProvider { 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 4c10298f1..a7493460b 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 @@ -24,7 +24,7 @@ public class DAGNodeRefReceiver extends AbstractClusterWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(DAGNodeRefPersistence.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(DAGNodeRefPersistence.Role.INSTANCE).create(this); } @Override 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 679dc55b4..2ea786e1f 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 @@ -24,8 +24,8 @@ public class TraceSegmentReceiver extends AbstractClusterWorker { @Override public void preStart() throws ProviderNotFountException { - getClusterContext().findProvider(ApplicationMain.Role.INSTANCE).create(getClusterContext(), getSelfContext()); - getClusterContext().findProvider(ApplicationRefMain.Role.INSTANCE).create(getClusterContext(), getSelfContext()); + getClusterContext().findProvider(ApplicationMain.Role.INSTANCE).create(this); + getClusterContext().findProvider(ApplicationRefMain.Role.INSTANCE).create(this); } @Override