1. cluster and actor module amalgamated into single one
2. created workerref class use for cover akka actorref class
This commit is contained in:
parent
fda13399d8
commit
3b8109b393
|
|
@ -5,7 +5,6 @@
|
|||
<modules>
|
||||
<module>skywalking-collector-cluster</module>
|
||||
<module>skywalking-collector-worker</module>
|
||||
<module>skywalking-collector-actor</module>
|
||||
</modules>
|
||||
<parent>
|
||||
<artifactId>skywalking</artifactId>
|
||||
|
|
|
|||
|
|
@ -1,21 +0,0 @@
|
|||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<project xmlns="http://maven.apache.org/POM/4.0.0"
|
||||
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
|
||||
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
|
||||
<parent>
|
||||
<artifactId>skywalking-collector</artifactId>
|
||||
<groupId>com.a.eye</groupId>
|
||||
<version>3.0-2017</version>
|
||||
</parent>
|
||||
<modelVersion>4.0.0</modelVersion>
|
||||
|
||||
<artifactId>skywalking-collector-actor</artifactId>
|
||||
|
||||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>com.a.eye</groupId>
|
||||
<artifactId>skywalking-collector-cluster</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
</project>
|
||||
|
|
@ -1,24 +0,0 @@
|
|||
package com.a.eye.skywalking.collector.actor.selector;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* The <code>WorkerSelector</code> should be implemented
|
||||
* by any class whose instances are intended to provide select a {@link ActorRef} from a {@link ActorRef} list.
|
||||
* <p></p>
|
||||
* Actually, the <code>ActorRef</code> is designed to provide a routing ability in the collector cluster.
|
||||
*
|
||||
* @author wusheng
|
||||
*/
|
||||
public interface WorkerSelector<T> {
|
||||
/**
|
||||
* 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<ActorRef> members, T message);
|
||||
}
|
||||
|
|
@ -9,6 +9,7 @@ import com.a.eye.skywalking.collector.actor.selector.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;
|
||||
|
||||
/**
|
||||
|
|
@ -42,7 +43,7 @@ public abstract class AbstractWorker<T> extends UntypedActor {
|
|||
}
|
||||
|
||||
public void tell(AbstractWorkerProvider targetWorkerProvider, WorkerSelector selector, T message) throws Throwable {
|
||||
List<ActorRef> avaibleWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName());
|
||||
List<WorkerRef> avaibleWorks = WorkersRefCenter.INSTANCE.availableWorks(targetWorkerProvider.roleName());
|
||||
selector.select(avaibleWorks, message).tell(message, getSelf());
|
||||
}
|
||||
|
||||
|
|
@ -2,7 +2,6 @@ package com.a.eye.skywalking.collector.actor;
|
|||
|
||||
import akka.actor.ActorSystem;
|
||||
import akka.actor.Props;
|
||||
import com.a.eye.skywalking.api.util.StringUtil;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
|
|
@ -1,7 +1,5 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorSystem;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
|
|
@ -0,0 +1,32 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class WorkerRef {
|
||||
final ActorRef actorRef;
|
||||
|
||||
public WorkerRef(ActorRef actorRef) {
|
||||
this.actorRef = actorRef;
|
||||
}
|
||||
|
||||
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();
|
||||
}
|
||||
}
|
||||
|
|
@ -1,12 +1,13 @@
|
|||
package com.a.eye.skywalking.collector.actor.selector;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import com.a.eye.skywalking.collector.actor.WorkerRef;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* The <code>RollingSelector</code> is a simple implementation of {@link WorkerSelector}.
|
||||
* It choose {@link ActorRef} nearly random, by round-robin.
|
||||
* It choose {@link WorkerRef} nearly random, by round-robin.
|
||||
*
|
||||
* @author wusheng
|
||||
*/
|
||||
|
|
@ -19,14 +20,14 @@ public enum RollingSelector implements WorkerSelector<Object> {
|
|||
private int index = 0;
|
||||
|
||||
/**
|
||||
* Use round-robin to select {@link ActorRef}.
|
||||
* Use round-robin to select {@link WorkerRef}.
|
||||
*
|
||||
* @param members given {@link ActorRef} list, which size is greater than 0;
|
||||
* @param members given {@link WorkerRef} list, which size is greater than 0;
|
||||
* @param message the {@link AbstractWorker} is going to send.
|
||||
* @return the selected {@link ActorRef}
|
||||
* @return the selected {@link WorkerRef}
|
||||
*/
|
||||
@Override
|
||||
public ActorRef select(List<ActorRef> members, Object message) {
|
||||
public WorkerRef select(List<WorkerRef> members, Object message) {
|
||||
int size = members.size();
|
||||
index++;
|
||||
int selectIndex = Math.abs(index) % size;
|
||||
|
|
@ -0,0 +1,25 @@
|
|||
package com.a.eye.skywalking.collector.actor.selector;
|
||||
|
||||
import com.a.eye.skywalking.collector.actor.AbstractWorker;
|
||||
import com.a.eye.skywalking.collector.actor.WorkerRef;
|
||||
|
||||
import java.util.List;
|
||||
|
||||
/**
|
||||
* The <code>WorkerSelector</code> should be implemented
|
||||
* by any class whose instances are intended to provide select a {@link WorkerRef} from a {@link WorkerRef} list.
|
||||
* <p></p>
|
||||
* Actually, the <code>WorkerRef</code> is designed to provide a routing ability in the collector cluster.
|
||||
*
|
||||
* @author wusheng
|
||||
*/
|
||||
public interface WorkerSelector<T> {
|
||||
/**
|
||||
* select a {@link WorkerRef} from a {@link WorkerRef} list.
|
||||
*
|
||||
* @param members given {@link WorkerRef} list, which size is greater than 0;
|
||||
* @param message the {@link AbstractWorker} is going to send.
|
||||
* @return the selected {@link WorkerRef}
|
||||
*/
|
||||
WorkerRef select(List<WorkerRef> members, T message);
|
||||
}
|
||||
|
|
@ -1,6 +1,7 @@
|
|||
package com.a.eye.skywalking.collector.cluster;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
import com.a.eye.skywalking.collector.actor.WorkerRef;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.Collections;
|
||||
|
|
@ -18,38 +19,45 @@ import java.util.concurrent.ConcurrentHashMap;
|
|||
public enum WorkersRefCenter {
|
||||
INSTANCE;
|
||||
|
||||
private Map<String, List<ActorRef>> roleToActor = new ConcurrentHashMap();
|
||||
private Map<String, List<WorkerRef>> roleToActor = new ConcurrentHashMap();
|
||||
|
||||
private Map<ActorRef, String> actorToRole = new ConcurrentHashMap();
|
||||
private Map<WorkerRef, String> actorToRole = new ConcurrentHashMap();
|
||||
|
||||
public void register(ActorRef newRef, String workerRole) {
|
||||
// private Map<String, WorkerRef> pathToWorkerRef = new ConcurrentHashMap();
|
||||
|
||||
public void register(ActorRef newActorRef, String workerRole) {
|
||||
if (!roleToActor.containsKey(workerRole)) {
|
||||
List<ActorRef> actorList = Collections.synchronizedList(new ArrayList<ActorRef>());
|
||||
List<WorkerRef> actorList = Collections.synchronizedList(new ArrayList<WorkerRef>());
|
||||
roleToActor.putIfAbsent(workerRole, actorList);
|
||||
}
|
||||
roleToActor.get(workerRole).add(newRef);
|
||||
actorToRole.put(newRef, workerRole);
|
||||
|
||||
WorkerRef newWorkerRef = new WorkerRef(newActorRef);
|
||||
roleToActor.get(workerRole).add(newWorkerRef);
|
||||
actorToRole.put(newWorkerRef, workerRole);
|
||||
// pathToWorkerRef.put(newWorkerRef.path(), newWorkerRef);
|
||||
}
|
||||
|
||||
public void unregister(ActorRef newRef) {
|
||||
String workerRole = actorToRole.get(newRef);
|
||||
roleToActor.get(workerRole).remove(newRef);
|
||||
actorToRole.remove(newRef);
|
||||
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());
|
||||
}
|
||||
|
||||
/**
|
||||
* Get a copy all available {@link ActorRef} list, by the given worker role.
|
||||
* Get all available {@link WorkerRef} list, by the given worker role.
|
||||
*
|
||||
* @param workerRole the given role
|
||||
* @return available {@link ActorRef} list
|
||||
* @return available {@link WorkerRef} list
|
||||
* @throws NoAvailableWorkerException , when no available worker.
|
||||
*/
|
||||
public List<ActorRef> availableWorks(String workerRole) throws NoAvailableWorkerException {
|
||||
List<ActorRef> refs = roleToActor.get(workerRole);
|
||||
if(refs == null || refs.size() == 0){
|
||||
public List<WorkerRef> availableWorks(String workerRole) throws NoAvailableWorkerException {
|
||||
List<WorkerRef> refs = roleToActor.get(workerRole);
|
||||
if (refs == null || refs.size() == 0) {
|
||||
throw new NoAvailableWorkerException("role=" + workerRole + ", no available worker.");
|
||||
}
|
||||
List<ActorRef> availableList = new ArrayList<>(refs.size());
|
||||
availableList.addAll(refs);
|
||||
return Collections.unmodifiableList(availableList);
|
||||
return Collections.unmodifiableList(refs);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,9 +1,6 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.actor.ActorRef;
|
||||
import akka.actor.ActorSelection;
|
||||
import akka.actor.ActorSystem;
|
||||
import akka.testkit.JavaTestKit;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
|
|
@ -12,8 +9,6 @@ import org.junit.Test;
|
|||
* @author pengys5
|
||||
*/
|
||||
public class AbstractWorkerProviderTestCase {
|
||||
|
||||
|
||||
ActorSystem system;
|
||||
|
||||
@Before
|
||||
|
|
@ -28,36 +23,9 @@ public class AbstractWorkerProviderTestCase {
|
|||
system = null;
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testCreateWorkerWhenWorkNameIsNull() {
|
||||
AbstractWorkerProvider aWorkerProvider = new AbstractWorkerProvider() {
|
||||
@Override
|
||||
public String workerRole() {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
return Object.class;
|
||||
}
|
||||
|
||||
@Override
|
||||
public int workerNum() {
|
||||
return 1;
|
||||
}
|
||||
};
|
||||
|
||||
aWorkerProvider.createWorker(system);
|
||||
}
|
||||
|
||||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testCreateWorkerWhenWorkerClassIsNull() {
|
||||
AbstractWorkerProvider aWorkerProvider = new AbstractWorkerProvider() {
|
||||
@Override
|
||||
public String workerRole() {
|
||||
return "Test";
|
||||
}
|
||||
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
return Object.class;
|
||||
|
|
@ -75,10 +43,6 @@ public class AbstractWorkerProviderTestCase {
|
|||
@Test(expected = IllegalArgumentException.class)
|
||||
public void testCreateWorkerWhenWorkerNumLessThan_1() {
|
||||
AbstractWorkerProvider aWorkerProvider = new AbstractWorkerProvider() {
|
||||
@Override
|
||||
public String workerRole() {
|
||||
return "Test";
|
||||
}
|
||||
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
|
|
@ -1,7 +1,5 @@
|
|||
package com.a.eye.skywalking.collector.actor;
|
||||
|
||||
import akka.japi.Creator;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
|
|
@ -7,11 +7,6 @@ public class SpiTestWorkerFactory extends AbstractWorkerProvider {
|
|||
|
||||
public static final String WorkerRole = "SpiTestWorker";
|
||||
|
||||
@Override
|
||||
public String workerRole() {
|
||||
return WorkerRole;
|
||||
}
|
||||
|
||||
@Override
|
||||
public Class workerClass() {
|
||||
return SpiTestWorker.class;
|
||||
|
|
@ -5,13 +5,11 @@ import akka.testkit.JavaTestKit;
|
|||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Test;
|
||||
import org.mockito.Mockito;
|
||||
|
||||
/**
|
||||
* @author pengys5
|
||||
*/
|
||||
public class SpiTestWorkerFactoryTestCase {
|
||||
|
||||
ActorSystem system;
|
||||
|
||||
@Before
|
||||
|
|
@ -10,7 +10,6 @@ import org.junit.Test;
|
|||
* @author pengys5
|
||||
*/
|
||||
public class WorkersCreatorTestCase {
|
||||
|
||||
ActorSystem system;
|
||||
|
||||
@Before
|
||||
|
|
@ -4,6 +4,7 @@ import akka.actor.ActorRef;
|
|||
import akka.actor.ActorSystem;
|
||||
import akka.actor.Props;
|
||||
import akka.testkit.TestActorRef;
|
||||
import com.a.eye.skywalking.collector.actor.WorkerRef;
|
||||
import org.junit.After;
|
||||
import org.junit.Assert;
|
||||
import org.junit.Before;
|
||||
|
|
@ -47,14 +48,22 @@ public class WorkersRefCenterTestCase {
|
|||
WorkersRefCenter.INSTANCE.register(actorRef2, "WorkersListener");
|
||||
WorkersRefCenter.INSTANCE.register(actorRef3, "WorkersListener");
|
||||
|
||||
Map<ActorRef, String> actorToRole = (Map<ActorRef, String>) MemberModifier.field(WorkersRefCenter.class, "actorToRole").get(WorkersRefCenter.INSTANCE);
|
||||
Assert.assertEquals("WorkersListener", actorToRole.get(actorRef1));
|
||||
Assert.assertEquals("WorkersListener", actorToRole.get(actorRef2));
|
||||
Assert.assertEquals("WorkersListener", actorToRole.get(actorRef3));
|
||||
Map<WorkerRef, String> actorToRole = (Map<WorkerRef, String>) MemberModifier.field(WorkersRefCenter.class, "actorToRole").get(WorkersRefCenter.INSTANCE);
|
||||
|
||||
Map<String, List<ActorRef>> roleToActor = (Map<String, List<ActorRef>>) MemberModifier.field(WorkersRefCenter.class, "roleToActor").get(WorkersRefCenter.INSTANCE);
|
||||
ActorRef[] actorRefs = {actorRef1, actorRef2, actorRef3};
|
||||
Assert.assertArrayEquals(actorRefs, roleToActor.get("WorkersListener").toArray());
|
||||
for (Map.Entry<WorkerRef, String> entry : actorToRole.entrySet()) {
|
||||
WorkerRef workerRef = entry.getKey();
|
||||
if (workerRef.equals(actorRef1) || workerRef.equals(actorRef2) || workerRef.equals(actorRef3)) {
|
||||
Assert.assertEquals("WorkersListener", entry.getValue());
|
||||
} else {
|
||||
Assert.fail();
|
||||
}
|
||||
}
|
||||
|
||||
Map<String, List<WorkerRef>> roleToActor = (Map<String, List<WorkerRef>>) MemberModifier.field(WorkersRefCenter.class, "roleToActor").get(WorkersRefCenter.INSTANCE);
|
||||
List<WorkerRef> workerRefs = roleToActor.get("WorkersListener");
|
||||
Assert.assertEquals(actorRef1.path().toString(), workerRefs.get(0).path().toString());
|
||||
Assert.assertEquals(actorRef2.path().toString(), workerRefs.get(1).path().toString());
|
||||
Assert.assertEquals(actorRef3.path().toString(), workerRefs.get(2).path().toString());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -79,7 +88,7 @@ public class WorkersRefCenterTestCase {
|
|||
}
|
||||
|
||||
@Test
|
||||
public void testSizeOf(){
|
||||
public void testSizeOf() throws NoAvailableWorkerException {
|
||||
final Props props = Props.create(WorkersListener.class);
|
||||
final TestActorRef<WorkersListener> actorRef1 = TestActorRef.create(system, props, "WorkersListener1");
|
||||
final TestActorRef<WorkersListener> actorRef2 = TestActorRef.create(system, props, "WorkersListener2");
|
||||
|
|
@ -89,6 +98,6 @@ public class WorkersRefCenterTestCase {
|
|||
WorkersRefCenter.INSTANCE.register(actorRef2, "WorkersListener");
|
||||
WorkersRefCenter.INSTANCE.register(actorRef3, "WorkersListener");
|
||||
|
||||
Assert.assertEquals(3, WorkersRefCenter.INSTANCE.sizeOf("WorkersListener"));
|
||||
Assert.assertEquals(3, WorkersRefCenter.INSTANCE.availableWorks("WorkersListener").size());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -14,7 +14,7 @@
|
|||
<dependencies>
|
||||
<dependency>
|
||||
<groupId>com.a.eye</groupId>
|
||||
<artifactId>skywalking-collector-actor</artifactId>
|
||||
<artifactId>skywalking-collector-cluster</artifactId>
|
||||
<version>${project.version}</version>
|
||||
</dependency>
|
||||
</dependencies>
|
||||
|
|
|
|||
Loading…
Reference in New Issue