diff --git a/apm-collector/apm-collector-agentstream/pom.xml b/apm-collector/apm-collector-agentstream/pom.xml
index 8ef444a8a..8e268c685 100644
--- a/apm-collector/apm-collector-agentstream/pom.xml
+++ b/apm-collector/apm-collector-agentstream/pom.xml
@@ -15,7 +15,12 @@
org.skywalking
- apm-collector-core
+ apm-collector-stream
+ ${project.version}
+
+
+ org.skywalking
+ apm-collector-cluster
${project.version}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleDefine.java
similarity index 80%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleDefine.java
index 97e377ba6..239e57424 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleDefine.java
@@ -1,11 +1,13 @@
-package org.skywalking.apm.collector.core.agentstream;
+package org.skywalking.apm.collector.agentstream;
import java.util.Map;
+import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.cluster.ClusterDataInitializer;
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.DataInitializer;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleDefine;
@@ -24,7 +26,7 @@ public abstract class AgentStreamModuleDefine extends ModuleDefine {
server.initialize();
String key = ClusterDataInitializer.BASE_CATALOG + "." + name();
- ClusterModuleContext.WRITER.write(key, registration().buildValue());
+ ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getWriter().write(key, registration().buildValue());
} catch (ConfigParseException | ServerException e) {
throw new AgentStreamModuleException(e.getMessage(), e);
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleException.java
similarity index 86%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleException.java
index f46602bbd..e51b3f6f6 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleException.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.core.agentstream;
+package org.skywalking.apm.collector.agentstream;
import org.skywalking.apm.collector.core.module.ModuleException;
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleGroupDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleGroupDefine.java
new file mode 100644
index 000000000..ea730c5d6
--- /dev/null
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleGroupDefine.java
@@ -0,0 +1,26 @@
+package org.skywalking.apm.collector.agentstream;
+
+import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
+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 AgentStreamModuleGroupDefine implements ModuleGroupDefine {
+
+ public static final String GROUP_NAME = "agent_stream";
+
+ @Override public String name() {
+ return GROUP_NAME;
+ }
+
+ @Override public Context groupContext() {
+ return new ClusterModuleContext(GROUP_NAME);
+ }
+
+ @Override public ModuleInstaller moduleInstaller() {
+ return new AgentStreamModuleInstaller();
+ }
+}
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java
new file mode 100644
index 000000000..df6a5f168
--- /dev/null
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java
@@ -0,0 +1,51 @@
+package org.skywalking.apm.collector.agentstream;
+
+import java.util.Iterator;
+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.ClusterModuleContext;
+import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
+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.util.CollectionUtils;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
+
+/**
+ * @author pengys5
+ */
+public class AgentStreamModuleInstaller implements ModuleInstaller {
+
+ private final Logger logger = LoggerFactory.getLogger(AgentStreamModuleInstaller.class);
+
+ @Override public void install(Map moduleConfig,
+ Map moduleDefineMap) throws DefineException, ClientException {
+ logger.info("beginning cluster module install");
+
+ ModuleDefine moduleDefine = null;
+ if (CollectionUtils.isEmpty(moduleConfig)) {
+ logger.info("could not configure cluster module, use the default");
+ Iterator> moduleDefineEntry = moduleDefineMap.entrySet().iterator();
+ while (moduleDefineEntry.hasNext()) {
+ moduleDefine = moduleDefineEntry.next().getValue();
+ if (moduleDefine.defaultModule()) {
+ logger.info("module {} initialize", moduleDefine.getClass().getName());
+ moduleDefine.initialize(null);
+ break;
+ }
+ }
+ } else {
+ Map.Entry clusterConfigEntry = moduleConfig.entrySet().iterator().next();
+ moduleDefine = moduleDefineMap.get(clusterConfigEntry.getKey());
+ moduleDefine.initialize(clusterConfigEntry.getValue());
+ }
+
+ ClusterModuleContext context = new ClusterModuleContext(ClusterModuleGroupDefine.GROUP_NAME);
+ context.setWriter(((ClusterModuleDefine)moduleDefine).registrationWriter());
+
+ CollectorContextHelper.INSTANCE.putContext(context);
+ }
+}
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfig.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCConfig.java
similarity index 66%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfig.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCConfig.java
index 950775c50..d04feda9f 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfig.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCConfig.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.agent.stream.server.grpc;
+package org.skywalking.apm.collector.agentstream.grpc;
/**
* @author pengys5
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfigParser.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCConfigParser.java
similarity index 93%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfigParser.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCConfigParser.java
index 05fd728a7..35b9e288f 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfigParser.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCConfigParser.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.agent.stream.server.grpc;
+package org.skywalking.apm.collector.agentstream.grpc;
import java.util.Map;
import org.skywalking.apm.collector.core.config.ConfigParseException;
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCModuleDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleDefine.java
similarity index 74%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCModuleDefine.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleDefine.java
index 5efeb7a58..b41b062d5 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCModuleDefine.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleDefine.java
@@ -1,8 +1,8 @@
-package org.skywalking.apm.collector.agent.stream.server.grpc;
+package org.skywalking.apm.collector.agentstream.grpc;
-import org.skywalking.apm.collector.core.agentstream.AgentStreamModuleDefine;
+import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine;
+import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine;
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.server.Server;
import org.skywalking.apm.collector.server.grpc.GRPCServer;
@@ -12,8 +12,8 @@ import org.skywalking.apm.collector.server.grpc.GRPCServer;
*/
public class AgentStreamGRPCModuleDefine extends AgentStreamModuleDefine {
- @Override protected ModuleGroup group() {
- return ModuleGroup.AgentStream;
+ @Override protected String group() {
+ return AgentStreamModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCModuleRegistration.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleRegistration.java
similarity index 83%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCModuleRegistration.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleRegistration.java
index 347b29b4e..800ebba24 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCModuleRegistration.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleRegistration.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.agent.stream.server.grpc;
+package org.skywalking.apm.collector.agentstream.grpc;
import org.skywalking.apm.collector.core.module.ModuleRegistration;
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyConfig.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyConfig.java
similarity index 72%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyConfig.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyConfig.java
index 240d5ef1e..5ed3e24c5 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyConfig.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyConfig.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.agent.stream.server.jetty;
+package org.skywalking.apm.collector.agentstream.jetty;
/**
* @author pengys5
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyConfigParser.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyConfigParser.java
similarity index 94%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyConfigParser.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyConfigParser.java
index d50a463c2..11afc270b 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyConfigParser.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyConfigParser.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.agent.stream.server.jetty;
+package org.skywalking.apm.collector.agentstream.jetty;
import java.util.Map;
import org.skywalking.apm.collector.core.config.ConfigParseException;
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyModuleDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java
similarity index 75%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyModuleDefine.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java
index b23de9696..ac14352cb 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyModuleDefine.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java
@@ -1,8 +1,8 @@
-package org.skywalking.apm.collector.agent.stream.server.jetty;
+package org.skywalking.apm.collector.agentstream.jetty;
-import org.skywalking.apm.collector.core.agentstream.AgentStreamModuleDefine;
+import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine;
+import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine;
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.server.Server;
import org.skywalking.apm.collector.server.jetty.JettyServer;
@@ -12,8 +12,8 @@ import org.skywalking.apm.collector.server.jetty.JettyServer;
*/
public class AgentStreamJettyModuleDefine extends AgentStreamModuleDefine {
- @Override protected ModuleGroup group() {
- return ModuleGroup.AgentStream;
+ @Override protected String group() {
+ return AgentStreamModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyModuleRegistration.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleRegistration.java
similarity index 88%
rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyModuleRegistration.java
rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleRegistration.java
index 82ac5aa62..3968fb8f7 100644
--- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/jetty/AgentStreamJettyModuleRegistration.java
+++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleRegistration.java
@@ -1,4 +1,4 @@
-package org.skywalking.apm.collector.agent.stream.server.jetty;
+package org.skywalking.apm.collector.agentstream.jetty;
import com.google.gson.JsonObject;
import org.skywalking.apm.collector.core.module.ModuleRegistration;
diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/group.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/group.define
new file mode 100644
index 000000000..2c200d484
--- /dev/null
+++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/group.define
@@ -0,0 +1 @@
+org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine
\ No newline at end of file
diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/module.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/module.define
index 6891cf0e2..7103f4219 100644
--- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/module.define
+++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/module.define
@@ -1,2 +1,2 @@
-org.skywalking.apm.collector.agent.stream.server.grpc.AgentStreamGRPCModuleDefine
-org.skywalking.apm.collector.agent.stream.server.jetty.AgentStreamJettyModuleDefine
\ No newline at end of file
+org.skywalking.apm.collector.agentstream.grpc.AgentStreamGRPCModuleDefine
+org.skywalking.apm.collector.agentstream.jetty.AgentStreamJettyModuleDefine
\ No newline at end of file
diff --git a/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorBootStartUp.java b/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorBootStartUp.java
index 46acc6cc5..ef2dc90df 100644
--- a/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorBootStartUp.java
+++ b/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorBootStartUp.java
@@ -2,7 +2,6 @@ package org.skywalking.apm.collector.boot;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.config.ConfigException;
-import org.skywalking.apm.collector.core.framework.CollectorStarter;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorStarter.java b/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorStarter.java
similarity index 57%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorStarter.java
rename to apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorStarter.java
index 44534172d..a5d55fd25 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorStarter.java
+++ b/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorStarter.java
@@ -1,15 +1,16 @@
-package org.skywalking.apm.collector.core.framework;
+package org.skywalking.apm.collector.boot;
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.config.ConfigException;
+import org.skywalking.apm.collector.core.framework.DefineException;
+import org.skywalking.apm.collector.core.framework.Starter;
import org.skywalking.apm.collector.core.module.ModuleConfigLoader;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleDefineLoader;
-import org.skywalking.apm.collector.core.module.ModuleGroup;
import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
import org.skywalking.apm.collector.core.module.ModuleGroupDefineLoader;
-import org.skywalking.apm.collector.core.module.ModuleInstallerAdapter;
import org.skywalking.apm.collector.core.remote.SerializedDefineLoader;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -22,29 +23,23 @@ public class CollectorStarter implements Starter {
private final Logger logger = LoggerFactory.getLogger(CollectorStarter.class);
@Override public void start() throws ConfigException, DefineException, ClientException {
- Context context = new Context();
ModuleConfigLoader configLoader = new ModuleConfigLoader();
Map configuration = configLoader.load();
SerializedDefineLoader serializedDefineLoader = new SerializedDefineLoader();
serializedDefineLoader.load();
- ModuleDefineLoader defineLoader = new ModuleDefineLoader();
- Map> moduleDefineMap = defineLoader.load();
-
ModuleGroupDefineLoader groupDefineLoader = new ModuleGroupDefineLoader();
Map moduleGroupDefineMap = groupDefineLoader.load();
- ModuleInstallerAdapter moduleInstallerAdapter = new ModuleInstallerAdapter(ModuleGroup.Cluster);
- moduleInstallerAdapter.install(configuration.get(ModuleGroup.Cluster.name().toLowerCase()), moduleDefineMap.get(ModuleGroup.Cluster.name().toLowerCase()));
+ ModuleDefineLoader defineLoader = new ModuleDefineLoader();
+ Map> moduleDefineMap = defineLoader.load();
- ModuleGroup[] moduleGroups = ModuleGroup.values();
- for (ModuleGroup moduleGroup : moduleGroups) {
- if (!ModuleGroup.Cluster.equals(moduleGroup)) {
- moduleInstallerAdapter = new ModuleInstallerAdapter(moduleGroup);
- logger.info("module group {}, configuration {}", moduleGroup.name().toLowerCase(), configuration.get(moduleGroup.name().toLowerCase()));
- moduleInstallerAdapter.install(configuration.get(moduleGroup.name().toLowerCase()), moduleDefineMap.get(moduleGroup.name().toLowerCase()));
- }
+ moduleGroupDefineMap.get(ClusterModuleGroupDefine.GROUP_NAME).moduleInstaller().install(configuration.get(ClusterModuleGroupDefine.GROUP_NAME), moduleDefineMap.get(ClusterModuleGroupDefine.GROUP_NAME));
+ moduleGroupDefineMap.remove(ClusterModuleGroupDefine.GROUP_NAME);
+
+ for (ModuleGroupDefine moduleGroupDefine : moduleGroupDefineMap.values()) {
+ moduleGroupDefine.moduleInstaller().install(configuration.get(moduleGroupDefine.name()), moduleDefineMap.get(moduleGroupDefine.name()));
}
}
}
diff --git a/apm-collector/apm-collector-client/pom.xml b/apm-collector/apm-collector-client/pom.xml
index 1de50be53..ae71e6f14 100644
--- a/apm-collector/apm-collector-client/pom.xml
+++ b/apm-collector/apm-collector-client/pom.xml
@@ -32,6 +32,16 @@
org.apache.zookeeper
zookeeper
3.4.10
+
+
+ slf4j-api
+ org.slf4j
+
+
+ slf4j-log4j12
+ org.slf4j
+
+
\ No newline at end of file
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleGroupDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleGroupDefine.java
index 4f5246442..d51023d8a 100644
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleGroupDefine.java
+++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleGroupDefine.java
@@ -1,17 +1,26 @@
package org.skywalking.apm.collector.cluster;
+import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
+import org.skywalking.apm.collector.core.framework.Context;
import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
-import org.skywalking.apm.collector.core.module.ModuleInstallMode;
+import org.skywalking.apm.collector.core.module.ModuleInstaller;
/**
* @author pengys5
*/
public class ClusterModuleGroupDefine implements ModuleGroupDefine {
+
+ public static final String GROUP_NAME = "cluster";
+
@Override public String name() {
- return "cluster";
+ return GROUP_NAME;
}
- @Override public ModuleInstallMode mode() {
- return ModuleInstallMode.Single;
+ @Override public Context groupContext() {
+ return new ClusterModuleContext(GROUP_NAME);
+ }
+
+ @Override public ModuleInstaller moduleInstaller() {
+ return new ClusterModuleInstaller();
}
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleInstaller.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleInstaller.java
similarity index 77%
rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleInstaller.java
rename to apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleInstaller.java
index d632bace9..666366244 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleInstaller.java
+++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleInstaller.java
@@ -1,8 +1,11 @@
-package org.skywalking.apm.collector.core.cluster;
+package org.skywalking.apm.collector.cluster;
import java.util.Iterator;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
+import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
+import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
+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;
@@ -38,6 +41,10 @@ public class ClusterModuleInstaller implements ModuleInstaller {
moduleDefine = moduleDefineMap.get(clusterConfigEntry.getKey());
moduleDefine.initialize(clusterConfigEntry.getValue());
}
- ClusterModuleContext.WRITER = ((ClusterModuleDefine)moduleDefine).registrationWriter();
+
+ ClusterModuleContext context = new ClusterModuleContext(ClusterModuleGroupDefine.GROUP_NAME);
+ context.setWriter(((ClusterModuleDefine)moduleDefine).registrationWriter());
+
+ CollectorContextHelper.INSTANCE.putContext(context);
}
}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java
index cc7e9c71f..625d894fa 100644
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java
+++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java
@@ -1,25 +1,27 @@
package org.skywalking.apm.collector.cluster.redis;
import org.skywalking.apm.collector.client.redis.RedisClient;
+import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationWriter;
import org.skywalking.apm.collector.core.framework.DataInitializer;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
-import org.skywalking.apm.collector.core.module.ModuleGroup;
/**
* @author pengys5
*/
public class ClusterRedisModuleDefine extends ClusterModuleDefine {
- @Override public ModuleGroup group() {
- return ModuleGroup.Cluster;
+ public static final String MODULE_NAME = "redis";
+
+ @Override public String group() {
+ return ClusterModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
- return "redis";
+ return MODULE_NAME;
}
@Override public boolean defaultModule() {
@@ -38,11 +40,11 @@ public class ClusterRedisModuleDefine extends ClusterModuleDefine {
return new ClusterRedisDataInitializer();
}
- @Override protected ClusterModuleRegistrationWriter registrationWriter() {
+ @Override public ClusterModuleRegistrationWriter registrationWriter() {
return new ClusterRedisModuleRegistrationWriter(getClient());
}
- @Override protected ClusterModuleRegistrationReader registrationReader() {
+ @Override public ClusterModuleRegistrationReader registrationReader() {
return null;
}
}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java
index d7802e216..583439d7c 100644
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java
+++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java
@@ -1,25 +1,27 @@
package org.skywalking.apm.collector.cluster.standalone;
import org.skywalking.apm.collector.client.h2.H2Client;
+import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationWriter;
import org.skywalking.apm.collector.core.framework.DataInitializer;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
-import org.skywalking.apm.collector.core.module.ModuleGroup;
/**
* @author pengys5
*/
public class ClusterStandaloneModuleDefine extends ClusterModuleDefine {
- @Override public ModuleGroup group() {
- return ModuleGroup.Cluster;
+ public static final String MODULE_NAME = "standalone";
+
+ @Override public String group() {
+ return ClusterModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
- return "standalone";
+ return MODULE_NAME;
}
@Override public boolean defaultModule() {
@@ -38,11 +40,11 @@ public class ClusterStandaloneModuleDefine extends ClusterModuleDefine {
return new ClusterStandaloneDataInitializer();
}
- @Override protected ClusterModuleRegistrationWriter registrationWriter() {
+ @Override public ClusterModuleRegistrationWriter registrationWriter() {
return new ClusterStandaloneModuleRegistrationWriter(getClient());
}
- @Override protected ClusterModuleRegistrationReader registrationReader() {
+ @Override public ClusterModuleRegistrationReader registrationReader() {
return null;
}
}
diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java
index b555ac458..f9d570601 100644
--- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java
+++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java
@@ -1,25 +1,27 @@
package org.skywalking.apm.collector.cluster.zookeeper;
import org.skywalking.apm.collector.client.zookeeper.ZookeeperClient;
+import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.cluster.ClusterDataInitializer;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationWriter;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
-import org.skywalking.apm.collector.core.module.ModuleGroup;
/**
* @author pengys5
*/
public class ClusterZKModuleDefine extends ClusterModuleDefine {
- @Override protected ModuleGroup group() {
- return ModuleGroup.Cluster;
+ public static final String MODULE_NAME = "zookeeper";
+
+ @Override protected String group() {
+ return ClusterModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
- return "zookeeper";
+ return MODULE_NAME;
}
@Override public boolean defaultModule() {
@@ -38,11 +40,11 @@ public class ClusterZKModuleDefine extends ClusterModuleDefine {
return new ClusterZKDataInitializer();
}
- @Override protected ClusterModuleRegistrationWriter registrationWriter() {
+ @Override public ClusterModuleRegistrationWriter registrationWriter() {
return new ClusterZKModuleRegistrationWriter(getClient());
}
- @Override protected ClusterModuleRegistrationReader registrationReader() {
+ @Override public ClusterModuleRegistrationReader registrationReader() {
return null;
}
}
diff --git a/apm-collector/apm-collector-core/pom.xml b/apm-collector/apm-collector-core/pom.xml
index 5c21a41b7..8c3d9f630 100644
--- a/apm-collector/apm-collector-core/pom.xml
+++ b/apm-collector/apm-collector-core/pom.xml
@@ -18,11 +18,6 @@
snakeyaml
1.18
-
- ch.qos.logback
- logback-classic
- 1.2.3
-
com.google.code.gson
gson
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java
index 41a6a2c52..a104e24ea 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java
@@ -1,10 +1,33 @@
package org.skywalking.apm.collector.core.cluster;
+import org.skywalking.apm.collector.core.framework.Context;
+
/**
* @author pengys5
*/
-public class ClusterModuleContext {
- public static ClusterModuleRegistrationWriter WRITER;
+public class ClusterModuleContext extends Context {
- public static ClusterModuleRegistrationReader READER;
+ public ClusterModuleContext(String groupName) {
+ super(groupName);
+ }
+
+ private ClusterModuleRegistrationWriter writer;
+
+ private ClusterModuleRegistrationReader reader;
+
+ public ClusterModuleRegistrationWriter getWriter() {
+ return writer;
+ }
+
+ public void setWriter(ClusterModuleRegistrationWriter writer) {
+ this.writer = writer;
+ }
+
+ public ClusterModuleRegistrationReader getReader() {
+ return reader;
+ }
+
+ public void setReader(ClusterModuleRegistrationReader reader) {
+ this.reader = reader;
+ }
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java
index 42e4016c3..1d788d8ec 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleDefine.java
@@ -38,7 +38,7 @@ public abstract class ClusterModuleDefine extends ModuleDefine {
throw new UnsupportedOperationException("Cluster module do not need module registration.");
}
- protected abstract ClusterModuleRegistrationWriter registrationWriter();
+ public abstract ClusterModuleRegistrationWriter registrationWriter();
- protected abstract ClusterModuleRegistrationReader registrationReader();
+ public abstract ClusterModuleRegistrationReader registrationReader();
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorContextHelper.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorContextHelper.java
new file mode 100644
index 000000000..b0d6e96cf
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/CollectorContextHelper.java
@@ -0,0 +1,25 @@
+package org.skywalking.apm.collector.core.framework;
+
+import java.util.LinkedHashMap;
+import java.util.Map;
+
+/**
+ * @author pengys5
+ */
+public enum CollectorContextHelper {
+ INSTANCE;
+
+ private Map contexts = new LinkedHashMap();
+
+ public Context getContext(String moduleGroupName) {
+ return contexts.get(moduleGroupName);
+ }
+
+ public void putContext(Context context) {
+ if (contexts.containsKey(context.getGroupName())) {
+ throw new UnsupportedOperationException("This module context was put, do not allow put a new one");
+ } else {
+ contexts.put(context.getGroupName(), context);
+ }
+ }
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Context.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Context.java
index db34bb29a..9be638427 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Context.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Context.java
@@ -3,5 +3,14 @@ package org.skywalking.apm.collector.core.framework;
/**
* @author pengys5
*/
-public class Context {
+public abstract class Context {
+ private final String groupName;
+
+ public Context(String groupName) {
+ this.groupName = groupName;
+ }
+
+ public final String getGroupName() {
+ return groupName;
+ }
}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Executor.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Executor.java
new file mode 100644
index 000000000..b09717518
--- /dev/null
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/Executor.java
@@ -0,0 +1,8 @@
+package org.skywalking.apm.collector.core.framework;
+
+/**
+ * @author pengys5
+ */
+public interface Executor {
+ void execute(Object message);
+}
diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigLoader.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigLoader.java
index 55c8a0e78..8e89ef152 100644
--- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigLoader.java
+++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigLoader.java
@@ -1,6 +1,7 @@
package org.skywalking.apm.collector.core.module;
import java.io.FileNotFoundException;
+import java.io.FileReader;
import java.util.Map;
import org.skywalking.apm.collector.core.config.ConfigLoader;
import org.skywalking.apm.collector.core.util.ResourceUtils;
@@ -18,7 +19,13 @@ public class ModuleConfigLoader implements ConfigLoader