diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java index c7e2b9054..fb1bdaa5e 100644 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java @@ -1,17 +1,20 @@ package com.a.eye.skywalking.collector.actor; +import akka.actor.ActorRef; import akka.actor.UntypedActor; import akka.cluster.ClusterEvent; import akka.cluster.Member; import akka.cluster.MemberStatus; -import com.a.eye.skywalking.collector.actor.router.WorkerRouter; +import com.a.eye.skywalking.collector.actor.router.WorkerSelector; import com.a.eye.skywalking.collector.cluster.WorkerListenerMessage; import com.a.eye.skywalking.collector.cluster.WorkersListener; +import com.a.eye.skywalking.collector.cluster.WorkersRefCenter; +import java.util.List; /** * @author pengys5 */ -public abstract class AbstractWorker extends UntypedActor { +public abstract class AbstractWorker extends UntypedActor { final String workerRole; @@ -38,8 +41,14 @@ public abstract class AbstractWorker extends UntypedActor { } } +<<<<<<< HEAD protected void tell(String workerRole, WorkerRouter router, Object message) throws Throwable { router.find(workerRole).tell(message, getSelf()); +======= + public void tell(AbstractWorkerProvider targetWorkerProvider, WorkerSelector selector, T message) throws Throwable { + List avaibleWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName()); + selector.select(avaibleWorks, message).tell(message, getSelf()); +>>>>>>> e93c6a65c448419bb87c0777b3c42dfc425d533b } void register(Member member) { diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java index c567050aa..026e67dc4 100644 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorkerProvider.java @@ -9,16 +9,11 @@ import com.a.eye.skywalking.api.util.StringUtil; */ public abstract class AbstractWorkerProvider { - public abstract String workerRole(); - public abstract Class workerClass(); public abstract int workerNum(); public void createWorker(ActorSystem system) { - if (StringUtil.isEmpty(workerRole())) { - throw new IllegalArgumentException("cannot createWorker() with nothing obtained from workerRole()"); - } if (workerClass() == null) { throw new IllegalArgumentException("cannot createWorker() with nothing obtained from workerClass()"); } @@ -27,7 +22,11 @@ public abstract class AbstractWorkerProvider { } for (int i = 1; i <= workerNum(); i++) { - system.actorOf(Props.create(workerClass(), workerRole()), workerRole() + "_" + i); + system.actorOf(Props.create(workerClass(), roleName()), roleName() + "_" + i); } } + + protected String roleName(){ + return workerClass().getSimpleName(); + } } diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/RandomRouter.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/RandomRouter.java deleted file mode 100644 index 38f1877c3..000000000 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/RandomRouter.java +++ /dev/null @@ -1,20 +0,0 @@ -package com.a.eye.skywalking.collector.actor.router; - -import akka.actor.ActorRef; -import com.a.eye.skywalking.collector.cluster.WorkersRefCenter; - -import java.util.List; -import java.util.Random; - -/** - * @author pengys5 - */ -public class RandomRouter implements WorkerRouter { - - @Override - public ActorRef find(String workerRole) { - int workerNum = WorkersRefCenter.INSTANCE.sizeOf(workerRole); - Random random = new Random(workerNum); - return WorkersRefCenter.INSTANCE.find(workerRole, random.nextInt()); - } -} diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/RollingSelector.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/RollingSelector.java new file mode 100644 index 000000000..0352d4a26 --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/RollingSelector.java @@ -0,0 +1,35 @@ +package com.a.eye.skywalking.collector.actor.router; + +import akka.actor.ActorRef; +import com.a.eye.skywalking.collector.actor.AbstractWorker; +import java.util.List; + +/** + * The RollingSelector is a simple implementation of {@link WorkerSelector}. + * It choose {@link ActorRef} nearly random, by round-robin. + * + * @author wusheng + */ +public enum RollingSelector implements WorkerSelector { + INSTANCE; + + /** + * A simple round variable. + */ + private int index = 0; + + /** + * Use round-robin to select {@link ActorRef}. + * + * @param members given {@link ActorRef} list, which size is greater than 0; + * @param message the {@link AbstractWorker} is going to send. + * @return the selected {@link ActorRef} + */ + @Override + public ActorRef select(List members, Object message) { + int size = members.size(); + index++; + int selectIndex = Math.abs(index) % size; + return members.get(selectIndex); + } +} diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/WorkerRouter.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/WorkerRouter.java deleted file mode 100644 index 952131886..000000000 --- a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/WorkerRouter.java +++ /dev/null @@ -1,12 +0,0 @@ -package com.a.eye.skywalking.collector.actor.router; - -import akka.actor.ActorRef; - -import java.util.List; - -/** - * @author wusheng - */ -public interface WorkerRouter { - ActorRef find(String workerRole); -} diff --git a/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/WorkerSelector.java b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/WorkerSelector.java new file mode 100644 index 000000000..bde52a9c0 --- /dev/null +++ b/skywalking-collector/skywalking-collector-actor/src/main/java/com/a/eye/skywalking/collector/actor/router/WorkerSelector.java @@ -0,0 +1,24 @@ +package com.a.eye.skywalking.collector.actor.router; + +import akka.actor.ActorRef; +import com.a.eye.skywalking.collector.actor.AbstractWorker; +import java.util.List; + +/** + * The WorkerSelector should be implemented + * by any class whose instances are intended to provide select a {@link ActorRef} from a {@link ActorRef} list. + *

+ * Actually, the ActorRef is designed to provide a routing ability in the collector cluster. + * + * @author wusheng + */ +public interface WorkerSelector { + /** + * select a {@link ActorRef} from a {@link ActorRef} list. + * + * @param members given {@link ActorRef} list, which size is greater than 0; + * @param message the {@link AbstractWorker} is going to send. + * @return the selected {@link ActorRef} + */ + ActorRef select(List members, T message); +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/NoAvailableWorkerException.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/NoAvailableWorkerException.java new file mode 100644 index 000000000..08bf4f051 --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/NoAvailableWorkerException.java @@ -0,0 +1,15 @@ +package com.a.eye.skywalking.collector.cluster; + +/** + * The NoAvailableWorkerException represents no available member, + * when the {@link WorkersRefCenter#availableWorks(String)} try to get the list. + * + * Most likely, in the cluster, these is no active worker of the particular role. + * + * @author wusheng + */ +public class NoAvailableWorkerException extends Exception { + public NoAvailableWorkerException(String message){ + super(message); + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java index 65b26c51f..41d52b577 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersRefCenter.java @@ -37,11 +37,19 @@ public enum WorkersRefCenter { actorToRole.remove(newRef); } - public ActorRef find(String workerRole, int sequence) { - return roleToActor.get(workerRole).get(sequence); + /** + * Get a copy all available {@link ActorRef} list, by the given worker role. + * @param workerRole the given role + * @return available {@link ActorRef} list + * @throws NoAvailableWorkerException , when no available worker. + */ + public List availableWorks(String workerRole) throws NoAvailableWorkerException { + List refs = roleToActor.get(workerRole); + if(refs == null || refs.size() == 0){ + throw new NoAvailableWorkerException("role=" + workerRole + ", no available worker."); + } + List availableList = new ArrayList<>(refs.size()); + availableList.addAll(refs); + return Collections.unmodifiableList(availableList); } - - public int sizeOf(String workerRole) { - return roleToActor.get(workerRole).size(); - } -} \ No newline at end of file +}