diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java index 207075aba..249a152d0 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorker.java @@ -11,18 +11,45 @@ import org.apache.logging.log4j.LogManager; import org.apache.logging.log4j.Logger; /** + * The AbstractClusterWorker should be implemented by any class whose instances + * are intended to provide receive remote message that message will transfer across different jvm + * which running in same server or another address server. + *

+ * Usually the implemented class used to receive persistence data or aggregation the metric. + * * @author pengys5 + * @since feature3.0 */ public abstract class AbstractClusterWorker extends AbstractWorker { + /** + * Constructs a AbstractClusterWorker with the worker role and context. + * + * @param role The responsibility of worker in cluster, more than one workers can have + * same responsibility which use to provide load balancing ability. + * @param clusterContext See {@link ClusterWorkerContext} + * @param selfContext See {@link LocalWorkerContext} + */ protected AbstractClusterWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { super(role, clusterContext, selfContext); } + /** + * Receive message + * + * @param message The persistence data or metric data. + * @throws Exception The Exception happen in {@link #onWork(Object)} + */ final public void allocateJob(Object message) throws Exception { onWork(message); } + /** + * The data process logic in this method. + * + * @param message Cast the message object to a expect subclass. + * @throws Exception Don't handle the exception, throw it. + */ protected abstract void onWork(Object message) throws Exception; static class WorkerWithAkka extends UntypedActor { diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java index 421d44112..1820e5e6d 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractClusterWorkerProvider.java @@ -4,17 +4,36 @@ import akka.actor.ActorRef; import akka.actor.Props; /** + * The AbstractClusterWorkerProvider should be implemented by any class whose instances + * are intended to provide create the class instance whose implemented {@link AbstractClusterWorker}. + *

+ * * @author pengys5 + * @since feature3.0 */ public abstract class AbstractClusterWorkerProvider extends AbstractWorkerProvider { + /** + * Create how many worker instance of {@link AbstractClusterWorker} in one jvm. + * + * @return The worker instance number. + */ public abstract int workerNum(); + /** + * Create the worker instance into akka system, the akka system will control the cluster worker life cycle. + * + * @param localContext Not used, will be null. + * @return The created worker reference. See {@link ClusterWorkerRef} + * @throws IllegalArgumentException Not used. + * @throws ProviderNotFoundException This worker instance attempted to find a provider which use to create another worker + * instance, when the worker provider not find then Throw this Exception. + */ @Override final public WorkerRef onCreate(LocalWorkerContext localContext) throws IllegalArgumentException, ProviderNotFoundException { int num = ClusterWorkerRefCounter.INSTANCE.incrementAndGet(role()); - T clusterWorker = (T) workerInstance(getClusterContext()); + T clusterWorker = workerInstance(getClusterContext()); clusterWorker.preStart(); ActorRef actorRef = getClusterContext().getAkkaSystem().actorOf(Props.create(AbstractClusterWorker.WorkerWithAkka.class, clusterWorker), role().roleName() + "_" + num); diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java index eb919e57e..38bcd139a 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/AbstractLocalAsyncWorker.java @@ -6,34 +6,72 @@ import com.lmax.disruptor.EventHandler; import com.lmax.disruptor.RingBuffer; /** + * The AbstractLocalAsyncWorker should be implemented by any class whose instances + * are intended to provide receive asynchronous message in same jvm. + * * @author pengys5 + * @since feature3.0 */ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker { + /** + * Constructs a AbstractLocalAsyncWorker with the worker role and context. + * + * @param role The responsibility of worker in cluster, more than one workers can have + * same responsibility which use to provide load balancing ability. + * @param clusterContext See {@link ClusterWorkerContext} + * @param selfContext See {@link LocalWorkerContext} + */ public AbstractLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) { super(role, clusterContext, selfContext); } + /** + * The asynchronous worker always use to persistence data into db, this is the end of the streaming, + * so usually no need to create the next worker instance at the time of this worker instance create. + * + * @throws ProviderNotFoundException When worker provider not found, it will be throw this exception. + */ @Override public void preStart() throws ProviderNotFoundException { } - final public void allocateJob(Object request) throws Exception { - onWork(request); + /** + * Receive message + * + * @param message The persistence data or metric data. + * @throws Exception The Exception happen in {@link #onWork(Object)} + */ + final public void allocateJob(Object message) throws Exception { + onWork(message); } - protected abstract void onWork(Object request) throws Exception; + /** + * The data process logic in this method. + * + * @param message Cast the message object to a expect subclass. + * @throws Exception Don't handle the exception, throw it. + */ + protected abstract void onWork(Object message) throws Exception; static class WorkerWithDisruptor implements EventHandler { private RingBuffer ringBuffer; private AbstractLocalAsyncWorker asyncWorker; - public WorkerWithDisruptor(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) { + WorkerWithDisruptor(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) { this.ringBuffer = ringBuffer; this.asyncWorker = asyncWorker; } + /** + * Receive the message from disruptor, when message in disruptor is empty, then send the cached data + * to the next workers. + * + * @param event published to the {@link RingBuffer} + * @param sequence of the event being processed + * @param endOfBatch flag to indicate if this is the last event in a batch from the {@link RingBuffer} + */ public void onEvent(MessageHolder event, long sequence, boolean endOfBatch) { try { Object message = event.getMessage(); @@ -48,6 +86,12 @@ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker { } } + /** + * Push the message into disruptor ring buffer. + * + * @param message of the data to process. + * @throws Exception not used. + */ public void tell(Object message) throws Exception { long sequence = ringBuffer.next(); try { diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java index c3f336f83..a5fb46af4 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/AbstractHashMessage.java @@ -5,9 +5,10 @@ package com.a.eye.skywalking.collector.actor.selector; * are intended to provide send message with {@link HashCodeSelector}. *

* Usually the implemented class used to persistence data to database - * or aggregation the metric, + * or aggregation the metric. * * @author pengys5 + * @since feature3.0 */ public abstract class AbstractHashMessage { private int hashCode; @@ -16,7 +17,7 @@ public abstract class AbstractHashMessage { this.hashCode = key.hashCode(); } - protected int getHashCode() { + int getHashCode() { return hashCode; } } diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java index aef766ab1..6e8b52c8a 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/HashCodeSelector.java @@ -11,6 +11,7 @@ import java.util.List; * message to same {@link WorkerRef}. Usually, use to database operate which avoid dirty data. * * @author pengys5 + * @since feature3.0 */ public class HashCodeSelector implements WorkerSelector { diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java index e56fed414..6e529ac3e 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/RollingSelector.java @@ -10,6 +10,7 @@ import java.util.List; * It choose {@link WorkerRef} nearly random, by round-robin. * * @author pengys5 + * @since feature3.0 */ public class RollingSelector implements WorkerSelector { diff --git a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java index 733c9f23f..8e491fdd7 100644 --- a/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java +++ b/skywalking-collector/skywalking-collector-cluster/src/main/java/com/a/eye/skywalking/collector/actor/selector/WorkerSelector.java @@ -12,6 +12,7 @@ import java.util.List; * Actually, the WorkerRef is designed to provide a routing ability in the collector cluster * * @author pengys5 + * @since feature3.0 */ public interface WorkerSelector {