diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleInstaller.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMCommonModuleInstaller.java similarity index 56% rename from apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleInstaller.java rename to apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMCommonModuleInstaller.java index 5eac7dc7e..634e94d16 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleInstaller.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMCommonModuleInstaller.java @@ -1,12 +1,13 @@ package org.skywalking.apm.collector.agentjvm; +import java.util.List; import org.skywalking.apm.collector.core.framework.Context; -import org.skywalking.apm.collector.core.module.MultipleModuleInstaller; +import org.skywalking.apm.collector.core.module.MultipleCommonModuleInstaller; /** * @author pengys5 */ -public class AgentJVMModuleInstaller extends MultipleModuleInstaller { +public class AgentJVMCommonModuleInstaller extends MultipleCommonModuleInstaller { @Override public String groupName() { return AgentJVMModuleGroupDefine.GROUP_NAME; @@ -15,4 +16,8 @@ public class AgentJVMModuleInstaller extends MultipleModuleInstaller { @Override public Context moduleContext() { return new AgentJVMModuleContext(groupName()); } + + @Override public List dependenceModules() { + return null; + } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleDefine.java index aca1bbaae..4b5c26f96 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleDefine.java @@ -1,7 +1,6 @@ package org.skywalking.apm.collector.agentjvm; import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.client.DataMonitor; import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; import org.skywalking.apm.collector.core.module.ModuleDefine; @@ -10,7 +9,7 @@ import org.skywalking.apm.collector.core.module.ModuleDefine; */ public abstract class AgentJVMModuleDefine extends ModuleDefine implements ClusterDataListenerDefine { - @Override protected final Client createClient(DataMonitor dataMonitor) { + @Override protected final Client createClient() { throw new UnsupportedOperationException(""); } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleGroupDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleGroupDefine.java index 89b942290..9faab6c30 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleGroupDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleGroupDefine.java @@ -11,6 +11,12 @@ public class AgentJVMModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "agent_jvm"; + private final AgentJVMCommonModuleInstaller installer; + + public AgentJVMModuleGroupDefine() { + installer = new AgentJVMCommonModuleInstaller(); + } + @Override public String name() { return GROUP_NAME; } @@ -20,6 +26,6 @@ public class AgentJVMModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new AgentJVMModuleInstaller(); + return installer; } } diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCDataListener.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCDataListener.java index b60afc4fd..45b21b59a 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCDataListener.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCDataListener.java @@ -19,7 +19,7 @@ public class AgentJVMGRPCDataListener extends ClusterDataListener { } - @Override public void serverQuitNotify() { + @Override public void serverQuitNotify(String serverAddress) { } } diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleInstaller.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterCommonModuleInstaller.java similarity index 57% rename from apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleInstaller.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterCommonModuleInstaller.java index 5567fe9ba..2d8850c61 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleInstaller.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterCommonModuleInstaller.java @@ -1,12 +1,13 @@ package org.skywalking.apm.collector.agentregister; +import java.util.List; import org.skywalking.apm.collector.core.framework.Context; -import org.skywalking.apm.collector.core.module.MultipleModuleInstaller; +import org.skywalking.apm.collector.core.module.MultipleCommonModuleInstaller; /** * @author pengys5 */ -public class AgentRegisterModuleInstaller extends MultipleModuleInstaller { +public class AgentRegisterCommonModuleInstaller extends MultipleCommonModuleInstaller { @Override public String groupName() { return AgentRegisterModuleGroupDefine.GROUP_NAME; @@ -15,4 +16,8 @@ public class AgentRegisterModuleInstaller extends MultipleModuleInstaller { @Override public Context moduleContext() { return new AgentRegisterModuleContext(groupName()); } + + @Override public List dependenceModules() { + return null; + } } diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleDefine.java index 166721fb0..fd627cd88 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleDefine.java @@ -1,7 +1,6 @@ package org.skywalking.apm.collector.agentregister; import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.client.DataMonitor; import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; import org.skywalking.apm.collector.core.module.ModuleDefine; @@ -14,7 +13,7 @@ public abstract class AgentRegisterModuleDefine extends ModuleDefine implements } - @Override protected final Client createClient(DataMonitor dataMonitor) { + @Override protected final Client createClient() { throw new UnsupportedOperationException(""); } diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleGroupDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleGroupDefine.java index d59aa5b85..0aa99ade7 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleGroupDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleGroupDefine.java @@ -11,6 +11,12 @@ public class AgentRegisterModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "agent_register"; + private final AgentRegisterCommonModuleInstaller installer; + + public AgentRegisterModuleGroupDefine() { + installer = new AgentRegisterCommonModuleInstaller(); + } + @Override public String name() { return GROUP_NAME; } @@ -20,6 +26,6 @@ public class AgentRegisterModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new AgentRegisterModuleInstaller(); + return installer; } } diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCDataListener.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCDataListener.java index a8ad38c14..23667b5a7 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCDataListener.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCDataListener.java @@ -19,7 +19,7 @@ public class AgentRegisterGRPCDataListener extends ClusterDataListener { } - @Override public void serverQuitNotify() { + @Override public void serverQuitNotify(String serverAddress) { } } diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java index d56a04930..aabcfa8a9 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java @@ -19,7 +19,7 @@ public class AgentRegisterJettyDataListener extends ClusterDataListener { } - @Override public void serverQuitNotify() { + @Override public void serverQuitNotify(String serverAddress) { } } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleInstaller.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerCommonModuleInstaller.java similarity index 56% rename from apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleInstaller.java rename to apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerCommonModuleInstaller.java index 351a8655f..452c7cdfa 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleInstaller.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerCommonModuleInstaller.java @@ -1,12 +1,13 @@ package org.skywalking.apm.collector.agentserver; +import java.util.List; import org.skywalking.apm.collector.core.framework.Context; -import org.skywalking.apm.collector.core.module.MultipleModuleInstaller; +import org.skywalking.apm.collector.core.module.MultipleCommonModuleInstaller; /** * @author pengys5 */ -public class AgentServerModuleInstaller extends MultipleModuleInstaller { +public class AgentServerCommonModuleInstaller extends MultipleCommonModuleInstaller { @Override public String groupName() { return AgentServerModuleGroupDefine.GROUP_NAME; @@ -15,4 +16,8 @@ public class AgentServerModuleInstaller extends MultipleModuleInstaller { @Override public Context moduleContext() { return new AgentServerModuleContext(groupName()); } + + @Override public List dependenceModules() { + return null; + } } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleDefine.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleDefine.java index 84b75c0fc..8dc87b767 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleDefine.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleDefine.java @@ -1,7 +1,6 @@ package org.skywalking.apm.collector.agentserver; import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.client.DataMonitor; import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; import org.skywalking.apm.collector.core.module.ModuleDefine; @@ -14,7 +13,7 @@ public abstract class AgentServerModuleDefine extends ModuleDefine implements Cl } - @Override protected final Client createClient(DataMonitor dataMonitor) { + @Override protected final Client createClient() { throw new UnsupportedOperationException(""); } } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleGroupDefine.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleGroupDefine.java index 9539be419..5334db123 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleGroupDefine.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/AgentServerModuleGroupDefine.java @@ -10,6 +10,11 @@ import org.skywalking.apm.collector.core.module.ModuleInstaller; public class AgentServerModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "agent_server"; + private final AgentServerCommonModuleInstaller installer; + + public AgentServerModuleGroupDefine() { + installer = new AgentServerCommonModuleInstaller(); + } @Override public String name() { return GROUP_NAME; @@ -20,6 +25,6 @@ public class AgentServerModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new AgentServerModuleInstaller(); + return installer; } } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/AgentServerJettyDataListener.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/AgentServerJettyDataListener.java index 9c4f25468..225bc5b69 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/AgentServerJettyDataListener.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/AgentServerJettyDataListener.java @@ -17,7 +17,7 @@ public class AgentServerJettyDataListener extends ClusterDataListener { } - @Override public void serverQuitNotify() { + @Override public void serverQuitNotify(String serverAddress) { } } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java index 652af12c3..f6b05b028 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamGRPCServerHandler.java @@ -5,8 +5,6 @@ import com.google.gson.JsonElement; import java.util.Set; import javax.servlet.http.HttpServletRequest; import org.skywalking.apm.collector.agentstream.grpc.AgentStreamGRPCDataListener; -import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine; -import org.skywalking.apm.collector.core.cluster.ClusterModuleContext; import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; @@ -22,10 +20,10 @@ public class AgentStreamGRPCServerHandler extends JettyHandler { } @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { - ClusterModuleRegistrationReader reader = ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getReader(); + ClusterModuleRegistrationReader reader = CollectorContextHelper.INSTANCE.getClusterModuleContext().getReader(); Set servers = reader.read(AgentStreamGRPCDataListener.PATH); JsonArray serverArray = new JsonArray(); - servers.forEach(server -> serverArray.add(server)); + servers.forEach(serverArray::add); return serverArray; } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java index 40a1e339a..ea0f003be 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/AgentStreamJettyServerHandler.java @@ -5,8 +5,6 @@ import com.google.gson.JsonElement; import java.util.Set; import javax.servlet.http.HttpServletRequest; import org.skywalking.apm.collector.agentstream.jetty.AgentStreamJettyDataListener; -import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine; -import org.skywalking.apm.collector.core.cluster.ClusterModuleContext; import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; @@ -22,12 +20,10 @@ public class AgentStreamJettyServerHandler extends JettyHandler { } @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { - ClusterModuleRegistrationReader reader = ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getReader(); + ClusterModuleRegistrationReader reader = CollectorContextHelper.INSTANCE.getClusterModuleContext().getReader(); Set servers = reader.read(AgentStreamJettyDataListener.PATH); JsonArray serverArray = new JsonArray(); - servers.forEach(server -> { - serverArray.add(server); - }); + servers.forEach(serverArray::add); return serverArray; } diff --git a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/UIJettyServerHandler.java b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/UIJettyServerHandler.java index 9c81388bc..99378443a 100644 --- a/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/UIJettyServerHandler.java +++ b/apm-collector/apm-collector-agentserver/src/main/java/org/skywalking/apm/collector/agentserver/jetty/handler/UIJettyServerHandler.java @@ -4,8 +4,6 @@ import com.google.gson.JsonArray; import com.google.gson.JsonElement; import java.util.Set; import javax.servlet.http.HttpServletRequest; -import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine; -import org.skywalking.apm.collector.core.cluster.ClusterModuleContext; import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; @@ -22,12 +20,10 @@ public class UIJettyServerHandler extends JettyHandler { } @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { - ClusterModuleRegistrationReader reader = ((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getReader(); + ClusterModuleRegistrationReader reader = CollectorContextHelper.INSTANCE.getClusterModuleContext().getReader(); Set servers = reader.read(UIJettyDataListener.PATH); JsonArray serverArray = new JsonArray(); - servers.forEach(server -> { - serverArray.add(server); - }); + servers.forEach(serverArray::add); return serverArray; } 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/AgentStreamCommonModuleInstaller.java similarity index 66% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java rename to apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamCommonModuleInstaller.java index 073aa0c78..2f75107bc 100644 --- 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/AgentStreamCommonModuleInstaller.java @@ -1,16 +1,18 @@ package org.skywalking.apm.collector.agentstream; +import java.util.List; import org.skywalking.apm.collector.agentstream.worker.storage.PersistenceTimer; +import org.skywalking.apm.collector.core.client.ClientException; import org.skywalking.apm.collector.core.config.ConfigException; import org.skywalking.apm.collector.core.framework.Context; import org.skywalking.apm.collector.core.framework.DefineException; -import org.skywalking.apm.collector.core.module.MultipleModuleInstaller; +import org.skywalking.apm.collector.core.module.MultipleCommonModuleInstaller; import org.skywalking.apm.collector.core.server.ServerException; /** * @author pengys5 */ -public class AgentStreamModuleInstaller extends MultipleModuleInstaller { +public class AgentStreamCommonModuleInstaller extends MultipleCommonModuleInstaller { @Override public String groupName() { return AgentStreamModuleGroupDefine.GROUP_NAME; @@ -20,7 +22,11 @@ public class AgentStreamModuleInstaller extends MultipleModuleInstaller { return new AgentStreamModuleContext(groupName()); } - @Override public void install() throws DefineException, ConfigException, ServerException { + @Override public List dependenceModules() { + return null; + } + + @Override public void install() throws DefineException, ConfigException, ServerException, ClientException { super.install(); new PersistenceTimer().start(); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleDefine.java index 762e6c525..953091639 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleDefine.java @@ -1,7 +1,6 @@ package org.skywalking.apm.collector.agentstream; import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.client.DataMonitor; import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; import org.skywalking.apm.collector.core.module.ModuleDefine; @@ -10,7 +9,7 @@ import org.skywalking.apm.collector.core.module.ModuleDefine; */ public abstract class AgentStreamModuleDefine extends ModuleDefine implements ClusterDataListenerDefine { - @Override protected final Client createClient(DataMonitor dataMonitor) { + @Override protected final Client createClient() { throw new UnsupportedOperationException(""); } 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 index 216e2ef2e..9ba5517a6 100644 --- 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 @@ -10,6 +10,11 @@ import org.skywalking.apm.collector.core.module.ModuleInstaller; public class AgentStreamModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "agent_stream"; + private final AgentStreamCommonModuleInstaller installer; + + public AgentStreamModuleGroupDefine() { + installer = new AgentStreamCommonModuleInstaller(); + } @Override public String name() { return GROUP_NAME; @@ -20,6 +25,6 @@ public class AgentStreamModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new AgentStreamModuleInstaller(); + return installer; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCDataListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCDataListener.java index 12b28c8a0..a4903fd5d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCDataListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCDataListener.java @@ -19,7 +19,7 @@ public class AgentStreamGRPCDataListener extends ClusterDataListener { } - @Override public void serverQuitNotify() { + @Override public void serverQuitNotify(String serverAddress) { } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyDataListener.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyDataListener.java index 4b2774dd4..408c72194 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyDataListener.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyDataListener.java @@ -19,7 +19,7 @@ public class AgentStreamJettyDataListener extends ClusterDataListener { } - @Override public void serverQuitNotify() { + @Override public void serverQuitNotify(String serverAddress) { } } 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 ef2dc90df..ac8703252 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 @@ -1,8 +1,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.DefineException; +import org.skywalking.apm.collector.core.CollectorException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -13,9 +11,10 @@ public class CollectorBootStartUp { private static final Logger logger = LoggerFactory.getLogger(CollectorBootStartUp.class); - public static void main(String[] args) throws ConfigException, DefineException, ClientException { + public static void main(String[] args) throws CollectorException { logger.info("collector starting..."); CollectorStarter starter = new CollectorStarter(); starter.start(); + logger.info("collector start successful."); } } diff --git a/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorStarter.java b/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorStarter.java index 5e1752bcb..3f7c12922 100644 --- a/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorStarter.java +++ b/apm-collector/apm-collector-boot/src/main/java/org/skywalking/apm/collector/boot/CollectorStarter.java @@ -1,9 +1,7 @@ package org.skywalking.apm.collector.boot; import java.util.Map; -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.CollectorException; import org.skywalking.apm.collector.core.framework.Starter; import org.skywalking.apm.collector.core.module.ModuleConfigLoader; import org.skywalking.apm.collector.core.module.ModuleDefine; @@ -12,6 +10,7 @@ import org.skywalking.apm.collector.core.module.ModuleGroupDefine; import org.skywalking.apm.collector.core.module.ModuleGroupDefineLoader; import org.skywalking.apm.collector.core.server.ServerException; import org.skywalking.apm.collector.core.server.ServerHolder; +import org.skywalking.apm.collector.core.util.CollectionUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -21,24 +20,26 @@ import org.slf4j.LoggerFactory; public class CollectorStarter implements Starter { private final Logger logger = LoggerFactory.getLogger(CollectorStarter.class); + private Map moduleGroupDefineMap; - @Override public void start() throws ConfigException, DefineException, ClientException { + @Override public void start() throws CollectorException { ModuleConfigLoader configLoader = new ModuleConfigLoader(); Map configuration = configLoader.load(); ModuleGroupDefineLoader groupDefineLoader = new ModuleGroupDefineLoader(); - Map moduleGroupDefineMap = groupDefineLoader.load(); + moduleGroupDefineMap = groupDefineLoader.load(); ModuleDefineLoader defineLoader = new ModuleDefineLoader(); Map> moduleDefineMap = defineLoader.load(); ServerHolder serverHolder = new ServerHolder(); -// moduleGroupDefineMap.get(ClusterModuleGroupDefine.GROUP_NAME).moduleInstaller().install(configuration.get(ClusterModuleGroupDefine.GROUP_NAME), moduleDefineMap.get(ClusterModuleGroupDefine.GROUP_NAME), serverHolder); -// moduleGroupDefineMap.remove(ClusterModuleGroupDefine.GROUP_NAME); - for (ModuleGroupDefine moduleGroupDefine : moduleGroupDefineMap.values()) { moduleGroupDefine.moduleInstaller().injectConfiguration(configuration.get(moduleGroupDefine.name()), moduleDefineMap.get(moduleGroupDefine.name())); moduleGroupDefine.moduleInstaller().injectServerHolder(serverHolder); + moduleGroupDefine.moduleInstaller().preInstall(); + } + + for (ModuleGroupDefine moduleGroupDefine : moduleGroupDefineMap.values()) { moduleGroupDefine.moduleInstaller().install(); } @@ -49,5 +50,26 @@ public class CollectorStarter implements Starter { logger.error(e.getMessage(), e); } }); + + dependenceAfterInstall(); + } + + private void dependenceAfterInstall() throws CollectorException { + for (ModuleGroupDefine moduleGroupDefine : moduleGroupDefineMap.values()) { + moduleInstall(moduleGroupDefine); + } + } + + private void moduleInstall(ModuleGroupDefine moduleGroupDefine) throws CollectorException { + if (CollectionUtils.isNotEmpty(moduleGroupDefine.moduleInstaller().dependenceModules())) { + for (String groupName : moduleGroupDefine.moduleInstaller().dependenceModules()) { + moduleInstall(moduleGroupDefineMap.get(groupName)); + } + logger.info("after install module group: {}", moduleGroupDefine.name()); + moduleGroupDefine.moduleInstaller().afterInstall(); + } else { + logger.info("after install module group: {}", moduleGroupDefine.name()); + moduleGroupDefine.moduleInstaller().afterInstall(); + } } } diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java index b1f1284f8..6c5d2de09 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/elasticsearch/ElasticSearchClient.java @@ -64,6 +64,10 @@ public class ElasticSearchClient implements Client { } } + @Override public void shutdown() { + + } + private List parseClusterNodes(String nodes) { List pairsList = new LinkedList<>(); logger.info("elasticsearch cluster nodes: {}", nodes); diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClient.java index 19a077514..d9906c712 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/grpc/GRPCClient.java @@ -25,6 +25,10 @@ public class GRPCClient implements Client { channel = ManagedChannelBuilder.forAddress(host, port).usePlaintext(true).build(); } + @Override public void shutdown() { + channel.shutdownNow(); + } + public ManagedChannel getChannel() { return channel; } diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java index e4b686768..65960cbef 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java @@ -27,6 +27,10 @@ public class H2Client implements Client { } } + @Override public void shutdown() { + + } + public void execute(String sql) throws H2ClientException { Statement statement = null; try { diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/redis/RedisClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/redis/RedisClient.java index fffc7f6dd..da3940439 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/redis/RedisClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/redis/RedisClient.java @@ -23,6 +23,10 @@ public class RedisClient implements Client { jedis = new Jedis(host, port); } + @Override public void shutdown() { + + } + public void setex(String key, int seconds, String value) { jedis.setex(key, seconds, value); } diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClient.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClient.java index 89dc9174a..3077cc958 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClient.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/zookeeper/ZookeeperClient.java @@ -39,6 +39,10 @@ public class ZookeeperClient implements Client { } } + @Override public void shutdown() { + + } + public void create(final String path, byte data[], List acl, CreateMode createMode) throws ZookeeperClientException { try { diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleDefine.java index dc1880881..d6fd565b2 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleDefine.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleDefine.java @@ -1,6 +1,7 @@ package org.skywalking.apm.collector.cluster; import java.util.List; +import org.skywalking.apm.collector.core.CollectorException; import org.skywalking.apm.collector.core.client.Client; import org.skywalking.apm.collector.core.client.ClientException; import org.skywalking.apm.collector.core.client.DataMonitor; @@ -20,17 +21,15 @@ public abstract class ClusterModuleDefine extends ModuleDefine { public static final String BASE_CATALOG = "skywalking"; private Client client; - private DataMonitor dataMonitor; @Override protected void initializeOtherContext() { try { - dataMonitor = dataMonitor(); - client = createClient(dataMonitor); + client = createClient(); client.initialize(); - dataMonitor.setClient(client); - ClusterModuleRegistrationReader reader = registrationReader(dataMonitor); + dataMonitor().setClient(client); + ClusterModuleRegistrationReader reader = registrationReader(); - CollectorContextHelper.INSTANCE.getClusterModuleContext().setDataMonitor(dataMonitor); + CollectorContextHelper.INSTANCE.getClusterModuleContext().setDataMonitor(dataMonitor()); CollectorContextHelper.INSTANCE.getClusterModuleContext().setReader(reader); } catch (ClientException e) { throw new UnexpectedException(e.getMessage()); @@ -42,11 +41,11 @@ public abstract class ClusterModuleDefine extends ModuleDefine { } @Override public final Server server() { - throw new UnsupportedOperationException(""); + return null; } @Override public final List handlerList() { - throw new UnsupportedOperationException(""); + return null; } @Override protected final ModuleRegistration registration() { @@ -55,6 +54,9 @@ public abstract class ClusterModuleDefine extends ModuleDefine { public abstract DataMonitor dataMonitor(); - public abstract ClusterModuleRegistrationReader registrationReader(DataMonitor dataMonitor); + public abstract ClusterModuleRegistrationReader registrationReader(); + public void startMonitor() throws CollectorException { + dataMonitor().start(); + } } 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 d51023d8a..982a7128a 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 @@ -11,6 +11,11 @@ import org.skywalking.apm.collector.core.module.ModuleInstaller; public class ClusterModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "cluster"; + private final ClusterModuleInstaller installer; + + public ClusterModuleGroupDefine() { + installer = new ClusterModuleInstaller(); + } @Override public String name() { return GROUP_NAME; @@ -21,6 +26,6 @@ public class ClusterModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new ClusterModuleInstaller(); + return installer; } } diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleInstaller.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleInstaller.java index e7af136f1..529cf9547 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleInstaller.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/ClusterModuleInstaller.java @@ -1,5 +1,8 @@ package org.skywalking.apm.collector.cluster; +import java.util.LinkedList; +import java.util.List; +import org.skywalking.apm.collector.core.CollectorException; import org.skywalking.apm.collector.core.cluster.ClusterModuleContext; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.framework.Context; @@ -19,4 +22,14 @@ public class ClusterModuleInstaller extends SingleModuleInstaller { CollectorContextHelper.INSTANCE.putClusterContext(clusterModuleContext); return clusterModuleContext; } + + @Override public List dependenceModules() { + List dependenceModules = new LinkedList<>(); + dependenceModules.add("collector_inside"); + return dependenceModules; + } + + @Override public void onAfterInstall() throws CollectorException { + ((ClusterModuleDefine)getModuleDefine()).startMonitor(); + } } 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 b946cd709..5995b7b3c 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 @@ -35,11 +35,11 @@ public class ClusterRedisModuleDefine extends ClusterModuleDefine { return null; } - @Override protected Client createClient(DataMonitor dataMonitor) { + @Override protected Client createClient() { return new RedisClient(ClusterRedisConfig.HOST, ClusterRedisConfig.PORT); } - @Override public ClusterModuleRegistrationReader registrationReader(DataMonitor dataMonitor) { - return new ClusterRedisModuleRegistrationReader(dataMonitor); + @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 1c95e0642..d8572a4f4 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 @@ -35,11 +35,11 @@ public class ClusterStandaloneModuleDefine extends ClusterModuleDefine { return null; } - @Override protected Client createClient(DataMonitor dataMonitor) { + @Override protected Client createClient() { return new H2Client(); } - @Override public ClusterModuleRegistrationReader registrationReader(DataMonitor dataMonitor) { - return new ClusterStandaloneModuleRegistrationReader(dataMonitor); + @Override public ClusterModuleRegistrationReader registrationReader() { + return null; } } diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java index 0262005e2..7162adf97 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataMonitor.java @@ -53,9 +53,12 @@ public class ClusterZKDataMonitor implements DataMonitor, Watcher { String dataStr = new String(data); if (stat.getCzxid() == stat.getMzxid()) { logger.info("path children has been created, path: {}, data: {}", event.getPath() + "/" + serverPath, dataStr); + listeners.get(event.getPath()).addAddress(serverPath + dataStr); listeners.get(event.getPath()).serverJoinNotify(serverPath + dataStr); } else { logger.info("path children has been changed, path: {}, data: {}", event.getPath() + "/" + serverPath, dataStr); + listeners.get(event.getPath()).removeAddress(serverPath + dataStr); + listeners.get(event.getPath()).serverQuitNotify(serverPath + dataStr); } } } 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 c53676407..87fd223a4 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,6 +1,5 @@ package org.skywalking.apm.collector.cluster.zookeeper; -import org.apache.zookeeper.Watcher; import org.skywalking.apm.collector.client.zookeeper.ZookeeperClient; import org.skywalking.apm.collector.cluster.ClusterModuleDefine; import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine; @@ -15,6 +14,11 @@ import org.skywalking.apm.collector.core.module.ModuleConfigParser; public class ClusterZKModuleDefine extends ClusterModuleDefine { public static final String MODULE_NAME = "zookeeper"; + private final ClusterZKDataMonitor dataMonitor; + + public ClusterZKModuleDefine() { + dataMonitor = new ClusterZKDataMonitor(); + } @Override protected String group() { return ClusterModuleGroupDefine.GROUP_NAME; @@ -33,14 +37,14 @@ public class ClusterZKModuleDefine extends ClusterModuleDefine { } @Override public DataMonitor dataMonitor() { - return new ClusterZKDataMonitor(); + return dataMonitor; } - @Override protected Client createClient(DataMonitor dataMonitor) { - return new ZookeeperClient(ClusterZKConfig.HOST_PORT, ClusterZKConfig.SESSION_TIMEOUT, (Watcher)dataMonitor); + @Override protected Client createClient() { + return new ZookeeperClient(ClusterZKConfig.HOST_PORT, ClusterZKConfig.SESSION_TIMEOUT, dataMonitor); } - @Override public ClusterModuleRegistrationReader registrationReader(DataMonitor dataMonitor) { + @Override public ClusterModuleRegistrationReader registrationReader() { return new ClusterZKModuleRegistrationReader(dataMonitor); } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/Client.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/Client.java index 24b7f5357..192e8aa7f 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/Client.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/Client.java @@ -5,4 +5,6 @@ package org.skywalking.apm.collector.core.client; */ public interface Client { void initialize() throws ClientException; + + void shutdown(); } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/DataMonitor.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/DataMonitor.java index c8b1a5df0..e2a7c9efc 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/DataMonitor.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/client/DataMonitor.java @@ -7,7 +7,7 @@ import org.skywalking.apm.collector.core.module.ModuleRegistration; /** * @author pengys5 */ -public interface DataMonitor extends Starter{ +public interface DataMonitor extends Starter { void setClient(Client client); void addListener(ClusterDataListener listener, ModuleRegistration registration) throws ClientException; diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDataListener.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDataListener.java index 6613678c5..cc494980e 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDataListener.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterDataListener.java @@ -31,5 +31,5 @@ public abstract class ClusterDataListener implements Listener { public abstract void serverJoinNotify(String serverAddress); - public abstract void serverQuitNotify(); + public abstract void serverQuitNotify(String serverAddress); } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/CommonModuleInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/CommonModuleInstaller.java new file mode 100644 index 000000000..1e5e47a88 --- /dev/null +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/CommonModuleInstaller.java @@ -0,0 +1,37 @@ +package org.skywalking.apm.collector.core.module; + +import java.util.Map; +import org.skywalking.apm.collector.core.CollectorException; + +/** + * @author pengys5 + */ +public abstract class CommonModuleInstaller implements ModuleInstaller { + + private boolean isInstalled = false; + private Map moduleConfig; + private Map moduleDefineMap; + + @Override + public final void injectConfiguration(Map moduleConfig, Map moduleDefineMap) { + this.moduleConfig = moduleConfig; + this.moduleDefineMap = moduleDefineMap; + } + + protected final Map getModuleConfig() { + return moduleConfig; + } + + protected final Map getModuleDefineMap() { + return moduleDefineMap; + } + + public abstract void onAfterInstall() throws CollectorException; + + @Override public final void afterInstall() throws CollectorException { + if (!isInstalled) { + onAfterInstall(); + } + isInstalled = true; + } +} diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigContainer.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigContainer.java deleted file mode 100644 index 8c416f1d4..000000000 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigContainer.java +++ /dev/null @@ -1,26 +0,0 @@ -package org.skywalking.apm.collector.core.module; - -import java.util.Map; - -/** - * @author pengys5 - */ -public abstract class ModuleConfigContainer implements ModuleInstaller { - - private Map moduleConfig; - private Map moduleDefineMap; - - @Override - public final void injectConfiguration(Map moduleConfig, Map moduleDefineMap) { - this.moduleConfig = moduleConfig; - this.moduleDefineMap = moduleDefineMap; - } - - public final Map getModuleConfig() { - return moduleConfig; - } - - public final Map getModuleDefineMap() { - return moduleDefineMap; - } -} diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java index f06583209..b0706ecc8 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleDefine.java @@ -2,7 +2,6 @@ package org.skywalking.apm.collector.core.module; 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.framework.Define; import org.skywalking.apm.collector.core.framework.Handler; import org.skywalking.apm.collector.core.server.Server; @@ -18,7 +17,7 @@ public abstract class ModuleDefine implements Define { protected abstract ModuleConfigParser configParser(); - protected abstract Client createClient(DataMonitor dataMonitor); + protected abstract Client createClient(); protected abstract Server server(); diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleGroupDefineLoader.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleGroupDefineLoader.java index 6a4e1387c..09125d075 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleGroupDefineLoader.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleGroupDefineLoader.java @@ -1,5 +1,6 @@ package org.skywalking.apm.collector.core.module; +import java.util.Iterator; import java.util.LinkedHashMap; import java.util.Map; import org.skywalking.apm.collector.core.framework.DefineException; @@ -21,11 +22,11 @@ public class ModuleGroupDefineLoader implements Loader definitionLoader = DefinitionLoader.load(ModuleGroupDefine.class, definitionFile); - for (ModuleGroupDefine moduleGroupDefine : definitionLoader) { - logger.info("loaded group module definition class: {}", moduleGroupDefine.getClass().getName()); - - String groupName = moduleGroupDefine.name().toLowerCase(); - moduleGroupDefineMap.put(groupName, moduleGroupDefine); + Iterator defineIterator = definitionLoader.iterator(); + while (defineIterator.hasNext()) { + ModuleGroupDefine groupDefine = defineIterator.next(); + String groupName = groupDefine.name().toLowerCase(); + moduleGroupDefineMap.put(groupName, groupDefine); } return moduleGroupDefineMap; } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstaller.java index 563a30b00..b9418d445 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstaller.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstaller.java @@ -1,6 +1,8 @@ package org.skywalking.apm.collector.core.module; +import java.util.List; import java.util.Map; +import org.skywalking.apm.collector.core.CollectorException; import org.skywalking.apm.collector.core.client.ClientException; import org.skywalking.apm.collector.core.config.ConfigException; import org.skywalking.apm.collector.core.framework.Context; @@ -13,6 +15,8 @@ import org.skywalking.apm.collector.core.server.ServerHolder; */ public interface ModuleInstaller { + List dependenceModules(); + void injectServerHolder(ServerHolder serverHolder); String groupName(); @@ -24,4 +28,6 @@ public interface ModuleInstaller { void preInstall() throws DefineException, ConfigException, ServerException; void install() throws ClientException, DefineException, ConfigException, ServerException; + + void afterInstall() throws CollectorException; } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/MultipleModuleInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/MultipleCommonModuleInstaller.java similarity index 60% rename from apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/MultipleModuleInstaller.java rename to apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/MultipleCommonModuleInstaller.java index cca15bc04..1b1806df0 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/MultipleModuleInstaller.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/MultipleCommonModuleInstaller.java @@ -4,22 +4,26 @@ import java.util.Iterator; import java.util.LinkedList; import java.util.List; import java.util.Map; +import org.skywalking.apm.collector.core.CollectorException; +import org.skywalking.apm.collector.core.client.ClientException; +import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; import org.skywalking.apm.collector.core.config.ConfigException; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.framework.DefineException; import org.skywalking.apm.collector.core.server.ServerException; import org.skywalking.apm.collector.core.server.ServerHolder; +import org.skywalking.apm.collector.core.util.ObjectUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public abstract class MultipleModuleInstaller extends ModuleConfigContainer { +public abstract class MultipleCommonModuleInstaller extends CommonModuleInstaller { - private final Logger logger = LoggerFactory.getLogger(MultipleModuleInstaller.class); + private final Logger logger = LoggerFactory.getLogger(MultipleCommonModuleInstaller.class); - public MultipleModuleInstaller() { + public MultipleCommonModuleInstaller() { moduleDefines = new LinkedList<>(); } @@ -31,6 +35,7 @@ public abstract class MultipleModuleInstaller extends ModuleConfigContainer { } @Override public final void preInstall() throws DefineException, ConfigException, ServerException { + logger.info("install module group: {}", groupName()); Map moduleConfig = getModuleConfig(); Map moduleDefineMap = getModuleDefineMap(); @@ -44,12 +49,22 @@ public abstract class MultipleModuleInstaller extends ModuleConfigContainer { } } - @Override public void install() throws DefineException, ConfigException, ServerException { - preInstall(); - + @Override public void install() throws DefineException, ConfigException, ServerException, ClientException { CollectorContextHelper.INSTANCE.putContext(moduleContext()); - moduleDefines.forEach(moduleDefine -> { + for (ModuleDefine moduleDefine : moduleDefines) { moduleDefine.initializeOtherContext(); - }); + + if (moduleDefine instanceof ClusterDataListenerDefine) { + ClusterDataListenerDefine listenerDefine = (ClusterDataListenerDefine)moduleDefine; + if (ObjectUtils.isNotEmpty(listenerDefine.listener()) && ObjectUtils.isNotEmpty(moduleDefine.registration())) { + logger.info("add group: {}, module: {}, listener into cluster data monitor", moduleDefine.group(), moduleDefine.name()); + CollectorContextHelper.INSTANCE.getClusterModuleContext().getDataMonitor().addListener(listenerDefine.listener(), moduleDefine.registration()); + } + } + } + } + + @Override public void onAfterInstall() throws CollectorException { + } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/SingleModuleInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/SingleModuleInstaller.java index 9ee19d3fe..0fd55836c 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/SingleModuleInstaller.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/SingleModuleInstaller.java @@ -2,8 +2,10 @@ package org.skywalking.apm.collector.core.module; import java.util.Iterator; import java.util.Map; +import org.skywalking.apm.collector.core.CollectorException; 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.cluster.ClusterModuleException; import org.skywalking.apm.collector.core.config.ConfigException; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; @@ -11,13 +13,14 @@ import org.skywalking.apm.collector.core.framework.DefineException; import org.skywalking.apm.collector.core.server.ServerException; import org.skywalking.apm.collector.core.server.ServerHolder; import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.skywalking.apm.collector.core.util.ObjectUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; /** * @author pengys5 */ -public abstract class SingleModuleInstaller extends ModuleConfigContainer { +public abstract class SingleModuleInstaller extends CommonModuleInstaller { private final Logger logger = LoggerFactory.getLogger(SingleModuleInstaller.class); @@ -29,6 +32,7 @@ public abstract class SingleModuleInstaller extends ModuleConfigContainer { } @Override public final void preInstall() throws DefineException, ConfigException, ServerException { + logger.info("install module group: {}", groupName()); Map moduleConfig = getModuleConfig(); Map moduleDefineMap = getModuleDefineMap(); if (CollectionUtils.isNotEmpty(moduleConfig)) { @@ -45,17 +49,17 @@ public abstract class SingleModuleInstaller extends ModuleConfigContainer { } } else { logger.info("could not configure module, use the default"); - Iterator> moduleDefineEntry = moduleDefineMap.entrySet().iterator(); + Iterator> moduleDefineIterator = moduleDefineMap.entrySet().iterator(); boolean hasDefaultModule = false; - while (moduleDefineEntry.hasNext()) { - if (moduleDefineEntry.next().getValue().defaultModule()) { - logger.info("module {} initialize", moduleDefine.getClass().getName()); + while (moduleDefineIterator.hasNext()) { + Map.Entry moduleDefineEntry = moduleDefineIterator.next(); + if (moduleDefineEntry.getValue().defaultModule()) { if (hasDefaultModule) { throw new ClusterModuleException("single module, but configure multiple default module"); } - moduleDefine = moduleDefineEntry.next().getValue(); - moduleDefine.configParser().parse(null); + this.moduleDefine = moduleDefineEntry.getValue(); + this.moduleDefine.configParser().parse(null); hasDefaultModule = true; } } @@ -64,13 +68,21 @@ public abstract class SingleModuleInstaller extends ModuleConfigContainer { } @Override public void install() throws ClientException, DefineException, ConfigException, ServerException { - preInstall(); + if (!(moduleContext() instanceof ClusterModuleContext)) { + CollectorContextHelper.INSTANCE.putContext(moduleContext()); + } moduleDefine.initializeOtherContext(); - CollectorContextHelper.INSTANCE.putContext(moduleContext()); if (moduleDefine instanceof ClusterDataListenerDefine) { ClusterDataListenerDefine listenerDefine = (ClusterDataListenerDefine)moduleDefine; - CollectorContextHelper.INSTANCE.getClusterModuleContext().getDataMonitor().addListener(listenerDefine.listener(), moduleDefine.registration()); + if (ObjectUtils.isNotEmpty(listenerDefine.listener()) && ObjectUtils.isNotEmpty(moduleDefine.registration())) { + CollectorContextHelper.INSTANCE.getClusterModuleContext().getDataMonitor().addListener(listenerDefine.listener(), moduleDefine.registration()); + logger.info("add group: {}, module: {}, listener into cluster data monitor", moduleDefine.group(), moduleDefine.name()); + } } } + + protected ModuleDefine getModuleDefine() { + return moduleDefine; + } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java index 2653ab530..ec8e01238 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/server/ServerHolder.java @@ -4,6 +4,7 @@ import java.util.LinkedList; import java.util.List; import org.skywalking.apm.collector.core.framework.Handler; import org.skywalking.apm.collector.core.util.CollectionUtils; +import org.skywalking.apm.collector.core.util.ObjectUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -21,6 +22,10 @@ public class ServerHolder { } public void holdServer(Server newServer, List handlers) throws ServerException { + if (ObjectUtils.isEmpty(newServer) || CollectionUtils.isEmpty(handlers)) { + return; + } + boolean isNewServer = true; for (Server server : servers) { if (server.hostPort().equals(newServer.hostPort()) && server.serverClassify().equals(newServer.serverClassify())) { diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/DefinitionLoader.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/DefinitionLoader.java index 9030b1c2d..32b901072 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/DefinitionLoader.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/DefinitionLoader.java @@ -35,7 +35,6 @@ public class DefinitionLoader implements Iterable { @Override public final Iterator iterator() { logger.info("load definition file: {}", definitionFile.get()); - Properties properties = new Properties(); List definitionList = new LinkedList<>(); try { Enumeration urlEnumeration = this.getClass().getClassLoader().getResources(definitionFile.get()); @@ -43,6 +42,7 @@ public class DefinitionLoader implements Iterable { URL definitionFileURL = urlEnumeration.nextElement(); logger.info("definition file url: {}", definitionFileURL.getPath()); BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(definitionFileURL.openStream())); + Properties properties = new Properties(); properties.load(bufferedReader); Enumeration defineItem = properties.propertyNames(); @@ -52,7 +52,7 @@ public class DefinitionLoader implements Iterable { } } } catch (IOException e) { - e.printStackTrace(); + logger.error(e.getMessage(), e); } Iterator moduleDefineIterator = definitionList.iterator(); diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleDefine.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleDefine.java index 3b94c44f0..e7c4920d4 100644 --- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleDefine.java +++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/QueueModuleDefine.java @@ -4,7 +4,6 @@ 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.framework.Handler; -import org.skywalking.apm.collector.core.module.ModuleConfigParser; import org.skywalking.apm.collector.core.module.ModuleDefine; import org.skywalking.apm.collector.core.module.ModuleRegistration; import org.skywalking.apm.collector.core.server.Server; @@ -13,12 +12,9 @@ import org.skywalking.apm.collector.core.server.Server; * @author pengys5 */ public abstract class QueueModuleDefine extends ModuleDefine { - @Override protected final ModuleConfigParser configParser() { - throw new UnsupportedOperationException(""); - } - @Override protected final Client createClient(DataMonitor dataMonitor) { - throw new UnsupportedOperationException(""); + @Override protected Client createClient() { + return null; } @Override protected final ModuleRegistration registration() { @@ -26,10 +22,10 @@ public abstract class QueueModuleDefine extends ModuleDefine { } @Override protected final Server server() { - throw new UnsupportedOperationException(""); + return null; } @Override public final List handlerList() { - throw new UnsupportedOperationException(""); + return null; } } 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 index 9b40b7c2c..6f2648cc8 100644 --- 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 @@ -10,6 +10,11 @@ import org.skywalking.apm.collector.core.module.ModuleInstaller; public class QueueModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "queue"; + private final QueueModuleInstaller installer; + + public QueueModuleGroupDefine() { + installer = new QueueModuleInstaller(); + } @Override public String name() { return GROUP_NAME; @@ -20,6 +25,6 @@ public class QueueModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new QueueModuleInstaller(); + return installer; } } 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 index 1ed4c7996..054979197 100644 --- 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 @@ -1,13 +1,19 @@ package org.skywalking.apm.collector.queue; +import java.util.List; +import org.skywalking.apm.collector.core.CollectorException; import org.skywalking.apm.collector.core.client.ClientException; import org.skywalking.apm.collector.core.config.ConfigException; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.framework.Context; import org.skywalking.apm.collector.core.framework.DefineException; +import org.skywalking.apm.collector.core.framework.UnexpectedException; import org.skywalking.apm.collector.core.module.SingleModuleInstaller; import org.skywalking.apm.collector.core.server.ServerException; import org.skywalking.apm.collector.queue.datacarrier.DataCarrierQueueCreator; +import org.skywalking.apm.collector.queue.datacarrier.QueueDataCarrierModuleDefine; +import org.skywalking.apm.collector.queue.disruptor.DisruptorQueueCreator; +import org.skywalking.apm.collector.queue.disruptor.QueueDisruptorModuleDefine; /** * @author pengys5 @@ -22,8 +28,21 @@ public class QueueModuleInstaller extends SingleModuleInstaller { return new QueueModuleContext(groupName()); } + @Override public List dependenceModules() { + return null; + } + @Override public void install() throws ClientException, DefineException, ConfigException, ServerException { super.install(); - ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(groupName())).setQueueCreator(new DataCarrierQueueCreator()); + if (getModuleDefine() instanceof QueueDataCarrierModuleDefine) { + ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(groupName())).setQueueCreator(new DataCarrierQueueCreator()); + } else if (getModuleDefine() instanceof QueueDisruptorModuleDefine) { + ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(groupName())).setQueueCreator(new DisruptorQueueCreator()); + } else { + throw new UnexpectedException(""); + } + } + + @Override public void onAfterInstall() throws CollectorException { } } diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierQueueConfigParser.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierQueueConfigParser.java new file mode 100644 index 000000000..614ac2e7b --- /dev/null +++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/DataCarrierQueueConfigParser.java @@ -0,0 +1,15 @@ +package org.skywalking.apm.collector.queue.datacarrier; + +import java.util.Map; +import org.skywalking.apm.collector.core.config.ConfigParseException; +import org.skywalking.apm.collector.core.module.ModuleConfigParser; + +/** + * @author pengys5 + */ +public class DataCarrierQueueConfigParser implements ModuleConfigParser { + + @Override public void parse(Map config) throws ConfigParseException { + + } +} 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 12cc17463..384edb588 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 @@ -1,6 +1,7 @@ package org.skywalking.apm.collector.queue.datacarrier; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.core.module.ModuleConfigParser; import org.skywalking.apm.collector.queue.QueueModuleContext; import org.skywalking.apm.collector.queue.QueueModuleDefine; import org.skywalking.apm.collector.queue.QueueModuleGroupDefine; @@ -22,6 +23,10 @@ public class QueueDataCarrierModuleDefine extends QueueModuleDefine { return false; } + @Override protected ModuleConfigParser configParser() { + return new DataCarrierQueueConfigParser(); + } + @Override protected void initializeOtherContext() { ((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/QueueDisruptorConfigParser.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorQueueConfigParser.java similarity index 83% rename from apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorConfigParser.java rename to apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorQueueConfigParser.java index d1e52b4f9..f523401f6 100644 --- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorConfigParser.java +++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/DisruptorQueueConfigParser.java @@ -7,7 +7,9 @@ import org.skywalking.apm.collector.core.module.ModuleConfigParser; /** * @author pengys5 */ -public class QueueDisruptorConfigParser implements ModuleConfigParser { +public class DisruptorQueueConfigParser implements ModuleConfigParser { + @Override public void parse(Map config) throws ConfigParseException { + } } 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 4eeaf0de8..f086dfe54 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 @@ -1,6 +1,7 @@ package org.skywalking.apm.collector.queue.disruptor; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.core.module.ModuleConfigParser; import org.skywalking.apm.collector.queue.QueueModuleContext; import org.skywalking.apm.collector.queue.QueueModuleDefine; import org.skywalking.apm.collector.queue.QueueModuleGroupDefine; @@ -22,6 +23,10 @@ public class QueueDisruptorModuleDefine extends QueueModuleDefine { return true; } + @Override protected ModuleConfigParser configParser() { + return new DisruptorQueueConfigParser(); + } + @Override protected void initializeOtherContext() { ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(group())).setQueueCreator(new DisruptorQueueCreator()); } 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 index c088f8aa3..78d5fa23c 100644 --- 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 @@ -14,20 +14,16 @@ import org.skywalking.apm.collector.core.module.ModuleRegistration; import org.skywalking.apm.collector.core.server.Server; 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 protected void initializeOtherContext() { try { StorageModuleContext context = (StorageModuleContext)CollectorContextHelper.INSTANCE.getContext(StorageModuleGroupDefine.GROUP_NAME); - Client client = createClient(null); + Client client = createClient(); client.initialize(); context.setClient(client); injectClientIntoDAO(client); @@ -39,19 +35,19 @@ public abstract class StorageModuleDefine extends ModuleDefine implements Cluste } @Override public final List handlerList() { - throw new UnsupportedOperationException(""); + return null; } @Override protected final Server server() { - throw new UnsupportedOperationException(""); + return null; } @Override protected final ModuleRegistration registration() { - throw new UnsupportedOperationException(""); + return null; } @Override public final ClusterDataListener listener() { - throw new UnsupportedOperationException(""); + return null; } @Override public final boolean defaultModule() { 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 index c373f3a83..48e641d9f 100644 --- 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 @@ -10,6 +10,11 @@ import org.skywalking.apm.collector.core.module.ModuleInstaller; public class StorageModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "storage"; + private final StorageModuleInstaller installer; + + public StorageModuleGroupDefine() { + installer = new StorageModuleInstaller(); + } @Override public String name() { return GROUP_NAME; @@ -20,6 +25,6 @@ public class StorageModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new StorageModuleInstaller(); + return installer; } } 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 index 67e8d217b..a3b2d2794 100644 --- 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 @@ -1,5 +1,7 @@ package org.skywalking.apm.collector.storage; +import java.util.List; +import org.skywalking.apm.collector.core.CollectorException; import org.skywalking.apm.collector.core.framework.Context; import org.skywalking.apm.collector.core.module.SingleModuleInstaller; @@ -15,4 +17,12 @@ public class StorageModuleInstaller extends SingleModuleInstaller { @Override public Context moduleContext() { return new StorageModuleContext(groupName()); } + + @Override public List dependenceModules() { + return null; + } + + @Override public void onAfterInstall() throws CollectorException { + + } } 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 index 826ad4cb0..65d210d92 100644 --- 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 @@ -3,7 +3,6 @@ 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; @@ -33,7 +32,7 @@ public class StorageElasticSearchModuleDefine extends StorageModuleDefine { return new StorageElasticSearchConfigParser(); } - @Override protected Client createClient(DataMonitor dataMonitor) { + @Override protected Client createClient() { return new ElasticSearchClient(StorageElasticSearchConfig.CLUSTER_NAME, StorageElasticSearchConfig.CLUSTER_TRANSPORT_SNIFFER, StorageElasticSearchConfig.CLUSTER_NODES); } 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 index 84f9b31a7..99575e2b2 100644 --- 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 @@ -3,7 +3,6 @@ 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; @@ -33,7 +32,7 @@ public class StorageH2ModuleDefine extends StorageModuleDefine { return new StorageH2ConfigParser(); } - @Override protected Client createClient(DataMonitor dataMonitor) { + @Override protected Client createClient() { return new H2Client(); } 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 index 99b11bbbb..cb2142c2d 100644 --- 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 @@ -10,6 +10,11 @@ import org.skywalking.apm.collector.core.module.ModuleInstaller; public class StreamModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "collector_inside"; + private final StreamModuleInstaller installer; + + public StreamModuleGroupDefine() { + installer = new StreamModuleInstaller(); + } @Override public String name() { return GROUP_NAME; @@ -20,6 +25,6 @@ public class StreamModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new StreamModuleInstaller(); + return installer; } } 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 index 1c2a7fd21..da03133de 100644 --- 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 @@ -1,13 +1,12 @@ package org.skywalking.apm.collector.stream; +import java.util.LinkedList; import java.util.List; -import org.skywalking.apm.collector.core.client.ClientException; -import org.skywalking.apm.collector.core.config.ConfigException; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.core.framework.Context; import org.skywalking.apm.collector.core.framework.DefineException; import org.skywalking.apm.collector.core.module.SingleModuleInstaller; -import org.skywalking.apm.collector.core.server.ServerException; +import org.skywalking.apm.collector.queue.QueueModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorkerProvider; import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorkerProvider; import org.skywalking.apm.collector.stream.worker.ClusterWorkerContext; @@ -32,8 +31,13 @@ public class StreamModuleInstaller extends SingleModuleInstaller { return new StreamModuleContext(groupName()); } - @Override public void install() throws ClientException, DefineException, ConfigException, ServerException { - super.install(); + @Override public List dependenceModules() { + List dependenceModules = new LinkedList<>(); + dependenceModules.add(QueueModuleGroupDefine.GROUP_NAME); + return dependenceModules; + } + + @Override public void onAfterInstall() throws DefineException { initializeWorker((StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(groupName())); } 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 index 781e65d55..445a17e45 100644 --- 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 @@ -1,6 +1,8 @@ package org.skywalking.apm.collector.stream.grpc; import java.util.HashMap; +import java.util.LinkedList; +import java.util.List; import java.util.Map; import org.skywalking.apm.collector.client.grpc.GRPCClient; import org.skywalking.apm.collector.cluster.ClusterModuleDefine; @@ -9,6 +11,7 @@ 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; @@ -26,6 +29,7 @@ public class StreamGRPCDataListener extends ClusterDataListener { } private Map clients = new HashMap<>(); + private Map> remoteWorkerRefMap = new HashMap<>(); @Override public void serverJoinNotify(String serverAddress) { String selfAddress = StreamGRPCConfig.HOST + ":" + StreamGRPCConfig.PORT; @@ -50,7 +54,11 @@ public class StreamGRPCDataListener extends ClusterDataListener { } else { context.getClusterWorkerContext().getProviders().forEach(provider -> { logger.info("create remote worker reference, role: {}", provider.role().roleName()); - provider.create(client); + RemoteWorkerRef remoteWorkerRef = provider.create(client); + if (!remoteWorkerRefMap.containsKey(serverAddress)) { + remoteWorkerRefMap.put(selfAddress, new LinkedList<>()); + } + remoteWorkerRefMap.get(serverAddress).add(remoteWorkerRef); }); } } else { @@ -58,7 +66,17 @@ public class StreamGRPCDataListener extends ClusterDataListener { } } - @Override public void serverQuitNotify() { + @Override public void serverQuitNotify(String serverAddress) { + StreamModuleContext context = (StreamModuleContext)CollectorContextHelper.INSTANCE.getContext(StreamModuleGroupDefine.GROUP_NAME); + if (clients.containsKey(serverAddress)) { + clients.get(serverAddress).shutdown(); + clients.remove(serverAddress); + } + if (remoteWorkerRefMap.containsKey(serverAddress)) { + for (RemoteWorkerRef remoteWorkerRef : remoteWorkerRefMap.get(serverAddress)) { + context.getClusterWorkerContext().remove(remoteWorkerRef); + } + } } } 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 index 1e3897d51..388d01126 100644 --- 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 @@ -3,7 +3,6 @@ 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.Handler; import org.skywalking.apm.collector.core.module.ModuleConfigParser; @@ -33,7 +32,7 @@ public class StreamGRPCModuleDefine extends StreamModuleDefine { return new StreamGRPCConfigParser(); } - @Override protected Client createClient(DataMonitor dataMonitor) { + @Override protected Client createClient() { return null; } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleInstaller.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UICommonModuleInstaller.java similarity index 55% rename from apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleInstaller.java rename to apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UICommonModuleInstaller.java index 4308cf6c9..0466baf68 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleInstaller.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UICommonModuleInstaller.java @@ -1,12 +1,13 @@ package org.skywalking.apm.collector.ui; +import java.util.List; import org.skywalking.apm.collector.core.framework.Context; -import org.skywalking.apm.collector.core.module.MultipleModuleInstaller; +import org.skywalking.apm.collector.core.module.MultipleCommonModuleInstaller; /** * @author pengys5 */ -public class UIModuleInstaller extends MultipleModuleInstaller { +public class UICommonModuleInstaller extends MultipleCommonModuleInstaller { @Override public String groupName() { return UIModuleGroupDefine.GROUP_NAME; @@ -15,4 +16,8 @@ public class UIModuleInstaller extends MultipleModuleInstaller { @Override public Context moduleContext() { return new UIModuleContext(groupName()); } + + @Override public List dependenceModules() { + return null; + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleDefine.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleDefine.java index bc7643d7b..fca0c2a0f 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleDefine.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleDefine.java @@ -1,7 +1,6 @@ package org.skywalking.apm.collector.ui; import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.client.DataMonitor; import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; import org.skywalking.apm.collector.core.module.ModuleDefine; @@ -10,7 +9,7 @@ import org.skywalking.apm.collector.core.module.ModuleDefine; */ public abstract class UIModuleDefine extends ModuleDefine implements ClusterDataListenerDefine { - @Override protected final Client createClient(DataMonitor dataMonitor) { + @Override protected final Client createClient() { throw new UnsupportedOperationException(""); } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleGroupDefine.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleGroupDefine.java index 575cd4175..372174703 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleGroupDefine.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/UIModuleGroupDefine.java @@ -10,6 +10,11 @@ import org.skywalking.apm.collector.core.module.ModuleInstaller; public class UIModuleGroupDefine implements ModuleGroupDefine { public static final String GROUP_NAME = "ui"; + private final UICommonModuleInstaller installer; + + public UIModuleGroupDefine() { + installer = new UICommonModuleInstaller(); + } @Override public String name() { return GROUP_NAME; @@ -20,6 +25,6 @@ public class UIModuleGroupDefine implements ModuleGroupDefine { } @Override public ModuleInstaller moduleInstaller() { - return new UIModuleInstaller(); + return installer; } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyDataListener.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyDataListener.java index e3d4bffb2..d9fdd2c36 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyDataListener.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/jetty/UIJettyDataListener.java @@ -19,7 +19,6 @@ public class UIJettyDataListener extends ClusterDataListener { } - @Override public void serverQuitNotify() { - + @Override public void serverQuitNotify(String serverAddress) { } }