getColumnDefines() {
+ return columnDefines;
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java
index dd66a5ecb..22141d6a2 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java
@@ -1,5 +1,6 @@
package org.skywalking.apm.collector.core.util;
+import java.util.List;
import java.util.Map;
/**
@@ -10,4 +11,12 @@ public class CollectionUtils {
public static boolean isEmpty(Map map) {
return map == null || map.size() == 0;
}
+
+ public static boolean isEmpty(List list) {
+ return list == null || list.size() == 0;
+ }
+
+ public static boolean isNotEmpty(List list) {
+ return !isEmpty(list);
+ }
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java
index df264ab32..9300a103e 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java
@@ -7,4 +7,8 @@ public class ObjectUtils {
public static boolean isEmpty(Object obj) {
return obj == null;
}
+
+ public static boolean isNotEmpty(Object obj) {
+ return !isEmpty(obj);
+ }
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ResourceUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ResourceUtils.java
index 31daeec6f..257f85dd2 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ResourceUtils.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ResourceUtils.java
@@ -1,16 +1,21 @@
package org.skywalking.apm.collector.core.util;
+import java.io.File;
import java.io.FileNotFoundException;
import java.io.FileReader;
+import java.net.URL;
/**
* @author pengys5
*/
public class ResourceUtils {
- private static final String PATH = ResourceUtils.class.getResource("/").getPath();
-
public static FileReader read(String fileName) throws FileNotFoundException {
- return new FileReader(PATH + fileName);
+ URL url = ResourceUtils.class.getClassLoader().getResource(fileName);
+ if (url == null) {
+ throw new FileNotFoundException("file not found: " + fileName);
+ }
+ File file = new File(ResourceUtils.class.getClassLoader().getResource(fileName).getFile());
+ return new FileReader(file);
}
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java
index 83413c178..58d260d7e 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java
@@ -10,4 +10,8 @@ public class StringUtils {
public static boolean isEmpty(Object str) {
return str == null || EMPTY_STRING.equals(str);
}
+
+ public static boolean isNotEmpty(Object str) {
+ return !isEmpty(str);
+ }
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractClusterWorkerProvider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractClusterWorkerProvider.java
deleted file mode 100644
index 519689182..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractClusterWorkerProvider.java
+++ /dev/null
@@ -1,37 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-/**
- * The AbstractClusterWorkerProvider implementations represent providers,
- * which create instance of cluster workers whose implemented {@link AbstractClusterWorker}.
- *
- *
- * @author pengys5
- * @since v3.0-2017
- */
-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 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 ProviderNotFoundException {
- T clusterWorker = workerInstance(getClusterContext());
- clusterWorker.preStart();
-
- ClusterWorkerRef workerRef = new ClusterWorkerRef(role(), clusterWorker);
- getClusterContext().put(workerRef);
- return workerRef;
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorkerProvider.java
deleted file mode 100644
index 22b3e9a58..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorkerProvider.java
+++ /dev/null
@@ -1,28 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-import org.skywalking.apm.collector.core.queue.QueueEventHandler;
-import org.skywalking.apm.collector.core.queue.QueueModuleContext;
-
-/**
- * @author pengys5
- */
-public abstract class AbstractLocalAsyncWorkerProvider extends AbstractLocalWorkerProvider {
-
- public abstract int queueSize();
-
- @Override final public WorkerRef onCreate(
- LocalWorkerContext localContext) throws ProviderNotFoundException {
- T localAsyncWorker = workerInstance(getClusterContext());
- localAsyncWorker.preStart();
-
- QueueEventHandler queueEventHandler = QueueModuleContext.CREATOR.create(queueSize(), localAsyncWorker);
-
- LocalAsyncWorkerRef workerRef = new LocalAsyncWorkerRef(role(), queueEventHandler);
-
- if (localContext != null) {
- localContext.put(workerRef);
- }
-
- return workerRef;
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalSyncWorkerProvider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalSyncWorkerProvider.java
deleted file mode 100644
index 056b724f6..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalSyncWorkerProvider.java
+++ /dev/null
@@ -1,21 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-/**
- * @author pengys5
- */
-public abstract class AbstractLocalSyncWorkerProvider extends AbstractLocalWorkerProvider {
-
- @Override
- final public WorkerRef onCreate(
- LocalWorkerContext localContext) throws ProviderNotFoundException {
- T localSyncWorker = (T) workerInstance(getClusterContext());
- localSyncWorker.preStart();
-
- LocalSyncWorkerRef workerRef = new LocalSyncWorkerRef(role(), localSyncWorker);
-
- if (localContext != null) {
- localContext.put(workerRef);
- }
- return workerRef;
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractWorker.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractWorker.java
deleted file mode 100644
index cdfc839c7..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractWorker.java
+++ /dev/null
@@ -1,37 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-/**
- * @author pengys5
- */
-public abstract class AbstractWorker {
-
- private final LocalWorkerContext selfContext;
-
- private final Role role;
-
- private final ClusterWorkerContext clusterContext;
-
- public AbstractWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- this.role = role;
- this.clusterContext = clusterContext;
- this.selfContext = selfContext;
- }
-
- public abstract void preStart() throws ProviderNotFoundException;
-
- final public LookUp getSelfContext() {
- return selfContext;
- }
-
- final public LookUp getClusterContext() {
- return clusterContext;
- }
-
- final public Role getRole() {
- return role;
- }
-
- final public static AbstractWorker noOwner() {
- return null;
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractWorkerProvider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractWorkerProvider.java
deleted file mode 100644
index bc87c19e4..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractWorkerProvider.java
+++ /dev/null
@@ -1,36 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-/**
- * @author pengys5
- */
-public abstract class AbstractWorkerProvider implements Provider {
-
- private ClusterWorkerContext clusterContext;
-
- public abstract Role role();
-
- public abstract T workerInstance(ClusterWorkerContext clusterContext);
-
- public abstract WorkerRef onCreate(
- LocalWorkerContext localContext) throws ProviderNotFoundException;
-
- final public void setClusterContext(ClusterWorkerContext clusterContext) {
- this.clusterContext = clusterContext;
- }
-
- final protected ClusterWorkerContext getClusterContext() {
- return clusterContext;
- }
-
- final public WorkerRef create(
- AbstractWorker workerOwner) throws ProviderNotFoundException {
-
- if (workerOwner == null) {
- return onCreate(null);
- } else if (workerOwner.getSelfContext() instanceof LocalWorkerContext) {
- return onCreate((LocalWorkerContext)workerOwner.getSelfContext());
- } else {
- throw new IllegalArgumentException("the argument of workerOwner is Illegal");
- }
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/ClusterWorkerContext.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/ClusterWorkerContext.java
deleted file mode 100644
index ebd29da62..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/ClusterWorkerContext.java
+++ /dev/null
@@ -1,35 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-import org.slf4j.Logger;
-import org.slf4j.LoggerFactory;
-
-/**
- * @author pengys5
- */
-public class ClusterWorkerContext extends WorkerContext {
- private final Logger logger = LoggerFactory.getLogger(ClusterWorkerContext.class);
-
- private Map providers = new ConcurrentHashMap<>();
-
- @Override
- public AbstractWorkerProvider findProvider(Role role) throws ProviderNotFoundException {
- logger.debug("find role of %s provider from ClusterWorkerContext", role.roleName());
- if (providers.containsKey(role.roleName())) {
- return providers.get(role.roleName());
- } else {
- throw new ProviderNotFoundException("role=" + role.roleName() + ", no available provider.");
- }
- }
-
- @Override
- public void putProvider(AbstractWorkerProvider provider) throws UsedRoleNameException {
- logger.debug("put role of %s provider into ClusterWorkerContext", provider.role().roleName());
- if (providers.containsKey(provider.role().roleName())) {
- throw new UsedRoleNameException("provider with role=" + provider.role().roleName() + " duplicate each other.");
- } else {
- providers.put(provider.role().roleName(), provider);
- }
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/ClusterWorkerRef.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/ClusterWorkerRef.java
deleted file mode 100644
index dc2243f3a..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/ClusterWorkerRef.java
+++ /dev/null
@@ -1,19 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-/**
- * @author pengys5
- */
-public class ClusterWorkerRef extends WorkerRef {
-
- private AbstractClusterWorker clusterWorker;
-
- public ClusterWorkerRef(Role role, AbstractClusterWorker clusterWorker) {
- super(role);
- this.clusterWorker = clusterWorker;
- }
-
- @Override
- public void tell(Object message) throws WorkerInvokeException {
- clusterWorker.allocateJob(message);
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/LocalWorkerContext.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/LocalWorkerContext.java
deleted file mode 100644
index 3d70bbb56..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/LocalWorkerContext.java
+++ /dev/null
@@ -1,17 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-/**
- * @author pengys5
- */
-public class LocalWorkerContext extends WorkerContext {
-
- @Override
- final public AbstractWorkerProvider findProvider(Role role) throws ProviderNotFoundException {
- return null;
- }
-
- @Override
- final public void putProvider(AbstractWorkerProvider provider) throws UsedRoleNameException {
-
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/Provider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/Provider.java
deleted file mode 100644
index d9ccf9769..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/Provider.java
+++ /dev/null
@@ -1,9 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-/**
- * @author pengys5
- */
-public interface Provider {
-
- WorkerRef create(AbstractWorker workerOwner) throws ProviderNotFoundException;
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/Role.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/Role.java
deleted file mode 100644
index 8c59ff0ac..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/Role.java
+++ /dev/null
@@ -1,13 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-import org.skywalking.apm.collector.core.worker.selector.WorkerSelector;
-
-/**
- * @author pengys5
- */
-public interface Role {
-
- String roleName();
-
- WorkerSelector workerSelector();
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerContext.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerContext.java
deleted file mode 100644
index 78ff5e231..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerContext.java
+++ /dev/null
@@ -1,45 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-import java.util.ArrayList;
-import java.util.List;
-import java.util.Map;
-import java.util.concurrent.ConcurrentHashMap;
-
-/**
- * @author pengys5
- */
-public abstract class WorkerContext implements Context {
-
- private Map> roleWorkers;
-
- public WorkerContext() {
- this.roleWorkers = new ConcurrentHashMap<>();
- }
-
- private Map> getRoleWorkers() {
- return this.roleWorkers;
- }
-
- @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;
- } else {
- throw new WorkerNotFoundException("role=" + role.roleName() + ", no available worker.");
- }
- }
-
- @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) {
- getRoleWorkers().remove(workerRef.getRole().roleName());
- }
-}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerModuleInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerModuleInstaller.java
deleted file mode 100644
index 0b86aa90c..000000000
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/WorkerModuleInstaller.java
+++ /dev/null
@@ -1,25 +0,0 @@
-package org.skywalking.apm.collector.core.worker;
-
-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-core/src/main/resources/application-default.yml b/apm-collector/apm-collector-core/src/main/resources/application-default.yml
new file mode 100644
index 000000000..1e4e62547
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/resources/application-default.yml
@@ -0,0 +1,26 @@
+cluster:
+ zookeeper:
+ hostPort: localhost:2181
+ sessionTimeout: 100000
+# redis:
+# host: localhost
+# port: 6379
+queue:
+ disruptor: on
+ data_carrier: off
+agentstream:
+ grpc:
+ host: localhost
+ port: 1000
+ jetty:
+ host: localhost
+ port: 2000
+ context_path: /
+discovery:
+ grpc: localhost
+ port: 1000
+ui:
+ jetty:
+ host: localhost
+ port: 12800
+
diff --git a/apm-collector/apm-collector-core/src/main/resources/logback.xml b/apm-collector/apm-collector-core/src/main/resources/logback.xml
index b0df07af4..dc57f744d 100644
--- a/apm-collector/apm-collector-core/src/main/resources/logback.xml
+++ b/apm-collector/apm-collector-core/src/main/resources/logback.xml
@@ -2,11 +2,12 @@
- %d{HH:mm:ss.SSS} [%thread] %-5level %logger{36} - %msg%n
+ %d{HH:mm:ss.SSS} [%thread] %-5level %logger{70} - %msg%n
-
+
+
diff --git a/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/config/ModuleConfigLoaderTestCase.java b/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/config/ModuleConfigLoaderTestCase.java
deleted file mode 100644
index 85f802f23..000000000
--- a/apm-collector/apm-collector-core/src/test/java/org/skywalking/apm/collector/core/config/ModuleConfigLoaderTestCase.java
+++ /dev/null
@@ -1,18 +0,0 @@
-package org.skywalking.apm.collector.core.config;
-
-import java.io.FileNotFoundException;
-import org.junit.Test;
-import org.skywalking.apm.collector.core.module.ModuleConfigLoader;
-import org.skywalking.apm.collector.core.module.ModuleConfigLoaderException;
-
-/**
- * @author pengys5
- */
-public class ModuleConfigLoaderTestCase {
-
- @Test
- public void testLoad() throws ModuleConfigLoaderException {
- ModuleConfigLoader loader = new ModuleConfigLoader();
- loader.load();
- }
-}
diff --git a/apm-collector/apm-collector-discovery/pom.xml b/apm-collector/apm-collector-discovery/pom.xml
deleted file mode 100644
index 5aba21592..000000000
--- a/apm-collector/apm-collector-discovery/pom.xml
+++ /dev/null
@@ -1,14 +0,0 @@
-
-
-
- apm-collector
- org.skywalking
- 3.2-2017
-
- 4.0.0
-
- apm-collector-discovery
- jar
-
\ No newline at end of file
diff --git a/apm-collector/apm-collector-discovery/src/main/java/org/skywalking/apm/collector/discovery/DiscoveryJettyModuleDefine.java b/apm-collector/apm-collector-discovery/src/main/java/org/skywalking/apm/collector/discovery/DiscoveryJettyModuleDefine.java
deleted file mode 100644
index 3d696fa9e..000000000
--- a/apm-collector/apm-collector-discovery/src/main/java/org/skywalking/apm/collector/discovery/DiscoveryJettyModuleDefine.java
+++ /dev/null
@@ -1,7 +0,0 @@
-package org.skywalking.apm.collector.discovery;
-
-/**
- * @author pengys5
- */
-public class DiscoveryJettyModuleDefine {
-}
diff --git a/apm-collector/apm-collector-discovery/src/main/resources/META-INF/defines/module.define b/apm-collector/apm-collector-discovery/src/main/resources/META-INF/defines/module.define
deleted file mode 100644
index 427dbc6e4..000000000
--- a/apm-collector/apm-collector-discovery/src/main/resources/META-INF/defines/module.define
+++ /dev/null
@@ -1 +0,0 @@
-org.skywalking.apm.collector.discovery.DiscoveryJettyModuleDefine
\ No newline at end of file
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleContext.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleContext.java
new file mode 100644
index 000000000..cf9bf00cd
--- /dev/null
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleContext.java
@@ -0,0 +1,23 @@
+package org.skywalking.apm.collector.queue;
+
+import org.skywalking.apm.collector.core.framework.Context;
+import org.skywalking.apm.collector.core.queue.QueueCreator;
+
+/**
+ * @author pengys5
+ */
+public class QueueModuleContext extends Context {
+ private QueueCreator queueCreator;
+
+ public QueueModuleContext(String groupName) {
+ super(groupName);
+ }
+
+ public QueueCreator getQueueCreator() {
+ return queueCreator;
+ }
+
+ public void setQueueCreator(QueueCreator queueCreator) {
+ this.queueCreator = queueCreator;
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/queue/QueueModuleDefine.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleDefine.java
similarity index 73%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/queue/QueueModuleDefine.java
rename to apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleDefine.java
index 9b6bd447f..88bbae0c7 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/queue/QueueModuleDefine.java
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleDefine.java
@@ -1,7 +1,7 @@
-package org.skywalking.apm.collector.core.queue;
+package org.skywalking.apm.collector.queue;
import org.skywalking.apm.collector.core.client.Client;
-import org.skywalking.apm.collector.core.framework.DataInitializer;
+import org.skywalking.apm.collector.core.client.DataMonitor;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleRegistration;
@@ -15,11 +15,7 @@ public abstract class QueueModuleDefine extends ModuleDefine {
throw new UnsupportedOperationException("");
}
- @Override protected final Client createClient() {
- throw new UnsupportedOperationException("");
- }
-
- @Override protected final DataInitializer dataInitializer() {
+ @Override protected Client createClient(DataMonitor dataMonitor) {
throw new UnsupportedOperationException("");
}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleException.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleException.java
new file mode 100644
index 000000000..b22b24089
--- /dev/null
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleException.java
@@ -0,0 +1,16 @@
+package org.skywalking.apm.collector.queue;
+
+import org.skywalking.apm.collector.core.module.ModuleException;
+
+/**
+ * @author pengys5
+ */
+public class QueueModuleException extends ModuleException {
+ public QueueModuleException(String message) {
+ super(message);
+ }
+
+ public QueueModuleException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleGroupDefine.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleGroupDefine.java
new file mode 100644
index 000000000..9b40b7c2c
--- /dev/null
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleGroupDefine.java
@@ -0,0 +1,25 @@
+package org.skywalking.apm.collector.queue;
+
+import org.skywalking.apm.collector.core.framework.Context;
+import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
+import org.skywalking.apm.collector.core.module.ModuleInstaller;
+
+/**
+ * @author pengys5
+ */
+public class QueueModuleGroupDefine implements ModuleGroupDefine {
+
+ public static final String GROUP_NAME = "queue";
+
+ @Override public String name() {
+ return GROUP_NAME;
+ }
+
+ @Override public Context groupContext() {
+ return new QueueModuleContext(GROUP_NAME);
+ }
+
+ @Override public ModuleInstaller moduleInstaller() {
+ return new QueueModuleInstaller();
+ }
+}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleInstaller.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleInstaller.java
new file mode 100644
index 000000000..b0ed3238a
--- /dev/null
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleInstaller.java
@@ -0,0 +1,28 @@
+package org.skywalking.apm.collector.queue;
+
+import java.util.Map;
+import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
+import org.skywalking.apm.collector.core.framework.DefineException;
+import org.skywalking.apm.collector.core.module.ModuleDefine;
+import org.skywalking.apm.collector.core.module.SingleModuleInstaller;
+import org.skywalking.apm.collector.core.server.ServerHolder;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class QueueModuleInstaller extends SingleModuleInstaller {
+
+ private final Logger logger = LoggerFactory.getLogger(QueueModuleInstaller.class);
+
+ @Override public void install(Map moduleConfig,
+ Map moduleDefineMap, ServerHolder serverHolder) throws DefineException, ClientException {
+ logger.info("beginning queue module install");
+ QueueModuleContext context = new QueueModuleContext(QueueModuleGroupDefine.GROUP_NAME);
+ CollectorContextHelper.INSTANCE.putContext(context);
+
+ installSingle(moduleConfig, moduleDefineMap, serverHolder);
+ }
+}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierCreator.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierCreator.java
deleted file mode 100644
index 2402b83f4..000000000
--- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierCreator.java
+++ /dev/null
@@ -1,15 +0,0 @@
-package org.skywalking.apm.collector.queue.datacarrier;
-
-import org.skywalking.apm.collector.core.queue.Creator;
-import org.skywalking.apm.collector.core.queue.QueueEventHandler;
-import org.skywalking.apm.collector.core.worker.AbstractLocalAsyncWorker;
-
-/**
- * @author pengys5
- */
-public class DataCarrierCreator implements Creator {
-
- @Override public QueueEventHandler create(int queueSize, AbstractLocalAsyncWorker localAsyncWorker) {
- return null;
- }
-}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierQueueCreator.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierQueueCreator.java
new file mode 100644
index 000000000..e8593d6fc
--- /dev/null
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierQueueCreator.java
@@ -0,0 +1,15 @@
+package org.skywalking.apm.collector.queue.datacarrier;
+
+import org.skywalking.apm.collector.core.queue.QueueCreator;
+import org.skywalking.apm.collector.core.queue.QueueEventHandler;
+import org.skywalking.apm.collector.core.queue.QueueExecutor;
+
+/**
+ * @author pengys5
+ */
+public class DataCarrierQueueCreator implements QueueCreator {
+
+ @Override public QueueEventHandler create(int queueSize, QueueExecutor executor) {
+ return null;
+ }
+}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java
index 572557711..af5db748a 100644
--- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java
@@ -2,18 +2,20 @@ package org.skywalking.apm.collector.queue.datacarrier;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.framework.DefineException;
-import org.skywalking.apm.collector.core.module.ModuleGroup;
-import org.skywalking.apm.collector.core.queue.QueueModuleContext;
-import org.skywalking.apm.collector.core.queue.QueueModuleDefine;
+import org.skywalking.apm.collector.core.server.ServerHolder;
+import org.skywalking.apm.collector.queue.QueueModuleContext;
+import org.skywalking.apm.collector.queue.QueueModuleDefine;
+import org.skywalking.apm.collector.queue.QueueModuleGroupDefine;
/**
* @author pengys5
*/
public class QueueDataCarrierModuleDefine extends QueueModuleDefine {
- @Override protected ModuleGroup group() {
- return ModuleGroup.Queue;
+ @Override protected String group() {
+ return QueueModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
@@ -21,10 +23,11 @@ public class QueueDataCarrierModuleDefine extends QueueModuleDefine {
}
@Override public boolean defaultModule() {
- return true;
+ return false;
}
- @Override public final void initialize(Map config) throws DefineException, ClientException {
- QueueModuleContext.CREATOR = new DataCarrierCreator();
+ @Override
+ public final void initialize(Map config, ServerHolder serverHolder) throws DefineException, ClientException {
+ ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(group())).setQueueCreator(new DataCarrierQueueCreator());
}
}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorEventHandler.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorEventHandler.java
index 34090fa08..f25396531 100644
--- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorEventHandler.java
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorEventHandler.java
@@ -5,8 +5,7 @@ import com.lmax.disruptor.RingBuffer;
import org.skywalking.apm.collector.core.queue.EndOfBatchCommand;
import org.skywalking.apm.collector.core.queue.MessageHolder;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;
-import org.skywalking.apm.collector.core.worker.AbstractLocalAsyncWorker;
-import org.skywalking.apm.collector.core.worker.WorkerException;
+import org.skywalking.apm.collector.core.queue.QueueExecutor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -18,11 +17,11 @@ public class DisruptorEventHandler implements EventHandler, Queue
private final Logger logger = LoggerFactory.getLogger(DisruptorEventHandler.class);
private RingBuffer ringBuffer;
- private AbstractLocalAsyncWorker asyncWorker;
+ private QueueExecutor executor;
- DisruptorEventHandler(RingBuffer ringBuffer, AbstractLocalAsyncWorker asyncWorker) {
+ DisruptorEventHandler(RingBuffer ringBuffer, QueueExecutor executor) {
this.ringBuffer = ringBuffer;
- this.asyncWorker = asyncWorker;
+ this.executor = executor;
}
/**
@@ -34,16 +33,12 @@ public class DisruptorEventHandler implements EventHandler, Queue
* @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();
- event.reset();
+ Object message = event.getMessage();
+ event.reset();
- asyncWorker.allocateJob(message);
- if (endOfBatch) {
- asyncWorker.allocateJob(new EndOfBatchCommand());
- }
- } catch (WorkerException e) {
- logger.error(e.getMessage(), e);
+ executor.execute(message);
+ if (endOfBatch) {
+ executor.execute(new EndOfBatchCommand());
}
}
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorCreator.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorQueueCreator.java
similarity index 77%
rename from apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorCreator.java
rename to apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorQueueCreator.java
index 8be095b26..df027d0e6 100644
--- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorCreator.java
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorQueueCreator.java
@@ -2,18 +2,18 @@ package org.skywalking.apm.collector.queue.disruptor;
import com.lmax.disruptor.RingBuffer;
import com.lmax.disruptor.dsl.Disruptor;
-import org.skywalking.apm.collector.core.queue.Creator;
import org.skywalking.apm.collector.core.queue.DaemonThreadFactory;
import org.skywalking.apm.collector.core.queue.MessageHolder;
+import org.skywalking.apm.collector.core.queue.QueueCreator;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;
-import org.skywalking.apm.collector.core.worker.AbstractLocalAsyncWorker;
+import org.skywalking.apm.collector.core.queue.QueueExecutor;
/**
* @author pengys5
*/
-public class DisruptorCreator implements Creator {
+public class DisruptorQueueCreator implements QueueCreator {
- public QueueEventHandler create(int queueSize, AbstractLocalAsyncWorker localAsyncWorker) {
+ @Override public QueueEventHandler create(int queueSize, QueueExecutor executor) {
// Specify the size of the ring buffer, must be power of 2.
if (!((((queueSize - 1) & queueSize) == 0) && queueSize != 0)) {
throw new IllegalArgumentException("queue size must be power of 2");
@@ -23,7 +23,7 @@ public class DisruptorCreator implements Creator {
Disruptor disruptor = new Disruptor(MessageHolderFactory.INSTANCE, queueSize, DaemonThreadFactory.INSTANCE);
RingBuffer ringBuffer = disruptor.getRingBuffer();
- DisruptorEventHandler eventHandler = new DisruptorEventHandler(ringBuffer, localAsyncWorker);
+ DisruptorEventHandler eventHandler = new DisruptorEventHandler(ringBuffer, executor);
// Connect the handler
disruptor.handleEventsWith(eventHandler);
diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java
index cae0f1325..d74cb2dee 100644
--- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java
+++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java
@@ -2,18 +2,20 @@ package org.skywalking.apm.collector.queue.disruptor;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.framework.DefineException;
-import org.skywalking.apm.collector.core.module.ModuleGroup;
-import org.skywalking.apm.collector.core.queue.QueueModuleContext;
-import org.skywalking.apm.collector.core.queue.QueueModuleDefine;
+import org.skywalking.apm.collector.core.server.ServerHolder;
+import org.skywalking.apm.collector.queue.QueueModuleContext;
+import org.skywalking.apm.collector.queue.QueueModuleDefine;
+import org.skywalking.apm.collector.queue.QueueModuleGroupDefine;
/**
* @author pengys5
*/
public class QueueDisruptorModuleDefine extends QueueModuleDefine {
- @Override protected ModuleGroup group() {
- return ModuleGroup.Queue;
+ @Override protected String group() {
+ return QueueModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
@@ -24,7 +26,8 @@ public class QueueDisruptorModuleDefine extends QueueModuleDefine {
return true;
}
- @Override public final void initialize(Map config) throws DefineException, ClientException {
- QueueModuleContext.CREATOR = new DisruptorCreator();
+ @Override
+ public final void initialize(Map config, ServerHolder serverHolder) throws DefineException, ClientException {
+ ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(group())).setQueueCreator(new DisruptorQueueCreator());
}
}
diff --git a/apm-collector/apm-collector-queue/src/main/resources/META-INF/defines/group.define b/apm-collector/apm-collector-queue/src/main/resources/META-INF/defines/group.define
new file mode 100644
index 000000000..761c2b7e0
--- /dev/null
+++ b/apm-collector/apm-collector-queue/src/main/resources/META-INF/defines/group.define
@@ -0,0 +1 @@
+org.skywalking.apm.collector.queue.QueueModuleGroupDefine
\ No newline at end of file
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
deleted file mode 100644
index 30b74e893..000000000
--- a/apm-collector/apm-collector-remote/src/main/java/org/skywalking/apm/collector/remote/grpc/RemoteGRPCModuleDefine.java
+++ /dev/null
@@ -1,53 +0,0 @@
-package org.skywalking.apm.collector.remote.grpc;
-
-import java.util.Map;
-import org.skywalking.apm.collector.core.client.Client;
-import org.skywalking.apm.collector.core.client.ClientException;
-import org.skywalking.apm.collector.core.framework.DataInitializer;
-import org.skywalking.apm.collector.core.framework.DefineException;
-import org.skywalking.apm.collector.core.module.ModuleConfigParser;
-import org.skywalking.apm.collector.core.module.ModuleGroup;
-import org.skywalking.apm.collector.core.module.ModuleRegistration;
-import org.skywalking.apm.collector.core.remote.RemoteModuleDefine;
-import org.skywalking.apm.collector.core.server.Server;
-
-/**
- * @author pengys5
- */
-public class RemoteGRPCModuleDefine extends RemoteModuleDefine {
- @Override protected ModuleGroup group() {
- return ModuleGroup.Queue;
- }
-
- @Override public boolean defaultModule() {
- return false;
- }
-
- @Override protected ModuleConfigParser configParser() {
- return null;
- }
-
- @Override protected Client createClient() {
- return null;
- }
-
- @Override protected Server server() {
- return null;
- }
-
- @Override protected DataInitializer dataInitializer() {
- return null;
- }
-
- @Override protected ModuleRegistration registration() {
- return null;
- }
-
- @Override public void initialize(Map config) throws DefineException, ClientException {
-
- }
-
- @Override public String name() {
- return null;
- }
-}
diff --git a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto b/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto
deleted file mode 100644
index ae1da50f0..000000000
--- a/apm-collector/apm-collector-remote/src/main/proto/RemoteCommonService.proto
+++ /dev/null
@@ -1,18 +0,0 @@
-syntax = "proto3";
-
-option java_multiple_files = true;
-option java_package = "org.skywalking.apm.collector.remote.grpc.proto";
-
-service RemoteCommonService {
- rpc call (Message) returns (Empty) {
- }
-}
-
-message Message {
- string workerRole = 1;
- int32 objectId = 2;
- bytes objectBytes = 3; // the byte array of data object
-}
-
-message Empty {
-}
\ No newline at end of file
diff --git a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/grpc/GRPCHandler.java b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/grpc/GRPCHandler.java
new file mode 100644
index 000000000..7a6dd08bf
--- /dev/null
+++ b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/grpc/GRPCHandler.java
@@ -0,0 +1,9 @@
+package org.skywalking.apm.collector.server.grpc;
+
+import org.skywalking.apm.collector.core.framework.Handler;
+
+/**
+ * @author pengys5
+ */
+public interface GRPCHandler extends Handler {
+}
diff --git a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/grpc/GRPCServer.java b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/grpc/GRPCServer.java
index 309901ade..6f9472317 100644
--- a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/grpc/GRPCServer.java
+++ b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/grpc/GRPCServer.java
@@ -3,6 +3,7 @@ package org.skywalking.apm.collector.server.grpc;
import io.grpc.netty.NettyServerBuilder;
import java.io.IOException;
import java.net.InetSocketAddress;
+import org.skywalking.apm.collector.core.framework.Handler;
import org.skywalking.apm.collector.core.server.Server;
import org.skywalking.apm.collector.core.server.ServerException;
import org.slf4j.Logger;
@@ -17,27 +18,38 @@ public class GRPCServer implements Server {
private final String host;
private final int port;
+ private io.grpc.Server server;
+ private NettyServerBuilder nettyServerBuilder;
public GRPCServer(String host, int port) {
this.host = host;
this.port = port;
}
+ @Override public String hostPort() {
+ return host + ":" + port;
+ }
+
+ @Override public String serverClassify() {
+ return "Google-RPC";
+ }
+
@Override public void initialize() throws ServerException {
InetSocketAddress address = new InetSocketAddress(host, port);
- NettyServerBuilder nettyServerBuilder = NettyServerBuilder.forAddress(address);
- try {
- io.grpc.Server server = nettyServerBuilder.build().start();
- blockUntilShutdown(server);
- } catch (InterruptedException | IOException e) {
- throw new GRPCServerException(e.getMessage(), e);
- }
+ nettyServerBuilder = NettyServerBuilder.forAddress(address);
logger.info("Server started, host {} listening on {}", host, port);
}
- private void blockUntilShutdown(io.grpc.Server server) throws InterruptedException {
- if (server != null) {
- server.awaitTermination();
+ @Override public void start() throws ServerException {
+ try {
+ server = nettyServerBuilder.build();
+ server.start();
+ } catch (IOException e) {
+ throw new GRPCServerException(e.getMessage(), e);
}
}
+
+ @Override public void addHandler(Handler handler) {
+ nettyServerBuilder.addService((io.grpc.BindableService)handler);
+ }
}
diff --git a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/ArgumentsParseException.java b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/ArgumentsParseException.java
new file mode 100644
index 000000000..e7148b58d
--- /dev/null
+++ b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/ArgumentsParseException.java
@@ -0,0 +1,17 @@
+package org.skywalking.apm.collector.server.jetty;
+
+import org.skywalking.apm.collector.core.CollectorException;
+
+/**
+ * @author pengys5
+ */
+public class ArgumentsParseException extends CollectorException {
+
+ public ArgumentsParseException(String message) {
+ super(message);
+ }
+
+ public ArgumentsParseException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java
new file mode 100644
index 000000000..7b7c39bc5
--- /dev/null
+++ b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyHandler.java
@@ -0,0 +1,147 @@
+package org.skywalking.apm.collector.server.jetty;
+
+import com.google.gson.JsonElement;
+import java.io.IOException;
+import java.io.PrintWriter;
+import java.util.Enumeration;
+import javax.servlet.ServletConfig;
+import javax.servlet.ServletContext;
+import javax.servlet.ServletException;
+import javax.servlet.ServletRequest;
+import javax.servlet.ServletResponse;
+import javax.servlet.http.HttpServlet;
+import javax.servlet.http.HttpServletRequest;
+import javax.servlet.http.HttpServletResponse;
+import org.skywalking.apm.collector.core.framework.Handler;
+
+/**
+ * @author pengys5
+ */
+public abstract class JettyHandler extends HttpServlet implements Handler {
+
+ public abstract String pathSpec();
+
+ @Override
+ protected final void doGet(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
+ try {
+ reply(resp, doGet(req));
+ } catch (ArgumentsParseException e) {
+ replyError(resp, e.getMessage(), HttpServletResponse.SC_BAD_REQUEST);
+ }
+ }
+
+ protected abstract JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException;
+
+ @Override
+ protected final void doPost(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
+ super.doPost(req, resp);
+ }
+
+ @Override
+ protected final void doHead(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
+ super.doHead(req, resp);
+ }
+
+ @Override protected final long getLastModified(HttpServletRequest req) {
+ return super.getLastModified(req);
+ }
+
+ @Override
+ protected final void doPut(HttpServletRequest req, HttpServletResponse resp) throws ServletException, IOException {
+ super.doPut(req, resp);
+ }
+
+ @Override
+ protected final void doDelete(HttpServletRequest req,
+ HttpServletResponse resp) throws ServletException, IOException {
+ super.doDelete(req, resp);
+ }
+
+ @Override
+ protected final void doOptions(HttpServletRequest req,
+ HttpServletResponse resp) throws ServletException, IOException {
+ super.doOptions(req, resp);
+ }
+
+ @Override
+ protected final void doTrace(HttpServletRequest req,
+ HttpServletResponse resp) throws ServletException, IOException {
+ super.doTrace(req, resp);
+ }
+
+ @Override
+ protected final void service(HttpServletRequest req,
+ HttpServletResponse resp) throws ServletException, IOException {
+ super.service(req, resp);
+ }
+
+ @Override public final void service(ServletRequest req, ServletResponse res) throws ServletException, IOException {
+ super.service(req, res);
+ }
+
+ @Override public final void destroy() {
+ super.destroy();
+ }
+
+ @Override public final String getInitParameter(String name) {
+ return super.getInitParameter(name);
+ }
+
+ @Override public final Enumeration getInitParameterNames() {
+ return super.getInitParameterNames();
+ }
+
+ @Override public final ServletConfig getServletConfig() {
+ return super.getServletConfig();
+ }
+
+ @Override public final ServletContext getServletContext() {
+ return super.getServletContext();
+ }
+
+ @Override public final String getServletInfo() {
+ return super.getServletInfo();
+ }
+
+ @Override public final void init(ServletConfig config) throws ServletException {
+ super.init(config);
+ }
+
+ @Override public final void init() throws ServletException {
+ super.init();
+ }
+
+ @Override public final void log(String msg) {
+ super.log(msg);
+ }
+
+ @Override public final void log(String message, Throwable t) {
+ super.log(message, t);
+ }
+
+ @Override public final String getServletName() {
+ return super.getServletName();
+ }
+
+ private void reply(HttpServletResponse response, JsonElement resJson) throws IOException {
+ response.setContentType("text/json");
+ response.setCharacterEncoding("utf-8");
+ response.setStatus(HttpServletResponse.SC_OK);
+
+ PrintWriter out = response.getWriter();
+ out.print(resJson);
+ out.flush();
+ out.close();
+ }
+
+ private void replyError(HttpServletResponse response, String errorMessage, int status) throws IOException {
+ response.setContentType("text/plain");
+ response.setCharacterEncoding("utf-8");
+ response.setStatus(status);
+
+ PrintWriter out = response.getWriter();
+ out.print(errorMessage);
+ out.flush();
+ out.close();
+ }
+}
diff --git a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java
index 984e76269..441a3b16a 100644
--- a/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java
+++ b/apm-collector/apm-collector-server/src/main/java/org/skywalking/apm/collector/server/jetty/JettyServer.java
@@ -1,7 +1,10 @@
package org.skywalking.apm.collector.server.jetty;
import java.net.InetSocketAddress;
+import javax.servlet.http.HttpServlet;
import org.eclipse.jetty.servlet.ServletContextHandler;
+import org.eclipse.jetty.servlet.ServletHolder;
+import org.skywalking.apm.collector.core.framework.Handler;
import org.skywalking.apm.collector.core.server.Server;
import org.skywalking.apm.collector.core.server.ServerException;
import org.slf4j.Logger;
@@ -17,6 +20,8 @@ public class JettyServer implements Server {
private final String host;
private final int port;
private final String contextPath;
+ private org.eclipse.jetty.server.Server server;
+ private ServletContextHandler servletContextHandler;
public JettyServer(String host, int port, String contextPath) {
this.host = host;
@@ -24,14 +29,31 @@ public class JettyServer implements Server {
this.contextPath = contextPath;
}
- @Override public void initialize() throws ServerException {
- org.eclipse.jetty.server.Server server = new org.eclipse.jetty.server.Server(new InetSocketAddress(host, port));
+ @Override public String hostPort() {
+ return host + ":" + port;
+ }
- ServletContextHandler servletContextHandler = new ServletContextHandler(ServletContextHandler.NO_SESSIONS);
+ @Override public String serverClassify() {
+ return "Jetty";
+ }
+
+ @Override public void initialize() throws ServerException {
+ server = new org.eclipse.jetty.server.Server(new InetSocketAddress(host, port));
+
+ servletContextHandler = new ServletContextHandler(ServletContextHandler.NO_SESSIONS);
servletContextHandler.setContextPath(contextPath);
logger.info("http server root context path: {}", contextPath);
server.setHandler(servletContextHandler);
+ }
+
+ @Override public void addHandler(Handler handler) {
+ ServletHolder servletHolder = new ServletHolder();
+ servletHolder.setServlet((HttpServlet)handler);
+ servletContextHandler.addServlet(servletHolder, ((JettyHandler)handler).pathSpec());
+ }
+
+ @Override public void start() throws ServerException {
try {
server.start();
} catch (Exception e) {
diff --git a/apm-collector/apm-collector-storage/elasticsearch-storage/pom.xml b/apm-collector/apm-collector-storage/elasticsearch-storage/pom.xml
deleted file mode 100644
index 94b2517c0..000000000
--- a/apm-collector/apm-collector-storage/elasticsearch-storage/pom.xml
+++ /dev/null
@@ -1,14 +0,0 @@
-
-
-
- apm-collector-storage
- org.skywalking
- 3.2-2017
-
- 4.0.0
-
- elasticsearch-storage
- jar
-
\ No newline at end of file
diff --git a/apm-collector/apm-collector-storage/h2-storage/pom.xml b/apm-collector/apm-collector-storage/h2-storage/pom.xml
deleted file mode 100644
index 8bc63b3f0..000000000
--- a/apm-collector/apm-collector-storage/h2-storage/pom.xml
+++ /dev/null
@@ -1,14 +0,0 @@
-
-
-
- apm-collector-storage
- org.skywalking
- 3.2-2017
-
- 4.0.0
-
- h2-storage
- jar
-
\ No newline at end of file
diff --git a/apm-collector/apm-collector-storage/pom.xml b/apm-collector/apm-collector-storage/pom.xml
index 6ade3218b..144648e27 100644
--- a/apm-collector/apm-collector-storage/pom.xml
+++ b/apm-collector/apm-collector-storage/pom.xml
@@ -10,9 +10,18 @@
4.0.0
apm-collector-storage
- pom
-
- h2-storage
- elasticsearch-storage
-
+ jar
+
+
+
+ org.skywalking
+ apm-collector-core
+ ${project.version}
+
+
+ org.skywalking
+ apm-collector-cluster
+ ${project.version}
+
+
\ No newline at end of file
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleContext.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleContext.java
new file mode 100644
index 000000000..128a2fdf8
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleContext.java
@@ -0,0 +1,24 @@
+package org.skywalking.apm.collector.storage;
+
+import org.skywalking.apm.collector.core.client.Client;
+import org.skywalking.apm.collector.core.framework.Context;
+
+/**
+ * @author pengys5
+ */
+public class StorageModuleContext extends Context {
+
+ private Client client;
+
+ public StorageModuleContext(String groupName) {
+ super(groupName);
+ }
+
+ public Client getClient() {
+ return client;
+ }
+
+ public void setClient(Client client) {
+ this.client = client;
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleDefine.java
new file mode 100644
index 000000000..57ccb77e9
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleDefine.java
@@ -0,0 +1,63 @@
+package org.skywalking.apm.collector.storage;
+
+import java.util.Map;
+import org.skywalking.apm.collector.core.client.Client;
+import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.cluster.ClusterDataListener;
+import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine;
+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.module.ModuleDefine;
+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.core.storage.StorageException;
+import org.skywalking.apm.collector.core.storage.StorageInstaller;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public abstract class StorageModuleDefine extends ModuleDefine implements ClusterDataListenerDefine {
+
+ private final Logger logger = LoggerFactory.getLogger(StorageModuleDefine.class);
+
+ @Override
+ public final void initialize(Map config, ServerHolder serverHolder) throws DefineException, ClientException {
+ try {
+ configParser().parse(config);
+
+ StorageModuleContext context = (StorageModuleContext)CollectorContextHelper.INSTANCE.getContext(StorageModuleGroupDefine.GROUP_NAME);
+ Client client = createClient(null);
+ client.initialize();
+ context.setClient(client);
+ injectClientIntoDAO(client);
+
+ storageInstaller().install(client);
+ } catch (ConfigParseException | StorageException e) {
+ throw new StorageModuleException(e.getMessage(), e);
+ }
+ }
+
+ @Override protected final Server server() {
+ throw new UnsupportedOperationException("");
+ }
+
+ @Override protected final ModuleRegistration registration() {
+ throw new UnsupportedOperationException("");
+ }
+
+ @Override public final ClusterDataListener listener() {
+ throw new UnsupportedOperationException("");
+ }
+
+ @Override public final boolean defaultModule() {
+ return true;
+ }
+
+ public abstract StorageInstaller storageInstaller();
+
+ public abstract void injectClientIntoDAO(Client client) throws DefineException;
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleException.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleException.java
new file mode 100644
index 000000000..4136c9fa4
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleException.java
@@ -0,0 +1,16 @@
+package org.skywalking.apm.collector.storage;
+
+import org.skywalking.apm.collector.core.module.ModuleException;
+
+/**
+ * @author pengys5
+ */
+public class StorageModuleException extends ModuleException {
+ public StorageModuleException(String message) {
+ super(message);
+ }
+
+ public StorageModuleException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleGroupDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleGroupDefine.java
new file mode 100644
index 000000000..c373f3a83
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleGroupDefine.java
@@ -0,0 +1,25 @@
+package org.skywalking.apm.collector.storage;
+
+import org.skywalking.apm.collector.core.framework.Context;
+import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
+import org.skywalking.apm.collector.core.module.ModuleInstaller;
+
+/**
+ * @author pengys5
+ */
+public class StorageModuleGroupDefine implements ModuleGroupDefine {
+
+ public static final String GROUP_NAME = "storage";
+
+ @Override public String name() {
+ return GROUP_NAME;
+ }
+
+ @Override public Context groupContext() {
+ return new StorageModuleContext(GROUP_NAME);
+ }
+
+ @Override public ModuleInstaller moduleInstaller() {
+ return new StorageModuleInstaller();
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleInstaller.java
new file mode 100644
index 000000000..f15b9861c
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/StorageModuleInstaller.java
@@ -0,0 +1,29 @@
+package org.skywalking.apm.collector.storage;
+
+import java.util.Map;
+import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
+import org.skywalking.apm.collector.core.framework.DefineException;
+import org.skywalking.apm.collector.core.module.ModuleDefine;
+import org.skywalking.apm.collector.core.module.SingleModuleInstaller;
+import org.skywalking.apm.collector.core.server.ServerHolder;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class StorageModuleInstaller extends SingleModuleInstaller {
+
+ private final Logger logger = LoggerFactory.getLogger(StorageModuleInstaller.class);
+
+ @Override public void install(Map moduleConfig,
+ Map moduleDefineMap, ServerHolder serverHolder) throws DefineException, ClientException {
+ logger.info("beginning storage module install");
+
+ StorageModuleContext context = new StorageModuleContext(StorageModuleGroupDefine.GROUP_NAME);
+ CollectorContextHelper.INSTANCE.putContext(context);
+
+ installSingle(moduleConfig, moduleDefineMap, serverHolder);
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/DAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/DAO.java
new file mode 100644
index 000000000..1ddeb9445
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/DAO.java
@@ -0,0 +1,18 @@
+package org.skywalking.apm.collector.storage.dao;
+
+import org.skywalking.apm.collector.core.client.Client;
+
+/**
+ * @author pengys5
+ */
+public abstract class DAO {
+ private C client;
+
+ public final C getClient() {
+ return client;
+ }
+
+ public final void setClient(C client) {
+ this.client = client;
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/DAOContainer.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/DAOContainer.java
new file mode 100644
index 000000000..a4f884239
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/DAOContainer.java
@@ -0,0 +1,21 @@
+package org.skywalking.apm.collector.storage.dao;
+
+import java.util.HashMap;
+import java.util.Map;
+
+/**
+ * @author pengys5
+ */
+public enum DAOContainer {
+ INSTANCE;
+
+ private Map daos = new HashMap<>();
+
+ public void put(String interfaceName, DAO dao) {
+ daos.put(interfaceName, dao);
+ }
+
+ public DAO get(String interfaceName) {
+ return daos.get(interfaceName);
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/IBatchDAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/IBatchDAO.java
new file mode 100644
index 000000000..afcac1bf1
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/dao/IBatchDAO.java
@@ -0,0 +1,10 @@
+package org.skywalking.apm.collector.storage.dao;
+
+import java.util.List;
+
+/**
+ * @author pengys5
+ */
+public interface IBatchDAO {
+ void batchPersistence(List> batchCollection);
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/ElasticSearchStorageException.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/ElasticSearchStorageException.java
new file mode 100644
index 000000000..be95df896
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/ElasticSearchStorageException.java
@@ -0,0 +1,16 @@
+package org.skywalking.apm.collector.storage.elasticsearch;
+
+import org.skywalking.apm.collector.core.storage.StorageException;
+
+/**
+ * @author pengys5
+ */
+public class ElasticSearchStorageException extends StorageException {
+ public ElasticSearchStorageException(String message) {
+ super(message);
+ }
+
+ public ElasticSearchStorageException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfig.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfig.java
new file mode 100644
index 000000000..3526fc98b
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfig.java
@@ -0,0 +1,10 @@
+package org.skywalking.apm.collector.storage.elasticsearch;
+
+/**
+ * @author pengys5
+ */
+public class StorageElasticSearchConfig {
+ public static String CLUSTER_NAME;
+ public static Boolean CLUSTER_TRANSPORT_SNIFFER;
+ public static String CLUSTER_NODES;
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfigParser.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfigParser.java
new file mode 100644
index 000000000..6653fdc22
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfigParser.java
@@ -0,0 +1,29 @@
+package org.skywalking.apm.collector.storage.elasticsearch;
+
+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 StorageElasticSearchConfigParser implements ModuleConfigParser {
+
+ private static final String CLUSTER_NAME = "cluster_name";
+ private static final String CLUSTER_TRANSPORT_SNIFFER = "cluster_transport_sniffer";
+ private static final String CLUSTER_NODES = "cluster_nodes";
+
+ @Override public void parse(Map config) throws ConfigParseException {
+ if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(CLUSTER_NAME))) {
+ StorageElasticSearchConfig.CLUSTER_NAME = (String)config.get(CLUSTER_NAME);
+ }
+ if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(CLUSTER_TRANSPORT_SNIFFER))) {
+ StorageElasticSearchConfig.CLUSTER_TRANSPORT_SNIFFER = (Boolean)config.get(CLUSTER_TRANSPORT_SNIFFER);
+ }
+ if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(CLUSTER_NODES))) {
+ StorageElasticSearchConfig.CLUSTER_NODES = (String)config.get(CLUSTER_NODES);
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchModuleDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchModuleDefine.java
new file mode 100644
index 000000000..826ad4cb0
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchModuleDefine.java
@@ -0,0 +1,53 @@
+package org.skywalking.apm.collector.storage.elasticsearch;
+
+import java.util.List;
+import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
+import org.skywalking.apm.collector.core.client.Client;
+import org.skywalking.apm.collector.core.client.DataMonitor;
+import org.skywalking.apm.collector.core.framework.DefineException;
+import org.skywalking.apm.collector.core.module.ModuleConfigParser;
+import org.skywalking.apm.collector.core.storage.StorageInstaller;
+import org.skywalking.apm.collector.storage.StorageModuleDefine;
+import org.skywalking.apm.collector.storage.StorageModuleGroupDefine;
+import org.skywalking.apm.collector.storage.dao.DAOContainer;
+import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO;
+import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAODefineLoader;
+import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchStorageInstaller;
+
+/**
+ * @author pengys5
+ */
+public class StorageElasticSearchModuleDefine extends StorageModuleDefine {
+
+ public static final String MODULE_NAME = "elasticsearch";
+
+ @Override protected String group() {
+ return StorageModuleGroupDefine.GROUP_NAME;
+ }
+
+ @Override public String name() {
+ return MODULE_NAME;
+ }
+
+ @Override protected ModuleConfigParser configParser() {
+ return new StorageElasticSearchConfigParser();
+ }
+
+ @Override protected Client createClient(DataMonitor dataMonitor) {
+ return new ElasticSearchClient(StorageElasticSearchConfig.CLUSTER_NAME, StorageElasticSearchConfig.CLUSTER_TRANSPORT_SNIFFER, StorageElasticSearchConfig.CLUSTER_NODES);
+ }
+
+ @Override public StorageInstaller storageInstaller() {
+ return new ElasticSearchStorageInstaller();
+ }
+
+ @Override public void injectClientIntoDAO(Client client) throws DefineException {
+ EsDAODefineLoader loader = new EsDAODefineLoader();
+ List esDAOs = loader.load();
+ esDAOs.forEach(esDAO -> {
+ esDAO.setClient((ElasticSearchClient)client);
+ String interFaceName = esDAO.getClass().getInterfaces()[0].getName();
+ DAOContainer.INSTANCE.put(interFaceName, esDAO);
+ });
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java
new file mode 100644
index 000000000..74c06e634
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/BatchEsDAO.java
@@ -0,0 +1,35 @@
+package org.skywalking.apm.collector.storage.elasticsearch.dao;
+
+import java.util.List;
+import org.elasticsearch.action.bulk.BulkRequestBuilder;
+import org.elasticsearch.action.bulk.BulkResponse;
+import org.elasticsearch.action.index.IndexRequestBuilder;
+import org.skywalking.apm.collector.core.util.CollectionUtils;
+import org.skywalking.apm.collector.storage.dao.IBatchDAO;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class BatchEsDAO extends EsDAO implements IBatchDAO {
+
+ private final Logger logger = LoggerFactory.getLogger(BatchEsDAO.class);
+
+ @Override public void batchPersistence(List> batchCollection) {
+ BulkRequestBuilder bulkRequest = getClient().prepareBulk();
+
+ logger.info("bulk data size: {}", batchCollection.size());
+ if (CollectionUtils.isNotEmpty(batchCollection)) {
+ for (int i = 0; i < batchCollection.size(); i++) {
+ IndexRequestBuilder builder = (IndexRequestBuilder)batchCollection.get(i);
+ bulkRequest.add(builder);
+ }
+
+ BulkResponse bulkResponse = bulkRequest.execute().actionGet();
+ if (bulkResponse.hasFailures()) {
+ logger.error(bulkResponse.buildFailureMessage());
+ }
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAO.java
new file mode 100644
index 000000000..f6eb18a85
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAO.java
@@ -0,0 +1,55 @@
+package org.skywalking.apm.collector.storage.elasticsearch.dao;
+
+import org.elasticsearch.action.search.SearchRequestBuilder;
+import org.elasticsearch.action.search.SearchResponse;
+import org.elasticsearch.search.aggregations.AggregationBuilders;
+import org.elasticsearch.search.aggregations.metrics.max.Max;
+import org.elasticsearch.search.aggregations.metrics.max.MaxAggregationBuilder;
+import org.elasticsearch.search.aggregations.metrics.min.Min;
+import org.elasticsearch.search.aggregations.metrics.min.MinAggregationBuilder;
+import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
+import org.skywalking.apm.collector.storage.dao.DAO;
+
+/**
+ * @author pengys5
+ */
+public abstract class EsDAO extends DAO {
+
+ public final int getMaxId(String indexName, String columnName) {
+ ElasticSearchClient client = getClient();
+ SearchRequestBuilder searchRequestBuilder = client.prepareSearch(indexName);
+ searchRequestBuilder.setTypes("type");
+ searchRequestBuilder.setSize(0);
+ MaxAggregationBuilder aggregation = AggregationBuilders.max("agg").field(columnName);
+ searchRequestBuilder.addAggregation(aggregation);
+
+ SearchResponse searchResponse = searchRequestBuilder.execute().actionGet();
+ Max agg = searchResponse.getAggregations().get("agg");
+
+ int id = (int)agg.getValue();
+ if (id == Integer.MAX_VALUE || id == Integer.MIN_VALUE) {
+ return 0;
+ } else {
+ return id;
+ }
+ }
+
+ public final int getMinId(String indexName, String columnName) {
+ ElasticSearchClient client = getClient();
+ SearchRequestBuilder searchRequestBuilder = client.prepareSearch(indexName);
+ searchRequestBuilder.setTypes("type");
+ searchRequestBuilder.setSize(0);
+ MinAggregationBuilder aggregation = AggregationBuilders.min("agg").field(columnName);
+ searchRequestBuilder.addAggregation(aggregation);
+
+ SearchResponse searchResponse = searchRequestBuilder.execute().actionGet();
+ Min agg = searchResponse.getAggregations().get("agg");
+
+ int id = (int)agg.getValue();
+ if (id == Integer.MAX_VALUE || id == Integer.MIN_VALUE) {
+ return 0;
+ } else {
+ return id;
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAODefineLoader.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAODefineLoader.java
new file mode 100644
index 000000000..de5213bff
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAODefineLoader.java
@@ -0,0 +1,30 @@
+package org.skywalking.apm.collector.storage.elasticsearch.dao;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.skywalking.apm.collector.core.framework.DefineException;
+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 EsDAODefineLoader implements Loader> {
+
+ private final Logger logger = LoggerFactory.getLogger(EsDAODefineLoader.class);
+
+ @Override public List load() throws DefineException {
+ List esDAOs = new ArrayList<>();
+
+ EsDAODefinitionFile definitionFile = new EsDAODefinitionFile();
+ logger.info("elasticsearch dao definition file name: {}", definitionFile.fileName());
+ DefinitionLoader definitionLoader = DefinitionLoader.load(EsDAO.class, definitionFile);
+ for (EsDAO dao : definitionLoader) {
+ logger.info("loaded elasticsearch dao definition class: {}", dao.getClass().getName());
+ esDAOs.add(dao);
+ }
+ return esDAOs;
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAODefinitionFile.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAODefinitionFile.java
new file mode 100644
index 000000000..08791d866
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/dao/EsDAODefinitionFile.java
@@ -0,0 +1,13 @@
+package org.skywalking.apm.collector.storage.elasticsearch.dao;
+
+import org.skywalking.apm.collector.core.framework.DefinitionFile;
+
+/**
+ * @author pengys5
+ */
+public class EsDAODefinitionFile extends DefinitionFile {
+
+ @Override protected String fileName() {
+ return "es_dao.define";
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchColumnDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchColumnDefine.java
new file mode 100644
index 000000000..883480a44
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchColumnDefine.java
@@ -0,0 +1,16 @@
+package org.skywalking.apm.collector.storage.elasticsearch.define;
+
+import org.skywalking.apm.collector.core.storage.ColumnDefine;
+
+/**
+ * @author pengys5
+ */
+public class ElasticSearchColumnDefine extends ColumnDefine {
+ public ElasticSearchColumnDefine(String name, String type) {
+ super(name, type);
+ }
+
+ public enum Type {
+ Binary, Boolean, Date, Keyword, Long, Integer
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java
new file mode 100644
index 000000000..25dd33d61
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java
@@ -0,0 +1,98 @@
+package org.skywalking.apm.collector.storage.elasticsearch.define;
+
+import java.io.IOException;
+import java.util.List;
+import org.elasticsearch.common.settings.Settings;
+import org.elasticsearch.common.xcontent.XContentBuilder;
+import org.elasticsearch.common.xcontent.XContentFactory;
+import org.elasticsearch.index.IndexNotFoundException;
+import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient;
+import org.skywalking.apm.collector.core.client.Client;
+import org.skywalking.apm.collector.core.storage.ColumnDefine;
+import org.skywalking.apm.collector.core.storage.StorageInstaller;
+import org.skywalking.apm.collector.core.storage.TableDefine;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class ElasticSearchStorageInstaller extends StorageInstaller {
+
+ private final Logger logger = LoggerFactory.getLogger(ElasticSearchStorageInstaller.class);
+
+ @Override protected void defineFilter(List tableDefines) {
+ int size = tableDefines.size();
+ for (int i = size - 1; i >= 0; i--) {
+ if (!(tableDefines.get(i) instanceof ElasticSearchTableDefine)) {
+ tableDefines.remove(i);
+ }
+ }
+ }
+
+ @Override protected boolean createTable(Client client, TableDefine tableDefine) {
+ ElasticSearchClient esClient = (ElasticSearchClient)client;
+ ElasticSearchTableDefine esTableDefine = (ElasticSearchTableDefine)tableDefine;
+ // settings
+ String settingSource = "";
+ // mapping
+ XContentBuilder mappingBuilder = null;
+ try {
+ XContentBuilder settingsBuilder = createSettingBuilder(esTableDefine);
+ settingSource = settingsBuilder.string();
+ mappingBuilder = createMappingBuilder(esTableDefine);
+ logger.info("mapping builder str: {}", mappingBuilder.string());
+ } catch (Exception e) {
+ logger.error("create {} index mapping builder error", esTableDefine.getName());
+ }
+ Settings settings = Settings.builder().loadFromSource(settingSource).build();
+
+ boolean isAcknowledged = esClient.createIndex(esTableDefine.getName(), esTableDefine.type(), settings, mappingBuilder);
+ logger.info("create {} index with type of {} finished, isAcknowledged: {}", esTableDefine.getName(), esTableDefine.type(), isAcknowledged);
+ return isAcknowledged;
+ }
+
+ private XContentBuilder createSettingBuilder(ElasticSearchTableDefine tableDefine) throws IOException {
+ return XContentFactory.jsonBuilder()
+ .startObject()
+ .field("index.number_of_shards", tableDefine.numberOfShards())
+ .field("index.number_of_replicas", tableDefine.numberOfReplicas())
+ .field("index.refresh_interval", String.valueOf(tableDefine.refreshInterval()) + "s")
+ .endObject();
+ }
+
+ private XContentBuilder createMappingBuilder(ElasticSearchTableDefine tableDefine) throws IOException {
+ XContentBuilder mappingBuilder = XContentFactory.jsonBuilder()
+ .startObject()
+ .startObject("properties");
+
+ for (ColumnDefine columnDefine : tableDefine.getColumnDefines()) {
+ ElasticSearchColumnDefine elasticSearchColumnDefine = (ElasticSearchColumnDefine)columnDefine;
+ mappingBuilder
+ .startObject(elasticSearchColumnDefine.getName())
+ .field("type", elasticSearchColumnDefine.getType().toLowerCase())
+ .endObject();
+ }
+
+ mappingBuilder
+ .endObject()
+ .endObject();
+ logger.debug("create elasticsearch index: {}", mappingBuilder.string());
+ return mappingBuilder;
+ }
+
+ @Override protected boolean deleteIndex(Client client, TableDefine tableDefine) {
+ ElasticSearchClient esClient = (ElasticSearchClient)client;
+ try {
+ return esClient.deleteIndex(tableDefine.getName());
+ } catch (IndexNotFoundException e) {
+ logger.info("{} index not found", tableDefine.getName());
+ }
+ return false;
+ }
+
+ @Override protected boolean isExists(Client client, TableDefine tableDefine) {
+ ElasticSearchClient esClient = (ElasticSearchClient)client;
+ return esClient.isExistsIndex(tableDefine.getName());
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchTableDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchTableDefine.java
new file mode 100644
index 000000000..c9f7241f4
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchTableDefine.java
@@ -0,0 +1,23 @@
+package org.skywalking.apm.collector.storage.elasticsearch.define;
+
+import org.skywalking.apm.collector.core.storage.TableDefine;
+
+/**
+ * @author pengys5
+ */
+public abstract class ElasticSearchTableDefine extends TableDefine {
+
+ public ElasticSearchTableDefine(String name) {
+ super(name);
+ }
+
+ public final String type() {
+ return "type";
+ }
+
+ public abstract int refreshInterval();
+
+ public abstract int numberOfShards();
+
+ public abstract int numberOfReplicas();
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/StorageH2ConfigParser.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/StorageH2ConfigParser.java
new file mode 100644
index 000000000..72f7b561f
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/StorageH2ConfigParser.java
@@ -0,0 +1,14 @@
+package org.skywalking.apm.collector.storage.h2;
+
+import java.util.Map;
+import org.skywalking.apm.collector.core.config.ConfigParseException;
+import org.skywalking.apm.collector.core.module.ModuleConfigParser;
+
+/**
+ * @author pengys5
+ */
+public class StorageH2ConfigParser implements ModuleConfigParser {
+
+ @Override public void parse(Map config) throws ConfigParseException {
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/StorageH2ModuleDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/StorageH2ModuleDefine.java
new file mode 100644
index 000000000..84f9b31a7
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/StorageH2ModuleDefine.java
@@ -0,0 +1,53 @@
+package org.skywalking.apm.collector.storage.h2;
+
+import java.util.List;
+import org.skywalking.apm.collector.client.h2.H2Client;
+import org.skywalking.apm.collector.core.client.Client;
+import org.skywalking.apm.collector.core.client.DataMonitor;
+import org.skywalking.apm.collector.core.framework.DefineException;
+import org.skywalking.apm.collector.core.module.ModuleConfigParser;
+import org.skywalking.apm.collector.core.storage.StorageInstaller;
+import org.skywalking.apm.collector.storage.StorageModuleDefine;
+import org.skywalking.apm.collector.storage.StorageModuleGroupDefine;
+import org.skywalking.apm.collector.storage.dao.DAOContainer;
+import org.skywalking.apm.collector.storage.h2.dao.H2DAO;
+import org.skywalking.apm.collector.storage.h2.dao.H2DAODefineLoader;
+import org.skywalking.apm.collector.storage.h2.define.H2StorageInstaller;
+
+/**
+ * @author pengys5
+ */
+public class StorageH2ModuleDefine extends StorageModuleDefine {
+
+ public static final String MODULE_NAME = "h2";
+
+ @Override protected String group() {
+ return StorageModuleGroupDefine.GROUP_NAME;
+ }
+
+ @Override public String name() {
+ return MODULE_NAME;
+ }
+
+ @Override protected ModuleConfigParser configParser() {
+ return new StorageH2ConfigParser();
+ }
+
+ @Override protected Client createClient(DataMonitor dataMonitor) {
+ return new H2Client();
+ }
+
+ @Override public StorageInstaller storageInstaller() {
+ return new H2StorageInstaller();
+ }
+
+ @Override public void injectClientIntoDAO(Client client) throws DefineException {
+ H2DAODefineLoader loader = new H2DAODefineLoader();
+ List h2DAOs = loader.load();
+ h2DAOs.forEach(h2DAO -> {
+ h2DAO.setClient((H2Client)client);
+ String interFaceName = h2DAO.getClass().getInterfaces()[0].getName();
+ DAOContainer.INSTANCE.put(interFaceName, h2DAO);
+ });
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java
new file mode 100644
index 000000000..27b69f6d4
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/BatchH2DAO.java
@@ -0,0 +1,14 @@
+package org.skywalking.apm.collector.storage.h2.dao;
+
+import java.util.List;
+import org.skywalking.apm.collector.storage.dao.IBatchDAO;
+
+/**
+ * @author pengys5
+ */
+public class BatchH2DAO extends H2DAO implements IBatchDAO {
+
+ @Override public void batchPersistence(List> batchCollection) {
+
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAO.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAO.java
new file mode 100644
index 000000000..a9c9ddb5a
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAO.java
@@ -0,0 +1,10 @@
+package org.skywalking.apm.collector.storage.h2.dao;
+
+import org.skywalking.apm.collector.client.h2.H2Client;
+import org.skywalking.apm.collector.storage.dao.DAO;
+
+/**
+ * @author pengys5
+ */
+public abstract class H2DAO extends DAO {
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAODefineLoader.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAODefineLoader.java
new file mode 100644
index 000000000..733eb3bdd
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAODefineLoader.java
@@ -0,0 +1,30 @@
+package org.skywalking.apm.collector.storage.h2.dao;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.skywalking.apm.collector.core.framework.DefineException;
+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 H2DAODefineLoader implements Loader> {
+
+ private final Logger logger = LoggerFactory.getLogger(H2DAODefineLoader.class);
+
+ @Override public List load() throws DefineException {
+ List h2DAOs = new ArrayList<>();
+
+ H2DAODefinitionFile definitionFile = new H2DAODefinitionFile();
+ logger.info("h2 dao definition file name: {}", definitionFile.fileName());
+ DefinitionLoader definitionLoader = DefinitionLoader.load(H2DAO.class, definitionFile);
+ for (H2DAO dao : definitionLoader) {
+ logger.info("loaded h2 dao definition class: {}", dao.getClass().getName());
+ h2DAOs.add(dao);
+ }
+ return h2DAOs;
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAODefinitionFile.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAODefinitionFile.java
new file mode 100644
index 000000000..72d37883f
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/dao/H2DAODefinitionFile.java
@@ -0,0 +1,13 @@
+package org.skywalking.apm.collector.storage.h2.dao;
+
+import org.skywalking.apm.collector.core.framework.DefinitionFile;
+
+/**
+ * @author pengys5
+ */
+public class H2DAODefinitionFile extends DefinitionFile {
+
+ @Override protected String fileName() {
+ return "h2_dao.define";
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2ColumnDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2ColumnDefine.java
new file mode 100644
index 000000000..ca7cf49f1
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2ColumnDefine.java
@@ -0,0 +1,17 @@
+package org.skywalking.apm.collector.storage.h2.define;
+
+import org.skywalking.apm.collector.core.storage.ColumnDefine;
+
+/**
+ * @author pengys5
+ */
+public class H2ColumnDefine extends ColumnDefine {
+
+ public H2ColumnDefine(String name, String type) {
+ super(name, type);
+ }
+
+ public enum Type {
+ Boolean, Varchar, Int, Bigint, BINARY
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java
new file mode 100644
index 000000000..d27a8ca74
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2StorageInstaller.java
@@ -0,0 +1,58 @@
+package org.skywalking.apm.collector.storage.h2.define;
+
+import java.util.List;
+import org.skywalking.apm.collector.client.h2.H2Client;
+import org.skywalking.apm.collector.client.h2.H2ClientException;
+import org.skywalking.apm.collector.core.client.Client;
+import org.skywalking.apm.collector.core.storage.StorageException;
+import org.skywalking.apm.collector.core.storage.StorageInstallException;
+import org.skywalking.apm.collector.core.storage.StorageInstaller;
+import org.skywalking.apm.collector.core.storage.TableDefine;
+
+/**
+ * @author pengys5
+ */
+public class H2StorageInstaller extends StorageInstaller {
+
+ @Override protected void defineFilter(List tableDefines) {
+ int size = tableDefines.size();
+ for (int i = size - 1; i >= 0; i--) {
+ if (!(tableDefines.get(i) instanceof H2TableDefine)) {
+ tableDefines.remove(i);
+ }
+ }
+ }
+
+ @Override protected boolean isExists(Client client, TableDefine tableDefine) throws StorageException {
+ return false;
+ }
+
+ @Override protected boolean deleteIndex(Client client, TableDefine tableDefine) throws StorageException {
+ return false;
+ }
+
+ @Override protected boolean createTable(Client client, TableDefine tableDefine) throws StorageException {
+ H2Client h2Client = (H2Client)client;
+ H2TableDefine h2TableDefine = (H2TableDefine)tableDefine;
+
+ StringBuilder sqlBuilder = new StringBuilder();
+ sqlBuilder.append("CREATE TABLE ").append(h2TableDefine.getName()).append(" (");
+
+ h2TableDefine.getColumnDefines().forEach(columnDefine -> {
+ H2ColumnDefine h2ColumnDefine = (H2ColumnDefine)columnDefine;
+ if (h2ColumnDefine.getType().equals(H2ColumnDefine.Type.Varchar.name())) {
+ sqlBuilder.append(h2ColumnDefine.getName()).append(" ").append(h2ColumnDefine.getType()).append("(255)");
+ } else {
+ sqlBuilder.append(h2ColumnDefine.getName()).append(" ").append(h2ColumnDefine.getType());
+ }
+ });
+
+ sqlBuilder.append(")");
+ try {
+ h2Client.execute(sqlBuilder.toString());
+ } catch (H2ClientException e) {
+ throw new StorageInstallException(e.getMessage(), e);
+ }
+ return true;
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2TableDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2TableDefine.java
new file mode 100644
index 000000000..1854980dd
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/h2/define/H2TableDefine.java
@@ -0,0 +1,13 @@
+package org.skywalking.apm.collector.storage.h2.define;
+
+import org.skywalking.apm.collector.core.storage.TableDefine;
+
+/**
+ * @author pengys5
+ */
+public abstract class H2TableDefine extends TableDefine {
+
+ public H2TableDefine(String name) {
+ super(name);
+ }
+}
diff --git a/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/es_dao.define
new file mode 100644
index 000000000..1fcbc9a03
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/es_dao.define
@@ -0,0 +1 @@
+org.skywalking.apm.collector.storage.elasticsearch.dao.BatchEsDAO
\ No newline at end of file
diff --git a/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/group.define b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/group.define
new file mode 100644
index 000000000..8627ade6b
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/group.define
@@ -0,0 +1 @@
+org.skywalking.apm.collector.storage.StorageModuleGroupDefine
\ No newline at end of file
diff --git a/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/h2_dao.define
new file mode 100644
index 000000000..de143481d
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/h2_dao.define
@@ -0,0 +1 @@
+org.skywalking.apm.collector.storage.h2.dao.BatchH2DAO
\ No newline at end of file
diff --git a/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/module.define b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/module.define
new file mode 100644
index 000000000..a5844ab39
--- /dev/null
+++ b/apm-collector/apm-collector-storage/src/main/resources/META-INF/defines/module.define
@@ -0,0 +1,2 @@
+org.skywalking.apm.collector.storage.elasticsearch.StorageElasticSearchModuleDefine
+org.skywalking.apm.collector.storage.h2.StorageH2ModuleDefine
\ No newline at end of file
diff --git a/apm-collector/apm-collector-remote/pom.xml b/apm-collector/apm-collector-stream/pom.xml
similarity index 65%
rename from apm-collector/apm-collector-remote/pom.xml
rename to apm-collector/apm-collector-stream/pom.xml
index 81c877d81..cc73ad3ca 100644
--- a/apm-collector/apm-collector-remote/pom.xml
+++ b/apm-collector/apm-collector-stream/pom.xml
@@ -9,15 +9,9 @@
4.0.0
- apm-collector-remote
+ apm-collector-stream
jar
-
- UTF-8
- 1.4.0
- 4.1.12.Final
-
-
org.skywalking
@@ -26,43 +20,18 @@
org.skywalking
- apm-collector-server
+ apm-collector-cluster
${project.version}
- io.grpc
- grpc-netty
- ${grpc.version}
-
-
- io.netty
- netty-codec-http2
-
-
- io.netty
- netty-handler-proxy
-
-
+ org.skywalking
+ apm-collector-queue
+ ${project.version}
- io.grpc
- grpc-protobuf
- ${grpc.version}
-
-
- io.grpc
- grpc-stub
- ${grpc.version}
-
-
- io.netty
- netty-codec-http2
- ${netty.version}
-
-
- io.netty
- netty-handler-proxy
- ${netty.version}
+ org.skywalking
+ apm-collector-server
+ ${project.version}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleContext.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleContext.java
new file mode 100644
index 000000000..a68f89341
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleContext.java
@@ -0,0 +1,37 @@
+package org.skywalking.apm.collector.stream;
+
+import java.util.HashMap;
+import java.util.Map;
+import org.skywalking.apm.collector.core.framework.Context;
+import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext;
+import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine;
+
+/**
+ * @author pengys5
+ */
+public class StreamModuleContext extends Context {
+
+ private Map dataDefineMap;
+ private ClusterWorkerContext clusterWorkerContext;
+
+ public StreamModuleContext(String groupName) {
+ super(groupName);
+ dataDefineMap = new HashMap<>();
+ }
+
+ public void putAllDataDefine(Map dataDefineMap) {
+ this.dataDefineMap.putAll(dataDefineMap);
+ }
+
+ public DataDefine getDataDefine(int dataDefineId) {
+ return this.dataDefineMap.get(dataDefineId);
+ }
+
+ public ClusterWorkerContext getClusterWorkerContext() {
+ return clusterWorkerContext;
+ }
+
+ public void setClusterWorkerContext(ClusterWorkerContext clusterWorkerContext) {
+ this.clusterWorkerContext = clusterWorkerContext;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleDefine.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleDefine.java
new file mode 100644
index 000000000..89dad56db
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleDefine.java
@@ -0,0 +1,45 @@
+package org.skywalking.apm.collector.stream;
+
+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 StreamModuleDefine extends ModuleDefine implements ClusterDataListenerDefine {
+
+ private final Logger logger = LoggerFactory.getLogger(StreamModuleDefine.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 StreamModuleException(e.getMessage(), e);
+ }
+ }
+
+ @Override public final boolean defaultModule() {
+ return true;
+ }
+
+ public abstract List handlerList() throws DefineException;
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleException.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleException.java
new file mode 100644
index 000000000..81249ad74
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleException.java
@@ -0,0 +1,17 @@
+package org.skywalking.apm.collector.stream;
+
+import org.skywalking.apm.collector.core.module.ModuleException;
+
+/**
+ * @author pengys5
+ */
+public class StreamModuleException extends ModuleException {
+
+ public StreamModuleException(String message) {
+ super(message);
+ }
+
+ public StreamModuleException(String message, Throwable cause) {
+ super(message, cause);
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleGroupDefine.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleGroupDefine.java
new file mode 100644
index 000000000..e320a817d
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleGroupDefine.java
@@ -0,0 +1,25 @@
+package org.skywalking.apm.collector.stream;
+
+import org.skywalking.apm.collector.core.framework.Context;
+import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
+import org.skywalking.apm.collector.core.module.ModuleInstaller;
+
+/**
+ * @author pengys5
+ */
+public class StreamModuleGroupDefine implements ModuleGroupDefine {
+
+ public static final String GROUP_NAME = "stream";
+
+ @Override public String name() {
+ return GROUP_NAME;
+ }
+
+ @Override public Context groupContext() {
+ return new StreamModuleContext(GROUP_NAME);
+ }
+
+ @Override public ModuleInstaller moduleInstaller() {
+ return new StreamModuleInstaller();
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleInstaller.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleInstaller.java
new file mode 100644
index 000000000..908eecbb0
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/StreamModuleInstaller.java
@@ -0,0 +1,76 @@
+package org.skywalking.apm.collector.stream;
+
+import java.util.Iterator;
+import java.util.List;
+import java.util.Map;
+import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
+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.skywalking.apm.collector.core.server.ServerHolder;
+import org.skywalking.apm.collector.core.util.ObjectUtils;
+import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider;
+import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorkerProvider;
+import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext;
+import org.skywalking.apm.collector.stream.worker.LocalAsyncWorkerProviderDefineLoader;
+import org.skywalking.apm.collector.stream.worker.ProviderNotFoundException;
+import org.skywalking.apm.collector.stream.worker.RemoteWorkerProviderDefineLoader;
+import org.skywalking.apm.collector.stream.worker.impl.data.DataDefine;
+import org.skywalking.apm.collector.stream.worker.impl.data.DataDefineLoader;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class StreamModuleInstaller implements ModuleInstaller {
+
+ private final Logger logger = LoggerFactory.getLogger(StreamModuleInstaller.class);
+
+ @Override public void install(Map moduleConfig, Map moduleDefineMap,
+ ServerHolder serverHolder) throws DefineException, ClientException {
+ logger.info("beginning stream module install");
+ StreamModuleContext context = new StreamModuleContext(StreamModuleGroupDefine.GROUP_NAME);
+ CollectorContextHelper.INSTANCE.putContext(context);
+
+ DataDefineLoader dataDefineLoader = new DataDefineLoader();
+ Map dataDefineMap = dataDefineLoader.load();
+ context.putAllDataDefine(dataDefineMap);
+
+ initializeWorker(context);
+
+ logger.info("could not configure cluster module, use the default");
+ Iterator> moduleDefineEntry = moduleDefineMap.entrySet().iterator();
+ while (moduleDefineEntry.hasNext()) {
+ ModuleDefine moduleDefine = moduleDefineEntry.next().getValue();
+ logger.info("module {} initialize", moduleDefine.getClass().getName());
+ moduleDefine.initialize((ObjectUtils.isNotEmpty(moduleConfig) && moduleConfig.containsKey(moduleDefine.name())) ? moduleConfig.get(moduleDefine.name()) : null, serverHolder);
+ }
+ }
+
+ private void initializeWorker(StreamModuleContext context) throws DefineException {
+ ClusterWorkerContext clusterWorkerContext = new ClusterWorkerContext();
+ context.setClusterWorkerContext(clusterWorkerContext);
+
+ LocalAsyncWorkerProviderDefineLoader localAsyncProviderLoader = new LocalAsyncWorkerProviderDefineLoader();
+ RemoteWorkerProviderDefineLoader remoteProviderLoader = new RemoteWorkerProviderDefineLoader();
+ try {
+ List localAsyncProviders = localAsyncProviderLoader.load();
+ for (AbstractLocalAsyncWorkerProvider provider : localAsyncProviders) {
+ provider.setClusterContext(clusterWorkerContext);
+ provider.create();
+ clusterWorkerContext.putRole(provider.role());
+ }
+
+ List remoteProviders = remoteProviderLoader.load();
+ for (AbstractRemoteWorkerProvider provider : remoteProviders) {
+ provider.setClusterContext(clusterWorkerContext);
+ clusterWorkerContext.putRole(provider.role());
+ clusterWorkerContext.putProvider(provider);
+ }
+ } catch (ProviderNotFoundException e) {
+ logger.error(e.getMessage(), e);
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCConfig.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCConfig.java
new file mode 100644
index 000000000..5d05cbd57
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCConfig.java
@@ -0,0 +1,9 @@
+package org.skywalking.apm.collector.stream.grpc;
+
+/**
+ * @author pengys5
+ */
+public class StreamGRPCConfig {
+ public static String HOST;
+ public static int PORT;
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCConfigParser.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCConfigParser.java
new file mode 100644
index 000000000..621cdbf72
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCConfigParser.java
@@ -0,0 +1,30 @@
+package org.skywalking.apm.collector.stream.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 StreamGRPCConfigParser 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))) {
+ StreamGRPCConfig.HOST = "localhost";
+ } else {
+ StreamGRPCConfig.HOST = (String)config.get(HOST);
+ }
+
+ if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(PORT))) {
+ StreamGRPCConfig.PORT = 11800;
+ } else {
+ StreamGRPCConfig.PORT = (Integer)config.get(PORT);
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCDataListener.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCDataListener.java
new file mode 100644
index 000000000..89f7e64a3
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCDataListener.java
@@ -0,0 +1,74 @@
+package org.skywalking.apm.collector.stream.grpc;
+
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import org.skywalking.apm.collector.client.grpc.GRPCClient;
+import org.skywalking.apm.collector.cluster.ClusterModuleDefine;
+import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.cluster.ClusterDataListener;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
+import org.skywalking.apm.collector.stream.StreamModuleContext;
+import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
+import org.skywalking.apm.collector.stream.worker.RemoteWorkerRef;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class StreamGRPCDataListener extends ClusterDataListener {
+
+ private final Logger logger = LoggerFactory.getLogger(StreamGRPCDataListener.class);
+
+ public static final String PATH = ClusterModuleDefine.BASE_CATALOG + "." + StreamModuleGroupDefine.GROUP_NAME + "." + StreamGRPCModuleDefine.MODULE_NAME;
+
+ @Override public String path() {
+ return PATH;
+ }
+
+ private Map clients = new HashMap<>();
+ private Map workerRefs = new HashMap<>();
+
+ @Override public void addressChangedNotify() {
+ String selfAddress = StreamGRPCConfig.HOST + ":" + StreamGRPCConfig.PORT;
+
+ StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME);
+
+ List addresses = getAddresses();
+ clients.keySet().forEach(address -> {
+ if (!addresses.contains(address)) {
+ context.getClusterWorkerContext().remove(workerRefs.get(address));
+ workerRefs.remove(address);
+ }
+ });
+
+ for (String address : addresses) {
+ if (!clients.containsKey(address)) {
+ logger.debug("new address: {}, create this address remote worker reference", address);
+ String[] hostPort = address.split(":");
+ GRPCClient client = new GRPCClient(hostPort[0], Integer.valueOf(hostPort[1]));
+ try {
+ client.initialize();
+ } catch (ClientException e) {
+ e.printStackTrace();
+ }
+ clients.put(address, client);
+
+ if (selfAddress.equals(address)) {
+ context.getClusterWorkerContext().getProviders().forEach(provider -> {
+ logger.debug("create remote worker self reference, role: {}", provider.role().roleName());
+ provider.create();
+ });
+ } else {
+ context.getClusterWorkerContext().getProviders().forEach(provider -> {
+ logger.debug("create remote worker reference, role: {}", provider.role().roleName());
+ RemoteWorkerRef workerRef = provider.create(client);
+ });
+ }
+ } else {
+ logger.debug("address: {} had remote worker reference, ignore", address);
+ }
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCModuleDefine.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCModuleDefine.java
new file mode 100644
index 000000000..3a8166476
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCModuleDefine.java
@@ -0,0 +1,58 @@
+package org.skywalking.apm.collector.stream.grpc;
+
+import java.util.ArrayList;
+import java.util.List;
+import org.skywalking.apm.collector.core.client.Client;
+import org.skywalking.apm.collector.core.client.DataMonitor;
+import org.skywalking.apm.collector.core.cluster.ClusterDataListener;
+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.server.grpc.GRPCServer;
+import org.skywalking.apm.collector.stream.StreamModuleDefine;
+import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
+import org.skywalking.apm.collector.stream.grpc.handler.RemoteCommonServiceHandler;
+
+/**
+ * @author pengys5
+ */
+public class StreamGRPCModuleDefine extends StreamModuleDefine {
+
+ public static final String MODULE_NAME = "stream";
+
+ @Override public String name() {
+ return MODULE_NAME;
+ }
+
+ @Override protected String group() {
+ return StreamModuleGroupDefine.GROUP_NAME;
+ }
+
+ @Override protected ModuleConfigParser configParser() {
+ return new StreamGRPCConfigParser();
+ }
+
+ @Override protected Client createClient(DataMonitor dataMonitor) {
+ return null;
+ }
+
+ @Override protected Server server() {
+ return new GRPCServer(StreamGRPCConfig.HOST, StreamGRPCConfig.PORT);
+ }
+
+ @Override protected ModuleRegistration registration() {
+ return new StreamGRPCModuleRegistration();
+ }
+
+ @Override public ClusterDataListener listener() {
+ return new StreamGRPCDataListener();
+ }
+
+ @Override public List handlerList() throws DefineException {
+ List handlers = new ArrayList<>();
+ handlers.add(new RemoteCommonServiceHandler());
+ return handlers;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCModuleRegistration.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCModuleRegistration.java
new file mode 100644
index 000000000..1c74d45c3
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/StreamGRPCModuleRegistration.java
@@ -0,0 +1,13 @@
+package org.skywalking.apm.collector.stream.grpc;
+
+import org.skywalking.apm.collector.core.module.ModuleRegistration;
+
+/**
+ * @author pengys5
+ */
+public class StreamGRPCModuleRegistration extends ModuleRegistration {
+
+ @Override public Value buildValue() {
+ return new Value(StreamGRPCConfig.HOST, StreamGRPCConfig.PORT, null);
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java
new file mode 100644
index 000000000..86b5dc2b0
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/grpc/handler/RemoteCommonServiceHandler.java
@@ -0,0 +1,38 @@
+package org.skywalking.apm.collector.stream.grpc.handler;
+
+import io.grpc.stub.StreamObserver;
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
+import org.skywalking.apm.collector.remote.grpc.proto.Empty;
+import org.skywalking.apm.collector.remote.grpc.proto.RemoteCommonServiceGrpc;
+import org.skywalking.apm.collector.remote.grpc.proto.RemoteData;
+import org.skywalking.apm.collector.remote.grpc.proto.RemoteMessage;
+import org.skywalking.apm.collector.server.grpc.GRPCHandler;
+import org.skywalking.apm.collector.stream.StreamModuleContext;
+import org.skywalking.apm.collector.stream.StreamModuleGroupDefine;
+import org.skywalking.apm.collector.stream.worker.Role;
+import org.skywalking.apm.collector.stream.worker.WorkerInvokeException;
+import org.skywalking.apm.collector.stream.worker.WorkerNotFoundException;
+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(RemoteMessage request, StreamObserver responseObserver) {
+ String roleName = request.getWorkerRole();
+ RemoteData remoteData = request.getRemoteData();
+
+ StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME);
+ Role role = context.getClusterWorkerContext().getRole(roleName);
+ Object object = role.dataDefine().deserialize(remoteData);
+ try {
+ context.getClusterWorkerContext().lookupInSide(roleName).tell(object);
+ } catch (WorkerNotFoundException | WorkerInvokeException e) {
+ logger.error(e.getMessage(), e);
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorker.java
similarity index 73%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorker.java
rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorker.java
index bba13892b..e43561bee 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorker.java
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorker.java
@@ -1,4 +1,6 @@
-package org.skywalking.apm.collector.core.worker;
+package org.skywalking.apm.collector.stream.worker;
+
+import org.skywalking.apm.collector.core.queue.QueueExecutor;
/**
* The AbstractLocalAsyncWorker implementations represent workers,
@@ -7,7 +9,7 @@ package org.skywalking.apm.collector.core.worker;
* @author pengys5
* @since v3.0-2017
*/
-public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker {
+public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker implements QueueExecutor {
/**
* Construct an AbstractLocalAsyncWorker with the worker role and context.
@@ -15,10 +17,9 @@ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker {
* @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);
+ public AbstractLocalAsyncWorker(Role role, ClusterWorkerContext clusterContext) {
+ super(role, clusterContext);
}
/**
@@ -40,12 +41,4 @@ public abstract class AbstractLocalAsyncWorker extends AbstractLocalWorker {
final public void allocateJob(Object message) throws WorkerException {
onWork(message);
}
-
- /**
- * The data process logic in this method.
- *
- * @param message Cast the message object to a expect subclass.
- * @throws WorkerException Don't handle the exception, throw it.
- */
- protected abstract void onWork(Object message) throws WorkerException;
}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorkerProvider.java
new file mode 100644
index 000000000..094b960cb
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalAsyncWorkerProvider.java
@@ -0,0 +1,34 @@
+package org.skywalking.apm.collector.stream.worker;
+
+import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
+import org.skywalking.apm.collector.core.queue.QueueCreator;
+import org.skywalking.apm.collector.core.queue.QueueEventHandler;
+import org.skywalking.apm.collector.core.queue.QueueExecutor;
+import org.skywalking.apm.collector.queue.QueueModuleContext;
+import org.skywalking.apm.collector.queue.QueueModuleGroupDefine;
+import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorker;
+import org.skywalking.apm.collector.stream.worker.impl.PersistenceWorkerContainer;
+
+/**
+ * @author pengys5
+ */
+public abstract class AbstractLocalAsyncWorkerProvider extends AbstractLocalWorkerProvider {
+
+ public abstract int queueSize();
+
+ @Override final public WorkerRef create() throws ProviderNotFoundException {
+ T localAsyncWorker = workerInstance(getClusterContext());
+ localAsyncWorker.preStart();
+
+ if (localAsyncWorker instanceof PersistenceWorker) {
+ PersistenceWorkerContainer.INSTANCE.addWorker((PersistenceWorker)localAsyncWorker);
+ }
+
+ QueueCreator queueCreator = ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(QueueModuleGroupDefine.GROUP_NAME)).getQueueCreator();
+ QueueEventHandler queueEventHandler = queueCreator.create(queueSize(), localAsyncWorker);
+
+ LocalAsyncWorkerRef workerRef = new LocalAsyncWorkerRef(role(), queueEventHandler);
+ getClusterContext().put(workerRef);
+ return workerRef;
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalSyncWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorker.java
similarity index 90%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalSyncWorker.java
rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorker.java
index 0638c4556..af38f5f50 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalSyncWorker.java
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorker.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.core.worker;
+package org.skywalking.apm.collector.stream.worker;
/**
* The AbstractLocalSyncWorker defines workers who receive data from jvm inside call and response in real
@@ -8,8 +8,8 @@ package org.skywalking.apm.collector.core.worker;
* @since v3.0-2017
*/
public abstract class AbstractLocalSyncWorker extends AbstractLocalWorker {
- public AbstractLocalSyncWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- super(role, clusterContext, selfContext);
+ public AbstractLocalSyncWorker(Role role, ClusterWorkerContext clusterContext) {
+ super(role, clusterContext);
}
/**
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorkerProvider.java
new file mode 100644
index 000000000..d3a542a61
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalSyncWorkerProvider.java
@@ -0,0 +1,15 @@
+package org.skywalking.apm.collector.stream.worker;
+
+/**
+ * @author pengys5
+ */
+public abstract class AbstractLocalSyncWorkerProvider extends AbstractLocalWorkerProvider {
+
+ @Override final public WorkerRef create() throws ProviderNotFoundException {
+ T localSyncWorker = workerInstance(getClusterContext());
+ localSyncWorker.preStart();
+
+ LocalSyncWorkerRef workerRef = new LocalSyncWorkerRef(role(), localSyncWorker);
+ return workerRef;
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorker.java
similarity index 52%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalWorker.java
rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorker.java
index f370cdac2..c92a23fc7 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalWorker.java
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorker.java
@@ -1,10 +1,10 @@
-package org.skywalking.apm.collector.core.worker;
+package org.skywalking.apm.collector.stream.worker;
/**
* @author pengys5
*/
public abstract class AbstractLocalWorker extends AbstractWorker {
- public AbstractLocalWorker(Role role, ClusterWorkerContext clusterContext, LocalWorkerContext selfContext) {
- super(role, clusterContext, selfContext);
+ public AbstractLocalWorker(Role role, ClusterWorkerContext clusterContext) {
+ super(role, clusterContext);
}
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorkerProvider.java
similarity index 73%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalWorkerProvider.java
rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorkerProvider.java
index 1f7116be4..8cc96b71f 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalWorkerProvider.java
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractLocalWorkerProvider.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.core.worker;
+package org.skywalking.apm.collector.stream.worker;
/**
* @author pengys5
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractClusterWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorker.java
similarity index 52%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractClusterWorker.java
rename to apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorker.java
index 56cfc50f4..8dab25277 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractClusterWorker.java
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorker.java
@@ -1,7 +1,7 @@
-package org.skywalking.apm.collector.core.worker;
+package org.skywalking.apm.collector.stream.worker;
/**
- * The AbstractClusterWorker implementations represent workers,
+ * The AbstractRemoteWorker implementations represent workers,
* which receive remote messages.
*
* Usually, the implementations are doing persistent, or aggregate works.
@@ -9,18 +9,17 @@ package org.skywalking.apm.collector.core.worker;
* @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) {
- super(role, clusterContext, selfContext);
+ protected AbstractRemoteWorker(Role role, ClusterWorkerContext clusterContext) {
+ super(role, clusterContext);
}
/**
@@ -36,12 +35,4 @@ public abstract class AbstractClusterWorker extends AbstractWorker {
throw new WorkerInvokeException(e.getMessage(), e.getCause());
}
}
-
- /**
- * This method use for message receiver to analyse message.
- *
- * @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 WorkerException;
}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorkerProvider.java
new file mode 100644
index 000000000..33b97df63
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractRemoteWorkerProvider.java
@@ -0,0 +1,34 @@
+package org.skywalking.apm.collector.stream.worker;
+
+import org.skywalking.apm.collector.client.grpc.GRPCClient;
+
+/**
+ * The AbstractRemoteWorkerProvider implementations represent providers,
+ * which create instance of cluster workers whose implemented {@link AbstractRemoteWorker}.
+ *
+ *
+ * @author pengys5
+ * @since v3.0-2017
+ */
+public abstract class AbstractRemoteWorkerProvider extends AbstractWorkerProvider {
+
+ /**
+ * Create the worker instance into akka system, the akka system will control the cluster worker life cycle.
+ *
+ * @return The created worker reference. See {@link RemoteWorkerRef}
+ * @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 create() {
+ T clusterWorker = workerInstance(getClusterContext());
+ RemoteWorkerRef workerRef = new RemoteWorkerRef(role(), clusterWorker);
+ getClusterContext().put(workerRef);
+ return workerRef;
+ }
+
+ public final RemoteWorkerRef create(GRPCClient client) {
+ RemoteWorkerRef workerRef = new RemoteWorkerRef(role(), client);
+ getClusterContext().put(workerRef);
+ return workerRef;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorker.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorker.java
new file mode 100644
index 000000000..42d69415a
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorker.java
@@ -0,0 +1,48 @@
+package org.skywalking.apm.collector.stream.worker;
+
+import org.skywalking.apm.collector.core.framework.Executor;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public abstract class AbstractWorker implements Executor {
+
+ private final Logger logger = LoggerFactory.getLogger(AbstractWorker.class);
+
+ private final Role role;
+
+ private final ClusterWorkerContext clusterContext;
+
+ public AbstractWorker(Role role, ClusterWorkerContext clusterContext) {
+ this.role = role;
+ this.clusterContext = clusterContext;
+ }
+
+ @Override public final void execute(Object message) {
+ try {
+ onWork(message);
+ } catch (WorkerException e) {
+ logger.error(e.getMessage(), e);
+ }
+ }
+
+ /**
+ * The data process logic in this method.
+ *
+ * @param message Cast the message object to a expect subclass.
+ * @throws WorkerException Don't handle the exception, throw it.
+ */
+ protected abstract void onWork(Object message) throws WorkerException;
+
+ public abstract void preStart() throws ProviderNotFoundException;
+
+ final public ClusterWorkerContext getClusterContext() {
+ return clusterContext;
+ }
+
+ final public Role getRole() {
+ return role;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorkerProvider.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorkerProvider.java
new file mode 100644
index 000000000..142a371e5
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/AbstractWorkerProvider.java
@@ -0,0 +1,21 @@
+package org.skywalking.apm.collector.stream.worker;
+
+/**
+ * @author pengys5
+ */
+public abstract class AbstractWorkerProvider implements Provider {
+
+ private ClusterWorkerContext clusterContext;
+
+ public abstract Role role();
+
+ public abstract T workerInstance(ClusterWorkerContext clusterContext);
+
+ final public void setClusterContext(ClusterWorkerContext clusterContext) {
+ this.clusterContext = clusterContext;
+ }
+
+ final protected ClusterWorkerContext getClusterContext() {
+ return clusterContext;
+ }
+}
diff --git a/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/ClusterWorkerContext.java b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/ClusterWorkerContext.java
new file mode 100644
index 000000000..f6783dce1
--- /dev/null
+++ b/apm-collector/apm-collector-stream/src/main/java/org/skywalking/apm/collector/stream/worker/ClusterWorkerContext.java
@@ -0,0 +1,21 @@
+package org.skywalking.apm.collector.stream.worker;
+
+import java.util.ArrayList;
+import java.util.List;
+
+/**
+ * @author pengys5
+ */
+public class ClusterWorkerContext extends WorkerContext {
+
+ private List