1. Add stream module: The stream base framework

2. Add context help for every module to find another module’s context.
This commit is contained in:
pengys5 2017-07-18 14:49:47 +08:00
parent 04fbe180f3
commit 640440d581
97 changed files with 578 additions and 250 deletions

View File

@ -15,7 +15,12 @@
<dependencies>
<dependency>
<groupId>org.skywalking</groupId>
<artifactId>apm-collector-core</artifactId>
<artifactId>apm-collector-stream</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.skywalking</groupId>
<artifactId>apm-collector-cluster</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>

View File

@ -1,11 +1,13 @@
package org.skywalking.apm.collector.core.agentstream;
package org.skywalking.apm.collector.agentstream;
import java.util.Map;
import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.cluster.ClusterDataInitializer;
import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
import org.skywalking.apm.collector.core.config.ConfigParseException;
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.framework.DataInitializer;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleDefine;
@ -24,7 +26,7 @@ public abstract class AgentStreamModuleDefine extends ModuleDefine {
server.initialize();
String key = ClusterDataInitializer.BASE_CATALOG + "." + name();
ClusterModuleContext.WRITER.write(key, registration().buildValue());
((ClusterModuleContext)CollectorContextHelper.INSTANCE.getContext(ClusterModuleGroupDefine.GROUP_NAME)).getWriter().write(key, registration().buildValue());
} catch (ConfigParseException | ServerException e) {
throw new AgentStreamModuleException(e.getMessage(), e);
}

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.agentstream;
package org.skywalking.apm.collector.agentstream;
import org.skywalking.apm.collector.core.module.ModuleException;

View File

@ -0,0 +1,26 @@
package org.skywalking.apm.collector.agentstream;
import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
import org.skywalking.apm.collector.core.framework.Context;
import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
import org.skywalking.apm.collector.core.module.ModuleInstaller;
/**
* @author pengys5
*/
public class AgentStreamModuleGroupDefine implements ModuleGroupDefine {
public static final String GROUP_NAME = "agent_stream";
@Override public String name() {
return GROUP_NAME;
}
@Override public Context groupContext() {
return new ClusterModuleContext(GROUP_NAME);
}
@Override public ModuleInstaller moduleInstaller() {
return new AgentStreamModuleInstaller();
}
}

View File

@ -0,0 +1,51 @@
package org.skywalking.apm.collector.agentstream;
import java.util.Iterator;
import java.util.Map;
import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleInstaller;
import org.skywalking.apm.collector.core.util.CollectionUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* @author pengys5
*/
public class AgentStreamModuleInstaller implements ModuleInstaller {
private final Logger logger = LoggerFactory.getLogger(AgentStreamModuleInstaller.class);
@Override public void install(Map<String, Map> moduleConfig,
Map<String, ModuleDefine> moduleDefineMap) throws DefineException, ClientException {
logger.info("beginning cluster module install");
ModuleDefine moduleDefine = null;
if (CollectionUtils.isEmpty(moduleConfig)) {
logger.info("could not configure cluster module, use the default");
Iterator<Map.Entry<String, ModuleDefine>> moduleDefineEntry = moduleDefineMap.entrySet().iterator();
while (moduleDefineEntry.hasNext()) {
moduleDefine = moduleDefineEntry.next().getValue();
if (moduleDefine.defaultModule()) {
logger.info("module {} initialize", moduleDefine.getClass().getName());
moduleDefine.initialize(null);
break;
}
}
} else {
Map.Entry<String, Map> clusterConfigEntry = moduleConfig.entrySet().iterator().next();
moduleDefine = moduleDefineMap.get(clusterConfigEntry.getKey());
moduleDefine.initialize(clusterConfigEntry.getValue());
}
ClusterModuleContext context = new ClusterModuleContext(ClusterModuleGroupDefine.GROUP_NAME);
context.setWriter(((ClusterModuleDefine)moduleDefine).registrationWriter());
CollectorContextHelper.INSTANCE.putContext(context);
}
}

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.agent.stream.server.grpc;
package org.skywalking.apm.collector.agentstream.grpc;
import java.util.Map;
import org.skywalking.apm.collector.core.config.ConfigParseException;

View File

