add comment for the module of skywalking-collector-cluster
This commit is contained in:
parent
b6388251af
commit
cf5c0a2091
|
|
@ -1,6 +1,5 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
import akka.actor.UntypedActor;
|
||||
import akka.cluster.ClusterEvent;
|
||||
import akka.cluster.Member;
|
||||
|
|
@ -13,18 +12,46 @@ import com.a.eye.skywalking.collector.cluster.WorkersRefCenter;
|
|||
import java.util.List;
|
||||
|
||||
/**
|
||||
* Abstract implementation of the {@link akka.actor.UntypedActor} that represents an
|
||||
* analysis unit. <code>AbstractWorker</code> implementation process the message in
|
||||
* {@link #receive(Object)} method.
|
||||
*
|
||||
* <p>
|
||||
* Subclasses must implement the abstract {@link #receive(Object)} method to process message.
|
||||
* Subclasses forbid to override the {@link #onReceive(Object)} method.
|
||||
* <p>
|
||||
* Here is an example on how to create and use an {@link AbstractWorker}:
|
||||
* <p>
|
||||
* {{{
|
||||
* public class SampleWorker extends AbstractWorker {
|
||||
*
|
||||
* @Override
|
||||
* public void receive(Object message) throws Throwable {
|
||||
* if (message.equals("Tell Next")) {
|
||||
* Object sendMessage = new Object();
|
||||
* tell(new NextSampleWorkerFactory(), RollingSelector.INSTANCE, sendMessage);
|
||||
* }
|
||||
* }
|
||||
* }
|
||||
* }}}
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AbstractWorker<T> extends UntypedActor {
|
||||
|
||||
final String workerRole;
|
||||
|
||||
public AbstractWorker(String workerRole) {
|
||||
this.workerRole = workerRole;
|
||||
}
|
||||
|
||||
public abstract void receive(Object message);
|
||||
/**
|
||||
* 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;
|
||||
|
||||
/**
|
||||
* Listening {@link ClusterEvent.MemberUp} and {@link ClusterEvent.CurrentClusterState}
|
||||
* cluster event, when event send from the member of {@link WorkersListener} then tell
|
||||
* the sender to register self.
|
||||
*/
|
||||
@Override
|
||||
public void onReceive(Object message) throws Throwable {
|
||||
if (message instanceof ClusterEvent.CurrentClusterState) {
|
||||
|
|
@ -42,14 +69,28 @@ public abstract class AbstractWorker<T> extends UntypedActor {
|
|||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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());
|
||||
}
|
||||
|
||||
/**
|
||||
* When member role is {@link WorkersListener#WorkName} then Select actor from context
|
||||
* and send register message to {@link WorkersListener}
|
||||
*
|
||||
* @param member is the new created or restart worker
|
||||
*/
|
||||
void register(Member member) {
|
||||
if (member.getRoles().equals(WorkersListener.WorkName)) {
|
||||
WorkerListenerMessage.RegisterMessage registerMessage = new WorkerListenerMessage.RegisterMessage(workerRole);
|
||||
WorkerListenerMessage.RegisterMessage registerMessage = new WorkerListenerMessage.RegisterMessage(getClass().getSimpleName());
|
||||
getContext().actorSelection(member.address() + "/user/" + WorkersListener.WorkName).tell(registerMessage, getSelf());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -4,6 +4,24 @@ import akka.actor.ActorSystem;
|
|||
import akka.actor.Props;
|
||||
|
||||
/**
|
||||
* The <code>AbstractWorkerProvider</code> should be implemented by any class whose
|
||||
* instances are intended to provide create instance of the {@link AbstractWorker}.
|
||||
* <p>
|
||||
* Here is an example on how to create and use an {@link AbstractWorkerProvider}:
|
||||
* <p>
|
||||
* {{{
|
||||
* public class SampleWorkerFactory extends AbstractWorkerProvider {
|
||||
*
|
||||
* @Override public Class workerClass() {
|
||||
* return SampleWorker.class;
|
||||
* }
|
||||
*
|
||||
* @Override public int workerNum() {
|
||||
* return Config.SampleWorkerNum;
|
||||
* }
|
||||
* }
|
||||
* }}}
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public abstract class AbstractWorkerProvider {
|
||||
|
|
@ -12,6 +30,11 @@ public abstract class AbstractWorkerProvider {
|
|||
|
||||
public abstract int workerNum();
|
||||
|
||||
/**
|
||||
* Use {@link ActorSystem} to Create worker instance with the {@link #workerClass()} method returned class.
|
||||
*
|
||||
* @param system is a akka {@link ActorSystem} instance.
|
||||
*/
|
||||
public void createWorker(ActorSystem system) {
|
||||
if (workerClass() == null) {
|
||||
throw new IllegalArgumentException("cannot createWorker() with nothing obtained from workerClass()");
|
||||
|
|
@ -25,7 +48,12 @@ public abstract class AbstractWorkerProvider {
|
|||
}
|
||||
}
|
||||
|
||||
protected String roleName(){
|
||||
/**
|
||||
* Use {@link #workerClass()} method returned class's simple name as a role name.
|
||||
*
|
||||
* @return is role of Worker
|
||||
*/
|
||||
protected String roleName() {
|
||||
return workerClass().getSimpleName();
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -3,30 +3,25 @@ package com.a.eye.skywalking.collector.actor;
|
|||
import akka.actor.ActorRef;
|
||||
|
||||
/**
|
||||
* The Worker reference
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public class WorkerRef {
|
||||
final ActorRef actorRef;
|
||||
|
||||
public WorkerRef(ActorRef actorRef) {
|
||||
final String workerRole;
|
||||
|
||||
public WorkerRef(ActorRef actorRef, String workerRole) {
|
||||
this.actorRef = actorRef;
|
||||
this.workerRole = workerRole;
|
||||
}
|
||||
|
||||
void tell(Object message, ActorRef actorRef) {
|
||||
actorRef.tell(message, actorRef);
|
||||
}
|
||||
|
||||
public String path(){
|
||||
return actorRef.path().toString();
|
||||
}
|
||||
|
||||
@Override
|
||||
public boolean equals(Object obj) {
|
||||
return actorRef.equals(obj);
|
||||
}
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return actorRef.toString();
|
||||
public String getWorkerRole() {
|
||||
return workerRole;
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,11 +5,19 @@ import akka.actor.ActorSystem;
|
|||
import java.util.ServiceLoader;
|
||||
|
||||
/**
|
||||
* <code>WorkersCreator</code> is a util that use Java Spi to create
|
||||
* workers by META-INF config file.
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public enum WorkersCreator {
|
||||
INSTANCE;
|
||||
|
||||
/**
|
||||
* create worker to use Java Spi.
|
||||
*
|
||||
* @param system is create by akka {@link ActorSystem}
|
||||
*/
|
||||
public void boot(ActorSystem system) {
|
||||
ServiceLoader<AbstractWorkerProvider> serviceLoader = ServiceLoader.load(AbstractWorkerProvider.class);
|
||||
for (AbstractWorkerProvider provider : serviceLoader) {
|
||||
|
|
|
|||
|
|
@ -1,6 +1,16 @@
|
|||
package com.a.eye.skywalking.collector.cluster;
|
||||
|
||||
import akka.actor.ActorSystem;
|
||||
|
||||
/**
|
||||
* A static class contains some config values of cluster.
|
||||
* {@link Cluster.Current#hostname} is a ip address of server which start this process.
|
||||
* {@link Cluster.Current#port} is a port of server use to bind
|
||||
* {@link Cluster.Current#roles} is a roles of workers that use to create workers which
|
||||
* has those role in this process.
|
||||
* {@link Cluster#nodes} is a nodes which cluster have.
|
||||
* {@link Cluster#appname} is a name of {@link ActorSystem} in cluster.
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ClusterConfig {
|
||||
|
|
|
|||
|
|
@ -9,6 +9,14 @@ import java.io.InputStream;
|
|||
import java.util.Properties;
|
||||
|
||||
/**
|
||||
* <code>ClusterConfigInitializer</code> Contains static methods for setting
|
||||
* {@link ClusterConfig} attributes value.
|
||||
*
|
||||
* <p>
|
||||
* The priority of value setting is
|
||||
* system property -> collector.config -> {@link ClusterConfig} default value
|
||||
* <p>
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public class ClusterConfigInitializer {
|
||||
|
|
@ -17,6 +25,11 @@ public class ClusterConfigInitializer {
|
|||
|
||||
public static final String ConfigFileName = "collector.config";
|
||||
|
||||
/**
|
||||
* Read config file to setting {@link ClusterConfig} then get system property to overwrite it.
|
||||
*
|
||||
* @param configFileName is the config file name, the file format is key-value pairs
|
||||
*/
|
||||
public static void initialize(String configFileName) {
|
||||
InputStream configFileStream = ClusterConfigInitializer.class.getResourceAsStream("/" + configFileName);
|
||||
|
||||
|
|
|
|||
|
|
@ -3,6 +3,10 @@ package com.a.eye.skywalking.collector.cluster;
|
|||
import java.io.Serializable;
|
||||
|
||||
/**
|
||||
* <code>WorkerListenerMessage</code> is a message just for the worker
|
||||
* implementation of the {@link com.a.eye.skywalking.collector.actor.AbstractWorker}
|
||||
* to register.
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public class WorkerListenerMessage {
|
||||
|
|
|
|||
|
|
@ -5,6 +5,15 @@ import akka.actor.Terminated;
|
|||
import akka.actor.UntypedActor;
|
||||
|
||||
/**
|
||||
* <code>WorkersListener</code> listening the register message from workers
|
||||
* implementation of the {@link com.a.eye.skywalking.collector.actor.AbstractWorker}
|
||||
* and terminated message from akka cluster.
|
||||
* <p>
|
||||
* when listened register message then begin to watch the state for this worker
|
||||
* and register to {@link WorkersRefCenter}.
|
||||
* <p>
|
||||
* when listened terminate message then unregister from {@link WorkersRefCenter}.
|
||||
*
|
||||
* @author pengys5
|
||||
*/
|
||||
public class WorkersListener extends UntypedActor {
|
||||
|
|
|
|||
|
|
@ -19,31 +19,25 @@ import java.util.concurrent.ConcurrentHashMap;
|
|||
public enum WorkersRefCenter {
|
||||
INSTANCE;
|
||||
|
||||
private Map<String, List<WorkerRef>> roleToActor = new ConcurrentHashMap();
|
||||
private Map<String, List<WorkerRef>> roleToWorkerRef = new ConcurrentHashMap();
|
||||
|
||||
private Map<WorkerRef, String> actorToRole = new ConcurrentHashMap();
|
||||
|
||||
// private Map<String, WorkerRef> pathToWorkerRef = new ConcurrentHashMap();
|
||||
private Map<ActorRef, WorkerRef> actorRefToWorkerRef = new ConcurrentHashMap<>();
|
||||
|
||||
public void register(ActorRef newActorRef, String workerRole) {
|
||||
if (!roleToActor.containsKey(workerRole)) {
|
||||
if (!roleToWorkerRef.containsKey(workerRole)) {
|
||||
List<WorkerRef> actorList = Collections.synchronizedList(new ArrayList<WorkerRef>());
|
||||
roleToActor.putIfAbsent(workerRole, actorList);
|
||||
roleToWorkerRef.putIfAbsent(workerRole, actorList);
|
||||
}
|
||||
|
||||
WorkerRef newWorkerRef = new WorkerRef(newActorRef);
|
||||
roleToActor.get(workerRole).add(newWorkerRef);
|
||||
actorToRole.put(newWorkerRef, workerRole);
|
||||
// pathToWorkerRef.put(newWorkerRef.path(), newWorkerRef);
|
||||
WorkerRef newWorkerRef = new WorkerRef(newActorRef, workerRole);
|
||||
roleToWorkerRef.get(workerRole).add(newWorkerRef);
|
||||
actorRefToWorkerRef.put(newActorRef, newWorkerRef);
|
||||
}
|
||||
|
||||
public void unregister(ActorRef newActorRef) {
|
||||
String workerRole = actorToRole.get(newActorRef.path());
|
||||
// WorkerRef workerRef = pathToWorkerRef.get(newActorRef.path());
|
||||
|
||||
roleToActor.get(workerRole).remove(newActorRef);
|
||||
actorToRole.remove(newActorRef);
|
||||
// pathToWorkerRef.remove(newActorRef.path());
|
||||
public void unregister(ActorRef oldActorRef) {
|
||||
WorkerRef oldWorkerRef = actorRefToWorkerRef.get(oldActorRef);
|
||||
roleToWorkerRef.get(oldWorkerRef.getWorkerRole()).remove(oldWorkerRef);
|
||||
actorRefToWorkerRef.remove(oldActorRef);
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -54,7 +48,7 @@ public enum WorkersRefCenter {
|
|||
* @throws NoAvailableWorkerException , when no available worker.
|
||||
*/
|
||||
public List<WorkerRef> availableWorks(String workerRole) throws NoAvailableWorkerException {
|
||||
List<WorkerRef> refs = roleToActor.get(workerRole);
|
||||
List<WorkerRef> refs = roleToWorkerRef.get(workerRole);
|
||||
if (refs == null || refs.size() == 0) {
|
||||
throw new NoAvailableWorkerException("role=" + workerRole + ", no available worker.");
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,20 +1,21 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.selector.RollingSelector;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class SpiTestWorker extends AbstractWorker {
|
||||
|
||||
public SpiTestWorker(String workerRole) {
|
||||
super(workerRole);
|
||||
}
|
||||
|
||||
@Override
|
||||
public void receive(Object message) {
|
||||
public void receive(Object message) throws Throwable {
|
||||
if (message.equals("Test1")) {
|
||||
getSender().tell("Yes", getSelf());
|
||||
} else if (message.equals("Test2")) {
|
||||
getSender().tell("No", getSelf());
|
||||
} else if (message.equals("Test3")) {
|
||||
Object sendMessage = new Object();
|
||||
tell(new SpiTestWorkerFactory(), RollingSelector.INSTANCE, sendMessage);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
@ -5,8 +5,6 @@ package com.a.eye.skywalking.collector.actor;
|
|||
*/
|
||||
public class SpiTestWorkerFactory extends AbstractWorkerProvider {
|
||||
|
||||
public static final String WorkerRole = "SpiTestWorker";
|
||||
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
return SpiTestWorker.class;
|
||||
|
|
|
|||
|
|
@ -6,6 +6,7 @@ import akka.actor.Props;
|
|||
import akka.actor.Terminated;
|
||||
import akka.pattern.Patterns;
|
||||
import akka.testkit.TestActorRef;
|
||||
import com.a.eye.skywalking.collector.actor.WorkerRef;
|
||||
import org.junit.After;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
|
|
@ -23,10 +24,19 @@ import java.util.concurrent.ConcurrentHashMap;
|
|||
public class WorkerListenerTestCase {
|
||||
|
||||
ActorSystem system;
|
||||
TestActorRef<WorkersListener> senderActorRef;
|
||||
TestActorRef<WorkersListener> receiveactorRef;
|
||||
|
||||
@Before
|
||||
public void createSystem() {
|
||||
public void initData() {
|
||||
system = ActorSystem.create();
|
||||
|
||||
final Props props = Props.create(WorkersListener.class);
|
||||
senderActorRef = TestActorRef.create(system, props, "WorkersListenerSender");
|
||||
receiveactorRef = TestActorRef.create(system, props, "WorkersListenerReceive");
|
||||
|
||||
WorkerListenerMessage.RegisterMessage message = new WorkerListenerMessage.RegisterMessage("WorkersListener");
|
||||
receiveactorRef.tell(message, senderActorRef);
|
||||
}
|
||||
|
||||
@After
|
||||
|
|
@ -34,44 +44,32 @@ public class WorkerListenerTestCase {
|
|||
system.terminate();
|
||||
system.awaitTermination();
|
||||
system = null;
|
||||
MemberModifier.field(WorkersRefCenter.class, "actorToRole").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
|
||||
MemberModifier.field(WorkersRefCenter.class, "roleToActor").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
|
||||
MemberModifier.field(WorkersRefCenter.class, "roleToWorkerRef").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
|
||||
MemberModifier.field(WorkersRefCenter.class, "actorRefToWorkerRef").set(WorkersRefCenter.INSTANCE, new ConcurrentHashMap());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testRegister() throws IllegalAccessException {
|
||||
final Props props = Props.create(WorkersListener.class);
|
||||
final TestActorRef<WorkersListener> senderActorRef = TestActorRef.create(system, props, "WorkersListenerSender");
|
||||
final TestActorRef<WorkersListener> receiveactorRef = TestActorRef.create(system, props, "WorkersListenerReceive");
|
||||
Map<ActorRef, WorkerRef> actorRefToWorkerRef = (Map<ActorRef, WorkerRef>) MemberModifier.field(WorkersRefCenter.class, "actorRefToWorkerRef").get(WorkersRefCenter.INSTANCE);
|
||||
ActorRef senderRefInWorkerRef = (ActorRef) MemberModifier.field(WorkerRef.class, "actorRef").get(actorRefToWorkerRef.get(senderActorRef));
|
||||
Assert.assertEquals(senderActorRef, senderRefInWorkerRef);
|
||||
|
||||
WorkerListenerMessage.RegisterMessage message = new WorkerListenerMessage.RegisterMessage("WorkersListener");
|
||||
receiveactorRef.tell(message, senderActorRef);
|
||||
|
||||
Map<ActorRef, String> actorToRole = (Map<ActorRef, String>) MemberModifier.field(WorkersRefCenter.class, "actorToRole").get(WorkersRefCenter.INSTANCE);
|
||||
Assert.assertEquals("WorkersListener", actorToRole.get(senderActorRef));
|
||||
|
||||
Map<String, List<ActorRef>> roleToActor = (Map<String, List<ActorRef>>) MemberModifier.field(WorkersRefCenter.class, "roleToActor").get(WorkersRefCenter.INSTANCE);
|
||||
ActorRef[] actorRefs = {senderActorRef};
|
||||
Assert.assertArrayEquals(actorRefs, roleToActor.get("WorkersListener").toArray());
|
||||
Map<String, List<WorkerRef>> roleToWorkerRef = (Map<String, List<WorkerRef>>) MemberModifier.field(WorkersRefCenter.class, "roleToWorkerRef").get(WorkersRefCenter.INSTANCE);
|
||||
WorkerRef workerRef = roleToWorkerRef.get("WorkersListener").get(0);
|
||||
senderRefInWorkerRef = (ActorRef) MemberModifier.field(WorkerRef.class, "actorRef").get(workerRef);
|
||||
Assert.assertEquals(senderActorRef, senderRefInWorkerRef);
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testTerminated() throws IllegalAccessException {
|
||||
final Props props = Props.create(WorkersListener.class);
|
||||
final TestActorRef<WorkersListener> senderActorRef = TestActorRef.create(system, props, "WorkersListenerSender");
|
||||
final TestActorRef<WorkersListener> receiveactorRef = TestActorRef.create(system, props, "WorkersListenerReceive");
|
||||
|
||||
WorkerListenerMessage.RegisterMessage message = new WorkerListenerMessage.RegisterMessage("WorkersListener");
|
||||
receiveactorRef.tell(message, senderActorRef);
|
||||
|
||||
senderActorRef.stop();
|
||||
|
||||
Map<ActorRef, String> actorToRole = (Map<ActorRef, String>) MemberModifier.field(WorkersRefCenter.class, "actorToRole").get(WorkersRefCenter.INSTANCE);
|
||||
Assert.assertEquals(null, actorToRole.get(senderActorRef));
|
||||
Map<ActorRef, WorkerRef> actorRefToWorkerRef = (Map<ActorRef, WorkerRef>) MemberModifier.field(WorkersRefCenter.class, "actorRefToWorkerRef").get(WorkersRefCenter.INSTANCE);
|
||||
Assert.assertEquals(null, actorRefToWorkerRef.get(senderActorRef));
|
||||
|
||||
Map<String, List<ActorRef>> roleToActor = (Map<String, List<ActorRef>>) MemberModifier.field(WorkersRefCenter.class, "roleToActor").get(WorkersRefCenter.INSTANCE);
|
||||
Map<String, List<WorkerRef>> roleToWorkerRef = (Map<String, List<WorkerRef>>) MemberModifier.field(WorkersRefCenter.class, "roleToWorkerRef").get(WorkersRefCenter.INSTANCE);
|
||||
ActorRef[] actorRefs = {};
|
||||
Assert.assertArrayEquals(actorRefs, roleToActor.get("WorkersListener").toArray());
|
||||
Assert.assertArrayEquals(actorRefs, roleToWorkerRef.get("WorkersListener").toArray());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
|
|||
Loading…
Reference in New Issue