AbstractWorker implementation process the message in
- * {@link #receive(Object)} method.
- * - *
- * Subclasses must implement the abstract {@link #receive(Object)} method to process message. - * Subclasses forbid to override the {@link #onReceive(Object)} method. - *
- * Here is an example on how to create and use an {@link AbstractWorker}: - *
- * {{{
- * public class SampleWorker extends AbstractWorker {
- *
* @author pengys5
- * @Override public void receive(Object message) throws Throwable {
- * if (message.equals("Tell Next")) {
- * Object sendMessage = new Object();
- * tell(new NextSampleWorkerFactory(), RollingSelector.INSTANCE, sendMessage);
- * }
- * }
- * }
- * }}}
*/
-public abstract class AbstractWorker extends UntypedActor {
+public abstract class AbstractWorker {
- private Logger logger = LogManager.getFormatterLogger(AbstractWorker.class);
+ private final LocalWorkerContext selfContext = new LocalWorkerContext();
- private Cluster cluster = Cluster.get(getContext().system());
+ private final Role role;
- @Override
- public void preStart() throws Exception {
- cluster.subscribe(getSelf(), ClusterEvent.MemberUp.class);
+ private final ClusterWorkerContext clusterContext;
+
+ public AbstractWorker(Role role, ClusterWorkerContext clusterContext) {
+ this.role = role;
+ this.clusterContext = clusterContext;
}
- @Override
- public void postStop() throws Exception {
- cluster.unsubscribe(getSelf());
+ public abstract void preStart() throws Exception;
+
+ public abstract void work(Object message) throws Exception;
+
+ final public LocalWorkerContext getSelfContext() {
+ return selfContext;
}
- /**
- * 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) {
- ClusterEvent.CurrentClusterState state = (ClusterEvent.CurrentClusterState) message;
- for (Member member : state.getMembers()) {
- if (member.status().equals(MemberStatus.up())) {
- register(member);
- }
- }
- } else if (message instanceof ClusterEvent.MemberUp) {
- ClusterEvent.MemberUp memberUp = (ClusterEvent.MemberUp) message;
- logger.info("receive ClusterEvent.MemberUp message, address: %s", memberUp.member().address().toString());
- register(memberUp.member());
- } else {
- logger.debug("worker class: %s, message class: %s", this.getClass().getName(), message.getClass().getName());
- receive(message);
- }
+ final public ClusterWorkerContext getClusterContext() {
+ return clusterContext;
}
- /**
- * 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, Object message) throws Throwable {
- List
- * Here is an example on how to create and use an {@link AbstractWorkerProvider}:
- *
- * {{{
- * public class SampleWorkerFactory extends AbstractWorkerProvider {
- *
* @author pengys5
- * @Override public Class workerClass() {
- * return SampleWorker.class;
- * }
- * @Override public int workerNum() {
- * return Config.SampleWorkerNum;
- * }
- * }
- * }}}
- *
*/
-public abstract class AbstractWorkerProvider {
+public abstract class AbstractWorkerProviderAbstractWorkerProvider should be implemented by any class whose
- * instances are intended to provide create instance of the {@link AbstractWorker}.
- * The {@link WorkersCreator} use java service loader to load provider implementer,
- * so you should config the service file.
- * WorkersCreator is a util that use Java Spi to create
- * workers by META-INF config file.
- *
- * @author pengys5
- */
-public enum WorkersCreator {
- INSTANCE;
-
- private Logger logger = LogManager.getFormatterLogger(WorkersCreator.class);
-
- /**
- * create worker to use Java Spi.
- *
- * @param system is create by akka {@link ActorSystem}
- */
- public void boot(ActorSystem system) {
- system.actorOf(Props.create(WorkersListener.class), WorkersListener.WorkName);
-
- ServiceLoaderHashCodeSelector is a simple implementation of {@link WorkerSelector}.
- * It choose {@link WorkerRef} by message hashcode, so it can use to send the same hashcode
- * message to same {@link WorkerRef}. Usually, use to database operate which avoid dirty data.
- *
* @author pengys5
*/
-public enum HashCodeSelector implements WorkerSelector {
- INSTANCE;
+public class HashCodeSelector implements WorkerSelectorRollingSelector is a simple implementation of {@link WorkerSelector}.
- * It choose {@link WorkerRef} nearly random, by round-robin.
- *
- * @author wusheng
+ * @author pengys5
*/
-public enum RollingSelector implements WorkerSelector {
- INSTANCE;
+public class RollingSelector implements WorkerSelectorWorkerSelector should be implemented
- * by any class whose instances are intended to provide select a {@link WorkerRef} from a {@link WorkerRef} list.
- *
- * Actually, the WorkerRef is designed to provide a routing ability in the collector cluster.
- *
- * @author wusheng
+ * @author pengys5
*/
-public interface WorkerSelector {
- /**
- * 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(ListWorkersListener listening the register message from workers
* implementation of the {@link com.a.eye.skywalking.collector.actor.AbstractWorker}
@@ -28,6 +34,14 @@ public class WorkersListener extends UntypedActor {
private Cluster cluster = Cluster.get(getContext().system());
+ private MapWorkersRefCenter 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 WorkersRefCenter {
- INSTANCE;
-
- private Map