lastData);
+}
diff --git a/apm-collector/apm-collector-remote/pom.xml b/apm-collector/apm-collector-remote/pom.xml
index 81c877d81..f20ebd870 100644
--- a/apm-collector/apm-collector-remote/pom.xml
+++ b/apm-collector/apm-collector-remote/pom.xml
@@ -29,6 +29,11 @@
apm-collector-server
${project.version}
+
+ org.skywalking
+ apm-collector-cluster
+ ${project.version}
+
io.grpc
grpc-netty
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java
index ad6875615..87796d38c 100644
--- a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleDefine.java
@@ -1,9 +1,41 @@
package org.skywalking.apm.collector.remote;
+import java.util.List;
+import java.util.Map;
+import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
+import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine;
+import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
+import org.skywalking.apm.collector.core.config.ConfigParseException;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
+import org.skywalking.apm.collector.core.framework.DefineException;
+import org.skywalking.apm.collector.core.framework.Handler;
import org.skywalking.apm.collector.core.module.ModuleDefine;
+import org.skywalking.apm.collector.core.server.Server;
+import org.skywalking.apm.collector.core.server.ServerException;
+import org.skywalking.apm.collector.core.server.ServerHolder;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
/**
* @author pengys5
*/
-public abstract class RemoteModuleDefine extends ModuleDefine {
+public abstract class RemoteModuleDefine extends ModuleDefine implements ClusterDataListenerDefine {
+
+ private final Logger logger = LoggerFactory.getLogger(RemoteModuleDefine.class);
+
+ @Override
+ public final void initialize(Map config, ServerHolder serverHolder) throws DefineException, ClientException {
+ try {
+ configParser().parse(config);
+ Server server = server();
+ serverHolder.holdServer(server, handlerList());
+
+ ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getDataMonitor().addListener(listener(), registration());
+ } catch (ConfigParseException | ServerException e) {
+ throw new RemoteModuleException(e.getMessage(), e);
+ }
+ }
+
+ public abstract List handlerList() throws DefineException;
}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleException.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleException.java
new file mode 100644
index 000000000..87af6e986
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/RemoteModuleException.java
@@ -0,0 +1,17 @@
+package org.skywalking.apm.collector.remote;
+
+import org.skywalking.apm.collector.core.module.ModuleException;
+
+/**
+ * @author pengys5
+ */
+public class RemoteModuleException extends ModuleException {
+
+ public RemoteModuleException(String message) {
+ super(message);
+ }
+
+ public RemoteModuleException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfig.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfig.java
new file mode 100644
index 000000000..6b465917a
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfig.java
@@ -0,0 +1,9 @@
+package org.skywalking.apm.collector.remote.grpc;
+
+/**
+ * @author pengys5
+ */
+public class RemoteGRPCConfig {
+ public static String HOST;
+ public static int PORT;
+}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfigParser.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfigParser.java
new file mode 100644
index 000000000..73beb8a49
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCConfigParser.java
@@ -0,0 +1,27 @@
+package org.skywalking.apm.collector.remote.grpc;
+
+import java.util.Map;
+import org.skywalking.apm.collector.core.config.ConfigParseException;
+import org.skywalking.apm.collector.core.module.ModuleConfigParser;
+import org.skywalking.apm.collector.core.util.ObjectUtils;
+import org.skywalking.apm.collector.core.util.StringUtils;
+
+/**
+ * @author pengys5
+ */
+public class RemoteGRPCConfigParser implements ModuleConfigParser {
+
+ private static final String HOST = "host";
+ private static final String PORT = "port";
+
+ @Override public void parse(Map config) throws ConfigParseException {
+ if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(HOST))) {
+ RemoteGRPCConfig.HOST = "localhost";
+ }
+ if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(PORT))) {
+ RemoteGRPCConfig.PORT = 11800;
+ } else {
+ RemoteGRPCConfig.PORT = (Integer)config.get(PORT);
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCDataListener.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCDataListener.java
new file mode 100644
index 000000000..7d2e96c84
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCDataListener.java
@@ -0,0 +1,17 @@
+package org.skywalking.apm.collector.remote.grpc;
+
+import org.skywalking.apm.collector.cluster.ClusterModuleDefine;
+import org.skywalking.apm.collector.core.cluster.ClusterDataListener;
+import org.skywalking.apm.collector.remote.RemoteModuleGroupDefine;
+
+/**
+ * @author pengys5
+ */
+public class RemoteGRPCDataListener extends ClusterDataListener {
+
+ public static final String PATH = ClusterModuleDefine.BASE_CATALOG + "." + RemoteModuleGroupDefine.GROUP_NAME + "." + RemoteGRPCModuleDefine.MODULE_NAME;
+
+ @Override public String path() {
+ return PATH;
+ }
+}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java
index 877952cff..22b884e43 100644
--- a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java
@@ -1,31 +1,42 @@
package org.skywalking.apm.collector.remote.grpc;
-import java.util.Map;
+import java.util.List;
import org.skywalking.apm.collector.core.client.Client;
-import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.client.DataMonitor;
+import org.skywalking.apm.collector.core.cluster.ClusterDataListener;
+import org.skywalking.apm.collector.core.config.ConfigException;
import org.skywalking.apm.collector.core.framework.DefineException;
+import org.skywalking.apm.collector.core.framework.Handler;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleRegistration;
import org.skywalking.apm.collector.core.server.Server;
-import org.skywalking.apm.collector.core.server.ServerHolder;
import org.skywalking.apm.collector.remote.RemoteModuleDefine;
import org.skywalking.apm.collector.remote.RemoteModuleGroupDefine;
+import org.skywalking.apm.collector.remote.grpc.handler.RemoteHandlerDefineException;
+import org.skywalking.apm.collector.remote.grpc.handler.RemoteHandlerDefineLoader;
+import org.skywalking.apm.collector.server.grpc.GRPCServer;
/**
* @author pengys5
*/
public class RemoteGRPCModuleDefine extends RemoteModuleDefine {
+
+ public static final String MODULE_NAME = "remote";
+
+ @Override public String name() {
+ return MODULE_NAME;
+ }
+
@Override protected String group() {
return RemoteModuleGroupDefine.GROUP_NAME;
}
@Override public boolean defaultModule() {
- return false;
+ return true;
}
@Override protected ModuleConfigParser configParser() {
- return null;
+ return new RemoteGRPCConfigParser();
}
@Override protected Client createClient(DataMonitor dataMonitor) {
@@ -33,18 +44,25 @@ public class RemoteGRPCModuleDefine extends RemoteModuleDefine {
}
@Override protected Server server() {
- return null;
+ return new GRPCServer(RemoteGRPCConfig.HOST, RemoteGRPCConfig.PORT);
}
@Override protected ModuleRegistration registration() {
- return null;
+ return new RemoteGRPCModuleRegistration();
}
- @Override public void initialize(Map config, ServerHolder serverHolder) throws DefineException, ClientException {
-
+ @Override public ClusterDataListener listener() {
+ return new RemoteGRPCDataListener();
}
- @Override public String name() {
- return null;
+ @Override public List handlerList() throws DefineException {
+ RemoteHandlerDefineLoader loader = new RemoteHandlerDefineLoader();
+ List handlers = null;
+ try {
+ handlers = loader.load();
+ } catch (ConfigException e) {
+ throw new RemoteHandlerDefineException(e.getMessage(), e);
+ }
+ return handlers;
}
}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleRegistration.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleRegistration.java
new file mode 100644
index 000000000..4f5a371d2
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleRegistration.java
@@ -0,0 +1,13 @@
+package org.skywalking.apm.collector.remote.grpc;
+
+import org.skywalking.apm.collector.core.module.ModuleRegistration;
+
+/**
+ * @author pengys5
+ */
+public class RemoteGRPCModuleRegistration extends ModuleRegistration {
+
+ @Override public Value buildValue() {
+ return new Value(RemoteGRPCConfig.HOST, RemoteGRPCConfig.PORT, null);
+ }
+}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineException.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineException.java
new file mode 100644
index 000000000..7a5c6063b
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineException.java
@@ -0,0 +1,17 @@
+package org.skywalking.apm.collector.remote.grpc.handler;
+
+import org.skywalking.apm.collector.core.framework.DefineException;
+
+/**
+ * @author pengys5
+ */
+public class RemoteHandlerDefineException extends DefineException {
+
+ public RemoteHandlerDefineException(String message) {
+ super(message);
+ }
+
+ public RemoteHandlerDefineException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineLoader.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineLoader.java
new file mode 100644
index 000000000..b125bc090
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefineLoader.java
@@ -0,0 +1,30 @@
+package org.skywalking.apm.collector.remote.grpc.handler;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.skywalking.apm.collector.core.config.ConfigException;
+import org.skywalking.apm.collector.core.framework.Handler;
+import org.skywalking.apm.collector.core.framework.Loader;
+import org.skywalking.apm.collector.core.util.DefinitionLoader;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class RemoteHandlerDefineLoader implements Loader> {
+
+ private final Logger logger = LoggerFactory.getLogger(RemoteHandlerDefineLoader.class);
+
+ @Override public List load() throws ConfigException {
+ List handlers = new ArrayList<>();
+
+ RemoteHandlerDefinitionFile definitionFile = new RemoteHandlerDefinitionFile();
+ DefinitionLoader definitionLoader = DefinitionLoader.load(Handler.class, definitionFile);
+ for (Handler handler : definitionLoader) {
+ logger.info("loaded remote handler definition class: {}", handler.getClass().getName());
+ handlers.add(handler);
+ }
+ return handlers;
+ }
+}
diff --git a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefinitionFile.java b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefinitionFile.java
new file mode 100644
index 000000000..4231229c4
--- /dev/null
+++ b/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/handler/RemoteHandlerDefinitionFile.java
@@ -0,0 +1,13 @@
+package org.skywalking.apm.collector.remote.grpc.handler;
+
+import org.skywalking.apm.collector.core.framework.DefinitionFile;
+
+/**
+ * @author pengys5
+ */
+public class RemoteHandlerDefinitionFile extends DefinitionFile {
+
+ @Override protected String fileName() {
+ return "remote_handler.define";
+ }
+}
diff --git a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto
index ae1da50f0..bf0f2d8bc 100644
--- a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto
+++ b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto
@@ -10,8 +10,8 @@ service RemoteCommonService {
message Message {
string workerRole = 1;
- int32 objectId = 2;
- bytes objectBytes = 3; // the byte array of data object
+ int32 dataDefineId = 2;
+ bytes dataBytes = 3;
}
message Empty {
diff --git a/apm-collector/apm-collector-stream/pom.xml b/apm-collector/apm-collector-stream/pom.xml
index e9be3a610..9ace78e53 100644
--- a/apm-collector/apm-collector-stream/pom.xml
+++ b/apm-collector/apm-collector-stream/pom.xml
@@ -23,5 +23,10 @@
apm-collector-queue
${project.version}
+
+ org.skywalking
+ apm-collector-remote
+ ${project.version}
+
\ No newline at end of file
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorker.java
similarity index 79%
rename from apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorker.java
rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorker.java
index cc2991b4d..20c552bee 100644
--- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorker.java
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorker.java
@@ -1,7 +1,7 @@
package org.skywalking.apm.collector.stream;
/**
- * The AbstractClusterWorker implementations represent workers,
+ * The AbstractRemoteWorker implementations represent workers,
* which receive remote messages.
*
* Usually, the implementations are doing persistent, or aggregate works.
@@ -9,17 +9,17 @@ package org.skywalking.apm.collector.stream;
* @author pengys5
* @since v3.0-2017
*/
-public abstract class AbstractClusterWorker extends AbstractWorker {
+public abstract class AbstractRemoteWorker extends AbstractWorker {
/**
- * Construct an AbstractClusterWorker with the worker role and context.
+ * Construct an AbstractRemoteWorker with the worker role and context.
*
* @param role If multi-workers are for load balance, they should be more likely called worker instance. Meaning,
* each worker have multi instances.
* @param clusterContext See {@link ClusterWorkerContext}
* @param selfContext See {@link LocalWorkerContext}
*/
- protected AbstractClusterWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ protected AbstractRemoteWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
super(role, clusterContext, selfContext);
}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java
similarity index 66%
rename from apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorkerProvider.java
rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java
index 3e1fd1337..3774be6cb 100644
--- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractClusterWorkerProvider.java
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/AbstractRemoteWorkerProvider.java
@@ -1,17 +1,17 @@
package org.skywalking.apm.collector.stream;
/**
- * The AbstractClusterWorkerProvider implementations represent providers,
- * which create instance of cluster workers whose implemented {@link AbstractClusterWorker}.
+ * The AbstractRemoteWorkerProvider implementations represent providers,
+ * which create instance of cluster workers whose implemented {@link AbstractRemoteWorker}.
*
*
* @author pengys5
* @since v3.0-2017
*/
-public abstract class AbstractClusterWorkerProvider extends AbstractWorkerProvider {
+public abstract class AbstractRemoteWorkerProvider extends AbstractWorkerProvider {
/**
- * Create how many worker instance of {@link AbstractClusterWorker} in one jvm.
+ * Create how many worker instance of {@link AbstractRemoteWorker} in one jvm.
*
* @return The worker instance number.
*/
@@ -21,7 +21,7 @@ public abstract class AbstractClusterWorkerProvider dataDefineMap;
+
private Map> roleWorkers;
public WorkerContext() {
@@ -20,8 +23,11 @@ public abstract class WorkerContext implements Context {
return this.roleWorkers;
}
- @Override
- final public WorkerRefs lookup(Role role) throws WorkerNotFoundException {
+ public final DataDefine getDataDefine(int defineId) {
+ return dataDefineMap.get(defineId);
+ }
+
+ @Override final public WorkerRefs lookup(Role role) throws WorkerNotFoundException {
if (getRoleWorkers().containsKey(role.roleName())) {
WorkerRefs refs = new WorkerRefs(getRoleWorkers().get(role.roleName()), role.workerSelector());
return refs;
@@ -30,16 +36,14 @@ public abstract class WorkerContext implements Context {
}
}
- @Override
- final public void put(WorkerRef workerRef) {
+ @Override final public void put(WorkerRef workerRef) {
if (!getRoleWorkers().containsKey(workerRef.getRole().roleName())) {
getRoleWorkers().putIfAbsent(workerRef.getRole().roleName(), new ArrayList());
}
getRoleWorkers().get(workerRef.getRole().roleName()).add(workerRef);
}
- @Override
- final public void remove(WorkerRef workerRef) {
+ @Override final public void remove(WorkerRef workerRef) {
getRoleWorkers().remove(workerRef.getRole().roleName());
}
}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/WorkerModuleInstaller.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/WorkerModuleInstaller.java
deleted file mode 100644
index 5dd6095ca..000000000
--- a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/WorkerModuleInstaller.java
+++ /dev/null
@@ -1,25 +0,0 @@
-package org.skywalking.apm.collector.stream;
-
-import java.util.Map;
-import org.skywalking.apm.collector.core.client.ClientException;
-import org.skywalking.apm.collector.core.framework.DefineException;
-import org.skywalking.apm.collector.core.module.ModuleDefine;
-import org.skywalking.apm.collector.core.module.ModuleInstaller;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-/**
- * @author pengys5
- */
-public class WorkerModuleInstaller implements ModuleInstaller {
-
- private final Logger logger = LoggerFactory.getLogger(WorkerModuleInstaller.class);
-
- @Override public void install(Map moduleConfig,
- Map moduleDefineMap) throws DefineException, ClientException {
- logger.info("beginning worker module install");
- Map.Entry workerConfigEntry = moduleConfig.entrySet().iterator().next();
- ModuleDefine moduleDefine = moduleDefineMap.get(workerConfigEntry.getKey());
- moduleDefine.initialize(workerConfigEntry.getValue());
- }
-}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/AggregationWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/AggregationWorker.java
new file mode 100644
index 000000000..bd892035c
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/AggregationWorker.java
@@ -0,0 +1,55 @@
+package org.skywalking.apm.collector.stream.impl;
+
+import org.skywalking.apm.collector.core.queue.EndOfBatchCommand;
+import org.skywalking.apm.collector.stream.AbstractLocalAsyncWorker;
+import org.skywalking.apm.collector.stream.ClusterWorkerContext;
+import org.skywalking.apm.collector.stream.LocalWorkerContext;
+import org.skywalking.apm.collector.stream.ProviderNotFoundException;
+import org.skywalking.apm.collector.stream.Role;
+import org.skywalking.apm.collector.stream.WorkerException;
+import org.skywalking.apm.collector.stream.impl.data.Data;
+import org.skywalking.apm.collector.stream.impl.data.DataCache;
+
+/**
+ * @author pengys5
+ */
+public abstract class AggregationWorker extends AbstractLocalAsyncWorker {
+
+ private DataCache dataCache;
+
+ public AggregationWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ super(role, clusterContext, selfContext);
+ dataCache = new DataCache();
+ }
+
+ private int messageNum;
+
+ @Override public void preStart() throws ProviderNotFoundException {
+ super.preStart();
+ }
+
+ @Override protected final void onWork(Object message) throws WorkerException {
+ if (message instanceof EndOfBatchCommand) {
+ sendToNext();
+ } else {
+ messageNum++;
+ aggregate(message);
+
+ if (messageNum >= 100) {
+ sendToNext();
+ messageNum = 0;
+ }
+ }
+ }
+
+ protected abstract void sendToNext();
+
+ protected final void aggregate(Object message) {
+ Data data = (Data)message;
+ if (dataCache.containsKey(data.id())) {
+ getClusterContext().getDataDefine(data.getDefineId()).mergeData(data, dataCache.get(data.id()));
+ } else {
+ dataCache.put(data.id(), data);
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/Const.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/Const.java
new file mode 100644
index 000000000..1850a2ee2
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/Const.java
@@ -0,0 +1,13 @@
+package org.skywalking.apm.collector.stream.impl;
+
+/**
+ * @author pengys5
+ */
+public class Const {
+ public static final String ID_SPLIT = "..-..";
+ public static final String IDS_SPLIT = "\\.\\.-\\.\\.";
+ public static final String PEERS_FRONT_SPLIT = "[";
+ public static final String PEERS_BEHIND_SPLIT = "]";
+ public static final String USER_CODE = "User";
+ public static final String RESULT = "result";
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/GRPCRemoteWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/GRPCRemoteWorker.java
new file mode 100644
index 000000000..3bd9a33fa
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/GRPCRemoteWorker.java
@@ -0,0 +1,26 @@
+package org.skywalking.apm.collector.stream.impl;
+
+import org.skywalking.apm.collector.stream.AbstractRemoteWorker;
+import org.skywalking.apm.collector.stream.ClusterWorkerContext;
+import org.skywalking.apm.collector.stream.LocalWorkerContext;
+import org.skywalking.apm.collector.stream.ProviderNotFoundException;
+import org.skywalking.apm.collector.stream.Role;
+import org.skywalking.apm.collector.stream.WorkerException;
+
+/**
+ * @author pengys5
+ */
+public class GRPCRemoteWorker extends AbstractRemoteWorker {
+
+ protected GRPCRemoteWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
+ super(role, clusterContext, selfContext);
+ }
+
+ @Override public void preStart() throws ProviderNotFoundException {
+
+ }
+
+ @Override protected final void onWork(Object message) throws WorkerException {
+
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/RemoteCommonServiceHandler.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/RemoteCommonServiceHandler.java
new file mode 100644
index 000000000..ea560ef12
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/RemoteCommonServiceHandler.java
@@ -0,0 +1,24 @@
+package org.skywalking.apm.collector.stream.impl;
+
+import com.google.protobuf.ByteString;
+import io.grpc.stub.StreamObserver;
+import org.skywalking.apm.collector.remote.grpc.proto.Empty;
+import org.skywalking.apm.collector.remote.grpc.proto.Message;
+import org.skywalking.apm.collector.remote.grpc.proto.RemoteCommonServiceGrpc;
+import org.skywalking.apm.collector.server.grpc.GRPCHandler;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class RemoteCommonServiceHandler extends RemoteCommonServiceGrpc.RemoteCommonServiceImplBase implements GRPCHandler {
+
+ private final Logger logger = LoggerFactory.getLogger(RemoteCommonServiceHandler.class);
+
+ @Override public void call(Message request, StreamObserver responseObserver) {
+ String workerRole = request.getWorkerRole();
+ int dataDefineId = request.getDataDefineId();
+ ByteString bytesData = request.getDataBytes();
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Attribute.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Attribute.java
new file mode 100644
index 000000000..3075517de
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Attribute.java
@@ -0,0 +1,28 @@
+package org.skywalking.apm.collector.stream.impl.data;
+
+/**
+ * @author pengys5
+ */
+public class Attribute {
+ private final String name;
+ private final AttributeType type;
+ private final Operation operation;
+
+ public Attribute(String name, AttributeType type, Operation operation) {
+ this.name = name;
+ this.type = type;
+ this.operation = operation;
+ }
+
+ public String getName() {
+ return name;
+ }
+
+ public AttributeType getType() {
+ return type;
+ }
+
+ public Operation getOperation() {
+ return operation;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/AttributeType.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/AttributeType.java
new file mode 100644
index 000000000..84e317ec0
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/AttributeType.java
@@ -0,0 +1,8 @@
+package org.skywalking.apm.collector.stream.impl.data;
+
+/**
+ * @author pengys5
+ */
+public enum AttributeType {
+ STRING, LONG, FLOAT
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Data.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Data.java
new file mode 100644
index 000000000..be8d4eb21
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/Data.java
@@ -0,0 +1,50 @@
+package org.skywalking.apm.collector.stream.impl.data;
+
+/**
+ * @author pengys5
+ */
+public class Data {
+ private int defineId;
+ private String[] dataStrings;
+ private Long[] dataLongs;
+ private Float[] dataFloats;
+
+ public Data(int defineId, int stringCapacity, int longCapacity, int floatCapacity) {
+ this.defineId = defineId;
+ this.dataStrings = new String[stringCapacity];
+ this.dataLongs = new Long[longCapacity];
+ this.dataFloats = new Float[floatCapacity];
+ }
+
+ public void setDataString(int position, String value) {
+ dataStrings[position] = value;
+ }
+
+ public void setDataLong(int position, Long value) {
+ dataLongs[position] = value;
+ }
+
+ public void setDataFloat(int position, Float value) {
+ dataFloats[position] = value;
+ }
+
+ public String getDataString(int position) {
+ return dataStrings[position];
+ }
+
+ public Long getDataLong(int position) {
+ return dataLongs[position];
+ }
+
+ public Float getDataFloat(int position) {
+ return dataFloats[position];
+ }
+
+ public String id() {
+ return dataStrings[0];
+ }
+
+ public int getDefineId() {
+ return defineId;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCache.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCache.java
new file mode 100644
index 000000000..db51cf91d
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCache.java
@@ -0,0 +1,30 @@
+package org.skywalking.apm.collector.stream.impl.data;
+
+/**
+ * @author pengys5
+ */
+public class DataCache extends Window {
+
+ private DataCollection lockedDataCollection;
+
+ public boolean containsKey(String id) {
+ return lockedDataCollection.containsKey(id);
+ }
+
+ public Data get(String id) {
+ return lockedDataCollection.get(id);
+ }
+
+ public void put(String id, Data data) {
+ lockedDataCollection.put(id, data);
+ }
+
+ public void hold() {
+ lockedDataCollection = getCurrentAndHold();
+ }
+
+ public void release() {
+ lockedDataCollection.release();
+ lockedDataCollection = null;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCollection.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCollection.java
new file mode 100644
index 000000000..725dcca71
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataCollection.java
@@ -0,0 +1,53 @@
+package org.skywalking.apm.collector.stream.impl.data;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * @author pengys5
+ */
+public class DataCollection {
+ private Map data;
+ private volatile boolean isHold;
+
+ public DataCollection() {
+ this.data = new HashMap<>();
+ this.isHold = false;
+ }
+
+ public void release() {
+ isHold = false;
+ }
+
+ public void hold() {
+ isHold = true;
+ }
+
+ public boolean isHolding() {
+ return isHold;
+ }
+
+ public boolean containsKey(String key) {
+ return data.containsKey(key);
+ }
+
+ public void put(String key, Data value) {
+ data.put(key, value);
+ }
+
+ public Data get(String key) {
+ return data.get(key);
+ }
+
+ public int size() {
+ return data.size();
+ }
+
+ public void clear() {
+ data.clear();
+ }
+
+ public Map asMap() {
+ return data;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefine.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefine.java
new file mode 100644
index 000000000..5571ac264
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefine.java
@@ -0,0 +1,78 @@
+package org.skywalking.apm.collector.stream.impl.data;
+
+/**
+ * @author pengys5
+ */
+public abstract class DataDefine {
+ private Attribute[] attributes;
+ private int stringCapacity;
+ private int longCapacity;
+ private int floatCapacity;
+
+ public DataDefine() {
+ stringCapacity = 0;
+ longCapacity = 0;
+ floatCapacity = 0;
+ }
+
+ public final void initial() {
+ for (Attribute attribute : attributes) {
+ if (AttributeType.STRING.equals(attribute.getType())) {
+ stringCapacity++;
+ } else if (AttributeType.LONG.equals(attribute.getType())) {
+ longCapacity++;
+ } else if (AttributeType.FLOAT.equals(attribute.getType())) {
+ floatCapacity++;
+ }
+ }
+ }
+
+ public final void addAttribute(int position, Attribute attribute) {
+ attributes[position] = attribute;
+ }
+
+ public final void define() {
+ attributes = new Attribute[initialCapacity()];
+ }
+
+ protected abstract int defineId();
+
+ protected abstract int initialCapacity();
+
+ protected abstract void attributeDefine();
+
+ public int getStringCapacity() {
+ return stringCapacity;
+ }
+
+ public int getLongCapacity() {
+ return longCapacity;
+ }
+
+ public int getFloatCapacity() {
+ return floatCapacity;
+ }
+
+ public Data build() {
+ return new Data(defineId(), getStringCapacity(), getLongCapacity(), getFloatCapacity());
+ }
+
+ public void mergeData(Data newData, Data oldData) {
+ int stringPosition = 0;
+ int longPosition = 0;
+ int floatPosition = 0;
+ for (int i = 0; i < initialCapacity(); i++) {
+ Attribute attribute = attributes[i];
+ if (AttributeType.STRING.equals(attribute.getType())) {
+ attribute.getOperation().operate(newData.getDataString(stringPosition), oldData.getDataString(stringPosition));
+ stringPosition++;
+ } else if (AttributeType.LONG.equals(attribute.getType())) {
+ attribute.getOperation().operate(newData.getDataLong(longPosition), oldData.getDataLong(longPosition));
+ longPosition++;
+ } else if (AttributeType.FLOAT.equals(attribute.getType())) {
+ attribute.getOperation().operate(newData.getDataFloat(floatPosition), oldData.getDataFloat(floatPosition));
+ floatPosition++;
+ }
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefineLoader.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefineLoader.java
new file mode 100644
index 000000000..8ce454440
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/impl/data/DataDefineLoader.java
@@ -0,0 +1,29 @@
+package org.skywalking.apm.collector.stream.impl.data;
+
+import java.util.HashMap;
+import java.util.Map;
+import org.skywalking.apm.collector.core.config.ConfigException;
+import org.skywalking.apm.collector.core.framework.Loader;
+import org.skywalking.apm.collector.core.util.DefinitionLoader;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class DataDefineLoader implements Loader