diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/consumer/TraceConsumerActor.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/consumer/TraceConsumerActor.java index e02c326da..bc15fcaae 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/consumer/TraceConsumerActor.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/consumer/TraceConsumerActor.java @@ -2,7 +2,7 @@ package com.a.eye.skywalking.collector.cluster.consumer; import akka.cluster.ClusterEvent; import com.a.eye.skywalking.collector.cluster.Const; -import com.a.eye.skywalking.collector.cluster.message.ActorRegisteMessage; +import com.a.eye.skywalking.collector.cluster.message.ActorRegisterMessage; import com.a.eye.skywalking.collector.cluster.message.TraceMessages.TransformationJob; import com.a.eye.skywalking.collector.cluster.message.TraceMessages.TransformationResult; import akka.actor.UntypedActor; @@ -56,8 +56,8 @@ public class TraceConsumerActor extends UntypedActor { System.out.println("register"); if (member.hasRole(Const.Trace_Producer_Role)) { System.out.println("register: " + Const.Trace_Producer_Role); - ActorRegisteMessage.RegisteMessage registeMessage = new ActorRegisteMessage.RegisteMessage(Const.Trace_Consumer_Role, ""); - getContext().actorSelection(member.address() + Const.Actor_Manager_Path).tell(registeMessage, getSelf()); + ActorRegisterMessage.RegisterMessage registerMessage = new ActorRegisterMessage.RegisterMessage(Const.Trace_Consumer_Role, ""); + getContext().actorSelection(member.address() + Const.Actor_Manager_Path).tell(registerMessage, getSelf()); } } -} \ No newline at end of file +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorCache.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorCache.java deleted file mode 100644 index 47f29d3bc..000000000 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorCache.java +++ /dev/null @@ -1,17 +0,0 @@ -package com.a.eye.skywalking.collector.cluster.manager; - -import akka.actor.ActorRef; - -import java.util.List; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; - -/** - * Created by Administrator on 2017/2/21 0021. - */ -public class ActorCache { - - public static Map> roleToActor = new ConcurrentHashMap(); - - public static Map actorToRole = new ConcurrentHashMap(); -} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorManagerActor.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorManagerActor.java index 0d3b65a17..cc16d5a24 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorManagerActor.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorManagerActor.java @@ -1,11 +1,8 @@ package com.a.eye.skywalking.collector.cluster.manager; -import akka.actor.ActorRef; import akka.actor.Terminated; import akka.actor.UntypedActor; -import com.a.eye.skywalking.collector.cluster.message.ActorRegisteMessage; - -import java.util.*; +import com.a.eye.skywalking.collector.cluster.message.ActorRegisterMessage; /** * Created by Administrator on 2017/2/21 0021. @@ -14,23 +11,17 @@ public class ActorManagerActor extends UntypedActor { @Override public void onReceive(Object message) throws Throwable { - if (message instanceof ActorRegisteMessage.RegisteMessage) { - System.out.println("RegisteMessage"); - ActorRegisteMessage.RegisteMessage regist = (ActorRegisteMessage.RegisteMessage) message; + if (message instanceof ActorRegisterMessage.RegisterMessage) { + System.out.println("RegisterMessage"); + ActorRegisterMessage.RegisterMessage regist = (ActorRegisterMessage.RegisterMessage) message; getContext().watch(getSender()); - if (!ActorCache.roleToActor.containsKey(regist.getRole())) { - List actorList = Collections.synchronizedList(new ArrayList()); - ActorCache.roleToActor.putIfAbsent(regist.getRole(), actorList); - } - getContext().watch(getSender()); - ActorCache.roleToActor.get(regist.getRole()).add(getSender()); - ActorCache.actorToRole.put(getSender(), regist.getRole()); + + ActorRefCenter.INSTANCE.register(getSender(), regist.getRole()); } else if (message instanceof Terminated) { System.out.println("Terminated"); Terminated terminated = (Terminated) message; - String role = ActorCache.actorToRole.get(terminated.getActor()); - ActorCache.roleToActor.get(role).remove(terminated.getActor()); - ActorCache.actorToRole.remove(terminated.getActor()); + + ActorRefCenter.INSTANCE.unregister(terminated.getActor()); } else { unhandled(message); } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorRefCenter.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorRefCenter.java new file mode 100644 index 000000000..440224f9a --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/ActorRefCenter.java @@ -0,0 +1,47 @@ +package com.a.eye.skywalking.collector.cluster.manager; + +import akka.actor.ActorRef; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; + +/** + * ActorRefCenter represent a cache center, + * store all {@link ActorRef}s, each of them represent a Akka Actor instance. + * All the Actors in this JVM, can find alive-actor in here, and send message. + * + * @author wusheng + */ +public enum ActorRefCenter { + INSTANCE; + + private Map> roleToActor = new ConcurrentHashMap(); + + private Map actorToRole = new ConcurrentHashMap(); + + public void register(ActorRef newRef, String name){ + if (!roleToActor.containsKey(name)) { + List actorList = Collections.synchronizedList(new ArrayList()); + roleToActor.putIfAbsent(name, actorList); + } + roleToActor.get(name).add(newRef); + actorToRole.put(newRef, name); + } + + public void unregister(ActorRef newRef){ + String role = actorToRole.get(newRef); + roleToActor.get(role).remove(newRef); + actorToRole.remove(newRef); + } + + public ActorRef find(String name, RefRouter router){ + return router.find(roleToActor.get(name)); + } + + public int sizeOf(String name){ + return roleToActor.get(name).size(); + } +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/RefRouter.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/RefRouter.java new file mode 100644 index 000000000..f25af6e30 --- /dev/null +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/manager/RefRouter.java @@ -0,0 +1,11 @@ +package com.a.eye.skywalking.collector.cluster.manager; + +import akka.actor.ActorRef; +import java.util.List; + +/** + * Created by wusheng on 2017/2/21. + */ +public interface RefRouter { + ActorRef find(List candidates); +} diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/message/ActorRegisteMessage.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/message/ActorRegisterMessage.java similarity index 83% rename from skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/message/ActorRegisteMessage.java rename to skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/message/ActorRegisterMessage.java index 431e7feb8..1b2cc11c9 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/message/ActorRegisteMessage.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/message/ActorRegisterMessage.java @@ -5,13 +5,13 @@ import java.io.Serializable; /** * Created by Administrator on 2017/2/21 0021. */ -public interface ActorRegisteMessage { +public interface ActorRegisterMessage { - public static class RegisteMessage implements Serializable { + public static class RegisterMessage implements Serializable { public final String role; public final String action; - public RegisteMessage(String role, String action) { + public RegisterMessage(String role, String action) { this.role = role; this.action = action; } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerActor.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerActor.java index e6913de02..73f6e63eb 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerActor.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/producer/TraceProducerActor.java @@ -1,18 +1,12 @@ package com.a.eye.skywalking.collector.cluster.producer; -import static com.a.eye.skywalking.collector.cluster.message.TraceMessages.BACKEND_REGISTRATION; - -import akka.actor.ActorRef; import com.a.eye.skywalking.collector.cluster.Const; -import com.a.eye.skywalking.collector.cluster.manager.ActorCache; +import com.a.eye.skywalking.collector.cluster.manager.ActorRefCenter; import com.a.eye.skywalking.collector.cluster.message.TraceMessages.JobFailed; import com.a.eye.skywalking.collector.cluster.message.TraceMessages.TransformationJob; -import akka.actor.Terminated; import akka.actor.UntypedActor; import org.springframework.context.annotation.Scope; -import java.util.List; - //#frontend //@Named("TraceProducerActor") @Scope("prototype") @@ -22,20 +16,21 @@ public class TraceProducerActor extends UntypedActor { @Override public void onReceive(Object message) { - List actorList = ActorCache.roleToActor.get(Const.Trace_Consumer_Role); - if (actorList == null) { + int actorSize = ActorRefCenter.INSTANCE.sizeOf(Const.Trace_Consumer_Role); + if (actorSize == 0) { System.out.println("actorList null"); } else { - System.out.println("size: " + actorList.size()); + System.out.println("sizeOf: " + actorSize); } - if ((message instanceof TransformationJob) && actorList == null) { - TransformationJob job = (TransformationJob) message; + if ((message instanceof TransformationJob) && actorSize == 0) { + TransformationJob job = (TransformationJob)message; getSender().tell(new JobFailed("Service unavailable, try again later", job), getSender()); } else if (message instanceof TransformationJob) { - TransformationJob job = (TransformationJob) message; + TransformationJob job = (TransformationJob)message; jobCounter++; - actorList.get(jobCounter % actorList.size()).forward(job, getContext()); + ActorRefCenter.INSTANCE.find(Const.Trace_Consumer_Role, + (candidates) -> candidates.get(jobCounter % candidates.size())); } else { unhandled(message); }