modify worker member
This commit is contained in:
parent
34a8925430
commit
c19158142e
|
|
@ -1,25 +0,0 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorSystem;
|
||||
import akka.actor.Props;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AbstractClusterWorkerProvider extends AbstractWorkerProvider<ActorSystem> {
|
||||
|
||||
@Override
|
||||
public void createWorker(ActorSystem system) {
|
||||
if (workerClass() == null) {
|
||||
throw new IllegalArgumentException("cannot createInstance() with nothing obtained from workerClass()");
|
||||
}
|
||||
if (workerNum() <= 0) {
|
||||
throw new IllegalArgumentException("cannot createInstance() with obtained from workerNum() must greater than 0");
|
||||
}
|
||||
|
||||
for (int i = 1; i <= workerNum(); i++) {
|
||||
system.actorOf(Props.create(workerClass()), roleName() + "_" + i);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
|
@ -1,28 +0,0 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.selector.WorkerSelector;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AbstractLocalWorker<T> implements Worker {
|
||||
|
||||
/**
|
||||
* Receive the message to analyse.
|
||||
*
|
||||
* @param message is the data send from the forward worker
|
||||
* @throws Throwable is the exception thrown by that worker implementation processing
|
||||
*/
|
||||
public abstract void receive(Object message) throws Throwable;
|
||||
|
||||
/**
|
||||
* Send analysed data to next Worker.
|
||||
*
|
||||
* @param targetWorkerProvider is the worker provider to create worker instance.
|
||||
* @param message is the data used to send to next worker.
|
||||
* @throws Throwable
|
||||
*/
|
||||
public void tell(AbstractLocalWorkerProvider targetWorkerProvider, T message) throws Throwable {
|
||||
LocalSystem.actorFor(targetWorkerProvider.getClass(), targetWorkerProvider.roleName());
|
||||
}
|
||||
}
|
||||
|
|
@ -1,28 +0,0 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorSystem;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AbstractLocalWorkerProvider extends AbstractWorkerProvider<LocalSystem> {
|
||||
|
||||
/**
|
||||
* Use {@link ActorSystem} to Create worker instance with the {@link #workerClass()} method returned class.
|
||||
*
|
||||
* @param system is a akka {@link ActorSystem} instance.
|
||||
*/
|
||||
@Override
|
||||
public void createWorker(LocalSystem system) {
|
||||
if (workerClass() == null) {
|
||||
throw new IllegalArgumentException("cannot createInstance() with nothing obtained from workerClass()");
|
||||
}
|
||||
if (workerNum() <= 0) {
|
||||
throw new IllegalArgumentException("cannot createInstance() with obtained from workerNum() must greater than 0");
|
||||
}
|
||||
|
||||
for (int i = 1; i <= workerNum(); i++) {
|
||||
LocalSystem.actorOf(getClass(), roleName());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,44 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
import com.a.eye.skywalking.collector.actor.selector.WorkerSelector;
|
||||
import com.a.eye.skywalking.collector.cluster.WorkersRefCenter;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AbstractMember<T> {
|
||||
|
||||
private ActorRef actorRef;
|
||||
|
||||
public ActorRef getSelf() {
|
||||
return actorRef;
|
||||
}
|
||||
|
||||
public void creatorRef(ActorRef actorRef) {
|
||||
this.actorRef = actorRef;
|
||||
}
|
||||
|
||||
/**
|
||||
* Receive the message to analyse.
|
||||
*
|
||||
* @param message is the data send from the forward worker
|
||||
* @throws Throwable is the exception thrown by that worker implementation processing
|
||||
*/
|
||||
public abstract void receive(Object message) throws Throwable;
|
||||
|
||||
/**
|
||||
* Send analysed data to next Worker.
|
||||
*
|
||||
* @param targetWorkerProvider is the worker provider to create worker instance.
|
||||
* @param selector is the selector to select a same role worker instance form cluster.
|
||||
* @param message is the data used to send to next worker.
|
||||
* @throws Throwable
|
||||
*/
|
||||
public void tell(AbstractWorkerProvider targetWorkerProvider, WorkerSelector selector, T message) throws Throwable {
|
||||
List<WorkerRef> availableWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName());
|
||||
selector.select(availableWorks, message).tell(message, getSelf());
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,28 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AbstractMemberProvider {
|
||||
public abstract Class memberClass();
|
||||
|
||||
public void createWorker(MemberSystem system, ActorRef actorRef) {
|
||||
if (memberClass() == null) {
|
||||
throw new IllegalArgumentException("cannot createInstance() with nothing obtained from memberClass()");
|
||||
}
|
||||
|
||||
AbstractMember member = system.memberOf(memberClass(), roleName());
|
||||
member.creatorRef(actorRef);
|
||||
}
|
||||
|
||||
/**
|
||||
* Use {@link #memberClass()} method returned class's simple name as a role name.
|
||||
*
|
||||
* @return is role of Worker
|
||||
*/
|
||||
protected String roleName() {
|
||||
return memberClass().getSimpleName();
|
||||
}
|
||||
}
|
||||
|
|
@ -35,7 +35,9 @@ import java.util.List;
|
|||
* }
|
||||
* }}}
|
||||
*/
|
||||
public abstract class AbstractWorker<T> extends UntypedActor implements Worker{
|
||||
public abstract class AbstractWorker<T> extends UntypedActor {
|
||||
|
||||
private MemberSystem memberSystem = new MemberSystem();
|
||||
|
||||
/**
|
||||
* Receive the message to analyse.
|
||||
|
|
@ -76,13 +78,8 @@ public abstract class AbstractWorker<T> extends UntypedActor implements Worker{
|
|||
* @throws Throwable
|
||||
*/
|
||||
public void tell(AbstractWorkerProvider targetWorkerProvider, WorkerSelector selector, T message) throws Throwable {
|
||||
if (targetWorkerProvider instanceof AbstractLocalWorkerProvider) {
|
||||
Worker worker = LocalSystem.actorFor(targetWorkerProvider.getClass(), targetWorkerProvider.roleName());
|
||||
worker.receive(message);
|
||||
} else if (targetWorkerProvider instanceof AbstractClusterWorkerProvider) {
|
||||
List<WorkerRef> availableWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName());
|
||||
selector.select(availableWorks, message).tell(message, getSelf());
|
||||
}
|
||||
List<WorkerRef> availableWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName());
|
||||
selector.select(availableWorks, message).tell(message, getSelf());
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -97,4 +94,8 @@ public abstract class AbstractWorker<T> extends UntypedActor implements Worker{
|
|||
getContext().actorSelection(member.address() + "/user/" + WorkersListener.WorkName).tell(registerMessage, getSelf());
|
||||
}
|
||||
}
|
||||
|
||||
public MemberSystem getMemberContext() {
|
||||
return memberSystem;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -31,7 +31,18 @@ public abstract class AbstractWorkerProvider<T> {
|
|||
|
||||
public abstract int workerNum();
|
||||
|
||||
public abstract void createWorker(T system);
|
||||
public void createWorker(ActorSystem system) {
|
||||
if (workerClass() == null) {
|
||||
throw new IllegalArgumentException("cannot createInstance() with nothing obtained from workerClass()");
|
||||
}
|
||||
if (workerNum() <= 0) {
|
||||
throw new IllegalArgumentException("cannot createInstance() with obtained from workerNum() must greater than 0");
|
||||
}
|
||||
|
||||
for (int i = 1; i <= workerNum(); i++) {
|
||||
system.actorOf(Props.create(workerClass()), roleName() + "_" + i);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Use {@link #workerClass()} method returned class's simple name as a role name.
|
||||
|
|
|
|||
|
|
@ -1,27 +0,0 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class LocalSystem {
|
||||
|
||||
private static Map<String, Worker> context = new HashMap();
|
||||
|
||||
public static void actorOf(Class clazz, String role) {
|
||||
try {
|
||||
Worker classInstance = (Worker) clazz.newInstance();
|
||||
context.put(clazz.getName() + "_" + role, classInstance);
|
||||
} catch (InstantiationException e) {
|
||||
e.printStackTrace();
|
||||
} catch (IllegalAccessException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
}
|
||||
|
||||
public static Worker actorFor(Class clazz, String role) {
|
||||
return context.get(clazz.getName() + "_" + role);
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,29 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class MemberSystem {
|
||||
|
||||
private Map<String, AbstractMember> memberMap = new HashMap();
|
||||
|
||||
public AbstractMember memberOf(Class clazz, String role) {
|
||||
try {
|
||||
AbstractMember member = (AbstractMember) clazz.newInstance();
|
||||
memberMap.put(role, member);
|
||||
return member;
|
||||
} catch (InstantiationException e) {
|
||||
e.printStackTrace();
|
||||
} catch (IllegalAccessException e) {
|
||||
e.printStackTrace();
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
public AbstractMember memberFor(String role) {
|
||||
return memberMap.get(role);
|
||||
}
|
||||
}
|
||||
|
|
@ -1,9 +0,0 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public interface Worker {
|
||||
|
||||
public void receive(Object message) throws Throwable;
|
||||
}
|
||||
|
|
@ -23,11 +23,5 @@ public enum WorkersCreator {
|
|||
for (AbstractClusterWorkerProvider provider : clusterServiceLoader) {
|
||||
provider.createWorker(system);
|
||||
}
|
||||
|
||||
LocalSystem localSystem = new LocalSystem();
|
||||
ServiceLoader<AbstractLocalWorkerProvider> localServiceLoader = ServiceLoader.load(AbstractLocalWorkerProvider.class);
|
||||
for (AbstractLocalWorkerProvider provider : localServiceLoader) {
|
||||
provider.createWorker(localSystem);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -0,0 +1,30 @@
|
|||
package com.a.eye.skywalking.collector.worker.application;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractMember;
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import com.a.eye.skywalking.collector.worker.application.member.ApplicationDiscoverFactory;
|
||||
import com.a.eye.skywalking.collector.worker.application.member.ApplicationDiscoverMember;
|
||||
import com.a.eye.skywalking.trace.TraceSegment;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationWorker extends AbstractWorker {
|
||||
|
||||
@Override
|
||||
public void preStart() throws Exception {
|
||||
ApplicationDiscoverFactory factory = new ApplicationDiscoverFactory();
|
||||
factory.createWorker(getMemberContext(), getSelf());
|
||||
|
||||
super.preStart();
|
||||
}
|
||||
|
||||
@Override
|
||||
public void receive(Object message) throws Throwable {
|
||||
if (message instanceof TraceSegment) {
|
||||
TraceSegment traceSegment = (TraceSegment) message;
|
||||
AbstractMember discoverMember = getMemberContext().memberFor(ApplicationDiscoverMember.class.getSimpleName());
|
||||
discoverMember.receive(traceSegment);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -1,15 +1,14 @@
|
|||
package com.a.eye.skywalking.collector.worker.metric;
|
||||
package com.a.eye.skywalking.collector.worker.application;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationDiscoverFactory extends AbstractWorkerProvider {
|
||||
|
||||
public class ApplicationWorkerFactory extends AbstractWorkerProvider {
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
return ApplicationDiscoverMetric.class;
|
||||
return ApplicationWorker.class;
|
||||
}
|
||||
|
||||
@Override
|
||||
|
|
@ -0,0 +1,14 @@
|
|||
package com.a.eye.skywalking.collector.worker.application.member;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractMemberProvider;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationDiscoverFactory extends AbstractMemberProvider {
|
||||
|
||||
@Override
|
||||
public Class memberClass() {
|
||||
return ApplicationDiscoverMember.class;
|
||||
}
|
||||
}
|
||||
|
|
@ -1,17 +1,17 @@
|
|||
package com.a.eye.skywalking.collector.worker.metric;
|
||||
package com.a.eye.skywalking.collector.worker.application.member;
|
||||
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import com.a.eye.skywalking.collector.actor.AbstractMember;
|
||||
import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
|
||||
import com.a.eye.skywalking.collector.worker.persistence.ApplicationMessage;
|
||||
import com.a.eye.skywalking.collector.worker.persistence.ApplicationPersistenceFactory;
|
||||
import com.a.eye.skywalking.collector.worker.application.persistence.ApplicationMessage;
|
||||
import com.a.eye.skywalking.collector.worker.application.persistence.ApplicationPersistenceFactory;
|
||||
import com.a.eye.skywalking.trace.TraceSegment;
|
||||
import com.a.eye.skywalking.trace.tag.Tags;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationDiscoverMetric extends AbstractWorker {
|
||||
public class ApplicationDiscoverMember extends AbstractMember {
|
||||
|
||||
@Override
|
||||
public void receive(Object message) throws Throwable {
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package com.a.eye.skywalking.collector.worker.persistence;
|
||||
package com.a.eye.skywalking.collector.worker.application.persistence;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,6 +1,7 @@
|
|||
package com.a.eye.skywalking.collector.worker.persistence;
|
||||
package com.a.eye.skywalking.collector.worker.application.persistence;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import com.a.eye.skywalking.collector.worker.persistence.PersistenceMessage;
|
||||
import com.a.eye.skywalking.collector.worker.persistence.PersistenceWorker;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
|
@ -8,7 +9,7 @@ import java.util.Map;
|
|||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationPersistence extends AbstractWorker<Object> {
|
||||
public class ApplicationPersistence extends PersistenceWorker<Object> {
|
||||
|
||||
private Map<String, ApplicationMessage> appData = new HashMap();
|
||||
|
||||
|
|
@ -1,4 +1,4 @@
|
|||
package com.a.eye.skywalking.collector.worker.persistence;
|
||||
package com.a.eye.skywalking.collector.worker.application.persistence;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
|
||||
|
||||
|
|
@ -0,0 +1,13 @@
|
|||
package com.a.eye.skywalking.collector.worker.applicationref;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationRefWorker extends AbstractWorker {
|
||||
@Override
|
||||
public void receive(Object message) throws Throwable {
|
||||
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,18 @@
|
|||
package com.a.eye.skywalking.collector.worker.applicationref;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ApplicationRefWorkerFactory extends AbstractWorkerProvider {
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
return ApplicationRefWorker.class;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int workerNum() {
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
|
|
@ -2,11 +2,9 @@ package com.a.eye.skywalking.collector.worker.persistence;
|
|||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import com.a.eye.skywalking.collector.worker.RecordCollection;
|
||||
import com.a.eye.skywalking.collector.worker.application.persistence.ApplicationMessage;
|
||||
import com.google.gson.JsonObject;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
|
|
|
|||
|
|
@ -0,0 +1,9 @@
|
|||
package com.a.eye.skywalking.collector.worker.persistence;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class PersistenceWorker<T> extends AbstractWorker<T> {
|
||||
}
|
||||
|
|
@ -0,0 +1,17 @@
|
|||
package com.a.eye.skywalking.collector.worker.receiver;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import com.a.eye.skywalking.trace.TraceSegment;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class TraceSegmentReceiver extends AbstractWorker {
|
||||
|
||||
@Override
|
||||
public void receive(Object message) throws Throwable {
|
||||
if (message instanceof TraceSegment) {
|
||||
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -0,0 +1,19 @@
|
|||
package com.a.eye.skywalking.collector.worker.receiver;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorkerProvider;
|
||||
import com.a.eye.skywalking.collector.worker.application.member.ApplicationDiscoverMember;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class TraceSegmentReceiverFactory extends AbstractWorkerProvider {
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
return ApplicationDiscoverMember.class;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int workerNum() {
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
Loading…
Reference in New Issue