@ -1,8 +1,8 @@
package org.skywalking.apm.collector.agent.stream.server.grpc;
package org.skywalking.apm.collector.agentstream.grpc;
import org.skywalking.apm.collector.core.agentstream.AgentStreamModuleDefine;
import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine;
import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleGroup;
import org.skywalking.apm.collector.core.module.ModuleRegistration;
import org.skywalking.apm.collector.core.server.Server;
import org.skywalking.apm.collector.server.grpc.GRPCServer;
@ -12,8 +12,8 @@ import org.skywalking.apm.collector.server.grpc.GRPCServer;
*/
public class AgentStreamGRPCModuleDefine extends AgentStreamModuleDefine {
@Override protected ModuleGroup group() {
return ModuleGroup.AgentStream;
@Override protected String group() {
return AgentStreamModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.agent.stream.server.grpc;
package org.skywalking.apm.collector.agentstream.grpc;
import org.skywalking.apm.collector.core.module.ModuleRegistration;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.agent.stream.server.jetty;
package org.skywalking.apm.collector.agentstream.jetty;
import java.util.Map;
import org.skywalking.apm.collector.core.config.ConfigParseException;

View File

@ -1,8 +1,8 @@
package org.skywalking.apm.collector.agent.stream.server.jetty;
package org.skywalking.apm.collector.agentstream.jetty;
import org.skywalking.apm.collector.core.agentstream.AgentStreamModuleDefine;
import org.skywalking.apm.collector.agentstream.AgentStreamModuleDefine;
import org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleGroup;
import org.skywalking.apm.collector.core.module.ModuleRegistration;
import org.skywalking.apm.collector.core.server.Server;
import org.skywalking.apm.collector.server.jetty.JettyServer;
@ -12,8 +12,8 @@ import org.skywalking.apm.collector.server.jetty.JettyServer;
*/
public class AgentStreamJettyModuleDefine extends AgentStreamModuleDefine {
@Override protected ModuleGroup group() {
return ModuleGroup.AgentStream;
@Override protected String group() {
return AgentStreamModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.agent.stream.server.jetty;
package org.skywalking.apm.collector.agentstream.jetty;
import com.google.gson.JsonObject;
import org.skywalking.apm.collector.core.module.ModuleRegistration;

View File

@ -0,0 +1 @@
org.skywalking.apm.collector.agentstream.AgentStreamModuleGroupDefine

View File

@ -1,2 +1,2 @@
org.skywalking.apm.collector.agent.stream.server.grpc.AgentStreamGRPCModuleDefine
org.skywalking.apm.collector.agent.stream.server.jetty.AgentStreamJettyModuleDefine
org.skywalking.apm.collector.agentstream.grpc.AgentStreamGRPCModuleDefine
org.skywalking.apm.collector.agentstream.jetty.AgentStreamJettyModuleDefine

View File

@ -2,7 +2,6 @@ package org.skywalking.apm.collector.boot;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.config.ConfigException;
import org.skywalking.apm.collector.core.framework.CollectorStarter;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

View File

@ -1,15 +1,16 @@
package org.skywalking.apm.collector.core.framework;
package org.skywalking.apm.collector.boot;
import java.util.Map;
import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.config.ConfigException;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.framework.Starter;
import org.skywalking.apm.collector.core.module.ModuleConfigLoader;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleDefineLoader;
import org.skywalking.apm.collector.core.module.ModuleGroup;
import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
import org.skywalking.apm.collector.core.module.ModuleGroupDefineLoader;
import org.skywalking.apm.collector.core.module.ModuleInstallerAdapter;
import org.skywalking.apm.collector.core.remote.SerializedDefineLoader;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -22,29 +23,23 @@ public class CollectorStarter implements Starter {
private final Logger logger = LoggerFactory.getLogger(CollectorStarter.class);
@Override public void start() throws ConfigException, DefineException, ClientException {
Context context = new Context();
ModuleConfigLoader configLoader = new ModuleConfigLoader();
Map<String, Map> configuration = configLoader.load();
SerializedDefineLoader serializedDefineLoader = new SerializedDefineLoader();
serializedDefineLoader.load();
ModuleDefineLoader defineLoader = new ModuleDefineLoader();
Map<String, Map<String, ModuleDefine>> moduleDefineMap = defineLoader.load();
ModuleGroupDefineLoader groupDefineLoader = new ModuleGroupDefineLoader();
Map<String, ModuleGroupDefine> moduleGroupDefineMap = groupDefineLoader.load();
ModuleInstallerAdapter moduleInstallerAdapter = new ModuleInstallerAdapter(ModuleGroup.Cluster);
moduleInstallerAdapter.install(configuration.get(ModuleGroup.Cluster.name().toLowerCase()), moduleDefineMap.get(ModuleGroup.Cluster.name().toLowerCase()));
ModuleDefineLoader defineLoader = new ModuleDefineLoader();
Map<String, Map<String, ModuleDefine>> moduleDefineMap = defineLoader.load();
ModuleGroup[] moduleGroups = ModuleGroup.values();
for (ModuleGroup moduleGroup : moduleGroups) {
if (!ModuleGroup.Cluster.equals(moduleGroup)) {
moduleInstallerAdapter = new ModuleInstallerAdapter(moduleGroup);
logger.info("module group {}, configuration {}", moduleGroup.name().toLowerCase(), configuration.get(moduleGroup.name().toLowerCase()));
moduleInstallerAdapter.install(configuration.get(moduleGroup.name().toLowerCase()), moduleDefineMap.get(moduleGroup.name().toLowerCase()));
}
moduleGroupDefineMap.get(ClusterModuleGroupDefine.GROUP_NAME).moduleInstaller().install(configuration.get(ClusterModuleGroupDefine.GROUP_NAME), moduleDefineMap.get(ClusterModuleGroupDefine.GROUP_NAME));
moduleGroupDefineMap.remove(ClusterModuleGroupDefine.GROUP_NAME);
for (ModuleGroupDefine moduleGroupDefine : moduleGroupDefineMap.values()) {
moduleGroupDefine.moduleInstaller().install(configuration.get(moduleGroupDefine.name()), moduleDefineMap.get(moduleGroupDefine.name()));
}
}
}

View File

@ -32,6 +32,16 @@
<groupId>org.apache.zookeeper</groupId>
<artifactId>zookeeper</artifactId>
<version>3.4.10</version>
<exclusions>
<exclusion>
<artifactId>slf4j-api</artifactId>
<groupId>org.slf4j</groupId>
</exclusion>
<exclusion>
<artifactId>slf4j-log4j12</artifactId>
<groupId>org.slf4j</groupId>
</exclusion>
</exclusions>
</dependency>
</dependencies>
</project>

View File

@ -1,17 +1,26 @@
package org.skywalking.apm.collector.cluster;
import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
import org.skywalking.apm.collector.core.framework.Context;
import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
import org.skywalking.apm.collector.core.module.ModuleInstallMode;
import org.skywalking.apm.collector.core.module.ModuleInstaller;
/**
* @author pengys5
*/
public class ClusterModuleGroupDefine implements ModuleGroupDefine {
public static final String GROUP_NAME = "cluster";
@Override public String name() {
return "cluster";
return GROUP_NAME;
}
@Override public ModuleInstallMode mode() {
return ModuleInstallMode.Single;
@Override public Context groupContext() {
return new ClusterModuleContext(GROUP_NAME);
}
@Override public ModuleInstaller moduleInstaller() {
return new ClusterModuleInstaller();
}
}

View File

@ -1,8 +1,11 @@
package org.skywalking.apm.collector.core.cluster;
package org.skywalking.apm.collector.cluster;
import java.util.Iterator;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.cluster.ClusterModuleContext;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleInstaller;
@ -38,6 +41,10 @@ public class ClusterModuleInstaller implements ModuleInstaller {
moduleDefine = moduleDefineMap.get(clusterConfigEntry.getKey());
moduleDefine.initialize(clusterConfigEntry.getValue());
}
ClusterModuleContext.WRITER = ((ClusterModuleDefine)moduleDefine).registrationWriter();
ClusterModuleContext context = new ClusterModuleContext(ClusterModuleGroupDefine.GROUP_NAME);
context.setWriter(((ClusterModuleDefine)moduleDefine).registrationWriter());
CollectorContextHelper.INSTANCE.putContext(context);
}
}

View File

@ -1,25 +1,27 @@
package org.skywalking.apm.collector.cluster.redis;
import org.skywalking.apm.collector.client.redis.RedisClient;
import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationWriter;
import org.skywalking.apm.collector.core.framework.DataInitializer;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleGroup;
/**
* @author pengys5
*/
public class ClusterRedisModuleDefine extends ClusterModuleDefine {
@Override public ModuleGroup group() {
return ModuleGroup.Cluster;
public static final String MODULE_NAME = "redis";
@Override public String group() {
return ClusterModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
return "redis";
return MODULE_NAME;
}
@Override public boolean defaultModule() {
@ -38,11 +40,11 @@ public class ClusterRedisModuleDefine extends ClusterModuleDefine {
return new ClusterRedisDataInitializer();
}
@Override protected ClusterModuleRegistrationWriter registrationWriter() {
@Override public ClusterModuleRegistrationWriter registrationWriter() {
return new ClusterRedisModuleRegistrationWriter(getClient());
}
@Override protected ClusterModuleRegistrationReader registrationReader() {
@Override public ClusterModuleRegistrationReader registrationReader() {
return null;
}
}

View File

@ -1,25 +1,27 @@
package org.skywalking.apm.collector.cluster.standalone;
import org.skywalking.apm.collector.client.h2.H2Client;
import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationWriter;
import org.skywalking.apm.collector.core.framework.DataInitializer;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleGroup;
/**
* @author pengys5
*/
public class ClusterStandaloneModuleDefine extends ClusterModuleDefine {
@Override public ModuleGroup group() {
return ModuleGroup.Cluster;
public static final String MODULE_NAME = "standalone";
@Override public String group() {
return ClusterModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
return "standalone";
return MODULE_NAME;
}
@Override public boolean defaultModule() {
@ -38,11 +40,11 @@ public class ClusterStandaloneModuleDefine extends ClusterModuleDefine {
return new ClusterStandaloneDataInitializer();
}
@Override protected ClusterModuleRegistrationWriter registrationWriter() {
@Override public ClusterModuleRegistrationWriter registrationWriter() {
return new ClusterStandaloneModuleRegistrationWriter(getClient());
}
@Override protected ClusterModuleRegistrationReader registrationReader() {
@Override public ClusterModuleRegistrationReader registrationReader() {
return null;
}
}

View File

@ -1,25 +1,27 @@
package org.skywalking.apm.collector.cluster.zookeeper;
import org.skywalking.apm.collector.client.zookeeper.ZookeeperClient;
import org.skywalking.apm.collector.cluster.ClusterModuleGroupDefine;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.cluster.ClusterDataInitializer;
import org.skywalking.apm.collector.core.cluster.ClusterModuleDefine;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationReader;
import org.skywalking.apm.collector.core.cluster.ClusterModuleRegistrationWriter;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleGroup;
/**
* @author pengys5
*/
public class ClusterZKModuleDefine extends ClusterModuleDefine {
@Override protected ModuleGroup group() {
return ModuleGroup.Cluster;
public static final String MODULE_NAME = "zookeeper";
@Override protected String group() {
return ClusterModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
return "zookeeper";
return MODULE_NAME;
}
@Override public boolean defaultModule() {
@ -38,11 +40,11 @@ public class ClusterZKModuleDefine extends ClusterModuleDefine {
return new ClusterZKDataInitializer();
}
@Override protected ClusterModuleRegistrationWriter registrationWriter() {
@Override public ClusterModuleRegistrationWriter registrationWriter() {
return new ClusterZKModuleRegistrationWriter(getClient());
}
@Override protected ClusterModuleRegistrationReader registrationReader() {
@Override public ClusterModuleRegistrationReader registrationReader() {
return null;
}
}

View File

@ -18,11 +18,6 @@
<artifactId>snakeyaml</artifactId>
<version>1.18</version>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
<version>1.2.3</version>
</dependency>
<dependency>
<groupId>com.google.code.gson</groupId>
<artifactId>gson</artifactId>

View File

@ -1,10 +1,33 @@
package org.skywalking.apm.collector.core.cluster;
import org.skywalking.apm.collector.core.framework.Context;
/**
* @author pengys5
*/
public class ClusterModuleContext {
public static ClusterModuleRegistrationWriter WRITER;
public class ClusterModuleContext extends Context {
public static ClusterModuleRegistrationReader READER;
public ClusterModuleContext(String groupName) {
super(groupName);
}
private ClusterModuleRegistrationWriter writer;
private ClusterModuleRegistrationReader reader;
public ClusterModuleRegistrationWriter getWriter() {
return writer;
}
public void setWriter(ClusterModuleRegistrationWriter writer) {
this.writer = writer;
}
public ClusterModuleRegistrationReader getReader() {
return reader;
}
public void setReader(ClusterModuleRegistrationReader reader) {
this.reader = reader;
}
}

View File

@ -38,7 +38,7 @@ public abstract class ClusterModuleDefine extends ModuleDefine {
throw new UnsupportedOperationException("Cluster module do not need module registration.");
}
protected abstract ClusterModuleRegistrationWriter registrationWriter();
public abstract ClusterModuleRegistrationWriter registrationWriter();
protected abstract ClusterModuleRegistrationReader registrationReader();
public abstract ClusterModuleRegistrationReader registrationReader();
}

View File

@ -0,0 +1,25 @@
package org.skywalking.apm.collector.core.framework;
import java.util.LinkedHashMap;
import java.util.Map;
/**
* @author pengys5
*/
public enum CollectorContextHelper {
INSTANCE;
private Map<String, Context> contexts = new LinkedHashMap();
public Context getContext(String moduleGroupName) {
return contexts.get(moduleGroupName);
}
public void putContext(Context context) {
if (contexts.containsKey(context.getGroupName())) {
throw new UnsupportedOperationException("This module context was put, do not allow put a new one");
} else {
contexts.put(context.getGroupName(), context);
}
}
}

View File

@ -3,5 +3,14 @@ package org.skywalking.apm.collector.core.framework;
/**
* @author pengys5
*/
public class Context {
public abstract class Context {
private final String groupName;
public Context(String groupName) {
this.groupName = groupName;
}
public final String getGroupName() {
return groupName;
}
}

View File

@ -0,0 +1,8 @@
package org.skywalking.apm.collector.core.framework;
/**
* @author pengys5
*/
public interface Executor {
void execute(Object message);
}

View File

@ -1,6 +1,7 @@
package org.skywalking.apm.collector.core.module;
import java.io.FileNotFoundException;
import java.io.FileReader;
import java.util.Map;
import org.skywalking.apm.collector.core.config.ConfigLoader;
import org.skywalking.apm.collector.core.util.ResourceUtils;
@ -18,7 +19,13 @@ public class ModuleConfigLoader implements ConfigLoader<Map<String, Map>> {
@Override public Map<String, Map> load() throws ModuleConfigLoaderException {
Yaml yaml = new Yaml();
try {
return (Map<String, Map>)yaml.load(ResourceUtils.read("application.yml"));
try {
FileReader applicationFileReader = ResourceUtils.read("application.yml");
return (Map<String, Map>)yaml.load(applicationFileReader);
} catch (FileNotFoundException e) {
logger.info("Could not found application.yml file, use default");
return (Map<String, Map>)yaml.load(ResourceUtils.read("application-default.yml"));
}
} catch (FileNotFoundException e) {
throw new ModuleConfigLoaderException(e.getMessage(), e);
}

View File

@ -10,7 +10,7 @@ import org.skywalking.apm.collector.core.server.Server;
*/
public abstract class ModuleDefine implements Define {
protected abstract ModuleGroup group();
protected abstract String group();
public abstract boolean defaultModule();

View File

@ -24,7 +24,7 @@ public class ModuleDefineLoader implements Loader<Map<String, Map<String, Module
for (ModuleDefine moduleDefine : definitionLoader) {
logger.info("loaded module definition class: {}", moduleDefine.getClass().getName());
String groupName = moduleDefine.group().name().toLowerCase();
String groupName = moduleDefine.group();
if (!moduleDefineMap.containsKey(groupName)) {
moduleDefineMap.put(groupName, new LinkedHashMap<>());
}

View File

@ -1,8 +0,0 @@
package org.skywalking.apm.collector.core.module;
/**
* @author pengys5
*/
public enum ModuleGroup {
Cluster, Queue, AgentStream
}

View File

@ -1,10 +1,14 @@
package org.skywalking.apm.collector.core.module;
import org.skywalking.apm.collector.core.framework.Context;
/**
* @author pengys5
*/
public interface ModuleGroupDefine {
String name();
ModuleInstallMode mode();
Context groupContext();
ModuleInstaller moduleInstaller();
}

View File

@ -1,8 +0,0 @@
package org.skywalking.apm.collector.core.module;
/**
* @author pengys5
*/
public enum ModuleInstallMode {
Single, Multiple
}

View File

@ -1,25 +0,0 @@
package org.skywalking.apm.collector.core.module;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.cluster.ClusterModuleInstaller;
import org.skywalking.apm.collector.core.framework.DefineException;
/**
* @author pengys5
*/
public class ModuleInstallerAdapter implements ModuleInstaller {
private ModuleInstaller moduleInstaller;
public ModuleInstallerAdapter(ModuleGroup moduleGroup) {
if (ModuleGroup.Cluster.equals(moduleGroup)) {
moduleInstaller = new ClusterModuleInstaller();
}
}
@Override public void install(Map<String, Map> moduleConfig,
Map<String, ModuleDefine> moduleDefineMap) throws DefineException, ClientException {
moduleInstaller.install(moduleConfig, moduleDefineMap);
}
}

View File

@ -1,10 +0,0 @@
package org.skywalking.apm.collector.core.queue;
import org.skywalking.apm.collector.core.worker.AbstractLocalAsyncWorker;
/**
* @author pengys5
*/
public interface Creator {
QueueEventHandler create(int queueSize, AbstractLocalAsyncWorker localAsyncWorker);
}

View File

@ -0,0 +1,8 @@
package org.skywalking.apm.collector.core.queue;
/**
* @author pengys5
*/
public interface QueueCreator {
QueueEventHandler create(int queueSize, QueueExecutor executor);
}

View File

@ -0,0 +1,9 @@
package org.skywalking.apm.collector.core.queue;
import org.skywalking.apm.collector.core.framework.Executor;
/**
* @author pengys5
*/
public interface QueueExecutor extends Executor {
}

View File

@ -1,8 +0,0 @@
package org.skywalking.apm.collector.core.queue;
/**
* @author pengys5
*/
public class QueueModuleContext {
public static Creator CREATOR;
}

View File

@ -1,16 +1,21 @@
package org.skywalking.apm.collector.core.util;
import java.io.File;
import java.io.FileNotFoundException;
import java.io.FileReader;
import java.net.URL;
/**
* @author pengys5
*/
public class ResourceUtils {
private static final String PATH = ResourceUtils.class.getResource("/").getPath();
public static FileReader read(String fileName) throws FileNotFoundException {
return new FileReader(PATH + fileName);
URL url = ResourceUtils.class.getClassLoader().getResource(fileName);
if (url == null) {
throw new FileNotFoundException("file not found: " + fileName);
}
File file = new File(ResourceUtils.class.getClassLoader().getResource(fileName).getFile());
return new FileReader(file);
}
}

View File

@ -1,13 +0,0 @@
package org.skywalking.apm.collector.core.worker;
import org.skywalking.apm.collector.core.worker.selector.WorkerSelector;
/**
* @author pengys5
*/
public interface Role {
String roleName();
WorkerSelector workerSelector();
}

View File

@ -0,0 +1,23 @@
package org.skywalking.apm.collector.queue;
import org.skywalking.apm.collector.core.framework.Context;
import org.skywalking.apm.collector.core.queue.QueueCreator;
/**
* @author pengys5
*/
public class QueueModuleContext extends Context {
private QueueCreator queueCreator;
public QueueModuleContext(String groupName) {
super(groupName);
}
public QueueCreator getQueueCreator() {
return queueCreator;
}
public void setQueueCreator(QueueCreator queueCreator) {
this.queueCreator = queueCreator;
}
}

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.queue;
package org.skywalking.apm.collector.queue;
import org.skywalking.apm.collector.core.client.Client;
import org.skywalking.apm.collector.core.framework.DataInitializer;

View File

@ -0,0 +1,16 @@
package org.skywalking.apm.collector.queue;
import org.skywalking.apm.collector.core.module.ModuleException;
/**
* @author pengys5
*/
public class QueueModuleException extends ModuleException {
public QueueModuleException(String message) {
super(message);
}
public QueueModuleException(String message, Throwable cause) {
super(message, cause);
}
}

View File

@ -0,0 +1,25 @@
package org.skywalking.apm.collector.queue;
import org.skywalking.apm.collector.core.framework.Context;
import org.skywalking.apm.collector.core.module.ModuleGroupDefine;
import org.skywalking.apm.collector.core.module.ModuleInstaller;
/**
* @author pengys5
*/
public class QueueModuleGroupDefine implements ModuleGroupDefine {
public static final String GROUP_NAME = "queue";
@Override public String name() {
return GROUP_NAME;
}
@Override public Context groupContext() {
return new QueueModuleContext(GROUP_NAME);
}
@Override public ModuleInstaller moduleInstaller() {
return new QueueModuleInstaller();
}
}

View File

@ -0,0 +1,23 @@
package org.skywalking.apm.collector.queue;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleInstaller;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* @author pengys5
*/
public class QueueModuleInstaller implements ModuleInstaller {
private final Logger logger = LoggerFactory.getLogger(QueueModuleInstaller.class);
@Override public void install(Map<String, Map> moduleConfig,
Map<String, ModuleDefine> moduleDefineMap) throws DefineException, ClientException {
logger.info("beginning cluster module install");
}
}

View File

@ -1,15 +0,0 @@
package org.skywalking.apm.collector.queue.datacarrier;
import org.skywalking.apm.collector.core.queue.Creator;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;
import org.skywalking.apm.collector.core.worker.AbstractLocalAsyncWorker;
/**
* @author pengys5
*/
public class DataCarrierCreator implements Creator {
@Override public QueueEventHandler create(int queueSize, AbstractLocalAsyncWorker localAsyncWorker) {
return null;
}
}

View File

@ -0,0 +1,15 @@
package org.skywalking.apm.collector.queue.datacarrier;
import org.skywalking.apm.collector.core.queue.QueueCreator;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;
import org.skywalking.apm.collector.core.queue.QueueExecutor;
/**
* @author pengys5
*/
public class DataCarrierQueueCreator implements QueueCreator {
@Override public QueueEventHandler create(int queueSize, QueueExecutor executor) {
return null;
}
}

View File

@ -2,18 +2,19 @@ package org.skywalking.apm.collector.queue.datacarrier;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleGroup;
import org.skywalking.apm.collector.core.queue.QueueModuleContext;
import org.skywalking.apm.collector.core.queue.QueueModuleDefine;
import org.skywalking.apm.collector.queue.QueueModuleContext;
import org.skywalking.apm.collector.queue.QueueModuleDefine;
import org.skywalking.apm.collector.queue.QueueModuleGroupDefine;
/**
* @author pengys5
*/
public class QueueDataCarrierModuleDefine extends QueueModuleDefine {
@Override protected ModuleGroup group() {
return ModuleGroup.Queue;
@Override protected String group() {
return QueueModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
@ -25,6 +26,6 @@ public class QueueDataCarrierModuleDefine extends QueueModuleDefine {
}
@Override public final void initialize(Map config) throws DefineException, ClientException {
QueueModuleContext.CREATOR = new DataCarrierCreator();
((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(group())).setQueueCreator(new DataCarrierQueueCreator());
}
}

View File

@ -5,8 +5,7 @@ import com.lmax.disruptor.RingBuffer;
import org.skywalking.apm.collector.core.queue.EndOfBatchCommand;
import org.skywalking.apm.collector.core.queue.MessageHolder;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;
import org.skywalking.apm.collector.core.worker.AbstractLocalAsyncWorker;
import org.skywalking.apm.collector.core.worker.WorkerException;
import org.skywalking.apm.collector.core.queue.QueueExecutor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@ -18,11 +17,11 @@ public class DisruptorEventHandler implements EventHandler<MessageHolder>, Queue
private final Logger logger = LoggerFactory.getLogger(DisruptorEventHandler.class);
private RingBuffer<MessageHolder> ringBuffer;
private AbstractLocalAsyncWorker asyncWorker;
private QueueExecutor executor;
DisruptorEventHandler(RingBuffer<MessageHolder> ringBuffer, AbstractLocalAsyncWorker asyncWorker) {
DisruptorEventHandler(RingBuffer<MessageHolder> ringBuffer, QueueExecutor executor) {
this.ringBuffer = ringBuffer;
this.asyncWorker = asyncWorker;
this.executor = executor;
}
/**
@ -34,16 +33,12 @@ public class DisruptorEventHandler implements EventHandler<MessageHolder>, Queue
* @param endOfBatch flag to indicate if this is the last event in a batch from the {@link RingBuffer}
*/
public void onEvent(MessageHolder event, long sequence, boolean endOfBatch) {
try {
Object message = event.getMessage();
event.reset();
Object message = event.getMessage();
event.reset();
asyncWorker.allocateJob(message);
if (endOfBatch) {
asyncWorker.allocateJob(new EndOfBatchCommand());
}
} catch (WorkerException e) {
logger.error(e.getMessage(), e);
executor.execute(message);
if (endOfBatch) {
executor.execute(new EndOfBatchCommand());
}
}

View File

@ -2,18 +2,18 @@ package org.skywalking.apm.collector.queue.disruptor;
import com.lmax.disruptor.RingBuffer;
import com.lmax.disruptor.dsl.Disruptor;
import org.skywalking.apm.collector.core.queue.Creator;
import org.skywalking.apm.collector.core.queue.DaemonThreadFactory;
import org.skywalking.apm.collector.core.queue.MessageHolder;
import org.skywalking.apm.collector.core.queue.QueueCreator;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;
import org.skywalking.apm.collector.core.worker.AbstractLocalAsyncWorker;
import org.skywalking.apm.collector.core.queue.QueueExecutor;
/**
* @author pengys5
*/
public class DisruptorCreator implements Creator {
public class DisruptorQueueCreator implements QueueCreator {
public QueueEventHandler create(int queueSize, AbstractLocalAsyncWorker localAsyncWorker) {
@Override public QueueEventHandler create(int queueSize, QueueExecutor executor) {
// Specify the size of the ring buffer, must be power of 2.
if (!((((queueSize - 1) & queueSize) == 0) && queueSize != 0)) {
throw new IllegalArgumentException("queue size must be power of 2");
@ -23,7 +23,7 @@ public class DisruptorCreator implements Creator {
Disruptor<MessageHolder> disruptor = new Disruptor(MessageHolderFactory.INSTANCE, queueSize, DaemonThreadFactory.INSTANCE);
RingBuffer<MessageHolder> ringBuffer = disruptor.getRingBuffer();
DisruptorEventHandler eventHandler = new DisruptorEventHandler(ringBuffer, localAsyncWorker);
DisruptorEventHandler eventHandler = new DisruptorEventHandler(ringBuffer, executor);
// Connect the handler
disruptor.handleEventsWith(eventHandler);

View File

@ -2,18 +2,19 @@ package org.skywalking.apm.collector.queue.disruptor;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleGroup;
import org.skywalking.apm.collector.core.queue.QueueModuleContext;
import org.skywalking.apm.collector.core.queue.QueueModuleDefine;
import org.skywalking.apm.collector.queue.QueueModuleContext;
import org.skywalking.apm.collector.queue.QueueModuleDefine;
import org.skywalking.apm.collector.queue.QueueModuleGroupDefine;
/**
* @author pengys5
*/
public class QueueDisruptorModuleDefine extends QueueModuleDefine {
@Override protected ModuleGroup group() {
return ModuleGroup.Queue;
@Override protected String group() {
return QueueModuleGroupDefine.GROUP_NAME;
}
@Override public String name() {
@ -25,6 +26,6 @@ public class QueueDisruptorModuleDefine extends QueueModuleDefine {
}
@Override public final void initialize(Map config) throws DefineException, ClientException {
QueueModuleContext.CREATOR = new DisruptorCreator();
((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(group())).setQueueCreator(new DisruptorQueueCreator());
}
}

View File

@ -0,0 +1 @@
org.skywalking.apm.collector.queue.QueueModuleGroupDefine

View File

@ -0,0 +1,12 @@
package org.skywalking.apm.collector.remote;
import org.skywalking.apm.collector.core.framework.Context;
/**
* @author pengys5
*/
public class RemoteModuleContext extends Context {
public RemoteModuleContext(String groupName) {
super(groupName);
}
}

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.remote;
package org.skywalking.apm.collector.remote;
import org.skywalking.apm.collector.core.module.ModuleDefine;

View File

@ -0,0 +1,25 @@
package org.skywalking.apm.collector.remote;
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 RemoteModuleGroupDefine implements ModuleGroupDefine {
public static final String GROUP_NAME = "remote";
@Override public String name() {
return GROUP_NAME;
}
@Override public Context groupContext() {
return new RemoteModuleContext(GROUP_NAME);
}
@Override public ModuleInstaller moduleInstaller() {
return new RemoteModuleInstaller();
}
}

View File

@ -0,0 +1,17 @@
package org.skywalking.apm.collector.remote;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleDefine;
import org.skywalking.apm.collector.core.module.ModuleInstaller;
/**
* @author pengys5
*/
public class RemoteModuleInstaller implements ModuleInstaller {
@Override public void install(Map<String, Map> moduleConfig,
Map<String, ModuleDefine> moduleDefineMap) throws DefineException, ClientException {
}
}

View File

@ -6,17 +6,17 @@ import org.skywalking.apm.collector.core.client.ClientException;
import org.skywalking.apm.collector.core.framework.DataInitializer;
import org.skywalking.apm.collector.core.framework.DefineException;
import org.skywalking.apm.collector.core.module.ModuleConfigParser;
import org.skywalking.apm.collector.core.module.ModuleGroup;
import org.skywalking.apm.collector.core.module.ModuleRegistration;
import org.skywalking.apm.collector.core.remote.RemoteModuleDefine;
import org.skywalking.apm.collector.core.server.Server;
import org.skywalking.apm.collector.remote.RemoteModuleDefine;
import org.skywalking.apm.collector.remote.RemoteModuleGroupDefine;
/**
* @author pengys5
*/
public class RemoteGRPCModuleDefine extends RemoteModuleDefine {
@Override protected ModuleGroup group() {
return ModuleGroup.Queue;
@Override protected String group() {
return RemoteModuleGroupDefine.GROUP_NAME;
}
@Override public boolean defaultModule() {

View File

@ -0,0 +1 @@
org.skywalking.apm.collector.remote.RemoteModuleGroupDefine

View File

@ -0,0 +1 @@
org.skywalking.apm.collector.remote.grpc.RemoteGRPCModuleDefine

View File

@ -0,0 +1,27 @@
<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="http://maven.apache.org/POM/4.0.0 http://maven.apache.org/xsd/maven-4.0.0.xsd">
<parent>
<artifactId>apm-collector</artifactId>
<groupId>org.skywalking</groupId>
<version>3.2-2017</version>
</parent>
<modelVersion>4.0.0</modelVersion>
<artifactId>apm-collector-stream</artifactId>
<packaging>jar</packaging>
<dependencies>
<dependency>
<groupId>org.skywalking</groupId>
<artifactId>apm-collector-core</artifactId>
<version>${project.version}</version>
</dependency>
<dependency>
<groupId>org.skywalking</groupId>
<artifactId>apm-collector-queue</artifactId>
<version>${project.version}</version>
</dependency>
</dependencies>
</project>

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* The <code>AbstractClusterWorker</code> implementations represent workers,

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* The <code>AbstractClusterWorkerProvider</code> implementations represent providers,

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* The <code>AbstractLocalAsyncWorker</code> implementations represent workers,

View File

@ -1,7 +1,10 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import org.skywalking.apm.collector.core.framework.CollectorContextHelper;
import org.skywalking.apm.collector.core.queue.QueueCreator;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;
import org.skywalking.apm.collector.core.queue.QueueModuleContext;
import org.skywalking.apm.collector.queue.QueueModuleContext;
import org.skywalking.apm.collector.queue.QueueModuleGroupDefine;
/**
* @author pengys5
@ -15,7 +18,8 @@ public abstract class AbstractLocalAsyncWorkerProvider<T extends AbstractLocalAs
T localAsyncWorker = workerInstance(getClusterContext());
localAsyncWorker.preStart();
QueueEventHandler queueEventHandler = QueueModuleContext.CREATOR.create(queueSize(), localAsyncWorker);
QueueCreator queueCreator = ((QueueModuleContext)CollectorContextHelper.INSTANCE.getContext(QueueModuleGroupDefine.GROUP_NAME)).getQueueCreator();
QueueEventHandler queueEventHandler = queueCreator.create(queueSize(), localAsyncWorker);
LocalAsyncWorkerRef workerRef = new LocalAsyncWorkerRef(role(), queueEventHandler);

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* The <code>AbstractLocalSyncWorker</code> defines workers who receive data from jvm inside call and response in real

View File

@ -1,9 +1,11 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import org.skywalking.apm.collector.core.framework.Executor;
/**
* @author pengys5
*/
public abstract class AbstractWorker {
public abstract class AbstractWorker implements Executor {
private final LocalWorkerContext selfContext;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* @author pengys5

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* @author pengys5

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import org.skywalking.apm.collector.core.queue.QueueEventHandler;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* @author pengys5

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* @author pengys5

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
public class ProviderNotFoundException extends Exception {
public ProviderNotFoundException(String message) {

View File

@ -0,0 +1,13 @@
package org.skywalking.apm.collector.stream;
import org.skywalking.apm.collector.stream.selector.WorkerSelector;
/**
* @author pengys5
*/
public interface Role {
String roleName();
WorkerSelector workerSelector();
}

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
public class UsedRoleNameException extends Exception {
public UsedRoleNameException(String message) {

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import java.util.ArrayList;
import java.util.List;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* Defines a general exception a worker can throw when it

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* This exception is raised when worker fails to process job during "call" or "ask"

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import java.util.Map;
import org.skywalking.apm.collector.core.client.ClientException;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
public class WorkerNotFoundException extends WorkerException {
public WorkerNotFoundException(String message) {

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
/**
* @author pengys5

View File

@ -1,7 +1,7 @@
package org.skywalking.apm.collector.core.worker;
package org.skywalking.apm.collector.stream;
import java.util.List;
import org.skywalking.apm.collector.core.worker.selector.WorkerSelector;
import org.skywalking.apm.collector.stream.selector.WorkerSelector;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

View File

@ -1,4 +1,4 @@
package org.skywalking.apm.collector.core.worker.selector;
package org.skywalking.apm.collector.stream.selector;
/**
* The <code>AbstractHashMessage</code> implementations represent aggregate message,

View File

@ -1,7 +1,7 @@
package org.skywalking.apm.collector.core.worker.selector;
package org.skywalking.apm.collector.stream.selector;
import java.util.List;
import org.skywalking.apm.collector.core.worker.WorkerRef;
import org.skywalking.apm.collector.stream.WorkerRef;
/**
* The <code>HashCodeSelector</code> is a simple implementation of {@link WorkerSelector}. It choose {@link WorkerRef}
@ -17,13 +17,13 @@ public class HashCodeSelector implements WorkerSelector<WorkerRef> {
* Use message hashcode to select {@link WorkerRef}.
*
* @param members given {@link WorkerRef} list, which size is greater than 0;
* @param message the {@link org.skywalking.apm.collector.core.worker.AbstractWorker} is going to send.
* @param message the {@link org.skywalking.apm.collector.stream.AbstractWorker} is going to send.
* @return the selected {@link WorkerRef}
*/
@Override
public WorkerRef select(List<WorkerRef> members, Object message) {
if (message instanceof AbstractHashMessage) {
AbstractHashMessage hashMessage = (AbstractHashMessage) message;
AbstractHashMessage hashMessage = (AbstractHashMessage)message;
int size = members.size();
int selectIndex = Math.abs(hashMessage.getHashCode()) % size;
return members.get(selectIndex);

View File

@ -1,7 +1,7 @@
package org.skywalking.apm.collector.core.worker.selector;
package org.skywalking.apm.collector.stream.selector;
import java.util.List;
import org.skywalking.apm.collector.core.worker.WorkerRef;
import org.skywalking.apm.collector.stream.WorkerRef;
/**
* The <code>RollingSelector</code> is a simple implementation of {@link WorkerSelector}.
@ -18,7 +18,7 @@ public class RollingSelector implements WorkerSelector<WorkerRef> {
* Use round-robin to select {@link WorkerRef}.
*
* @param members given {@link WorkerRef} list, which size is greater than 0;
* @param message message the {@link org.skywalking.apm.collector.core.worker.AbstractWorker} is going to send.
* @param message message the {@link org.skywalking.apm.collector.stream.AbstractWorker} is going to send.
* @return the selected {@link WorkerRef}
*/
@Override

View File

@ -1,7 +1,7 @@
package org.skywalking.apm.collector.core.worker.selector;
package org.skywalking.apm.collector.stream.selector;
import java.util.List;
import org.skywalking.apm.collector.core.worker.WorkerRef;
import org.skywalking.apm.collector.stream.WorkerRef;
/**
* The <code>WorkerSelector</code> should be implemented by any class whose instances
@ -18,7 +18,7 @@ public interface WorkerSelector<T extends WorkerRef> {
* select a {@link WorkerRef} from a {@link WorkerRef} list.
*
* @param members given {@link WorkerRef} list, which size is greater than 0;
* @param message the {@link org.skywalking.apm.collector.core.worker.AbstractWorker} is going to send.
* @param message the {@link org.skywalking.apm.collector.stream.AbstractWorker} is going to send.
* @return the selected {@link WorkerRef}
*/
T select(List<T> members, Object message);

View File

@ -14,6 +14,7 @@
<module>apm-collector-ui</module>
<module>apm-collector-boot</module>
<module>apm-collector-remote</module>
<module>apm-collector-stream</module>
</modules>
<parent>
<artifactId>apm</artifactId>
@ -24,10 +25,18 @@
<packaging>pom</packaging>
<properties>
<akka.version>2.4.17</akka.version>
<log4j.version>2.8.1</log4j.version>
</properties>
<dependencies>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-api</artifactId>
<version>1.7.25</version>
</dependency>
<dependency>
<groupId>ch.qos.logback</groupId>
<artifactId>logback-classic</artifactId>
<version>1.2.3</version>
</dependency>
</dependencies>
</project>