refactor ClusterContext and LocalContext implements Lookup interface that let the user of context just use lookup or findprovider method.
This commit is contained in:
parent
934aa6bb6c
commit
baf542a7a3
|
|
@ -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<AbstractClusterWorkerProvider> 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<AbstractLocalWorkerProvider> clusterServiceLoader = ServiceLoader.load(AbstractLocalWorkerProvider.class);
|
||||
for (AbstractLocalWorkerProvider provider : clusterServiceLoader) {
|
||||
provider.setClusterContext(clusterContext);
|
||||
clusterContext.putProvider(provider);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -11,16 +11,16 @@ public abstract class AbstractClusterWorkerProvider<T extends AbstractClusterWor
|
|||
public abstract int workerNum();
|
||||
|
||||
@Override
|
||||
final public WorkerRef onCreate(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException {
|
||||
final public WorkerRef onCreate(LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException {
|
||||
int num = ClusterWorkerRefCounter.INSTANCE.incrementAndGet(role());
|
||||
|
||||
T clusterWorker = (T) workerInstance(clusterContext);
|
||||
T clusterWorker = (T) workerInstance(getClusterContext());
|
||||
clusterWorker.preStart();
|
||||
|
||||
ActorRef actorRef = clusterContext.getAkkaSystem().actorOf(Props.create(AbstractClusterWorker.WorkerWithAkka.class, clusterWorker), role() + "_" + num);
|
||||
ActorRef actorRef = getClusterContext().getAkkaSystem().actorOf(Props.create(AbstractClusterWorker.WorkerWithAkka.class, clusterWorker), role() + "_" + num);
|
||||
|
||||
ClusterWorkerRef workerRef = new ClusterWorkerRef(actorRef, role());
|
||||
clusterContext.put(workerRef);
|
||||
getClusterContext().put(workerRef);
|
||||
return workerRef;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -14,8 +14,8 @@ public abstract class AbstractLocalAsyncWorkerProvider<T extends AbstractLocalAs
|
|||
public abstract int queueSize();
|
||||
|
||||
@Override
|
||||
final public WorkerRef onCreate(ClusterWorkerContext clusterContext, LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException {
|
||||
T localAsyncWorker = (T) workerInstance(clusterContext);
|
||||
final public WorkerRef onCreate(LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFountException {
|
||||
T localAsyncWorker = (T) workerInstance(getClusterContext());
|
||||
localAsyncWorker.preStart();
|
||||
|
||||
// Specify the size of the ring buffer, must be power of 2.
|
||||
|
|
@ -37,7 +37,11 @@ public abstract class AbstractLocalAsyncWorkerProvider<T extends AbstractLocalAs
|
|||
disruptor.start();
|
||||
|
||||
LocalAsyncWorkerRef workerRef = new LocalAsyncWorkerRef(role(), disruptorWorker);
|
||||
localContext.put(workerRef);
|
||||
|
||||
if (localContext != null) {
|
||||
localContext.put(workerRef);
|
||||
}
|
||||
|
||||
return workerRef;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -7,4 +7,10 @@ public abstract class AbstractLocalSyncWorker extends AbstractLocalWorker {
|
|||
public AbstractLocalSyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
|
||||
super(role, clusterContext, selfContext);
|
||||
}
|
||||
|
||||
@Override
|
||||
final public void work(Object message) throws Exception {
|
||||
}
|
||||
|
||||
public abstract Object onWork(Object message) throws Exception;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -6,12 +6,15 @@ package com.a.eye.skywalking.collector.actor;
|
|||
public abstract class AbstractLocalSyncWorkerProvider<T extends AbstractLocalSyncWorker> extends AbstractLocalWorkerProvider<T> {
|
||||
|
||||
@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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,16 +5,33 @@ package com.a.eye.skywalking.collector.actor;
|
|||
*/
|
||||
public abstract class AbstractWorkerProvider<T extends AbstractWorker> 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");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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);
|
||||
|
|
|
|||
|
|
@ -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;
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
}
|
||||
|
|
@ -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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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");
|
||||
|
|
|
|||
|
|
@ -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<TestLocalSyncWorker> {
|
||||
|
|
|
|||
|
|
@ -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<ApplicationMain> {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<ApplicationRefMain> {
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in New Issue