From c2a1ef4d6d6e46d51117e8b76f5a8cfd445a5d75 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Fri, 10 Mar 2017 22:35:52 +0800 Subject: [PATCH] add ClusterEvent.UnreachableMember in WorkersListener to unregister from WorkersRefCenter --- .../collector/actor/AbstractWorker.java | 15 ---------- .../skywalking/collector/actor/WorkerRef.java | 5 ++++ .../collector/cluster/WorkersListener.java | 14 ++++++++- .../collector/cluster/WorkersRefCenter.java | 29 +++++++++++++++---- .../src/main/resources/application.conf | 18 ++++++++++-- .../src/main/resources/collector.config | 2 +- 6 files changed, 59 insertions(+), 24 deletions(-) diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java index a483b0f74..758ae2ff8 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractWorker.java @@ -47,7 +47,6 @@ public abstract class AbstractWorker extends UntypedActor { @Override public void preStart() throws Exception { cluster.subscribe(getSelf(), ClusterEvent.MemberUp.class); - register(); } @Override @@ -71,7 +70,6 @@ public abstract class AbstractWorker extends UntypedActor { @Override public void onReceive(Object message) throws Throwable { if (message instanceof ClusterEvent.CurrentClusterState) { - logger.info("receive ClusterEvent.CurrentClusterState message"); ClusterEvent.CurrentClusterState state = (ClusterEvent.CurrentClusterState) message; for (Member member : state.getMembers()) { if (member.status().equals(MemberStatus.up())) { @@ -82,14 +80,6 @@ public abstract class AbstractWorker extends UntypedActor { ClusterEvent.MemberUp memberUp = (ClusterEvent.MemberUp) message; logger.info("receive ClusterEvent.MemberUp message, address: %s", memberUp.member().address().toString()); register(memberUp.member()); - } else if (message instanceof ClusterEvent.MemberEvent) { - System.out.println("other event: " + message.getClass().getSimpleName()); - } else if (message instanceof ClusterEvent.UnreachableMember) { - System.out.println("other event: " + message.getClass().getSimpleName()); - } else if (message instanceof ClusterEvent.MemberJoined) { - System.out.println("other event: " + message.getClass().getSimpleName()); - } else if (message instanceof ClusterEvent.ReachableMember) { - System.out.println("other event: " + message.getClass().getSimpleName()); } else { logger.debug("worker class: %s, message class: %s", this.getClass().getName(), message.getClass().getName()); receive(message); @@ -126,9 +116,4 @@ public abstract class AbstractWorker extends UntypedActor { getContext().actorSelection(member.address() + "/user/" + WorkersListener.WorkName).tell(registerMessage, getSelf()); } } - - void register() { - WorkerListenerMessage.RegisterMessage registerMessage = new WorkerListenerMessage.RegisterMessage(getClass().getSimpleName()); - getContext().actorSelection("/user/" + WorkersListener.WorkName).tell(registerMessage, getSelf()); - } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java index aaf424669..64a8a939b 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/WorkerRef.java @@ -1,5 +1,6 @@ package com.a.eye.skywalking.collector.actor; +import akka.actor.ActorPath; import akka.actor.ActorRef; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -26,6 +27,10 @@ public class WorkerRef { actorRef.tell(message, sender); } + public ActorPath path() { + return actorRef.path(); + } + public String getWorkerRole() { return workerRole; } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java index c1ecbcf8e..9e1c6fa88 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/cluster/WorkersListener.java @@ -3,6 +3,8 @@ package com.a.eye.skywalking.collector.cluster; import akka.actor.ActorRef; import akka.actor.Terminated; import akka.actor.UntypedActor; +import akka.cluster.Cluster; +import akka.cluster.ClusterEvent; import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; @@ -24,12 +26,19 @@ public class WorkersListener extends UntypedActor { public static final String WorkName = "WorkersListener"; + private Cluster cluster = Cluster.get(getContext().system()); + + @Override + public void preStart() throws Exception { + cluster.subscribe(getSelf(), ClusterEvent.UnreachableMember.class); + } + @Override public void onReceive(Object message) throws Throwable { if (message instanceof WorkerListenerMessage.RegisterMessage) { WorkerListenerMessage.RegisterMessage register = (WorkerListenerMessage.RegisterMessage) message; ActorRef sender = getSender(); - getContext().watch(sender); +// getContext().watch(sender); logger.info("register worker of role: %s, path: %s", register.getWorkRole(), sender.toString()); @@ -37,6 +46,9 @@ public class WorkersListener extends UntypedActor { } else if (message instanceof Terminated) { Terminated terminated = (Terminated) message; WorkersRefCenter.INSTANCE.unregister(terminated.getActor()); + } else if (message instanceof ClusterEvent.UnreachableMember) { + ClusterEvent.UnreachableMember unreachableMember = (ClusterEvent.UnreachableMember) message; + WorkersRefCenter.INSTANCE.unregister(unreachableMember.member().address()); } else { unhandled(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 99a209855..e601652e5 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 @@ -1,12 +1,10 @@ package com.a.eye.skywalking.collector.cluster; import akka.actor.ActorRef; +import akka.actor.Address; import com.a.eye.skywalking.collector.actor.WorkerRef; -import java.util.ArrayList; -import java.util.Collections; -import java.util.List; -import java.util.Map; +import java.util.*; import java.util.concurrent.ConcurrentHashMap; /** @@ -24,13 +22,13 @@ public enum WorkersRefCenter { private Map actorRefToWorkerRef = new ConcurrentHashMap<>(); public void register(ActorRef newActorRef, String workerRole) { - System.out.println("register: " + workerRole); if (!roleToWorkerRef.containsKey(workerRole)) { List actorList = Collections.synchronizedList(new ArrayList()); roleToWorkerRef.putIfAbsent(workerRole, actorList); } WorkerRef newWorkerRef = new WorkerRef(newActorRef, workerRole); + roleToWorkerRef.get(workerRole).add(newWorkerRef); actorRefToWorkerRef.put(newActorRef, newWorkerRef); } @@ -41,6 +39,27 @@ public enum WorkersRefCenter { actorRefToWorkerRef.remove(oldActorRef); } + public void unregister(Address address) { + Iterator actorRefToWorkerRefIterator = actorRefToWorkerRef.keySet().iterator(); + while (actorRefToWorkerRefIterator.hasNext()) { + if (address.equals(actorRefToWorkerRefIterator.next().path().address())) { + actorRefToWorkerRefIterator.remove(); + } + } + + Iterator>> roleToWorkerRefIterator = roleToWorkerRef.entrySet().iterator(); + while (roleToWorkerRefIterator.hasNext()) { + List workerRefList = roleToWorkerRefIterator.next().getValue(); + + Iterator workerRefIterator = workerRefList.iterator(); + while (workerRefIterator.hasNext()) { + if (workerRefIterator.next().path().address().equals(address)) { + workerRefIterator.remove(); + } + } + } + } + /** * Get all available {@link WorkerRef} list, by the given worker role. * diff --git a/skywalking-collector/skywalking-collector-worker/src/main/resources/application.conf b/skywalking-collector/skywalking-collector-worker/src/main/resources/application.conf index d40cceeaf..ab9bb7e3f 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/resources/application.conf +++ b/skywalking-collector/skywalking-collector-worker/src/main/resources/application.conf @@ -14,10 +14,24 @@ akka { "com.google.protobuf.Message" = proto "com.a.eye.skywalking.messages.ISerializable" = data "com.google.gson.JsonObject" = json -// "java.io.Serializable" = none + // "java.io.Serializable" = none } -// serialize-messages = on + // serialize-messages = on warn-about-java-serializer-usage = on } + + remote { + log-remote-lifecycle-events = off + + netty.tcp { + hostname = "127.0.0.1" + port = 1000 + } + } + + cluster { + auto-down-unreachable-after = off + metrics.enabled = off + } } \ No newline at end of file diff --git a/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config b/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config index dcc47d2fd..901420e88 100644 --- a/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config +++ b/skywalking-collector/skywalking-collector-worker/src/main/resources/collector.config @@ -1,4 +1,4 @@ cluster.current.hostname=127.0.0.1 cluster.current.port=1000 cluster.current.roles=[WorkersListener, TraceSegmentReceiver, NodeInstancePersistence] -cluster.nodes=["akka.tcp://CollectorSystem@127.0.0.1:1000", "akka.tcp://CollectorSystem@127.0.0.1:1002"] \ No newline at end of file +cluster.nodes=["akka.tcp://CollectorSystem@127.0.0.1:1000", "akka.tcp://CollectorSystem@127.0.0.1:1001", "akka.tcp://CollectorSystem@127.0.0.1:1002"] \ No newline at end of file