diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMCommonModuleInstaller.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMCommonModuleInstaller.java deleted file mode 100644 index 634e94d16..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMCommonModuleInstaller.java +++ /dev/null @@ -1,23 +0,0 @@ -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.MultipleCommonModuleInstaller; - -/** - * @author pengys5 - */ -public class AgentJVMCommonModuleInstaller extends MultipleCommonModuleInstaller { - - @Override public String groupName() { - return AgentJVMModuleGroupDefine.GROUP_NAME; - } - - @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/AgentJVMModuleContext.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleContext.java deleted file mode 100644 index 5152b68d6..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleContext.java +++ /dev/null @@ -1,12 +0,0 @@ -package org.skywalking.apm.collector.agentjvm; - -import org.skywalking.apm.collector.core.framework.Context; - -/** - * @author pengys5 - */ -public class AgentJVMModuleContext extends Context { - public AgentJVMModuleContext(String groupName) { - super(groupName); - } -} 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 deleted file mode 100644 index 4b5c26f96..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleDefine.java +++ /dev/null @@ -1,23 +0,0 @@ -package org.skywalking.apm.collector.agentjvm; - -import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; -import org.skywalking.apm.collector.core.module.ModuleDefine; - -/** - * @author pengys5 - */ -public abstract class AgentJVMModuleDefine extends ModuleDefine implements ClusterDataListenerDefine { - - @Override protected final Client createClient() { - throw new UnsupportedOperationException(""); - } - - @Override protected void initializeOtherContext() { - - } - - @Override public final boolean defaultModule() { - return true; - } -} diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleException.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleException.java deleted file mode 100644 index 6732f2636..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleException.java +++ /dev/null @@ -1,16 +0,0 @@ -package org.skywalking.apm.collector.agentjvm; - -import org.skywalking.apm.collector.core.module.ModuleException; - -/** - * @author pengys5 - */ -public class AgentJVMModuleException extends ModuleException { - public AgentJVMModuleException(String message) { - super(message); - } - - public AgentJVMModuleException(String message, Throwable cause) { - super(message, cause); - } -} 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 deleted file mode 100644 index 9faab6c30..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/AgentJVMModuleGroupDefine.java +++ /dev/null @@ -1,31 +0,0 @@ -package org.skywalking.apm.collector.agentjvm; - -import org.skywalking.apm.collector.core.framework.Context; -import org.skywalking.apm.collector.core.module.ModuleGroupDefine; -import org.skywalking.apm.collector.core.module.ModuleInstaller; - -/** - * @author pengys5 - */ -public class 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; - } - - @Override public Context groupContext() { - return new AgentJVMModuleContext(GROUP_NAME); - } - - @Override public ModuleInstaller moduleInstaller() { - return installer; - } -} diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCConfig.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCConfig.java deleted file mode 100644 index 1376814cd..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCConfig.java +++ /dev/null @@ -1,9 +0,0 @@ -package org.skywalking.apm.collector.agentjvm.grpc; - -/** - * @author pengys5 - */ -public class AgentJVMGRPCConfig { - public static String HOST; - public static int PORT; -} diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCConfigParser.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCConfigParser.java deleted file mode 100644 index f66ea0a8a..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCConfigParser.java +++ /dev/null @@ -1,30 +0,0 @@ -package org.skywalking.apm.collector.agentjvm.grpc; - -import java.util.Map; -import org.skywalking.apm.collector.core.config.ConfigParseException; -import org.skywalking.apm.collector.core.module.ModuleConfigParser; -import org.skywalking.apm.collector.core.util.ObjectUtils; -import org.skywalking.apm.collector.core.util.StringUtils; - -/** - * @author pengys5 - */ -public class AgentJVMGRPCConfigParser implements ModuleConfigParser { - - private static final String HOST = "host"; - private static final String PORT = "port"; - - @Override public void parse(Map config) throws ConfigParseException { - if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(HOST))) { - AgentJVMGRPCConfig.HOST = "localhost"; - } else { - AgentJVMGRPCConfig.HOST = (String)config.get(HOST); - } - - if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(PORT))) { - AgentJVMGRPCConfig.PORT = 11800; - } else { - AgentJVMGRPCConfig.PORT = (Integer)config.get(PORT); - } - } -} 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 deleted file mode 100644 index 45b21b59a..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCDataListener.java +++ /dev/null @@ -1,25 +0,0 @@ -package org.skywalking.apm.collector.agentjvm.grpc; - -import org.skywalking.apm.collector.agentjvm.AgentJVMModuleGroupDefine; -import org.skywalking.apm.collector.cluster.ClusterModuleDefine; -import org.skywalking.apm.collector.core.cluster.ClusterDataListener; - -/** - * @author pengys5 - */ -public class AgentJVMGRPCDataListener extends ClusterDataListener { - - public static final String PATH = ClusterModuleDefine.BASE_CATALOG + "." + AgentJVMModuleGroupDefine.GROUP_NAME + "." + AgentJVMGRPCModuleDefine.MODULE_NAME; - - @Override public String path() { - return PATH; - } - - @Override public void serverJoinNotify(String serverAddress) { - - } - - @Override public void serverQuitNotify(String serverAddress) { - - } -} diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCModuleDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCModuleDefine.java deleted file mode 100644 index 358adedfb..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCModuleDefine.java +++ /dev/null @@ -1,51 +0,0 @@ -package org.skywalking.apm.collector.agentjvm.grpc; - -import java.util.LinkedList; -import java.util.List; -import org.skywalking.apm.collector.agentjvm.AgentJVMModuleDefine; -import org.skywalking.apm.collector.agentjvm.AgentJVMModuleGroupDefine; -import org.skywalking.apm.collector.agentjvm.grpc.handler.JVMMetricsServiceHandler; -import org.skywalking.apm.collector.core.cluster.ClusterDataListener; -import org.skywalking.apm.collector.core.framework.Handler; -import org.skywalking.apm.collector.core.module.ModuleConfigParser; -import org.skywalking.apm.collector.core.module.ModuleRegistration; -import org.skywalking.apm.collector.core.server.Server; -import org.skywalking.apm.collector.server.grpc.GRPCServer; - -/** - * @author pengys5 - */ -public class AgentJVMGRPCModuleDefine extends AgentJVMModuleDefine { - - public static final String MODULE_NAME = "grpc"; - - @Override protected String group() { - return AgentJVMModuleGroupDefine.GROUP_NAME; - } - - @Override public String name() { - return MODULE_NAME; - } - - @Override protected ModuleConfigParser configParser() { - return new AgentJVMGRPCConfigParser(); - } - - @Override protected Server server() { - return new GRPCServer(AgentJVMGRPCConfig.HOST, AgentJVMGRPCConfig.PORT); - } - - @Override protected ModuleRegistration registration() { - return new AgentJVMGRPCModuleRegistration(); - } - - @Override public ClusterDataListener listener() { - return new AgentJVMGRPCDataListener(); - } - - @Override public List handlerList() { - List handlers = new LinkedList<>(); - handlers.add(new JVMMetricsServiceHandler()); - return handlers; - } -} diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCModuleRegistration.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCModuleRegistration.java deleted file mode 100644 index ca249dea7..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/grpc/AgentJVMGRPCModuleRegistration.java +++ /dev/null @@ -1,13 +0,0 @@ -package org.skywalking.apm.collector.agentjvm.grpc; - -import org.skywalking.apm.collector.core.module.ModuleRegistration; - -/** - * @author pengys5 - */ -public class AgentJVMGRPCModuleRegistration extends ModuleRegistration { - - @Override public Value buildValue() { - return new Value(AgentJVMGRPCConfig.HOST, AgentJVMGRPCConfig.PORT, null); - } -} diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/define/CpuMetricEsTableDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/define/CpuMetricEsTableDefine.java index cffcf2184..b08d94685 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/define/CpuMetricEsTableDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/cpu/define/CpuMetricEsTableDefine.java @@ -17,14 +17,6 @@ public class CpuMetricEsTableDefine extends ElasticSearchTableDefine { return 1; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(CpuMetricTable.COLUMN_USAGE_PERCENT, ElasticSearchColumnDefine.Type.Double.name())); diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java index 677bd4889..381ecde82 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/gc/define/GCMetricEsTableDefine.java @@ -17,14 +17,6 @@ public class GCMetricEsTableDefine extends ElasticSearchTableDefine { return 1; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(GCMetricTable.COLUMN_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(GCMetricTable.COLUMN_PHRASE, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/define/MemoryMetricEsTableDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/define/MemoryMetricEsTableDefine.java index 6d65b4499..8b347e5a4 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/define/MemoryMetricEsTableDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memory/define/MemoryMetricEsTableDefine.java @@ -17,14 +17,6 @@ public class MemoryMetricEsTableDefine extends ElasticSearchTableDefine { return 1; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(MemoryMetricTable.COLUMN_APPLICATION_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(MemoryMetricTable.COLUMN_IS_HEAP, ElasticSearchColumnDefine.Type.Boolean.name())); diff --git a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/define/MemoryPoolMetricEsTableDefine.java b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/define/MemoryPoolMetricEsTableDefine.java index 52c33ac4d..af14122dc 100644 --- a/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/define/MemoryPoolMetricEsTableDefine.java +++ b/apm-collector/apm-collector-agentjvm/src/main/java/org/skywalking/apm/collector/agentjvm/worker/memorypool/define/MemoryPoolMetricEsTableDefine.java @@ -17,14 +17,6 @@ public class MemoryPoolMetricEsTableDefine extends ElasticSearchTableDefine { return 1; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(MemoryPoolMetricTable.COLUMN_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(MemoryPoolMetricTable.COLUMN_POOL_TYPE, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/group.define b/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/group.define deleted file mode 100644 index 0b2741b98..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/group.define +++ /dev/null @@ -1 +0,0 @@ -org.skywalking.apm.collector.agentjvm.AgentJVMModuleGroupDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/module.define b/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/module.define deleted file mode 100644 index 78d6e4a41..000000000 --- a/apm-collector/apm-collector-agentjvm/src/main/resources/META-INF/defines/module.define +++ /dev/null @@ -1 +0,0 @@ -org.skywalking.apm.collector.agentjvm.grpc.AgentJVMGRPCModuleDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentregister/pom.xml b/apm-collector/apm-collector-agentregister/pom.xml index 258ffa456..da7f97c64 100644 --- a/apm-collector/apm-collector-agentregister/pom.xml +++ b/apm-collector/apm-collector-agentregister/pom.xml @@ -13,6 +13,11 @@ jar + + org.skywalking + apm-collector-stream + ${project.version} + org.skywalking apm-collector-cluster @@ -28,10 +33,5 @@ apm-network ${project.version} - - org.skywalking - apm-collector-agentstream - ${project.version} - \ No newline at end of file diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterCommonModuleInstaller.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterCommonModuleInstaller.java deleted file mode 100644 index 2d8850c61..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterCommonModuleInstaller.java +++ /dev/null @@ -1,23 +0,0 @@ -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.MultipleCommonModuleInstaller; - -/** - * @author pengys5 - */ -public class AgentRegisterCommonModuleInstaller extends MultipleCommonModuleInstaller { - - @Override public String groupName() { - return AgentRegisterModuleGroupDefine.GROUP_NAME; - } - - @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/AgentRegisterModuleContext.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleContext.java deleted file mode 100644 index 0bea80ed0..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleContext.java +++ /dev/null @@ -1,12 +0,0 @@ -package org.skywalking.apm.collector.agentregister; - -import org.skywalking.apm.collector.core.framework.Context; - -/** - * @author pengys5 - */ -public class AgentRegisterModuleContext extends Context { - public AgentRegisterModuleContext(String groupName) { - super(groupName); - } -} 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 deleted file mode 100644 index fd627cd88..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleDefine.java +++ /dev/null @@ -1,23 +0,0 @@ -package org.skywalking.apm.collector.agentregister; - -import org.skywalking.apm.collector.core.client.Client; -import org.skywalking.apm.collector.core.cluster.ClusterDataListenerDefine; -import org.skywalking.apm.collector.core.module.ModuleDefine; - -/** - * @author pengys5 - */ -public abstract class AgentRegisterModuleDefine extends ModuleDefine implements ClusterDataListenerDefine { - - @Override protected void initializeOtherContext() { - - } - - @Override protected final Client createClient() { - throw new UnsupportedOperationException(""); - } - - @Override public final boolean defaultModule() { - return true; - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleException.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleException.java deleted file mode 100644 index 061937a7e..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleException.java +++ /dev/null @@ -1,16 +0,0 @@ -package org.skywalking.apm.collector.agentregister; - -import org.skywalking.apm.collector.core.module.ModuleException; - -/** - * @author pengys5 - */ -public class AgentRegisterModuleException extends ModuleException { - public AgentRegisterModuleException(String message) { - super(message); - } - - public AgentRegisterModuleException(String message, Throwable cause) { - super(message, cause); - } -} 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 deleted file mode 100644 index 0aa99ade7..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/AgentRegisterModuleGroupDefine.java +++ /dev/null @@ -1,31 +0,0 @@ -package org.skywalking.apm.collector.agentregister; - -import org.skywalking.apm.collector.core.framework.Context; -import org.skywalking.apm.collector.core.module.ModuleGroupDefine; -import org.skywalking.apm.collector.core.module.ModuleInstaller; - -/** - * @author pengys5 - */ -public class 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; - } - - @Override public Context groupContext() { - return new AgentRegisterModuleContext(GROUP_NAME); - } - - @Override public ModuleInstaller moduleInstaller() { - return installer; - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java index 438b43d8b..95aa422f8 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/application/ApplicationIDService.java @@ -1,9 +1,9 @@ package org.skywalking.apm.collector.agentregister.application; -import org.skywalking.apm.collector.agentstream.worker.cache.ApplicationCache; -import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; -import org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterRemoteWorker; +import org.skywalking.apm.collector.agentregister.worker.application.ApplicationRegisterRemoteWorker; +import org.skywalking.apm.collector.agentregister.worker.cache.ApplicationCache; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; +import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; import org.skywalking.apm.collector.stream.StreamModuleContext; import org.skywalking.apm.collector.stream.StreamModuleGroupDefine; import org.skywalking.apm.collector.stream.worker.WorkerInvokeException; diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCConfig.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCConfig.java deleted file mode 100644 index 2d2ad6bc7..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCConfig.java +++ /dev/null @@ -1,9 +0,0 @@ -package org.skywalking.apm.collector.agentregister.grpc; - -/** - * @author pengys5 - */ -public class AgentRegisterGRPCConfig { - public static String HOST; - public static int PORT; -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCConfigParser.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCConfigParser.java deleted file mode 100644 index 90f0d39da..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCConfigParser.java +++ /dev/null @@ -1,30 +0,0 @@ -package org.skywalking.apm.collector.agentregister.grpc; - -import java.util.Map; -import org.skywalking.apm.collector.core.config.ConfigParseException; -import org.skywalking.apm.collector.core.module.ModuleConfigParser; -import org.skywalking.apm.collector.core.util.ObjectUtils; -import org.skywalking.apm.collector.core.util.StringUtils; - -/** - * @author pengys5 - */ -public class AgentRegisterGRPCConfigParser implements ModuleConfigParser { - - private static final String HOST = "host"; - private static final String PORT = "port"; - - @Override public void parse(Map config) throws ConfigParseException { - if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(HOST))) { - AgentRegisterGRPCConfig.HOST = "localhost"; - } else { - AgentRegisterGRPCConfig.HOST = (String)config.get(HOST); - } - - if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(PORT))) { - AgentRegisterGRPCConfig.PORT = 11800; - } else { - AgentRegisterGRPCConfig.PORT = (Integer)config.get(PORT); - } - } -} 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 deleted file mode 100644 index 23667b5a7..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCDataListener.java +++ /dev/null @@ -1,25 +0,0 @@ -package org.skywalking.apm.collector.agentregister.grpc; - -import org.skywalking.apm.collector.agentregister.AgentRegisterModuleGroupDefine; -import org.skywalking.apm.collector.cluster.ClusterModuleDefine; -import org.skywalking.apm.collector.core.cluster.ClusterDataListener; - -/** - * @author pengys5 - */ -public class AgentRegisterGRPCDataListener extends ClusterDataListener { - - public static final String PATH = ClusterModuleDefine.BASE_CATALOG + "." + AgentRegisterModuleGroupDefine.GROUP_NAME + "." + AgentRegisterGRPCModuleDefine.MODULE_NAME; - - @Override public String path() { - return PATH; - } - - @Override public void serverJoinNotify(String serverAddress) { - - } - - @Override public void serverQuitNotify(String serverAddress) { - - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCModuleDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCModuleDefine.java deleted file mode 100644 index 509009c37..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCModuleDefine.java +++ /dev/null @@ -1,55 +0,0 @@ -package org.skywalking.apm.collector.agentregister.grpc; - -import java.util.LinkedList; -import java.util.List; -import org.skywalking.apm.collector.agentregister.AgentRegisterModuleDefine; -import org.skywalking.apm.collector.agentregister.AgentRegisterModuleGroupDefine; -import org.skywalking.apm.collector.agentregister.grpc.handler.ApplicationRegisterServiceHandler; -import org.skywalking.apm.collector.agentregister.grpc.handler.InstanceDiscoveryServiceHandler; -import org.skywalking.apm.collector.agentregister.grpc.handler.ServiceNameDiscoveryServiceHandler; -import org.skywalking.apm.collector.core.cluster.ClusterDataListener; -import org.skywalking.apm.collector.core.framework.Handler; -import org.skywalking.apm.collector.core.module.ModuleConfigParser; -import org.skywalking.apm.collector.core.module.ModuleRegistration; -import org.skywalking.apm.collector.core.server.Server; -import org.skywalking.apm.collector.server.grpc.GRPCServer; - -/** - * @author pengys5 - */ -public class AgentRegisterGRPCModuleDefine extends AgentRegisterModuleDefine { - - public static final String MODULE_NAME = "grpc"; - - @Override protected String group() { - return AgentRegisterModuleGroupDefine.GROUP_NAME; - } - - @Override public String name() { - return MODULE_NAME; - } - - @Override protected ModuleConfigParser configParser() { - return new AgentRegisterGRPCConfigParser(); - } - - @Override protected Server server() { - return new GRPCServer(AgentRegisterGRPCConfig.HOST, AgentRegisterGRPCConfig.PORT); - } - - @Override protected ModuleRegistration registration() { - return new AgentRegisterGRPCModuleRegistration(); - } - - @Override public ClusterDataListener listener() { - return new AgentRegisterGRPCDataListener(); - } - - @Override public List handlerList() { - List handlers = new LinkedList<>(); - handlers.add(new ApplicationRegisterServiceHandler()); - handlers.add(new InstanceDiscoveryServiceHandler()); - handlers.add(new ServiceNameDiscoveryServiceHandler()); - return handlers; - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCModuleRegistration.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCModuleRegistration.java deleted file mode 100644 index bc97dba60..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/grpc/AgentRegisterGRPCModuleRegistration.java +++ /dev/null @@ -1,13 +0,0 @@ -package org.skywalking.apm.collector.agentregister.grpc; - -import org.skywalking.apm.collector.core.module.ModuleRegistration; - -/** - * @author pengys5 - */ -public class AgentRegisterGRPCModuleRegistration extends ModuleRegistration { - - @Override public Value buildValue() { - return new Value(AgentRegisterGRPCConfig.HOST, AgentRegisterGRPCConfig.PORT, null); - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java index 85ba57eac..0fdf1ba16 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/instance/InstanceIDService.java @@ -1,7 +1,7 @@ package org.skywalking.apm.collector.agentregister.instance; -import org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterRemoteWorker; -import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.IInstanceDAO; +import org.skywalking.apm.collector.agentregister.worker.instance.InstanceRegisterRemoteWorker; +import org.skywalking.apm.collector.agentregister.worker.instance.dao.IInstanceDAO; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfig.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfig.java deleted file mode 100644 index cf1e45207..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfig.java +++ /dev/null @@ -1,10 +0,0 @@ -package org.skywalking.apm.collector.agentregister.jetty; - -/** - * @author pengys5 - */ -public class AgentRegisterJettyConfig { - public static String HOST; - public static int PORT; - public static String CONTEXT_PATH; -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfigParser.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfigParser.java deleted file mode 100644 index 65a10bebe..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyConfigParser.java +++ /dev/null @@ -1,36 +0,0 @@ -package org.skywalking.apm.collector.agentregister.jetty; - -import java.util.Map; -import org.skywalking.apm.collector.core.config.ConfigParseException; -import org.skywalking.apm.collector.core.module.ModuleConfigParser; -import org.skywalking.apm.collector.core.util.ObjectUtils; -import org.skywalking.apm.collector.core.util.StringUtils; - -/** - * @author pengys5 - */ -public class AgentRegisterJettyConfigParser implements ModuleConfigParser { - - private static final String HOST = "host"; - private static final String PORT = "port"; - public static final String CONTEXT_PATH = "contextPath"; - - @Override public void parse(Map config) throws ConfigParseException { - AgentRegisterJettyConfig.CONTEXT_PATH = "/"; - - if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(HOST))) { - AgentRegisterJettyConfig.HOST = "localhost"; - } else { - AgentRegisterJettyConfig.HOST = (String)config.get(HOST); - } - - if (ObjectUtils.isEmpty(config) || StringUtils.isEmpty(config.get(PORT))) { - AgentRegisterJettyConfig.PORT = 12800; - } else { - AgentRegisterJettyConfig.PORT = (Integer)config.get(PORT); - } - if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(CONTEXT_PATH))) { - AgentRegisterJettyConfig.CONTEXT_PATH = (String)config.get(CONTEXT_PATH); - } - } -} 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 deleted file mode 100644 index aabcfa8a9..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyDataListener.java +++ /dev/null @@ -1,25 +0,0 @@ -package org.skywalking.apm.collector.agentregister.jetty; - -import org.skywalking.apm.collector.agentregister.AgentRegisterModuleGroupDefine; -import org.skywalking.apm.collector.cluster.ClusterModuleDefine; -import org.skywalking.apm.collector.core.cluster.ClusterDataListener; - -/** - * @author pengys5 - */ -public class AgentRegisterJettyDataListener extends ClusterDataListener { - - public static final String PATH = ClusterModuleDefine.BASE_CATALOG + "." + AgentRegisterModuleGroupDefine.GROUP_NAME + "." + AgentRegisterJettyModuleDefine.MODULE_NAME; - - @Override public String path() { - return PATH; - } - - @Override public void serverJoinNotify(String serverAddress) { - - } - - @Override public void serverQuitNotify(String serverAddress) { - - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleDefine.java deleted file mode 100644 index 413149d71..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleDefine.java +++ /dev/null @@ -1,53 +0,0 @@ -package org.skywalking.apm.collector.agentregister.jetty; - -import java.util.LinkedList; -import java.util.List; -import org.skywalking.apm.collector.agentregister.AgentRegisterModuleDefine; -import org.skywalking.apm.collector.agentregister.AgentRegisterModuleGroupDefine; -import org.skywalking.apm.collector.agentregister.jetty.handler.ApplicationRegisterServletHandler; -import org.skywalking.apm.collector.agentregister.jetty.handler.InstanceDiscoveryServletHandler; -import org.skywalking.apm.collector.core.cluster.ClusterDataListener; -import org.skywalking.apm.collector.core.framework.Handler; -import org.skywalking.apm.collector.core.module.ModuleConfigParser; -import org.skywalking.apm.collector.core.module.ModuleRegistration; -import org.skywalking.apm.collector.core.server.Server; -import org.skywalking.apm.collector.server.jetty.JettyServer; - -/** - * @author pengys5 - */ -public class AgentRegisterJettyModuleDefine extends AgentRegisterModuleDefine { - - public static final String MODULE_NAME = "jetty"; - - @Override protected String group() { - return AgentRegisterModuleGroupDefine.GROUP_NAME; - } - - @Override public String name() { - return MODULE_NAME; - } - - @Override protected ModuleConfigParser configParser() { - return new AgentRegisterJettyConfigParser(); - } - - @Override protected Server server() { - return new JettyServer(AgentRegisterJettyConfig.HOST, AgentRegisterJettyConfig.PORT, AgentRegisterJettyConfig.CONTEXT_PATH); - } - - @Override protected ModuleRegistration registration() { - return new AgentRegisterJettyModuleRegistration(); - } - - @Override public ClusterDataListener listener() { - return new AgentRegisterJettyDataListener(); - } - - @Override public List handlerList() { - List handlers = new LinkedList<>(); - handlers.add(new ApplicationRegisterServletHandler()); - handlers.add(new InstanceDiscoveryServletHandler()); - return handlers; - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleRegistration.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleRegistration.java deleted file mode 100644 index d6e2d260a..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/AgentRegisterJettyModuleRegistration.java +++ /dev/null @@ -1,13 +0,0 @@ -package org.skywalking.apm.collector.agentregister.jetty; - -import org.skywalking.apm.collector.core.module.ModuleRegistration; - -/** - * @author pengys5 - */ -public class AgentRegisterJettyModuleRegistration extends ModuleRegistration { - - @Override public Value buildValue() { - return new Value(AgentRegisterJettyConfig.HOST, AgentRegisterJettyConfig.PORT, AgentRegisterJettyConfig.CONTEXT_PATH); - } -} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/ServiceNameDiscoveryServiceHandler.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/ServiceNameDiscoveryServiceHandler.java new file mode 100644 index 000000000..91aa68a2f --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/jetty/handler/ServiceNameDiscoveryServiceHandler.java @@ -0,0 +1,59 @@ +package org.skywalking.apm.collector.agentregister.jetty.handler; + +import com.google.gson.Gson; +import com.google.gson.JsonArray; +import com.google.gson.JsonElement; +import com.google.gson.JsonObject; +import java.io.IOException; +import javax.servlet.http.HttpServletRequest; +import org.skywalking.apm.collector.agentregister.servicename.ServiceNameService; +import org.skywalking.apm.collector.server.jetty.ArgumentsParseException; +import org.skywalking.apm.collector.server.jetty.JettyHandler; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + +/** + * @author pengys5 + */ +public class ServiceNameDiscoveryServiceHandler extends JettyHandler { + + private final Logger logger = LoggerFactory.getLogger(ServiceNameDiscoveryServiceHandler.class); + + private ServiceNameService serviceNameService = new ServiceNameService(); + private Gson gson = new Gson(); + + @Override public String pathSpec() { + return "/servicename/discovery"; + } + + private static final String APPLICATION_ID = "ai"; + private static final String SERVICE_NAME = "sn"; + private static final String SERVICE_ID = "si"; + private static final String ELEMENT = "el"; + + @Override protected JsonElement doGet(HttpServletRequest req) throws ArgumentsParseException { + throw new UnsupportedOperationException(); + } + + @Override protected JsonElement doPost(HttpServletRequest req) throws ArgumentsParseException { + JsonArray responseArray = new JsonArray(); + try { + JsonArray services = gson.fromJson(req.getReader(), JsonArray.class); + for (JsonElement service : services) { + int applicationId = service.getAsJsonObject().get(APPLICATION_ID).getAsInt(); + String serviceName = service.getAsJsonObject().get(SERVICE_NAME).getAsString(); + + int serviceId = serviceNameService.getOrCreate(applicationId, serviceName); + if (serviceId != 0) { + JsonObject responseJson = new JsonObject(); + responseJson.addProperty(SERVICE_ID, serviceId); + responseJson.add(ELEMENT, service); + responseArray.add(responseJson); + } + } + } catch (IOException e) { + logger.error(e.getMessage(), e); + } + return responseArray; + } +} diff --git a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/servicename/ServiceNameService.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/servicename/ServiceNameService.java index 25f092d85..ae5a51ef9 100644 --- a/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/servicename/ServiceNameService.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/servicename/ServiceNameService.java @@ -1,8 +1,8 @@ package org.skywalking.apm.collector.agentregister.servicename; import org.skywalking.apm.collector.storage.define.register.ServiceNameDataDefine; -import org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterRemoteWorker; -import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO; +import org.skywalking.apm.collector.agentregister.worker.servicename.ServiceNameRegisterRemoteWorker; +import org.skywalking.apm.collector.agentregister.worker.servicename.dao.IServiceNameDAO; import org.skywalking.apm.collector.core.framework.CollectorContextHelper; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.stream.StreamModuleContext; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/IdAutoIncrement.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/IdAutoIncrement.java similarity index 88% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/IdAutoIncrement.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/IdAutoIncrement.java index d9864e4e7..808d53f65 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/IdAutoIncrement.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/IdAutoIncrement.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register; +package org.skywalking.apm.collector.agentregister.worker; /** * @author pengys5 diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationEsTableDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationEsTableDefine.java similarity index 79% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationEsTableDefine.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationEsTableDefine.java index fe070606f..968c3f84a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationEsTableDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationEsTableDefine.java @@ -1,8 +1,8 @@ -package org.skywalking.apm.collector.agentstream.worker.register.application; +package org.skywalking.apm.collector.agentregister.worker.application; +import org.skywalking.apm.collector.storage.define.register.ApplicationTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; -import org.skywalking.apm.collector.storage.define.register.ApplicationTable; /** * @author pengys5 @@ -17,14 +17,6 @@ public class ApplicationEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(ApplicationTable.COLUMN_APPLICATION_CODE, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(ApplicationTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationH2TableDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationH2TableDefine.java similarity index 89% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationH2TableDefine.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationH2TableDefine.java index 4eac88e63..cb3892cd2 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationH2TableDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationH2TableDefine.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.application; +package org.skywalking.apm.collector.agentregister.worker.application; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterRemoteWorker.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationRegisterRemoteWorker.java similarity index 96% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterRemoteWorker.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationRegisterRemoteWorker.java index 33564f8af..41f70e3c5 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterRemoteWorker.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationRegisterRemoteWorker.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.application; +package org.skywalking.apm.collector.agentregister.worker.application; import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorker; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterSerialWorker.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationRegisterSerialWorker.java similarity index 93% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterSerialWorker.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationRegisterSerialWorker.java index 79ef0c6ba..caa3a3af1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/ApplicationRegisterSerialWorker.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/ApplicationRegisterSerialWorker.java @@ -1,8 +1,8 @@ -package org.skywalking.apm.collector.agentstream.worker.register.application; +package org.skywalking.apm.collector.agentregister.worker.application; import org.skywalking.apm.collector.core.util.Const; -import org.skywalking.apm.collector.agentstream.worker.register.IdAutoIncrement; -import org.skywalking.apm.collector.agentstream.worker.register.application.dao.IApplicationDAO; +import org.skywalking.apm.collector.agentregister.worker.IdAutoIncrement; +import org.skywalking.apm.collector.agentregister.worker.application.dao.IApplicationDAO; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractLocalAsyncWorker; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/ApplicationEsDAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationEsDAO.java similarity index 97% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/ApplicationEsDAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationEsDAO.java index bbb03b629..76ad11394 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/ApplicationEsDAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationEsDAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.application.dao; +package org.skywalking.apm.collector.agentregister.worker.application.dao; import java.util.HashMap; import java.util.Map; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/ApplicationH2DAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationH2DAO.java similarity index 89% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/ApplicationH2DAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationH2DAO.java index 164949088..0081c0294 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/ApplicationH2DAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/ApplicationH2DAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.application.dao; +package org.skywalking.apm.collector.agentregister.worker.application.dao; import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; import org.skywalking.apm.collector.client.h2.H2Client; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/IApplicationDAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/IApplicationDAO.java similarity index 79% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/IApplicationDAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/IApplicationDAO.java index cc409b606..d85e8407d 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/application/dao/IApplicationDAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/application/dao/IApplicationDAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.application.dao; +package org.skywalking.apm.collector.agentregister.worker.application.dao; import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/cache/ApplicationCache.java similarity index 87% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/cache/ApplicationCache.java index a38cb3976..148991ce9 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ApplicationCache.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/cache/ApplicationCache.java @@ -1,8 +1,8 @@ -package org.skywalking.apm.collector.agentstream.worker.cache; +package org.skywalking.apm.collector.agentregister.worker.cache; import com.google.common.cache.Cache; import com.google.common.cache.CacheBuilder; -import org.skywalking.apm.collector.agentstream.worker.register.application.dao.IApplicationDAO; +import org.skywalking.apm.collector.agentregister.worker.application.dao.IApplicationDAO; import org.skywalking.apm.collector.storage.dao.DAOContainer; /** diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceEsTableDefine.java similarity index 86% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceEsTableDefine.java index 7bc301426..3a6c95209 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceEsTableDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceEsTableDefine.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.instance; +package org.skywalking.apm.collector.agentregister.worker.instance; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; @@ -17,14 +17,6 @@ public class InstanceEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(InstanceTable.COLUMN_AGENT_UUID, ElasticSearchColumnDefine.Type.Keyword.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceH2TableDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceH2TableDefine.java similarity index 93% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceH2TableDefine.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceH2TableDefine.java index ad55a7e21..1788de6b4 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceH2TableDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceH2TableDefine.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.instance; +package org.skywalking.apm.collector.agentregister.worker.instance; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceRegisterRemoteWorker.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceRegisterRemoteWorker.java similarity index 97% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceRegisterRemoteWorker.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceRegisterRemoteWorker.java index b35409300..022114af8 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceRegisterRemoteWorker.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceRegisterRemoteWorker.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.instance; +package org.skywalking.apm.collector.agentregister.worker.instance; import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorker; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceRegisterSerialWorker.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceRegisterSerialWorker.java similarity index 95% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceRegisterSerialWorker.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceRegisterSerialWorker.java index a0ad4209c..bea9b677e 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/InstanceRegisterSerialWorker.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/InstanceRegisterSerialWorker.java @@ -1,6 +1,6 @@ -package org.skywalking.apm.collector.agentstream.worker.register.instance; +package org.skywalking.apm.collector.agentregister.worker.instance; -import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.IInstanceDAO; +import org.skywalking.apm.collector.agentregister.worker.instance.dao.IInstanceDAO; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.storage.define.DataDefine; import org.skywalking.apm.collector.storage.define.register.ApplicationDataDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/IInstanceDAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/IInstanceDAO.java similarity index 84% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/IInstanceDAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/IInstanceDAO.java index fc318ae60..2f3b08c85 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/IInstanceDAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/IInstanceDAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.instance.dao; +package org.skywalking.apm.collector.agentregister.worker.instance.dao; import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceEsDAO.java similarity index 98% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceEsDAO.java index c2ffbdc0d..3de62b1ce 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceEsDAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceEsDAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.instance.dao; +package org.skywalking.apm.collector.agentregister.worker.instance.dao; import java.util.HashMap; import java.util.Map; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceH2DAO.java similarity index 90% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceH2DAO.java index 544513e74..c11920a7a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/instance/dao/InstanceH2DAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/instance/dao/InstanceH2DAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.instance.dao; +package org.skywalking.apm.collector.agentregister.worker.instance.dao; import org.skywalking.apm.collector.storage.define.register.InstanceDataDefine; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameEsTableDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameEsTableDefine.java similarity index 81% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameEsTableDefine.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameEsTableDefine.java index 16e4f5596..794e794a3 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameEsTableDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameEsTableDefine.java @@ -1,8 +1,8 @@ -package org.skywalking.apm.collector.agentstream.worker.register.servicename; +package org.skywalking.apm.collector.agentregister.worker.servicename; +import org.skywalking.apm.collector.storage.define.register.ServiceNameTable; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchColumnDefine; import org.skywalking.apm.collector.storage.elasticsearch.define.ElasticSearchTableDefine; -import org.skywalking.apm.collector.storage.define.register.ServiceNameTable; /** * @author pengys5 @@ -17,14 +17,6 @@ public class ServiceNameEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(ServiceNameTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(ServiceNameTable.COLUMN_SERVICE_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameH2TableDefine.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameH2TableDefine.java similarity index 90% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameH2TableDefine.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameH2TableDefine.java index 797db44f7..23cec8fc3 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameH2TableDefine.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameH2TableDefine.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.servicename; +package org.skywalking.apm.collector.agentregister.worker.servicename; import org.skywalking.apm.collector.storage.h2.define.H2ColumnDefine; import org.skywalking.apm.collector.storage.h2.define.H2TableDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterRemoteWorker.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameRegisterRemoteWorker.java similarity index 96% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterRemoteWorker.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameRegisterRemoteWorker.java index 9b4822be1..03e9b90b5 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterRemoteWorker.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameRegisterRemoteWorker.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.servicename; +package org.skywalking.apm.collector.agentregister.worker.servicename; import org.skywalking.apm.collector.storage.define.register.ServiceNameDataDefine; import org.skywalking.apm.collector.stream.worker.AbstractRemoteWorker; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameRegisterSerialWorker.java similarity index 93% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameRegisterSerialWorker.java index 3f9ba4b31..1c8e0c6c4 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/ServiceNameRegisterSerialWorker.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/ServiceNameRegisterSerialWorker.java @@ -1,7 +1,7 @@ -package org.skywalking.apm.collector.agentstream.worker.register.servicename; +package org.skywalking.apm.collector.agentregister.worker.servicename; -import org.skywalking.apm.collector.agentstream.worker.register.IdAutoIncrement; -import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO; +import org.skywalking.apm.collector.agentregister.worker.IdAutoIncrement; +import org.skywalking.apm.collector.agentregister.worker.servicename.dao.IServiceNameDAO; import org.skywalking.apm.collector.core.util.Const; import org.skywalking.apm.collector.storage.dao.DAOContainer; import org.skywalking.apm.collector.storage.define.DataDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/IServiceNameDAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/IServiceNameDAO.java similarity index 81% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/IServiceNameDAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/IServiceNameDAO.java index b02ddf180..c4b8635fa 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/IServiceNameDAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/IServiceNameDAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.servicename.dao; +package org.skywalking.apm.collector.agentregister.worker.servicename.dao; import org.skywalking.apm.collector.storage.define.register.ServiceNameDataDefine; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/ServiceNameEsDAO.java similarity index 97% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/ServiceNameEsDAO.java index a2e3c63d8..dbc3f6be9 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameEsDAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/ServiceNameEsDAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.servicename.dao; +package org.skywalking.apm.collector.agentregister.worker.servicename.dao; import java.util.HashMap; import java.util.Map; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameH2DAO.java b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/ServiceNameH2DAO.java similarity index 89% rename from apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameH2DAO.java rename to apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/ServiceNameH2DAO.java index c73b9687f..1c463e798 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/register/servicename/dao/ServiceNameH2DAO.java +++ b/apm-collector/apm-collector-agentregister/src/main/java/org/skywalking/apm/collector/agentregister/worker/servicename/dao/ServiceNameH2DAO.java @@ -1,4 +1,4 @@ -package org.skywalking.apm.collector.agentstream.worker.register.servicename.dao; +package org.skywalking.apm.collector.agentregister.worker.servicename.dao; import org.skywalking.apm.collector.storage.define.register.ServiceNameDataDefine; import org.skywalking.apm.collector.storage.h2.dao.H2DAO; diff --git a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/es_dao.define new file mode 100644 index 000000000..90d0e020d --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/es_dao.define @@ -0,0 +1,3 @@ +org.skywalking.apm.collector.agentregister.worker.application.dao.ApplicationEsDAO +org.skywalking.apm.collector.agentregister.worker.instance.dao.InstanceEsDAO +org.skywalking.apm.collector.agentregister.worker.servicename.dao.ServiceNameEsDAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/group.define b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/group.define deleted file mode 100644 index 0afe036bd..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/group.define +++ /dev/null @@ -1 +0,0 @@ -org.skywalking.apm.collector.agentregister.AgentRegisterModuleGroupDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/h2_dao.define new file mode 100644 index 000000000..f381c6958 --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/h2_dao.define @@ -0,0 +1,3 @@ +org.skywalking.apm.collector.agentregister.worker.application.dao.ApplicationH2DAO +org.skywalking.apm.collector.agentregister.worker.instance.dao.InstanceH2DAO +org.skywalking.apm.collector.agentregister.worker.servicename.dao.ServiceNameH2DAO \ No newline at end of file diff --git a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/local_worker_provider.define b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/local_worker_provider.define new file mode 100644 index 000000000..758eed4ba --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/local_worker_provider.define @@ -0,0 +1,3 @@ +org.skywalking.apm.collector.agentregister.worker.application.ApplicationRegisterSerialWorker$Factory +org.skywalking.apm.collector.agentregister.worker.instance.InstanceRegisterSerialWorker$Factory +org.skywalking.apm.collector.agentregister.worker.servicename.ServiceNameRegisterSerialWorker$Factory \ No newline at end of file diff --git a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/module.define b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/module.define deleted file mode 100644 index e826d4160..000000000 --- a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/module.define +++ /dev/null @@ -1,2 +0,0 @@ -org.skywalking.apm.collector.agentregister.grpc.AgentRegisterGRPCModuleDefine -org.skywalking.apm.collector.agentregister.jetty.AgentRegisterJettyModuleDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/storage.define b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/storage.define new file mode 100644 index 000000000..52807b047 --- /dev/null +++ b/apm-collector/apm-collector-agentregister/src/main/resources/META-INF/defines/storage.define @@ -0,0 +1,8 @@ +org.skywalking.apm.collector.agentregister.worker.application.ApplicationEsTableDefine +org.skywalking.apm.collector.agentregister.worker.application.ApplicationH2TableDefine + +org.skywalking.apm.collector.agentregister.worker.instance.InstanceEsTableDefine +org.skywalking.apm.collector.agentregister.worker.instance.InstanceH2TableDefine + +org.skywalking.apm.collector.agentregister.worker.servicename.ServiceNameEsTableDefine +org.skywalking.apm.collector.agentregister.worker.servicename.ServiceNameH2TableDefine \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/pom.xml b/apm-collector/apm-collector-agentstream/pom.xml index 7cf497179..b9f463f33 100644 --- a/apm-collector/apm-collector-agentstream/pom.xml +++ b/apm-collector/apm-collector-agentstream/pom.xml @@ -33,5 +33,15 @@ apm-collector-storage ${project.version} + + org.skywalking + apm-collector-agentjvm + ${project.version} + + + org.skywalking + apm-collector-agentregister + ${project.version} + \ No newline at end of file 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 9ba5517a6..9a72418cf 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,10 +10,10 @@ 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; + private final AgentStreamModuleInstaller installer; public AgentStreamModuleGroupDefine() { - installer = new AgentStreamCommonModuleInstaller(); + installer = new AgentStreamModuleInstaller(); } @Override public String name() { diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java new file mode 100644 index 000000000..a493b5fab --- /dev/null +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/AgentStreamModuleInstaller.java @@ -0,0 +1,33 @@ +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.MultipleCommonModuleInstaller; +import org.skywalking.apm.collector.core.server.ServerException; + +/** + * @author pengys5 + */ +public class AgentStreamModuleInstaller extends MultipleCommonModuleInstaller { + + @Override public String groupName() { + return AgentStreamModuleGroupDefine.GROUP_NAME; + } + + @Override public Context moduleContext() { + return new AgentStreamModuleContext(groupName()); + } + + @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/grpc/AgentStreamGRPCModuleDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleDefine.java index 1cc1bdfad..186281ae1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/grpc/AgentStreamGRPCModuleDefine.java @@ -2,6 +2,10 @@ package org.skywalking.apm.collector.agentstream.grpc; import java.util.LinkedList; import java.util.List; +import org.skywalking.apm.collector.agentjvm.grpc.handler.JVMMetricsServiceHandler; +import org.skywalking.apm.collector.agentregister.grpc.handler.ApplicationRegisterServiceHandler; +import org.skywalking.apm.collector.agentregister.grpc.handler.InstanceDiscoveryServiceHandler; +import org.skywalking.apm.collector.agentregister.grpc.handler.ServiceNameDiscoveryServiceHandler; import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine; import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine; import org.skywalking.apm.collector.agentstream.grpc.handler.TraceSegmentServiceHandler; @@ -46,6 +50,10 @@ public class AgentStreamGRPCModuleDefine extends AgentStreamModuleDefine { @Override public List handlerList() { List handlers = new LinkedList<>(); handlers.add(new TraceSegmentServiceHandler()); + handlers.add(new ApplicationRegisterServiceHandler()); + handlers.add(new InstanceDiscoveryServiceHandler()); + handlers.add(new ServiceNameDiscoveryServiceHandler()); + handlers.add(new JVMMetricsServiceHandler()); return handlers; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java index 87428c339..1f4dec378 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/jetty/AgentStreamJettyModuleDefine.java @@ -2,6 +2,9 @@ package org.skywalking.apm.collector.agentstream.jetty; import java.util.LinkedList; import java.util.List; +import org.skywalking.apm.collector.agentregister.jetty.handler.ApplicationRegisterServletHandler; +import org.skywalking.apm.collector.agentregister.jetty.handler.InstanceDiscoveryServletHandler; +import org.skywalking.apm.collector.agentregister.jetty.handler.ServiceNameDiscoveryServiceHandler; import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine; import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine; import org.skywalking.apm.collector.agentstream.jetty.handler.TraceSegmentServletHandler; @@ -46,6 +49,9 @@ public class AgentStreamJettyModuleDefine extends AgentStreamModuleDefine { @Override public List handlerList() { List handlers = new LinkedList<>(); handlers.add(new TraceSegmentServletHandler()); + handlers.add(new ApplicationRegisterServletHandler()); + handlers.add(new InstanceDiscoveryServletHandler()); + handlers.add(new ServiceNameDiscoveryServiceHandler()); return handlers; } } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java index e2808ef6f..a1982042a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/InstanceCache.java @@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.cache; import com.google.common.cache.Cache; import com.google.common.cache.CacheBuilder; -import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.IInstanceDAO; +import org.skywalking.apm.collector.agentregister.worker.instance.dao.IInstanceDAO; import org.skywalking.apm.collector.storage.dao.DAOContainer; /** diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java index 1dcc065b0..bb29dee65 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/cache/ServiceCache.java @@ -2,7 +2,7 @@ package org.skywalking.apm.collector.agentstream.worker.cache; import com.google.common.cache.Cache; import com.google.common.cache.CacheBuilder; -import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.IServiceNameDAO; +import org.skywalking.apm.collector.agentregister.worker.servicename.dao.IServiceNameDAO; import org.skywalking.apm.collector.core.util.Const; import org.skywalking.apm.collector.storage.dao.DAOContainer; diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/CacheSizeConfig.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/CacheSizeConfig.java deleted file mode 100644 index fe792ff4f..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/CacheSizeConfig.java +++ /dev/null @@ -1,17 +0,0 @@ -package org.skywalking.apm.collector.agentstream.worker.config; - -/** - * @author pengys5 - */ -public class CacheSizeConfig { - - public static class Cache { - public static class Analysis { - public static int SIZE = 1024; - } - - public static class Persistence { - public static int SIZE = 5000; - } - } -} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/WorkerConfig.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/WorkerConfig.java deleted file mode 100644 index 9c94a6ae8..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/config/WorkerConfig.java +++ /dev/null @@ -1,125 +0,0 @@ -package org.skywalking.apm.collector.agentstream.worker.config; - -/** - * @author pengys5 - */ -public class WorkerConfig { - - public static class WorkerNum { - public static class Node { - public static class NodeCompAgg { - public static int VALUE = 2; - } - - public static class NodeMappingDayAgg { - public static int VALUE = 2; - } - - public static class NodeMappingHourAgg { - public static int VALUE = 2; - } - - public static class NodeMappingMinuteAgg { - public static int VALUE = 2; - } - } - - public static class NodeRef { - public static class NodeRefDayAgg { - public static int VALUE = 2; - } - - public static class NodeRefHourAgg { - public static int VALUE = 2; - } - - public static class NodeRefMinuteAgg { - public static int VALUE = 2; - } - - public static class NodeRefResSumDayAgg { - public static int VALUE = 2; - } - - public static class NodeRefResSumHourAgg { - public static int VALUE = 2; - } - - public static class NodeRefResSumMinuteAgg { - public static int VALUE = 2; - } - } - - public static class GlobalTrace { - public static class GlobalTraceAgg { - public static int VALUE = 2; - } - } - } - - public static class Queue { - public static class GlobalTrace { - public static class GlobalTraceAnalysis { - public static int SIZE = 1024; - } - } - - public static class Segment { - public static class SegmentAnalysis { - public static int SIZE = 1024; - } - - public static class SegmentCostAnalysis { - public static int SIZE = 4096; - } - - public static class SegmentExceptionAnalysis { - public static int SIZE = 4096; - } - } - - public static class Node { - public static class NodeCompAnalysis { - public static int SIZE = 1024; - } - - public static class NodeMappingDayAnalysis { - public static int SIZE = 1024; - } - - public static class NodeMappingHourAnalysis { - public static int SIZE = 1024; - } - - public static class NodeMappingMinuteAnalysis { - public static int SIZE = 1024; - } - } - - public static class NodeRef { - public static class NodeRefDayAnalysis { - public static int SIZE = 1024; - } - - public static class NodeRefHourAnalysis { - public static int SIZE = 1024; - } - - public static class NodeRefMinuteAnalysis { - public static int SIZE = 1024; - } - - public static class NodeRefResSumDayAnalysis { - public static int SIZE = 1024; - } - - public static class NodeRefResSumHourAnalysis { - public static int SIZE = 1024; - } - - public static class NodeRefResSumMinuteAnalysis { - public static int SIZE = 1024; - } - } - } -} diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java index 8a2b244a2..759635417 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/global/define/GlobalTraceEsTableDefine.java @@ -17,14 +17,6 @@ public class GlobalTraceEsTableDefine extends ElasticSearchTableDefine { return 5; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(GlobalTraceTable.COLUMN_SEGMENT_ID, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(GlobalTraceTable.COLUMN_GLOBAL_TRACE_ID, ElasticSearchColumnDefine.Type.Keyword.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java index f0afdd86c..073280cdd 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/instance/performance/define/InstPerformanceEsTableDefine.java @@ -17,14 +17,6 @@ public class InstPerformanceEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(InstPerformanceTable.COLUMN_INSTANCE_ID, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java index 1ded3d4fe..610bc9866 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/component/define/NodeComponentEsTableDefine.java @@ -17,14 +17,6 @@ public class NodeComponentEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_COMPONENT_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(NodeComponentTable.COLUMN_COMPONENT_NAME, ElasticSearchColumnDefine.Type.Keyword.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java index e1027b392..460abfc23 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/node/mapping/define/NodeMappingEsTableDefine.java @@ -17,14 +17,6 @@ public class NodeMappingEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(NodeMappingTable.COLUMN_ADDRESS_ID, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/define/NodeReferenceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/define/NodeReferenceEsTableDefine.java index da96dc169..6d1bf8673 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/define/NodeReferenceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/noderef/define/NodeReferenceEsTableDefine.java @@ -17,14 +17,6 @@ public class NodeReferenceEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(NodeReferenceTable.COLUMN_FRONT_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(NodeReferenceTable.COLUMN_BEHIND_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java index 82f05c550..1e350ba07 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/cost/define/SegmentCostEsTableDefine.java @@ -17,14 +17,6 @@ public class SegmentCostEsTableDefine extends ElasticSearchTableDefine { return 5; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_SEGMENT_ID, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(SegmentCostTable.COLUMN_SERVICE_NAME, ElasticSearchColumnDefine.Type.Text.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentEsTableDefine.java index e32c7129f..ea9288ecd 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/segment/origin/define/SegmentEsTableDefine.java @@ -17,14 +17,6 @@ public class SegmentEsTableDefine extends ElasticSearchTableDefine { return 10; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(SegmentTable.COLUMN_DATA_BINARY, ElasticSearchColumnDefine.Type.Binary.name())); } diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java index cf2e43608..39de04b31 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/service/entry/define/ServiceEntryEsTableDefine.java @@ -17,14 +17,6 @@ public class ServiceEntryEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(ServiceEntryTable.COLUMN_APPLICATION_ID, ElasticSearchColumnDefine.Type.Integer.name())); addColumn(new ElasticSearchColumnDefine(ServiceEntryTable.COLUMN_ENTRY_SERVICE_ID, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/define/ServiceReferenceEsTableDefine.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/define/ServiceReferenceEsTableDefine.java index efbe857c2..676810afb 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/define/ServiceReferenceEsTableDefine.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agentstream/worker/serviceref/define/ServiceReferenceEsTableDefine.java @@ -17,14 +17,6 @@ public class ServiceReferenceEsTableDefine extends ElasticSearchTableDefine { return 2; } - @Override public int numberOfShards() { - return 2; - } - - @Override public int numberOfReplicas() { - return 0; - } - @Override public void initialize() { addColumn(new ElasticSearchColumnDefine(ServiceReferenceTable.COLUMN_AGG, ElasticSearchColumnDefine.Type.Keyword.name())); addColumn(new ElasticSearchColumnDefine(ServiceReferenceTable.COLUMN_ENTRY_SERVICE_ID, ElasticSearchColumnDefine.Type.Integer.name())); diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define index 57e7f95dd..5f55512c3 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/es_dao.define @@ -1,6 +1,3 @@ -org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationEsDAO -org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceEsDAO -org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameEsDAO org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentEsDAO org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.NodeMappingEsDAO org.skywalking.apm.collector.agentstream.worker.noderef.dao.NodeReferenceEsDAO diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define index 5d8eea1f2..b1064d603 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/h2_dao.define @@ -1,6 +1,3 @@ -org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationH2DAO -org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceH2DAO -org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameH2DAO org.skywalking.apm.collector.agentstream.worker.node.component.dao.NodeComponentH2DAO org.skywalking.apm.collector.agentstream.worker.node.mapping.dao.NodeMappingH2DAO org.skywalking.apm.collector.agentstream.worker.noderef.dao.NodeReferenceH2DAO diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define index 472b02738..638630fa1 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/local_worker_provider.define @@ -17,8 +17,4 @@ org.skywalking.apm.collector.agentstream.worker.segment.origin.SegmentPersistenc org.skywalking.apm.collector.agentstream.worker.segment.cost.SegmentCostPersistenceWorker$Factory org.skywalking.apm.collector.agentstream.worker.global.GlobalTracePersistenceWorker$Factory -org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterSerialWorker$Factory -org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterSerialWorker$Factory -org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterSerialWorker$Factory - org.skywalking.apm.collector.agentstream.worker.instance.performance.InstPerformancePersistenceWorker$Factory \ No newline at end of file diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define index da980634f..5608386d3 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/remote_worker_provider.define @@ -1,6 +1,6 @@ -org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationRegisterRemoteWorker$Factory -org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceRegisterRemoteWorker$Factory -org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameRegisterRemoteWorker$Factory +org.skywalking.apm.collector.agentregister.worker.application.ApplicationRegisterRemoteWorker$Factory +org.skywalking.apm.collector.agentregister.worker.instance.InstanceRegisterRemoteWorker$Factory +org.skywalking.apm.collector.agentregister.worker.servicename.ServiceNameRegisterRemoteWorker$Factory org.skywalking.apm.collector.agentstream.worker.node.component.NodeComponentRemoteWorker$Factory org.skywalking.apm.collector.agentstream.worker.node.mapping.NodeMappingRemoteWorker$Factory diff --git a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define index 3b03aea61..098e51e7a 100644 --- a/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define +++ b/apm-collector/apm-collector-agentstream/src/main/resources/META-INF/defines/storage.define @@ -7,15 +7,6 @@ org.skywalking.apm.collector.agentstream.worker.node.mapping.define.NodeMappingH org.skywalking.apm.collector.agentstream.worker.noderef.define.NodeReferenceEsTableDefine org.skywalking.apm.collector.agentstream.worker.noderef.define.NodeReferenceH2TableDefine -org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationEsTableDefine -org.skywalking.apm.collector.agentstream.worker.register.application.ApplicationH2TableDefine - -org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceEsTableDefine -org.skywalking.apm.collector.agentstream.worker.register.instance.InstanceH2TableDefine - -org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameEsTableDefine -org.skywalking.apm.collector.agentstream.worker.register.servicename.ServiceNameH2TableDefine - org.skywalking.apm.collector.agentstream.worker.segment.origin.define.SegmentEsTableDefine org.skywalking.apm.collector.agentstream.worker.segment.origin.define.SegmentH2TableDefine diff --git a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java index 23293258b..e184dcd8f 100644 --- a/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java +++ b/apm-collector/apm-collector-agentstream/src/test/java/org/skywalking/apm/collector/agentstream/mock/SegmentPost.java @@ -5,9 +5,9 @@ import com.google.gson.JsonElement; import com.google.gson.JsonObject; import java.io.IOException; import org.skywalking.apm.collector.agentstream.HttpClientTools; -import org.skywalking.apm.collector.agentstream.worker.register.application.dao.ApplicationEsDAO; -import org.skywalking.apm.collector.agentstream.worker.register.instance.dao.InstanceEsDAO; -import org.skywalking.apm.collector.agentstream.worker.register.servicename.dao.ServiceNameEsDAO; +import org.skywalking.apm.collector.agentregister.worker.application.dao.ApplicationEsDAO; +import org.skywalking.apm.collector.agentregister.worker.instance.dao.InstanceEsDAO; +import org.skywalking.apm.collector.agentregister.worker.servicename.dao.ServiceNameEsDAO; import org.skywalking.apm.collector.client.elasticsearch.ElasticSearchClient; import org.skywalking.apm.collector.core.CollectorException; import org.skywalking.apm.collector.core.util.TimeBucketUtils; diff --git a/apm-collector/apm-collector-boot/src/main/resources/application.yml b/apm-collector/apm-collector-boot/src/main/resources/application.yml index 0ff32ca04..c9e94f239 100644 --- a/apm-collector/apm-collector-boot/src/main/resources/application.yml +++ b/apm-collector/apm-collector-boot/src/main/resources/application.yml @@ -7,14 +7,6 @@ agent_server: host: localhost port: 10800 context_path: / -agent_register: - grpc: - host: localhost - port: 11800 - jetty: - host: localhost - port: 12800 - context_path: / agent_stream: grpc: host: localhost @@ -23,10 +15,6 @@ agent_stream: host: localhost port: 12800 context_path: / -agent_jvm: - grpc: - host: localhost - port: 11800 ui: jetty: host: localhost @@ -41,3 +29,5 @@ storage: cluster_name: CollectorDBCluster cluster_transport_sniffer: true cluster_nodes: localhost:9300 + index_shards_number: 2 + index_replicas_number: 0 diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfig.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfig.java index 3526fc98b..c05fea702 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfig.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfig.java @@ -7,4 +7,6 @@ public class StorageElasticSearchConfig { public static String CLUSTER_NAME; public static Boolean CLUSTER_TRANSPORT_SNIFFER; public static String CLUSTER_NODES; + public static Integer INDEX_SHARDS_NUMBER; + public static Integer INDEX_REPLICAS_NUMBER; } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfigParser.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfigParser.java index 6653fdc22..223f4ad00 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfigParser.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/StorageElasticSearchConfigParser.java @@ -14,6 +14,8 @@ public class StorageElasticSearchConfigParser implements ModuleConfigParser { private static final String CLUSTER_NAME = "cluster_name"; private static final String CLUSTER_TRANSPORT_SNIFFER = "cluster_transport_sniffer"; private static final String CLUSTER_NODES = "cluster_nodes"; + private static final String INDEX_SHARDS_NUMBER = "index_shards_number"; + private static final String INDEX_REPLICAS_NUMBER = "index_replicas_number"; @Override public void parse(Map config) throws ConfigParseException { if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(CLUSTER_NAME))) { @@ -25,5 +27,15 @@ public class StorageElasticSearchConfigParser implements ModuleConfigParser { if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(CLUSTER_NODES))) { StorageElasticSearchConfig.CLUSTER_NODES = (String)config.get(CLUSTER_NODES); } + if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(INDEX_SHARDS_NUMBER))) { + StorageElasticSearchConfig.INDEX_SHARDS_NUMBER = (Integer)config.get(INDEX_SHARDS_NUMBER); + } else { + StorageElasticSearchConfig.INDEX_SHARDS_NUMBER = 2; + } + if (ObjectUtils.isNotEmpty(config) && StringUtils.isNotEmpty(config.get(INDEX_REPLICAS_NUMBER))) { + StorageElasticSearchConfig.INDEX_REPLICAS_NUMBER = (Integer)config.get(INDEX_REPLICAS_NUMBER); + } else { + StorageElasticSearchConfig.INDEX_REPLICAS_NUMBER = 0; + } } } diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java index 05e3f2738..03c1805b9 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchStorageInstaller.java @@ -11,6 +11,7 @@ import org.skywalking.apm.collector.core.client.Client; import org.skywalking.apm.collector.core.storage.ColumnDefine; import org.skywalking.apm.collector.core.storage.StorageInstaller; import org.skywalking.apm.collector.core.storage.TableDefine; +import org.skywalking.apm.collector.storage.elasticsearch.StorageElasticSearchConfig; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -51,8 +52,8 @@ public class ElasticSearchStorageInstaller extends StorageInstaller { private Settings createSettingBuilder(ElasticSearchTableDefine tableDefine) { return Settings.builder() - .put("index.number_of_shards", tableDefine.numberOfShards()) - .put("index.number_of_replicas", tableDefine.numberOfReplicas()) + .put("index.number_of_shards", StorageElasticSearchConfig.INDEX_SHARDS_NUMBER) + .put("index.number_of_replicas", StorageElasticSearchConfig.INDEX_REPLICAS_NUMBER) .put("index.refresh_interval", String.valueOf(tableDefine.refreshInterval()) + "s") .put("analysis.analyzer.collector_analyzer.tokenizer", "collector_tokenizer") diff --git a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchTableDefine.java b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchTableDefine.java index c9f7241f4..92c6807a8 100644 --- a/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchTableDefine.java +++ b/apm-collector/apm-collector-storage/src/main/java/org/skywalking/apm/collector/storage/elasticsearch/define/ElasticSearchTableDefine.java @@ -16,8 +16,4 @@ public abstract class ElasticSearchTableDefine extends TableDefine { } public abstract int refreshInterval(); - - public abstract int numberOfShards(); - - public abstract int numberOfReplicas(); } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceReferenceEsDAO.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceReferenceEsDAO.java index b6947051c..7ac32ce63 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceReferenceEsDAO.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/dao/ServiceReferenceEsDAO.java @@ -2,6 +2,7 @@ package org.skywalking.apm.collector.ui.dao; import com.google.gson.JsonArray; import com.google.gson.JsonObject; +import java.util.Iterator; import java.util.LinkedHashMap; import java.util.Map; import org.elasticsearch.action.search.SearchRequestBuilder; @@ -14,6 +15,7 @@ import org.elasticsearch.search.aggregations.bucket.terms.Terms; import org.elasticsearch.search.aggregations.metrics.sum.Sum; import org.skywalking.apm.collector.core.util.ColumnNameUtils; import org.skywalking.apm.collector.core.util.Const; +import org.skywalking.apm.collector.core.util.ObjectUtils; import org.skywalking.apm.collector.core.util.StringUtils; import org.skywalking.apm.collector.storage.define.serviceref.ServiceReferenceTable; import org.skywalking.apm.collector.storage.elasticsearch.dao.EsDAO; @@ -112,13 +114,12 @@ public class ServiceReferenceEsDAO extends EsDAO implements IServiceReferenceDAO Map serviceReferenceMap = new LinkedHashMap<>(); - JsonArray serviceReferenceArray = new JsonArray(); SearchResponse searchResponse = searchRequestBuilder.get(); Terms frontServiceIdTerms = searchResponse.getAggregations().get(ServiceReferenceTable.COLUMN_FRONT_SERVICE_ID); for (Terms.Bucket frontServiceBucket : frontServiceIdTerms.getBuckets()) { int frontServiceId = frontServiceBucket.getKeyAsNumber().intValue(); if (frontServiceId != 0) { - parseSubAggregate(serviceReferenceMap, serviceReferenceArray, frontServiceBucket, frontServiceId); + parseSubAggregate(serviceReferenceMap, frontServiceBucket, frontServiceId); } } @@ -128,15 +129,23 @@ public class ServiceReferenceEsDAO extends EsDAO implements IServiceReferenceDAO if (StringUtils.isNotEmpty(frontServiceName)) { String[] serviceNames = frontServiceName.split(Const.ID_SPLIT); int frontServiceId = ServiceIdCache.getForUI(Integer.parseInt(serviceNames[0]), serviceNames[1]); - parseSubAggregate(serviceReferenceMap, serviceReferenceArray, frontServiceBucket, frontServiceId); + parseSubAggregate(serviceReferenceMap, frontServiceBucket, frontServiceId); } } - serviceReferenceMap.values().forEach(serviceReferenceArray::add); + JsonArray serviceReferenceArray = new JsonArray(); + JsonObject rootServiceReference = findRoot(serviceReferenceMap); + if (ObjectUtils.isNotEmpty(rootServiceReference)) { + String id = rootServiceReference.get(ColumnNameUtils.INSTANCE.rename(ServiceReferenceTable.COLUMN_FRONT_SERVICE_ID)) + Const.ID_SPLIT + rootServiceReference.get(ColumnNameUtils.INSTANCE.rename(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_ID)); + serviceReferenceMap.remove(id); + + int rootServiceId = rootServiceReference.get(ColumnNameUtils.INSTANCE.rename(ServiceReferenceTable.COLUMN_FRONT_SERVICE_ID)).getAsInt(); + sortAsTree(rootServiceId, serviceReferenceArray, serviceReferenceMap); + } return serviceReferenceArray; } - private void parseSubAggregate(Map serviceReferenceMap, JsonArray serviceReferenceArray, + private void parseSubAggregate(Map serviceReferenceMap, Terms.Bucket frontServiceBucket, int frontServiceId) { Terms behindServiceIdTerms = frontServiceBucket.getAggregations().get(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_ID); @@ -232,4 +241,29 @@ public class ServiceReferenceEsDAO extends EsDAO implements IServiceReferenceDAO long newValue = newReference.get(key).getAsLong(); oldReference.addProperty(key, oldValue + newValue); } + + private JsonObject findRoot(Map serviceReferenceMap) { + for (JsonObject serviceReference : serviceReferenceMap.values()) { + int behindServiceId = serviceReference.get(ColumnNameUtils.INSTANCE.rename(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_ID)).getAsInt(); + if (behindServiceId == 1) { + return serviceReference; + } + } + return null; + } + + private void sortAsTree(int serviceId, JsonArray serviceReferenceArray, + Map serviceReferenceMap) { + Iterator iterator = serviceReferenceMap.values().iterator(); + while (iterator.hasNext()) { + JsonObject serviceReference = iterator.next(); + int frontServiceId = serviceReference.get(ColumnNameUtils.INSTANCE.rename(ServiceReferenceTable.COLUMN_FRONT_SERVICE_ID)).getAsInt(); + if (serviceId == frontServiceId) { + serviceReferenceArray.add(serviceReference); + + int behindServiceId = serviceReference.get(ColumnNameUtils.INSTANCE.rename(ServiceReferenceTable.COLUMN_BEHIND_SERVICE_ID)).getAsInt(); + sortAsTree(behindServiceId, serviceReferenceArray, serviceReferenceMap); + } + } + } } diff --git a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java index a0434ed2a..a29e509cc 100644 --- a/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java +++ b/apm-collector/apm-collector-ui/src/main/java/org/skywalking/apm/collector/ui/service/InstanceHealthService.java @@ -28,11 +28,12 @@ public class InstanceHealthService { IInstanceDAO instanceDAO = (IInstanceDAO)DAOContainer.INSTANCE.get(IInstanceDAO.class.getName()); List instanceList = instanceDAO.getInstances(applicationId, halfHourBeforeTimeBucket); + JsonArray instances = new JsonArray(); + response.add("instances", instances); + instanceList.forEach(instance -> { - JsonArray instances = new JsonArray(); response.addProperty("applicationCode", ApplicationCache.getForUI(applicationId)); response.addProperty("applicationId", applicationId); - response.add("instances", instances); IInstPerformanceDAO instPerformanceDAO = (IInstPerformanceDAO)DAOContainer.INSTANCE.get(IInstPerformanceDAO.class.getName()); IInstPerformanceDAO.InstPerformance performance = instPerformanceDAO.get(timeBuckets, instance.getInstanceId());