From 38aeb92e5407d57730f6e153dc4a4308e4ef76f4 Mon Sep 17 00:00:00 2001 From: pengys5 <8082209@qq.com> Date: Mon, 17 Jul 2017 22:33:15 +0800 Subject: [PATCH] collector compile successful --- .../grpc/AgentStreamGRPCConfigParser.java | 4 +- .../grpc/impl/TraceSegmentServiceImpl.java | 63 ------------------- apm-collector/apm-collector-boot/pom.xml | 2 +- apm-collector/apm-collector-client/pom.xml | 15 +++++ .../apm/collector/client/h2/H2Client.java | 6 +- .../redis/ClusterRedisConfigParser.java | 6 +- .../redis/ClusterRedisModuleDefine.java | 5 ++ .../ClusterStandaloneModuleDefine.java | 5 ++ .../zookeeper/ClusterZKConfigParser.java | 4 +- .../zookeeper/ClusterZKDataInitializer.java | 10 +-- .../zookeeper/ClusterZKModuleDefine.java | 5 ++ .../agentstream/AgentStreamModuleDefine.java | 2 +- .../AgentStreamModuleException.java | 2 +- .../core/cluster/ClusterModuleContext.java | 4 +- .../core/cluster/ClusterModuleInstaller.java | 2 +- .../core/framework/DefinitionFile.java | 2 +- .../core/module/ModuleConfigParser.java | 1 - .../core/module/ModuleInstallerAdapter.java | 3 - .../core/queue/QueueModuleContext.java | 2 +- .../apm/collector/core/util/BytesUtils.java | 2 +- .../collector/core/util/CollectionUtils.java | 5 +- .../apm/collector/core/util/ObjectUtils.java | 4 +- .../apm/collector/core/util/StringUtils.java | 6 +- .../AbstractLocalAsyncWorkerProvider.java | 2 +- .../worker/selector/HashCodeSelector.java | 6 +- .../core/worker/selector/RollingSelector.java | 6 +- .../QueueDataCarrierModuleDefine.java | 2 +- .../disruptor/QueueDisruptorModuleDefine.java | 2 +- apm-collector/apm-collector-remote/pom.xml | 20 +++--- apm-collector/pom.xml | 1 - 30 files changed, 78 insertions(+), 121 deletions(-) delete mode 100644 apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/impl/TraceSegmentServiceImpl.java diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfigParser.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfigParser.java index 575d9134f..05fd728a7 100644 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfigParser.java +++ b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/AgentStreamGRPCConfigParser.java @@ -10,8 +10,8 @@ import org.skywalking.apm.collector.core.util.StringUtils; */ public class AgentStreamGRPCConfigParser implements ModuleConfigParser { - private final String HOST = "host"; - private final String PORT = "port"; + private static final String HOST = "host"; + private static final String PORT = "port"; @Override public void parse(Map config) throws ConfigParseException { AgentStreamGRPCConfig.HOST = (String)config.get(HOST); diff --git a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/impl/TraceSegmentServiceImpl.java b/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/impl/TraceSegmentServiceImpl.java deleted file mode 100644 index c8bc5916f..000000000 --- a/apm-collector/apm-collector-agentstream/src/main/java/org/skywalking/apm/collector/agent/stream/server/grpc/impl/TraceSegmentServiceImpl.java +++ /dev/null @@ -1,63 +0,0 @@ -package org.skywalking.apm.collector.agent.stream.server.grpc.impl; - -import io.grpc.stub.StreamObserver; -import org.apache.logging.log4j.LogManager; -import org.apache.logging.log4j.Logger; -import org.skywalking.apm.collector.actor.AbstractWorker; -import org.skywalking.apm.collector.actor.ClusterWorkerContext; -import org.skywalking.apm.collector.actor.ProviderNotFoundException; -import org.skywalking.apm.collector.actor.WorkerInvokeException; -import org.skywalking.apm.collector.actor.WorkerRef; -import org.skywalking.apm.collector.worker.grpcserver.WorkerCaller; -import org.skywalking.apm.network.proto.Downstream; -import org.skywalking.apm.network.proto.TraceSegmentServiceGrpc; -import org.skywalking.apm.network.proto.UpstreamSegment; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -/** - * @author pengys5 - */ -public class TraceSegmentServiceImpl extends TraceSegmentServiceGrpc.TraceSegmentServiceImplBase implements WorkerCaller { - - private final Logger logger = LoggerFactory.getLogger(TraceSegmentServiceImpl.class); - - private ClusterWorkerContext clusterWorkerContext; - private WorkerRef segmentReceiverWorkRef; - - @Override public void preStart() throws ProviderNotFoundException { - segmentReceiverWorkRef = clusterWorkerContext.findProvider(SegmentReceiver.WorkerRole.INSTANCE).create(AbstractWorker.noOwner()); - } - - @Override public StreamObserver collect(StreamObserver responseObserver) { - return new StreamObserver() { - @Override public void onNext(UpstreamSegment segment) { - if (logger.isDebugEnabled()) { - StringBuffer globalTraceIds = new StringBuffer(); - logger.debug("global trace ids count: %s", segment.getGlobalTraceIdsList().size()); - segment.getGlobalTraceIdsList().forEach(globalTraceId -> { - globalTraceIds.append(globalTraceId).append(","); - }); - logger.debug("receive segment, global trace ids: %s, segment byte size: %s", globalTraceIds, segment.getSegment().size()); - try { - segmentReceiverWorkRef.tell(segment); - } catch (WorkerInvokeException e) { - onError(e); - } - } - } - - @Override public void onError(Throwable throwable) { - logger.error(throwable.getMessage(), throwable); - } - - @Override public void onCompleted() { - responseObserver.onCompleted(); - } - }; - } - - @Override public void inject(ClusterWorkerContext clusterWorkerContext) { - this.clusterWorkerContext = clusterWorkerContext; - } -} diff --git a/apm-collector/apm-collector-boot/pom.xml b/apm-collector/apm-collector-boot/pom.xml index fa5a229a6..2c66d9874 100644 --- a/apm-collector/apm-collector-boot/pom.xml +++ b/apm-collector/apm-collector-boot/pom.xml @@ -15,7 +15,7 @@ org.skywalking - apm-collector-cluster-new + apm-collector-cluster ${project.version} diff --git a/apm-collector/apm-collector-client/pom.xml b/apm-collector/apm-collector-client/pom.xml index f5038e4c8..1de50be53 100644 --- a/apm-collector/apm-collector-client/pom.xml +++ b/apm-collector/apm-collector-client/pom.xml @@ -18,5 +18,20 @@ apm-collector-core ${project.version} + + com.h2database + h2 + 1.4.196 + + + redis.clients + jedis + 2.9.0 + + + org.apache.zookeeper + zookeeper + 3.4.10 + \ No newline at end of file diff --git a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java index a1b5acca9..e4b686768 100644 --- a/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java +++ b/apm-collector/apm-collector-client/src/main/java/org/skywalking/apm/collector/client/h2/H2Client.java @@ -6,12 +6,16 @@ import java.sql.ResultSet; import java.sql.SQLException; import java.sql.Statement; import org.skywalking.apm.collector.core.client.Client; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; /** * @author pengys5 */ public class H2Client implements Client { + private final Logger logger = LoggerFactory.getLogger(H2Client.class); + private Connection conn; @Override public void initialize() throws H2ClientException { @@ -40,7 +44,7 @@ public class H2Client implements Client { statement = conn.createStatement(); ResultSet rs = statement.executeQuery(sql); while (rs.next()) { - System.out.println(rs.getString("ADDRESS") + "," + rs.getString("DATA")); + logger.debug(rs.getString("ADDRESS") + "," + rs.getString("DATA")); } statement.closeOnCompletion(); } catch (SQLException e) { diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java index 63efb7680..943564ff3 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisConfigParser.java @@ -10,12 +10,12 @@ import org.skywalking.apm.collector.core.util.StringUtils; */ public class ClusterRedisConfigParser implements ModuleConfigParser { - private final String HOST = "host"; - private final String PORT = "port"; + private static final String HOST = "host"; + private static final String PORT = "port"; @Override public void parse(Map config) throws ConfigParseException { ClusterRedisConfig.HOST = (String)config.get(HOST); - ClusterRedisConfig.PORT = ((Integer)config.get(PORT)); + ClusterRedisConfig.PORT = (Integer)config.get(PORT); if (StringUtils.isEmpty(ClusterRedisConfig.HOST) || ClusterRedisConfig.PORT == 0) { throw new ConfigParseException(""); } diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java index 69afcc1a2..cc7e9c71f 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/redis/ClusterRedisModuleDefine.java @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.cluster.redis; import org.skywalking.apm.collector.client.redis.RedisClient; 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; @@ -40,4 +41,8 @@ public class ClusterRedisModuleDefine extends ClusterModuleDefine { @Override protected ClusterModuleRegistrationWriter registrationWriter() { return new ClusterRedisModuleRegistrationWriter(getClient()); } + + @Override protected ClusterModuleRegistrationReader registrationReader() { + return null; + } } diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java index baa27c40c..d7802e216 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/standalone/ClusterStandaloneModuleDefine.java @@ -3,6 +3,7 @@ package org.skywalking.apm.collector.cluster.standalone; import org.skywalking.apm.collector.client.h2.H2Client; 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; @@ -40,4 +41,8 @@ public class ClusterStandaloneModuleDefine extends ClusterModuleDefine { @Override protected ClusterModuleRegistrationWriter registrationWriter() { return new ClusterStandaloneModuleRegistrationWriter(getClient()); } + + @Override protected ClusterModuleRegistrationReader registrationReader() { + return null; + } } diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java index 82e5e86ab..9d0d62871 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKConfigParser.java @@ -10,8 +10,8 @@ import org.skywalking.apm.collector.core.util.StringUtils; */ public class ClusterZKConfigParser implements ModuleConfigParser { - private final String HOST_PORT = "hostPort"; - private final String SESSION_TIMEOUT = "sessionTimeout"; + private static final String HOST_PORT = "hostPort"; + private static final String SESSION_TIMEOUT = "sessionTimeout"; @Override public void parse(Map config) throws ConfigParseException { ClusterZKConfig.HOST_PORT = (String)config.get(HOST_PORT); diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java index 9ccb13c87..4f741245f 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKDataInitializer.java @@ -40,10 +40,10 @@ public class ClusterZKDataInitializer extends ClusterDataInitializer { pathBuilder.append("/").append(catalog); } - if (zkClient.exists(pathBuilder.toString(), false) == null) { - return false; - } else { - return true; - } +// if (zkClient.exists(pathBuilder.toString(), false) == null) { +// return false; +// } else { + return true; +// } } } diff --git a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java index 078a1a1b4..b555ac458 100644 --- a/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java +++ b/apm-collector/apm-collector-cluster/src/main/java/org/skywalking/apm/collector/cluster/zookeeper/ClusterZKModuleDefine.java @@ -4,6 +4,7 @@ import org.skywalking.apm.collector.client.zookeeper.ZookeeperClient; 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; @@ -40,4 +41,8 @@ public class ClusterZKModuleDefine extends ClusterModuleDefine { @Override protected ClusterModuleRegistrationWriter registrationWriter() { return new ClusterZKModuleRegistrationWriter(getClient()); } + + @Override protected ClusterModuleRegistrationReader registrationReader() { + return null; + } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java index e564c258a..97e377ba6 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleDefine.java @@ -24,7 +24,7 @@ public abstract class AgentStreamModuleDefine extends ModuleDefine { server.initialize(); String key = ClusterDataInitializer.BASE_CATALOG + "." + name(); - ClusterModuleContext.writer.write(key, registration().buildValue()); + ClusterModuleContext.WRITER.write(key, registration().buildValue()); } catch (ConfigParseException | ServerException e) { throw new AgentStreamModuleException(e.getMessage(), e); } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java index 11380f6b4..f46602bbd 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/agentstream/AgentStreamModuleException.java @@ -5,7 +5,7 @@ import org.skywalking.apm.collector.core.module.ModuleException; /** * @author pengys5 */ -public class AgentStreamModuleException extends ModuleException{ +public class AgentStreamModuleException extends ModuleException { public AgentStreamModuleException(String message) { super(message); } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java index 980873d4e..41a6a2c52 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleContext.java @@ -4,7 +4,7 @@ package org.skywalking.apm.collector.core.cluster; * @author pengys5 */ public class ClusterModuleContext { - public static ClusterModuleRegistrationWriter writer; + public static ClusterModuleRegistrationWriter WRITER; - public static ClusterModuleRegistrationReader reader; + public static ClusterModuleRegistrationReader READER; } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleInstaller.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleInstaller.java index 68470a624..d632bace9 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleInstaller.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/cluster/ClusterModuleInstaller.java @@ -38,6 +38,6 @@ public class ClusterModuleInstaller implements ModuleInstaller { moduleDefine = moduleDefineMap.get(clusterConfigEntry.getKey()); moduleDefine.initialize(clusterConfigEntry.getValue()); } - ClusterModuleContext.writer = ((ClusterModuleDefine)moduleDefine).registrationWriter(); + ClusterModuleContext.WRITER = ((ClusterModuleDefine)moduleDefine).registrationWriter(); } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/DefinitionFile.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/DefinitionFile.java index dae247d78..c1e4091d6 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/DefinitionFile.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/framework/DefinitionFile.java @@ -5,7 +5,7 @@ package org.skywalking.apm.collector.core.framework; */ public abstract class DefinitionFile { - private final String CATALOG = "META-INF/defines/"; + private static final String CATALOG = "META-INF/defines/"; protected abstract String fileName(); diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigParser.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigParser.java index 69330a7dd..1f2fbd941 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigParser.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleConfigParser.java @@ -1,7 +1,6 @@ package org.skywalking.apm.collector.core.module; import java.util.Map; -import org.skywalking.apm.collector.core.config.Config; import org.skywalking.apm.collector.core.config.ConfigParseException; /** diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstallerAdapter.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstallerAdapter.java index 5e9d1c27b..48c6c2d19 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstallerAdapter.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/module/ModuleInstallerAdapter.java @@ -4,7 +4,6 @@ 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; -import org.skywalking.apm.collector.core.worker.WorkerModuleInstaller; /** * @author pengys5 @@ -16,8 +15,6 @@ public class ModuleInstallerAdapter implements ModuleInstaller { public ModuleInstallerAdapter(ModuleGroup moduleGroup) { if (ModuleGroup.Cluster.equals(moduleGroup)) { moduleInstaller = new ClusterModuleInstaller(); - } else if (ModuleGroup.Worker.equals(moduleGroup)) { - moduleInstaller = new WorkerModuleInstaller(); } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/queue/QueueModuleContext.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/queue/QueueModuleContext.java index a1bf50f61..7825ca638 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/queue/QueueModuleContext.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/queue/QueueModuleContext.java @@ -4,5 +4,5 @@ package org.skywalking.apm.collector.core.queue; * @author pengys5 */ public class QueueModuleContext { - public static Creator creator; + public static Creator CREATOR; } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/BytesUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/BytesUtils.java index 9e58f5946..db9fd9bc4 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/BytesUtils.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/BytesUtils.java @@ -18,7 +18,7 @@ public class BytesUtils { long num = 0; for (int ix = 0; ix < 8; ++ix) { num <<= 8; - num |= (byteNum[ix] & 0xff); + num |= byteNum[ix] & 0xff; } return num; } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java index c959bc70c..dd66a5ecb 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/CollectionUtils.java @@ -1,6 +1,5 @@ package org.skywalking.apm.collector.core.util; -import com.sun.istack.internal.Nullable; import java.util.Map; /** @@ -8,7 +7,7 @@ import java.util.Map; */ public class CollectionUtils { - public static boolean isEmpty(@Nullable Map map) { - return (map == null || map.size() == 0); + public static boolean isEmpty(Map map) { + return map == null || map.size() == 0; } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java index f65b38d32..df264ab32 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/ObjectUtils.java @@ -1,12 +1,10 @@ package org.skywalking.apm.collector.core.util; -import com.sun.istack.internal.Nullable; - /** * @author pengys5 */ public class ObjectUtils { - public static boolean isEmpty(@Nullable Object obj) { + public static boolean isEmpty(Object obj) { return obj == null; } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java index c1b0ab3de..83413c178 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/util/StringUtils.java @@ -1,7 +1,5 @@ package org.skywalking.apm.collector.core.util; -import com.sun.istack.internal.Nullable; - /** * @author pengys5 */ @@ -9,7 +7,7 @@ public class StringUtils { public static final String EMPTY_STRING = ""; - public static boolean isEmpty(@Nullable Object str) { - return (str == null || EMPTY_STRING.equals(str)); + public static boolean isEmpty(Object str) { + return str == null || EMPTY_STRING.equals(str); } } diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorkerProvider.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorkerProvider.java index 6259e1ab7..22b3e9a58 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorkerProvider.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/AbstractLocalAsyncWorkerProvider.java @@ -15,7 +15,7 @@ public abstract class AbstractLocalAsyncWorkerProviderHashCodeSelector is a simple implementation of {@link WorkerSelector}. It choose {@link WorkerRef} @@ -19,7 +17,7 @@ public class HashCodeSelector implements WorkerSelector { * Use message hashcode to select {@link WorkerRef}. * * @param members given {@link WorkerRef} list, which size is greater than 0; - * @param message the {@link AbstractWorker} is going to send. + * @param message the {@link org.skywalking.apm.collector.core.worker.AbstractWorker} is going to send. * @return the selected {@link WorkerRef} */ @Override diff --git a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/selector/RollingSelector.java b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/selector/RollingSelector.java index a7f106b15..5bb133477 100644 --- a/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/selector/RollingSelector.java +++ b/apm-collector/apm-collector-core/src/main/java/org/skywalking/apm/collector/core/worker/selector/RollingSelector.java @@ -1,9 +1,7 @@ package org.skywalking.apm.collector.core.worker.selector; -import org.skywalking.apm.collector.actor.AbstractWorker; -import org.skywalking.apm.collector.actor.WorkerRef; - import java.util.List; +import org.skywalking.apm.collector.core.worker.WorkerRef; /** * The RollingSelector is a simple implementation of {@link WorkerSelector}. @@ -20,7 +18,7 @@ public class RollingSelector implements WorkerSelector { * Use round-robin to select {@link WorkerRef}. * * @param members given {@link WorkerRef} list, which size is greater than 0; - * @param message message the {@link AbstractWorker} is going to send. + * @param message message the {@link org.skywalking.apm.collector.core.worker.AbstractWorker} is going to send. * @return the selected {@link WorkerRef} */ @Override diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java index 1902143e1..572557711 100644 --- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java +++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/datacarrier/QueueDataCarrierModuleDefine.java @@ -25,6 +25,6 @@ public class QueueDataCarrierModuleDefine extends QueueModuleDefine { } @Override public final void initialize(Map config) throws DefineException, ClientException { - QueueModuleContext.creator = new DataCarrierCreator(); + QueueModuleContext.CREATOR = new DataCarrierCreator(); } } diff --git a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java index 9b6d48ae0..cae0f1325 100644 --- a/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java +++ b/apm-collector/apm-collector-queue/src/main/java/org/skywalking/apm/collector/queue/disruptor/QueueDisruptorModuleDefine.java @@ -25,6 +25,6 @@ public class QueueDisruptorModuleDefine extends QueueModuleDefine { } @Override public final void initialize(Map config) throws DefineException, ClientException { - QueueModuleContext.creator = new DisruptorCreator(); + QueueModuleContext.CREATOR = new DisruptorCreator(); } } diff --git a/apm-collector/apm-collector-remote/pom.xml b/apm-collector/apm-collector-remote/pom.xml index 8ca0b6141..81c877d81 100644 --- a/apm-collector/apm-collector-remote/pom.xml +++ b/apm-collector/apm-collector-remote/pom.xml @@ -19,16 +19,16 @@ - - - - - - - - - - + + org.skywalking + apm-collector-core + ${project.version} + + + org.skywalking + apm-collector-server + ${project.version} + io.grpc grpc-netty diff --git a/apm-collector/pom.xml b/apm-collector/pom.xml index 7f359e823..d7b882735 100644 --- a/apm-collector/pom.xml +++ b/apm-collector/pom.xml @@ -7,7 +7,6 @@ apm-collector-core apm-collector-queue apm-collector-storage - apm-collector-cluster apm-collector-client apm-collector-server apm-collector-discovery