diff --git a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/test/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/service/SegmentBase64Printer.java b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/test/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/service/SegmentBase64Printer.java index 13f7d5921..5f85a462d 100644 --- a/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/test/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/service/SegmentBase64Printer.java +++ b/apm-collector/apm-collector-analysis/analysis-segment-parser/segment-parser-provider/src/test/java/org/apache/skywalking/apm/collector/analysis/segment/parser/provider/service/SegmentBase64Printer.java @@ -19,13 +19,9 @@ package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.service; import com.google.protobuf.InvalidProtocolBufferException; -import java.util.Base64; -import java.util.List; -import org.apache.skywalking.apm.network.language.agent.SpanObject; -import org.apache.skywalking.apm.network.language.agent.TraceSegmentObject; -import org.apache.skywalking.apm.network.language.agent.UniqueId; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; +import java.util.*; +import org.apache.skywalking.apm.network.language.agent.*; +import org.slf4j.*; /** * @author peng-yongsheng @@ -35,8 +31,11 @@ public class SegmentBase64Printer { private static final Logger LOGGER = LoggerFactory.getLogger(SegmentBase64Printer.class); public static void main(String[] args) throws InvalidProtocolBufferException { - String segmentBase64 = "CgoKCJbf2NPCLBAQEiAQ////////////ARiV39jTwiwg2+7Y08IsMNQPWANgARIlCAEYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAIQARif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIAxACGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgEEAMYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAUQBBif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIBhAFGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgHEAYYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAgQBxif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicICRAIGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgKEAkYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAsQChif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIDBALGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgNEAwYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCA4QDRif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEvMCCA8QDhif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADggHIAhKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXIS8wIIEBAPGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAOCAcgCEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchLzAggREBAYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgA4IByAISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEvMCCBIQERif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADggHIAhKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXIS8wIIExASGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAOCAcgCEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchLzAggUEBMYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgA4IByAISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyGP7//////////wEgCA=="; - byte[] binarySegment = Base64.getDecoder().decode(segmentBase64); + if (args.length == 0) { + return; + } + + byte[] binarySegment = Base64.getDecoder().decode(args[0]); TraceSegmentObject segmentObject = TraceSegmentObject.parseFrom(binarySegment); UniqueId segmentId = segmentObject.getTraceSegmentId(); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModule.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModule.java index 565725d30..2b0860a0d 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModule.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModule.java @@ -19,12 +19,13 @@ package org.apache.skywalking.oap.server.core; import java.util.*; -import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper; -import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper; -import org.apache.skywalking.oap.server.core.source.SourceReceiver; import org.apache.skywalking.oap.server.core.remote.RemoteSenderService; +import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataClassGetter; import org.apache.skywalking.oap.server.core.remote.client.RemoteClientManager; import org.apache.skywalking.oap.server.core.server.*; +import org.apache.skywalking.oap.server.core.source.SourceReceiver; +import org.apache.skywalking.oap.server.core.storage.model.IModelGetter; +import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer; import org.apache.skywalking.oap.server.library.module.ModuleDefine; /** @@ -53,8 +54,9 @@ public class CoreModule extends ModuleDefine { } private void addInsideService(List classes) { - classes.add(IndicatorMapper.class); - classes.add(WorkerMapper.class); + classes.add(IModelGetter.class); + classes.add(StreamDataClassGetter.class); + classes.add(WorkerAnnotationContainer.class); classes.add(RemoteClientManager.class); classes.add(RemoteSenderService.class); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java index 936ca05eb..9b428f89f 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleProvider.java @@ -18,13 +18,17 @@ package org.apache.skywalking.oap.server.core; -import org.apache.skywalking.oap.server.core.analysis.indicator.define.*; -import org.apache.skywalking.oap.server.core.analysis.worker.define.*; +import java.io.IOException; +import org.apache.skywalking.oap.server.core.annotation.AnnotationScan; import org.apache.skywalking.oap.server.core.cluster.*; import org.apache.skywalking.oap.server.core.remote.*; +import org.apache.skywalking.oap.server.core.remote.annotation.*; import org.apache.skywalking.oap.server.core.remote.client.RemoteClientManager; import org.apache.skywalking.oap.server.core.server.*; import org.apache.skywalking.oap.server.core.source.*; +import org.apache.skywalking.oap.server.core.storage.annotation.StorageAnnotationListener; +import org.apache.skywalking.oap.server.core.storage.model.IModelGetter; +import org.apache.skywalking.oap.server.core.worker.annotation.*; import org.apache.skywalking.oap.server.library.module.*; import org.apache.skywalking.oap.server.library.server.ServerException; import org.apache.skywalking.oap.server.library.server.grpc.GRPCServer; @@ -41,14 +45,22 @@ public class CoreModuleProvider extends ModuleProvider { private final CoreModuleConfig moduleConfig; private GRPCServer grpcServer; private JettyServer jettyServer; - private final IndicatorMapper indicatorMapper; - private final WorkerMapper workerMapper; + private final AnnotationScan annotationScan; + private final StorageAnnotationListener storageAnnotationListener; + private final StreamAnnotationListener streamAnnotationListener; + private final WorkerAnnotationListener workerAnnotationListener; + private final StreamDataAnnotationContainer streamDataAnnotationContainer; + private final WorkerAnnotationContainer workerAnnotationContainer; public CoreModuleProvider() { super(); this.moduleConfig = new CoreModuleConfig(); - this.indicatorMapper = new IndicatorMapper(); - this.workerMapper = new WorkerMapper(); + this.annotationScan = new AnnotationScan(); + this.storageAnnotationListener = new StorageAnnotationListener(); + this.streamAnnotationListener = new StreamAnnotationListener(); + this.workerAnnotationListener = new WorkerAnnotationListener(); + this.streamDataAnnotationContainer = new StreamDataAnnotationContainer(); + this.workerAnnotationContainer = new WorkerAnnotationContainer(); } @Override public String name() { @@ -75,20 +87,27 @@ public class CoreModuleProvider extends ModuleProvider { this.registerServiceImplementation(SourceReceiver.class, new SourceReceiverImpl(getManager())); - this.registerServiceImplementation(IndicatorMapper.class, indicatorMapper); - this.registerServiceImplementation(WorkerMapper.class, workerMapper); + this.registerServiceImplementation(StreamDataClassGetter.class, streamDataAnnotationContainer); + this.registerServiceImplementation(WorkerAnnotationContainer.class, workerAnnotationContainer); this.registerServiceImplementation(RemoteClientManager.class, new RemoteClientManager(getManager())); this.registerServiceImplementation(RemoteSenderService.class, new RemoteSenderService(getManager())); + this.registerServiceImplementation(IModelGetter.class, storageAnnotationListener); + + annotationScan.registerListener(storageAnnotationListener); + annotationScan.registerListener(streamAnnotationListener); + annotationScan.registerListener(workerAnnotationListener); } @Override public void start() throws ModuleStartException { grpcServer.addHandler(new RemoteServiceHandler(getManager())); try { - indicatorMapper.load(); - workerMapper.load(getManager()); - } catch (IndicatorDefineLoadException | WorkerDefineLoadException e) { + annotationScan.scan(() -> { + streamDataAnnotationContainer.generate(streamAnnotationListener.getStreamClasses()); + workerAnnotationContainer.load(getManager(), workerAnnotationListener.getWorkerClasses()); + }); + } catch (WorkerDefineLoadException | IOException e) { throw new ModuleStartException(e.getMessage(), e); } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/MergeDataCollection.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/MergeDataCollection.java index 07aa88c50..bd53c39fa 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/MergeDataCollection.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/MergeDataCollection.java @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.core.analysis.data; import java.util.*; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; /** * @author peng-yongsheng diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCache.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCache.java index 5d4cf4366..ab40f7dea 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCache.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCache.java @@ -18,6 +18,8 @@ package org.apache.skywalking.oap.server.core.analysis.data; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; + /** * @author peng-yongsheng */ diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCollection.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCollection.java index 7eabe479f..b0498c3e4 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCollection.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/NonMergeDataCollection.java @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.core.analysis.data; import java.util.*; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; /** * @author peng-yongsheng diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointDispatcher.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointDispatcher.java index 22445a09e..b7aef09ad 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointDispatcher.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointDispatcher.java @@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.analysis.SourceDispatcher; -import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper; +import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer; import org.apache.skywalking.oap.server.core.source.Endpoint; import org.apache.skywalking.oap.server.library.module.ModuleManager; @@ -42,7 +42,7 @@ public class EndpointDispatcher implements SourceDispatcher { private void avg(Endpoint source) { if (avgAggregator == null) { - WorkerMapper workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class); + WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); avgAggregator = (EndpointLatencyAvgAggregateWorker)workerMapper.findInstanceByClass(EndpointLatencyAvgAggregateWorker.class); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgAggregateWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgAggregateWorker.java index f9177fe57..d5388c3bf 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgAggregateWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgAggregateWorker.java @@ -19,11 +19,13 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint; import org.apache.skywalking.oap.server.core.analysis.worker.AbstractAggregatorWorker; +import org.apache.skywalking.oap.server.core.worker.annotation.Worker; import org.apache.skywalking.oap.server.library.module.ModuleManager; /** * @author peng-yongsheng */ +@Worker public class EndpointLatencyAvgAggregateWorker extends AbstractAggregatorWorker { public EndpointLatencyAvgAggregateWorker(ModuleManager moduleManager) { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgIndicator.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgIndicator.java index 7054f3d8b..6872acb55 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgIndicator.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgIndicator.java @@ -21,15 +21,19 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint; import java.util.*; import lombok.*; import org.apache.skywalking.oap.server.core.analysis.indicator.*; +import org.apache.skywalking.oap.server.core.remote.annotation.StreamData; import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData; -import org.apache.skywalking.oap.server.core.storage.annotation.Column; +import org.apache.skywalking.oap.server.core.storage.annotation.*; /** * @author peng-yongsheng */ +@StreamData +@StorageEntity(name = EndpointLatencyAvgIndicator.NAME) public class EndpointLatencyAvgIndicator extends AvgIndicator { - private static final String NAME = "endpoint_latency_avg"; + public static final String NAME = "endpoint_latency_avg"; + private static final String ID = "id"; private static final String SERVICE_ID = "service_id"; private static final String SERVICE_INSTANCE_ID = "service_instance_id"; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgPersistentWorker.java index 288d568fc..99a823986 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgPersistentWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgPersistentWorker.java @@ -19,11 +19,13 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint; import org.apache.skywalking.oap.server.core.analysis.worker.AbstractPersistentWorker; +import org.apache.skywalking.oap.server.core.worker.annotation.Worker; import org.apache.skywalking.oap.server.library.module.ModuleManager; /** * @author peng-yongsheng */ +@Worker public class EndpointLatencyAvgPersistentWorker extends AbstractPersistentWorker { public EndpointLatencyAvgPersistentWorker(ModuleManager moduleManager) { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgRemoteWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgRemoteWorker.java index 47f75215a..f36483af4 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgRemoteWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgRemoteWorker.java @@ -20,11 +20,13 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint; import org.apache.skywalking.oap.server.core.analysis.worker.AbstractRemoteWorker; import org.apache.skywalking.oap.server.core.remote.selector.Selector; +import org.apache.skywalking.oap.server.core.worker.annotation.Worker; import org.apache.skywalking.oap.server.library.module.ModuleManager; /** * @author peng-yongsheng */ +@Worker public class EndpointLatencyAvgRemoteWorker extends AbstractRemoteWorker { public EndpointLatencyAvgRemoteWorker(ModuleManager moduleManager) { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/Indicator.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/Indicator.java index de932550e..439dee51d 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/Indicator.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/Indicator.java @@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.core.analysis.indicator; import java.util.Map; import lombok.*; -import org.apache.skywalking.oap.server.core.analysis.data.StreamData; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.storage.annotation.Column; /** diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorMapper.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorMapper.java deleted file mode 100644 index b60e3f779..000000000 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorMapper.java +++ /dev/null @@ -1,86 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ - -package org.apache.skywalking.oap.server.core.analysis.indicator.define; - -import java.io.*; -import java.net.URL; -import java.util.*; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.library.module.Service; -import org.slf4j.*; - -/** - * @author peng-yongsheng - */ -public class IndicatorMapper implements Service { - - private static final Logger logger = LoggerFactory.getLogger(IndicatorMapper.class); - - private int id = 0; - private final Map, Integer> classKeyMapping; - private final Map> idKeyMapping; - - public IndicatorMapper() { - this.classKeyMapping = new HashMap<>(); - this.idKeyMapping = new HashMap<>(); - } - - @SuppressWarnings(value = "unchecked") - public void load() throws IndicatorDefineLoadException { - try { - List indicatorClasses = new LinkedList<>(); - - Enumeration urlEnumeration = this.getClass().getClassLoader().getResources("META-INF/defines/indicator.def"); - while (urlEnumeration.hasMoreElements()) { - URL definitionFileURL = urlEnumeration.nextElement(); - logger.info("Load indicator definition file url: {}", definitionFileURL.getPath()); - BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(definitionFileURL.openStream())); - Properties properties = new Properties(); - properties.load(bufferedReader); - - Enumeration defineItem = properties.propertyNames(); - while (defineItem.hasMoreElements()) { - String fullNameClass = (String)defineItem.nextElement(); - indicatorClasses.add(fullNameClass); - } - } - - for (String indicatorClassName : indicatorClasses) { - Class indicatorClass = (Class)Class.forName(indicatorClassName); - id++; - classKeyMapping.put(indicatorClass, id); - idKeyMapping.put(id, indicatorClass); - } - } catch (IOException | ClassNotFoundException e) { - throw new IndicatorDefineLoadException(e.getMessage(), e); - } - } - - public int findIdByClass(Class indicatorClass) { - return classKeyMapping.get(indicatorClass); - } - - public Class findClassById(int id) { - return idKeyMapping.get(id); - } - - public Collection> indicatorClasses() { - return idKeyMapping.values(); - } -} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractAggregatorWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractAggregatorWorker.java index 1fadf906d..3e420c9b2 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractAggregatorWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractAggregatorWorker.java @@ -24,18 +24,19 @@ import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.analysis.data.*; import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper; +import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer; +import org.apache.skywalking.oap.server.core.worker.AbstractWorker; import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.slf4j.*; /** * @author peng-yongsheng */ -public abstract class AbstractAggregatorWorker extends Worker { +public abstract class AbstractAggregatorWorker extends AbstractWorker { private static final Logger logger = LoggerFactory.getLogger(AbstractAggregatorWorker.class); - private Worker worker; + private AbstractWorker worker; private final ModuleManager moduleManager; private final DataCarrier dataCarrier; private final MergeDataCache mergeDataCache; @@ -85,7 +86,7 @@ public abstract class AbstractAggregatorWorker extends private void onNext(INPUT data) { if (worker == null) { - WorkerMapper workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class); + WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); worker = workerMapper.findInstanceByClass(nextWorkerClass()); } worker.in(data); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractPersistentWorker.java index 59f15022a..a8bee9fbc 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractPersistentWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractPersistentWorker.java @@ -22,6 +22,7 @@ import java.util.*; import org.apache.skywalking.oap.server.core.analysis.data.*; import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; import org.apache.skywalking.oap.server.core.storage.*; +import org.apache.skywalking.oap.server.core.worker.AbstractWorker; import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.slf4j.*; @@ -30,7 +31,7 @@ import static java.util.Objects.nonNull; /** * @author peng-yongsheng */ -public abstract class AbstractPersistentWorker extends Worker { +public abstract class AbstractPersistentWorker extends AbstractWorker { private static final Logger logger = LoggerFactory.getLogger(AbstractPersistentWorker.class); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractRemoteWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractRemoteWorker.java index e3fb89d0c..fc9b48de6 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractRemoteWorker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractRemoteWorker.java @@ -20,22 +20,23 @@ package org.apache.skywalking.oap.server.core.analysis.worker; import org.apache.skywalking.oap.server.core.CoreModule; import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper; +import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer; import org.apache.skywalking.oap.server.core.remote.RemoteSenderService; import org.apache.skywalking.oap.server.core.remote.selector.Selector; +import org.apache.skywalking.oap.server.core.worker.AbstractWorker; import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.slf4j.*; /** * @author peng-yongsheng */ -public abstract class AbstractRemoteWorker extends Worker { +public abstract class AbstractRemoteWorker extends AbstractWorker { private static final Logger logger = LoggerFactory.getLogger(AbstractRemoteWorker.class); private final ModuleManager moduleManager; private RemoteSenderService remoteSender; - private WorkerMapper workerMapper; + private WorkerAnnotationContainer workerMapper; public AbstractRemoteWorker(ModuleManager moduleManager) { this.moduleManager = moduleManager; @@ -46,7 +47,7 @@ public abstract class AbstractRemoteWorker extends Work remoteSender = moduleManager.find(CoreModule.NAME).getService(RemoteSenderService.class); } if (workerMapper == null) { - workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class); + workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); } try { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/define/WorkerMapper.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/define/WorkerMapper.java deleted file mode 100644 index 5d81b9b46..000000000 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/define/WorkerMapper.java +++ /dev/null @@ -1,100 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ - -package org.apache.skywalking.oap.server.core.analysis.worker.define; - -import java.io.*; -import java.lang.reflect.Constructor; -import java.net.URL; -import java.util.*; -import org.apache.skywalking.oap.server.core.analysis.worker.Worker; -import org.apache.skywalking.oap.server.library.module.*; -import org.slf4j.*; - -/** - * @author peng-yongsheng - */ -public class WorkerMapper implements Service { - - private static final Logger logger = LoggerFactory.getLogger(WorkerMapper.class); - - private int id = 0; - private final Map, Integer> classKeyMapping; - private final Map> idKeyMapping; - private final Map, Worker> classKeyInstanceMapping; - private final Map idKeyInstanceMapping; - - public WorkerMapper() { - this.classKeyMapping = new HashMap<>(); - this.idKeyMapping = new HashMap<>(); - this.classKeyInstanceMapping = new HashMap<>(); - this.idKeyInstanceMapping = new HashMap<>(); - } - - @SuppressWarnings(value = "unchecked") - public void load(ModuleManager moduleManager) throws WorkerDefineLoadException { - try { - List workerClasses = new LinkedList<>(); - - Enumeration urlEnumeration = this.getClass().getClassLoader().getResources("META-INF/defines/worker.def"); - while (urlEnumeration.hasMoreElements()) { - URL definitionFileURL = urlEnumeration.nextElement(); - logger.info("Load worker definition file url: {}", definitionFileURL.getPath()); - BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(definitionFileURL.openStream())); - Properties properties = new Properties(); - properties.load(bufferedReader); - - Enumeration defineItem = properties.propertyNames(); - while (defineItem.hasMoreElements()) { - String fullNameClass = (String)defineItem.nextElement(); - workerClasses.add(fullNameClass); - } - } - - for (String workerClassName : workerClasses) { - Class workerClass = (Class)Class.forName(workerClassName); - id++; - classKeyMapping.put(workerClass, id); - idKeyMapping.put(id, workerClass); - - Constructor constructor = workerClass.getDeclaredConstructor(ModuleManager.class); - Worker worker = constructor.newInstance(moduleManager); - classKeyInstanceMapping.put(workerClass, worker); - idKeyInstanceMapping.put(id, worker); - } - } catch (Exception e) { - throw new WorkerDefineLoadException(e.getMessage(), e); - } - } - - public int findIdByClass(Class workerClass) { - return classKeyMapping.get(workerClass); - } - - public Class findClassById(int id) { - return idKeyMapping.get(id); - } - - public Worker findInstanceByClass(Class workerClass) { - return classKeyInstanceMapping.get(workerClass); - } - - public Worker findInstanceById(int id) { - return idKeyInstanceMapping.get(id); - } -} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorDefineLoadException.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/annotation/AnnotationListener.java similarity index 77% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorDefineLoadException.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/annotation/AnnotationListener.java index ca0521a97..a64db3160 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorDefineLoadException.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/annotation/AnnotationListener.java @@ -16,14 +16,16 @@ * */ -package org.apache.skywalking.oap.server.core.analysis.indicator.define; +package org.apache.skywalking.oap.server.core.annotation; + +import java.lang.annotation.Annotation; /** * @author peng-yongsheng */ -public class IndicatorDefineLoadException extends Exception { +public interface AnnotationListener { - public IndicatorDefineLoadException(String message, Throwable cause) { - super(message, cause); - } + Class annotation(); + + void notify(Class aClass); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/annotation/AnnotationScan.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/annotation/AnnotationScan.java new file mode 100644 index 000000000..94d67cd77 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/annotation/AnnotationScan.java @@ -0,0 +1,56 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.annotation; + +import com.google.common.collect.ImmutableSet; +import com.google.common.reflect.ClassPath; +import java.io.IOException; +import java.util.*; + +/** + * @author peng-yongsheng + */ +public class AnnotationScan { + + private final List listeners; + + public AnnotationScan() { + this.listeners = new LinkedList<>(); + } + + public void registerListener(AnnotationListener listener) { + listeners.add(listener); + } + + public void scan(Runnable callBack) throws IOException { + ClassPath classpath = ClassPath.from(this.getClass().getClassLoader()); + ImmutableSet classes = classpath.getTopLevelClassesRecursive("org.apache.skywalking"); + for (ClassPath.ClassInfo classInfo : classes) { + Class aClass = classInfo.load(); + + for (AnnotationListener listener : listeners) { + if (aClass.isAnnotationPresent(listener.annotation())) { + listener.notify(aClass); + } + } + } + + callBack.run(); + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteSenderService.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteSenderService.java index b2ce1c45e..d05fcb265 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteSenderService.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteSenderService.java @@ -19,7 +19,7 @@ package org.apache.skywalking.oap.server.core.remote; import org.apache.skywalking.oap.server.core.CoreModule; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.remote.client.*; import org.apache.skywalking.oap.server.core.remote.selector.*; import org.apache.skywalking.oap.server.library.module.*; @@ -41,20 +41,20 @@ public class RemoteSenderService implements Service { this.rollingSelector = new RollingSelector(); } - public void send(int nextWorkId, Indicator indicator, Selector selector) { + public void send(int nextWorkId, StreamData streamData, Selector selector) { RemoteClientManager clientManager = moduleManager.find(CoreModule.NAME).getService(RemoteClientManager.class); RemoteClient remoteClient; switch (selector) { case HashCode: - remoteClient = hashCodeSelector.select(clientManager.getRemoteClient(), indicator); - remoteClient.push(nextWorkId, indicator); + remoteClient = hashCodeSelector.select(clientManager.getRemoteClient(), streamData); + remoteClient.push(nextWorkId, streamData); case Rolling: - remoteClient = rollingSelector.select(clientManager.getRemoteClient(), indicator); - remoteClient.push(nextWorkId, indicator); + remoteClient = rollingSelector.select(clientManager.getRemoteClient(), streamData); + remoteClient.push(nextWorkId, streamData); case ForeverFirst: - remoteClient = foreverFirstSelector.select(clientManager.getRemoteClient(), indicator); - remoteClient.push(nextWorkId, indicator); + remoteClient = foreverFirstSelector.select(clientManager.getRemoteClient(), streamData); + remoteClient.push(nextWorkId, streamData); } } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteServiceHandler.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteServiceHandler.java index 892e9516d..9b94be896 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteServiceHandler.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/RemoteServiceHandler.java @@ -19,11 +19,12 @@ package org.apache.skywalking.oap.server.core.remote; import io.grpc.stub.StreamObserver; +import java.util.Objects; import org.apache.skywalking.oap.server.core.CoreModule; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper; -import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper; +import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataClassGetter; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.remote.grpc.proto.*; +import org.apache.skywalking.oap.server.core.worker.annotation.WorkerClassGetter; import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.apache.skywalking.oap.server.library.server.grpc.GRPCHandler; import org.slf4j.*; @@ -35,24 +36,33 @@ public class RemoteServiceHandler extends RemoteServiceGrpc.RemoteServiceImplBas private static final Logger logger = LoggerFactory.getLogger(RemoteServiceHandler.class); - private final IndicatorMapper indicatorMapper; - private final WorkerMapper workerMapper; + private final ModuleManager moduleManager; + private StreamDataClassGetter streamDataClassGetter; + private WorkerClassGetter workerClassGetter; public RemoteServiceHandler(ModuleManager moduleManager) { - this.indicatorMapper = moduleManager.find(CoreModule.NAME).getService(IndicatorMapper.class); - this.workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class); + this.moduleManager = moduleManager; } @Override public StreamObserver call(StreamObserver responseObserver) { + if (Objects.isNull(streamDataClassGetter)) { + streamDataClassGetter = moduleManager.find(CoreModule.NAME).getService(StreamDataClassGetter.class); + } + if (Objects.isNull(streamDataClassGetter)) { + workerClassGetter = moduleManager.find(CoreModule.NAME).getService(WorkerClassGetter.class); + } + return new StreamObserver() { @Override public void onNext(RemoteMessage message) { - int indicatorId = message.getIndicatorId(); + int streamDataId = message.getStreamDataId(); int nextWorkerId = message.getNextWorkerId(); RemoteData remoteData = message.getRemoteData(); - Class indicatorClass = indicatorMapper.findClassById(indicatorId); + Class streamDataClass = streamDataClassGetter.findClassById(streamDataId); try { - indicatorClass.newInstance().deserialize(remoteData); + StreamData streamData = streamDataClass.newInstance(); + streamData.deserialize(remoteData); + workerClassGetter.getClassById(nextWorkerId).newInstance().in(streamData); } catch (InstantiationException | IllegalAccessException e) { logger.warn(e.getMessage()); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamAnnotationListener.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamAnnotationListener.java new file mode 100644 index 000000000..43256b119 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamAnnotationListener.java @@ -0,0 +1,49 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.remote.annotation; + +import java.lang.annotation.Annotation; +import java.util.*; +import lombok.Getter; +import org.apache.skywalking.oap.server.core.annotation.AnnotationListener; +import org.slf4j.*; + +/** + * @author peng-yongsheng + */ +public class StreamAnnotationListener implements AnnotationListener { + + private static final Logger logger = LoggerFactory.getLogger(StreamAnnotationListener.class); + + @Getter private final List streamClasses; + + public StreamAnnotationListener() { + this.streamClasses = new LinkedList<>(); + } + + @Override public Class annotation() { + return StreamData.class; + } + + @Override public void notify(Class aClass) { + logger.info("The owner class of stream data annotation, class name: {}", aClass.getName()); + + streamClasses.add(aClass); + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamData.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamData.java new file mode 100644 index 000000000..13447c00e --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamData.java @@ -0,0 +1,29 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.remote.annotation; + +import java.lang.annotation.*; + +/** + * @author peng-yongsheng + */ +@Target(ElementType.TYPE) +@Retention(RetentionPolicy.RUNTIME) +public @interface StreamData { +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamDataAnnotationContainer.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamDataAnnotationContainer.java new file mode 100644 index 000000000..4b0cb0ab1 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamDataAnnotationContainer.java @@ -0,0 +1,59 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.remote.annotation; + +import java.util.*; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; +import org.slf4j.*; + +/** + * @author peng-yongsheng + */ +public class StreamDataAnnotationContainer implements StreamDataClassGetter { + + private static final Logger logger = LoggerFactory.getLogger(StreamDataAnnotationContainer.class); + + private int id = 0; + private final Map, Integer> classMap; + private final Map> idMap; + + public StreamDataAnnotationContainer() { + this.classMap = new HashMap<>(); + this.idMap = new HashMap<>(); + } + + @SuppressWarnings(value = "unchecked") + public synchronized void generate(List streamDataClasses) { + streamDataClasses.sort(Comparator.comparing(Class::getName)); + + for (Class streamDataClass : streamDataClasses) { + id++; + classMap.put(streamDataClass, id); + idMap.put(id, streamDataClass); + } + } + + public int findIdByClass(Class streamDataClass) { + return classMap.get(streamDataClass); + } + + @Override public Class findClassById(int id) { + return idMap.get(id); + } +} diff --git a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorMapperTestCase.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamDataClassGetter.java similarity index 65% rename from oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorMapperTestCase.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamDataClassGetter.java index 4b58eb567..10386c063 100644 --- a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/IndicatorMapperTestCase.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/annotation/StreamDataClassGetter.java @@ -16,21 +16,15 @@ * */ -package org.apache.skywalking.oap.server.core.analysis.indicator.define; +package org.apache.skywalking.oap.server.core.remote.annotation; -import org.junit.*; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; +import org.apache.skywalking.oap.server.library.module.Service; /** * @author peng-yongsheng */ -public class IndicatorMapperTestCase { +public interface StreamDataClassGetter extends Service { - @Test - public void test() throws IndicatorDefineLoadException { - IndicatorMapper mapper = new IndicatorMapper(); - mapper.load(); - - Assert.assertEquals(1, mapper.findIdByClass(TestAvgIndicator.class)); - Assert.assertEquals(TestAvgIndicator.class, mapper.findClassById(1)); - } + Class findClassById(int id); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java index 75a5952ad..0afc66a04 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/GRPCRemoteClient.java @@ -23,9 +23,9 @@ import java.util.List; import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; import org.apache.skywalking.apm.commons.datacarrier.buffer.BufferStrategy; import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.cluster.RemoteInstance; +import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer; import org.apache.skywalking.oap.server.core.remote.grpc.proto.*; import org.apache.skywalking.oap.server.library.client.grpc.GRPCClient; import org.slf4j.*; @@ -39,23 +39,23 @@ public class GRPCRemoteClient implements RemoteClient, Comparable carrier; - private final IndicatorMapper indicatorMapper; + private final StreamDataAnnotationContainer streamDataMapper; - public GRPCRemoteClient(IndicatorMapper indicatorMapper, RemoteInstance remoteInstance, int channelSize, + public GRPCRemoteClient(StreamDataAnnotationContainer streamDataMapper, RemoteInstance remoteInstance, int channelSize, int bufferSize) { - this.indicatorMapper = indicatorMapper; + this.streamDataMapper = streamDataMapper; this.client = new GRPCClient(remoteInstance.getHost(), remoteInstance.getPort()); this.carrier = new DataCarrier<>(channelSize, bufferSize); this.carrier.setBufferStrategy(BufferStrategy.BLOCKING); this.carrier.consume(new RemoteMessageConsumer(), 1); } - @Override public void push(int nextWorkerId, Indicator indicator) { - int indicatorId = indicatorMapper.findIdByClass(indicator.getClass()); + @Override public void push(int nextWorkerId, StreamData streamData) { + int streamDataId = streamDataMapper.findIdByClass(streamData.getClass()); RemoteMessage.Builder builder = RemoteMessage.newBuilder(); builder.setNextWorkerId(nextWorkerId); - builder.setIndicatorId(indicatorId); - builder.setRemoteData(indicator.serialize()); + builder.setStreamDataId(streamDataId); + builder.setRemoteData(streamData.serialize()); this.carrier.produce(builder.build()); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClient.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClient.java index 98d23f101..8172a28f8 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClient.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClient.java @@ -18,7 +18,7 @@ package org.apache.skywalking.oap.server.core.remote.client; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; /** * @author peng-yongsheng @@ -29,5 +29,5 @@ public interface RemoteClient { int getPort(); - void push(int nextWorkerId, Indicator indicator); + void push(int nextWorkerId, StreamData streamData); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClientManager.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClientManager.java index 4d2049bab..361c98257 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClientManager.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClientManager.java @@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.core.remote.client; import java.util.*; import java.util.concurrent.*; -import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper; +import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer; import org.apache.skywalking.oap.server.core.cluster.*; import org.apache.skywalking.oap.server.library.module.*; import org.slf4j.*; @@ -33,7 +33,7 @@ public class RemoteClientManager implements Service { private static final Logger logger = LoggerFactory.getLogger(RemoteClientManager.class); private final ModuleManager moduleManager; - private IndicatorMapper indicatorMapper; + private StreamDataAnnotationContainer indicatorMapper; private ClusterNodesQuery clusterNodesQuery; private final List clientsA; private final List clientsB; @@ -48,7 +48,7 @@ public class RemoteClientManager implements Service { public void start() { this.clusterNodesQuery = moduleManager.find(ClusterModule.NAME).getService(ClusterNodesQuery.class); - this.indicatorMapper = moduleManager.find(ClusterModule.NAME).getService(IndicatorMapper.class); + this.indicatorMapper = moduleManager.find(ClusterModule.NAME).getService(StreamDataAnnotationContainer.class); Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(this::refresh, 1, 2, TimeUnit.SECONDS); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/SelfRemoteClient.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/SelfRemoteClient.java index 109ff7798..9e0ba7320 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/SelfRemoteClient.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/SelfRemoteClient.java @@ -19,8 +19,8 @@ package org.apache.skywalking.oap.server.core.remote.client; import org.apache.skywalking.oap.server.core.CoreModule; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; +import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer; import org.apache.skywalking.oap.server.library.module.ModuleManager; /** @@ -46,8 +46,8 @@ public class SelfRemoteClient implements RemoteClient { return port; } - @Override public void push(int nextWorkerId, Indicator indicator) { - WorkerMapper workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class); - workerMapper.findInstanceById(nextWorkerId).in(indicator); + @Override public void push(int nextWorkerId, StreamData streamData) { + WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); + workerMapper.findInstanceById(nextWorkerId).in(streamData); } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/StreamData.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/data/StreamData.java similarity index 91% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/StreamData.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/data/StreamData.java index 57c75d449..e7721437e 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/data/StreamData.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/data/StreamData.java @@ -16,8 +16,9 @@ * */ -package org.apache.skywalking.oap.server.core.analysis.data; +package org.apache.skywalking.oap.server.core.remote.data; +import org.apache.skywalking.oap.server.core.analysis.data.*; import org.apache.skywalking.oap.server.core.remote.*; /** diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/ForeverFirstSelector.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/ForeverFirstSelector.java index e28f2031f..49e8b1fd5 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/ForeverFirstSelector.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/ForeverFirstSelector.java @@ -19,7 +19,7 @@ package org.apache.skywalking.oap.server.core.remote.selector; import java.util.List; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.remote.client.RemoteClient; import org.slf4j.*; @@ -30,7 +30,7 @@ public class ForeverFirstSelector implements RemoteClientSelector { private static final Logger logger = LoggerFactory.getLogger(ForeverFirstSelector.class); - @Override public RemoteClient select(List clients, Indicator indicator) { + @Override public RemoteClient select(List clients, StreamData streamData) { if (logger.isDebugEnabled()) { logger.debug("clients size: {}", clients.size()); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/HashCodeSelector.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/HashCodeSelector.java index 3d256b5ea..673034e91 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/HashCodeSelector.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/HashCodeSelector.java @@ -19,7 +19,7 @@ package org.apache.skywalking.oap.server.core.remote.selector; import java.util.List; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.remote.client.RemoteClient; /** @@ -27,9 +27,9 @@ import org.apache.skywalking.oap.server.core.remote.client.RemoteClient; */ public class HashCodeSelector implements RemoteClientSelector { - @Override public RemoteClient select(List clients, Indicator indicator) { + @Override public RemoteClient select(List clients, StreamData streamData) { int size = clients.size(); - int selectIndex = Math.abs(indicator.hashCode()) % size; + int selectIndex = Math.abs(streamData.hashCode()) % size; return clients.get(selectIndex); } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RemoteClientSelector.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RemoteClientSelector.java index 438cbad0b..6764742aa 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RemoteClientSelector.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RemoteClientSelector.java @@ -19,12 +19,12 @@ package org.apache.skywalking.oap.server.core.remote.selector; import java.util.List; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.remote.client.RemoteClient; /** * @author peng-yongsheng */ public interface RemoteClientSelector { - RemoteClient select(List clients, Indicator indicator); + RemoteClient select(List clients, StreamData streamData); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RollingSelector.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RollingSelector.java index c74e8c364..d48cfb168 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RollingSelector.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/selector/RollingSelector.java @@ -19,7 +19,7 @@ package org.apache.skywalking.oap.server.core.remote.selector; import java.util.List; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; import org.apache.skywalking.oap.server.core.remote.client.RemoteClient; /** @@ -29,7 +29,7 @@ public class RollingSelector implements RemoteClientSelector { private int index = 0; - @Override public RemoteClient select(List clients, Indicator indicator) { + @Override public RemoteClient select(List clients, StreamData streamData) { int size = clients.size(); index++; int selectIndex = Math.abs(index) % size; diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterLockDAO.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterLockDAO.java index 9c238f6ad..e64b15390 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterLockDAO.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterLockDAO.java @@ -24,7 +24,7 @@ import org.apache.skywalking.oap.server.core.source.Scope; * @author peng-yongsheng */ public interface IRegisterLockDAO extends DAO { - boolean tryLock(Scope scope, int timeout); + boolean tryLock(Scope scope); void releaseLock(Scope scope); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageInstaller.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageInstaller.java deleted file mode 100644 index 9e4545c63..000000000 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageInstaller.java +++ /dev/null @@ -1,81 +0,0 @@ -/* - * Licensed to the Apache Software Foundation (ASF) under one or more - * contributor license agreements. See the NOTICE file distributed with - * this work for additional information regarding copyright ownership. - * The ASF licenses this file to You under the Apache License, Version 2.0 - * (the "License"); you may not use this file except in compliance with - * the License. You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - * - */ - -package org.apache.skywalking.oap.server.core.storage; - -import java.util.*; -import org.apache.skywalking.oap.server.core.CoreModule; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper; -import org.apache.skywalking.oap.server.core.storage.annotation.ColumnAnnotationRetrieval; -import org.apache.skywalking.oap.server.core.storage.define.*; -import org.apache.skywalking.oap.server.library.client.Client; -import org.apache.skywalking.oap.server.library.module.ModuleManager; -import org.slf4j.*; - -/** - * @author peng-yongsheng - */ -public abstract class StorageInstaller { - - private static final Logger logger = LoggerFactory.getLogger(StorageInstaller.class); - - private final ModuleManager moduleManager; - private final ColumnAnnotationRetrieval annotationRetrieval; - - public StorageInstaller(ModuleManager moduleManager) { - this.moduleManager = moduleManager; - this.annotationRetrieval = new ColumnAnnotationRetrieval(); - } - - public final void install(Client client) throws StorageException { - IndicatorMapper indicatorMapper = moduleManager.find(CoreModule.NAME).getService(IndicatorMapper.class); - Collection> indicatorClasses = indicatorMapper.indicatorClasses(); - - Boolean debug = System.getProperty("debug") != null; - for (Class indicatorClass : indicatorClasses) { - List columnDefines = annotationRetrieval.retrieval(indicatorClass); - - String tableName; - try { - tableName = indicatorClass.newInstance().name(); - } catch (InstantiationException | IllegalAccessException e) { - throw new StorageException(e.getMessage()); - } - TableDefine tableDefine = new TableDefine(tableName, columnDefines); - - if (!isExists(client, tableDefine)) { - logger.info("table: {} not exists", tableDefine.getName()); - createTable(client, tableDefine); - } else if (debug) { - logger.info("table: {} exists", tableDefine.getName()); - deleteTable(client, tableDefine); - createTable(client, tableDefine); - } - columnCheck(client, tableDefine); - } - } - - protected abstract boolean isExists(Client client, TableDefine tableDefine) throws StorageException; - - protected abstract void columnCheck(Client client, TableDefine tableDefine) throws StorageException; - - protected abstract void deleteTable(Client client, TableDefine tableDefine) throws StorageException; - - protected abstract void createTable(Client client, TableDefine tableDefine) throws StorageException; -} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/ColumnAnnotationRetrieval.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageAnnotationListener.java similarity index 55% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/ColumnAnnotationRetrieval.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageAnnotationListener.java index cf6a3dc1a..894326cab 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/ColumnAnnotationRetrieval.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageAnnotationListener.java @@ -18,35 +18,48 @@ package org.apache.skywalking.oap.server.core.storage.annotation; +import java.lang.annotation.Annotation; import java.lang.reflect.Field; import java.util.*; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.storage.define.*; +import lombok.Getter; +import org.apache.skywalking.oap.server.core.annotation.AnnotationListener; +import org.apache.skywalking.oap.server.core.storage.model.*; import org.slf4j.*; /** * @author peng-yongsheng */ -public class ColumnAnnotationRetrieval { +public class StorageAnnotationListener implements AnnotationListener, IModelGetter { - private static final Logger logger = LoggerFactory.getLogger(ColumnAnnotationRetrieval.class); + private static final Logger logger = LoggerFactory.getLogger(StorageAnnotationListener.class); - public List retrieval(Class indicatorClass) { - if (logger.isDebugEnabled()) { - logger.debug("Retrieval column annotation from class {}", indicatorClass.getName()); - } - List columnDefines = new LinkedList<>(); - retrieval(indicatorClass, columnDefines); - return columnDefines; + @Getter private final List models; + + public StorageAnnotationListener() { + this.models = new LinkedList<>(); } - private void retrieval(Class clazz, List columnDefines) { + @Override public Class annotation() { + return StorageEntity.class; + } + + @Override public void notify(Class aClass) { + logger.info("The owner class of storage annotation, class name: {}", aClass.getName()); + + List modelColumns = new LinkedList<>(); + retrieval(aClass, modelColumns); + + StorageEntity annotation = (StorageEntity)aClass.getAnnotation(StorageEntity.class); + models.add(new Model(annotation.name(), modelColumns)); + } + + private void retrieval(Class clazz, List modelColumns) { Field[] fields = clazz.getDeclaredFields(); for (Field field : fields) { if (field.isAnnotationPresent(Column.class)) { Column column = field.getAnnotation(Column.class); - columnDefines.add(new ColumnDefine(new ColumnName(column.columnName(), column.columnName()), field.getType())); + modelColumns.add(new ModelColumn(new ColumnName(column.columnName(), column.columnName()), field.getType())); if (logger.isDebugEnabled()) { logger.debug("The field named {} with the {} type", column.columnName(), field.getType()); } @@ -54,7 +67,7 @@ public class ColumnAnnotationRetrieval { } if (Objects.nonNull(clazz.getSuperclass())) { - retrieval(clazz.getSuperclass(), columnDefines); + retrieval(clazz.getSuperclass(), modelColumns); } } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageEntity.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageEntity.java new file mode 100644 index 000000000..cfc7ce98f --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageEntity.java @@ -0,0 +1,30 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.storage.annotation; + +import java.lang.annotation.*; + +/** + * @author peng-yongsheng + */ +@Target(ElementType.TYPE) +@Retention(RetentionPolicy.RUNTIME) +public @interface StorageEntity { + String name(); +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnName.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ColumnName.java similarity index 95% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnName.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ColumnName.java index 8dc8882ec..be7e77564 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnName.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ColumnName.java @@ -16,7 +16,7 @@ * */ -package org.apache.skywalking.oap.server.core.storage.define; +package org.apache.skywalking.oap.server.core.storage.model; /** * @author peng-yongsheng diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnTypeMapping.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/DataTypeMapping.java similarity index 89% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnTypeMapping.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/DataTypeMapping.java index 3a6c3a87e..31fcc8a65 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnTypeMapping.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/DataTypeMapping.java @@ -16,12 +16,12 @@ * */ -package org.apache.skywalking.oap.server.core.storage.define; +package org.apache.skywalking.oap.server.core.storage.model; /** * @author peng-yongsheng */ -public interface ColumnTypeMapping { +public interface DataTypeMapping { String transform(Class type); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/IModelGetter.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/IModelGetter.java new file mode 100644 index 000000000..51469f5f1 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/IModelGetter.java @@ -0,0 +1,29 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.storage.model; + +import java.util.List; +import org.apache.skywalking.oap.server.library.module.Service; + +/** + * @author peng-yongsheng + */ +public interface IModelGetter extends Service { + List getModels(); +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/TableDefine.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/Model.java similarity index 66% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/TableDefine.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/Model.java index e0f59ffaf..b3f347895 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/TableDefine.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/Model.java @@ -16,27 +16,20 @@ * */ -package org.apache.skywalking.oap.server.core.storage.define; +package org.apache.skywalking.oap.server.core.storage.model; import java.util.List; +import lombok.Getter; /** * @author peng-yongsheng */ -public class TableDefine { - private final String name; - private final List columnDefines; +public class Model { + @Getter private final String name; + @Getter private final List columns; - public TableDefine(String name, List columnDefines) { + public Model(String name, List columns) { this.name = name; - this.columnDefines = columnDefines; - } - - public final String getName() { - return name; - } - - public final List getColumnDefines() { - return columnDefines; + this.columns = columns; } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnDefine.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ModelColumn.java similarity index 71% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnDefine.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ModelColumn.java index 9bf4a87e1..c75c862d0 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/define/ColumnDefine.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ModelColumn.java @@ -16,25 +16,19 @@ * */ -package org.apache.skywalking.oap.server.core.storage.define; +package org.apache.skywalking.oap.server.core.storage.model; + +import lombok.Getter; /** * @author peng-yongsheng */ -public class ColumnDefine { - private final ColumnName columnName; - private final Class type; +public class ModelColumn { + @Getter private final ColumnName columnName; + @Getter private final Class type; - public ColumnDefine(ColumnName columnName, Class type) { + public ModelColumn(ColumnName columnName, Class type) { this.columnName = columnName; this.type = type; } - - public final ColumnName getColumnName() { - return columnName; - } - - public final Class getType() { - return type; - } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ModelInstaller.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ModelInstaller.java new file mode 100644 index 000000000..e5ae1a33a --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/model/ModelInstaller.java @@ -0,0 +1,67 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.storage.model; + +import java.util.List; +import org.apache.skywalking.oap.server.core.CoreModule; +import org.apache.skywalking.oap.server.core.storage.StorageException; +import org.apache.skywalking.oap.server.library.client.Client; +import org.apache.skywalking.oap.server.library.module.ModuleManager; +import org.slf4j.*; + +/** + * @author peng-yongsheng + */ +public abstract class ModelInstaller { + + private static final Logger logger = LoggerFactory.getLogger(ModelInstaller.class); + + private final ModuleManager moduleManager; + + public ModelInstaller(ModuleManager moduleManager) { + this.moduleManager = moduleManager; + } + + public final void install(Client client) throws StorageException { + IModelGetter modelGetter = moduleManager.find(CoreModule.NAME).getService(IModelGetter.class); + List models = modelGetter.getModels(); + + Boolean debug = System.getProperty("debug") != null; + + for (Model model : models) { + if (!isExists(client, model)) { + logger.info("table: {} not exists", model.getName()); + createTable(client, model); + } else if (debug) { + logger.info("table: {} exists", model.getName()); + deleteTable(client, model); + createTable(client, model); + } + columnCheck(client, model); + } + } + + protected abstract boolean isExists(Client client, Model model) throws StorageException; + + protected abstract void columnCheck(Client client, Model model) throws StorageException; + + protected abstract void deleteTable(Client client, Model model) throws StorageException; + + protected abstract void createTable(Client client, Model model) throws StorageException; +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/AbstractWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/AbstractWorker.java new file mode 100644 index 000000000..fd73d4c2a --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/AbstractWorker.java @@ -0,0 +1,27 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.worker; + +/** + * @author peng-yongsheng + */ +public abstract class AbstractWorker { + + public abstract void in(INPUT input); +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/Worker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/Worker.java similarity index 78% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/Worker.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/Worker.java index 53010eedf..bc4c495fd 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/Worker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/Worker.java @@ -16,14 +16,14 @@ * */ -package org.apache.skywalking.oap.server.core.analysis.worker; +package org.apache.skywalking.oap.server.core.worker.annotation; -import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import java.lang.annotation.*; /** * @author peng-yongsheng */ -public abstract class Worker { - - public abstract void in(INPUT input); +@Target(ElementType.TYPE) +@Retention(RetentionPolicy.RUNTIME) +public @interface Worker { } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationContainer.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationContainer.java new file mode 100644 index 000000000..c06cabcf5 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationContainer.java @@ -0,0 +1,84 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.worker.annotation; + +import java.lang.reflect.Constructor; +import java.util.*; +import org.apache.skywalking.oap.server.core.worker.AbstractWorker; +import org.apache.skywalking.oap.server.library.module.ModuleManager; +import org.slf4j.*; + +/** + * @author peng-yongsheng + */ +public class WorkerAnnotationContainer implements WorkerClassGetter { + + private static final Logger logger = LoggerFactory.getLogger(WorkerAnnotationContainer.class); + + private int id = 0; + private final Map, Integer> classKeyMapping; + private final Map> idKeyMapping; + private final Map, AbstractWorker> classKeyInstanceMapping; + private final Map idKeyInstanceMapping; + + public WorkerAnnotationContainer() { + this.classKeyMapping = new HashMap<>(); + this.idKeyMapping = new HashMap<>(); + this.classKeyInstanceMapping = new HashMap<>(); + this.idKeyInstanceMapping = new HashMap<>(); + } + + @SuppressWarnings(value = "unchecked") + public void load(ModuleManager moduleManager, List workerClasses) throws WorkerDefineLoadException { + if (Objects.isNull(workerClasses)) { + return; + } + + try { + for (Class workerClass : workerClasses) { + id++; + classKeyMapping.put(workerClass, id); + idKeyMapping.put(id, workerClass); + + Constructor constructor = workerClass.getDeclaredConstructor(ModuleManager.class); + AbstractWorker worker = constructor.newInstance(moduleManager); + classKeyInstanceMapping.put(workerClass, worker); + idKeyInstanceMapping.put(id, worker); + } + } catch (Throwable t) { + throw new WorkerDefineLoadException(t.getMessage(), t); + } + } + + @Override public Class getClassById(int workerId) { + return idKeyMapping.get(id); + } + + public int findIdByClass(Class workerClass) { + return classKeyMapping.get(workerClass); + } + + public AbstractWorker findInstanceByClass(Class workerClass) { + return classKeyInstanceMapping.get(workerClass); + } + + public AbstractWorker findInstanceById(int id) { + return idKeyInstanceMapping.get(id); + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationListener.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationListener.java new file mode 100644 index 000000000..51f587517 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationListener.java @@ -0,0 +1,49 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.worker.annotation; + +import java.lang.annotation.Annotation; +import java.util.*; +import lombok.Getter; +import org.apache.skywalking.oap.server.core.annotation.AnnotationListener; +import org.slf4j.*; + +/** + * @author peng-yongsheng + */ +public class WorkerAnnotationListener implements AnnotationListener { + + private static final Logger logger = LoggerFactory.getLogger(WorkerAnnotationListener.class); + + @Getter private final List workerClasses; + + public WorkerAnnotationListener() { + this.workerClasses = new LinkedList<>(); + } + + @Override public Class annotation() { + return Worker.class; + } + + @Override public void notify(Class aClass) { + logger.info("The owner class of worker annotation, class name: {}", aClass.getName()); + + workerClasses.add(aClass); + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerClassGetter.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerClassGetter.java new file mode 100644 index 000000000..83b5b8d4a --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerClassGetter.java @@ -0,0 +1,29 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + * + */ + +package org.apache.skywalking.oap.server.core.worker.annotation; + +import org.apache.skywalking.oap.server.core.worker.AbstractWorker; +import org.apache.skywalking.oap.server.library.module.Service; + +/** + * @author peng-yongsheng + */ +public interface WorkerClassGetter extends Service { + Class getClassById(int workerId); +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/define/WorkerDefineLoadException.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerDefineLoadException.java similarity index 87% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/define/WorkerDefineLoadException.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerDefineLoadException.java index a24a0a7a7..1b24afa55 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/define/WorkerDefineLoadException.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerDefineLoadException.java @@ -16,12 +16,12 @@ * */ -package org.apache.skywalking.oap.server.core.analysis.worker.define; +package org.apache.skywalking.oap.server.core.worker.annotation; /** * @author peng-yongsheng */ -public class WorkerDefineLoadException extends Exception { +public class WorkerDefineLoadException extends RuntimeException { public WorkerDefineLoadException(String message, Throwable cause) { super(message, cause); diff --git a/oap-server/server-core/src/main/proto/RemoteService.proto b/oap-server/server-core/src/main/proto/RemoteService.proto index 410b0c31e..ddc9fa129 100644 --- a/oap-server/server-core/src/main/proto/RemoteService.proto +++ b/oap-server/server-core/src/main/proto/RemoteService.proto @@ -28,7 +28,7 @@ service RemoteService { message RemoteMessage { int32 nextWorkerId = 1; - int32 indicatorId = 2; + int32 streamDataId = 2; RemoteData remoteData = 3; } diff --git a/oap-server/server-core/src/main/resources/META-INF/defines/indicator.def b/oap-server/server-core/src/main/resources/META-INF/defines/indicator.def deleted file mode 100644 index ce6b6ddc3..000000000 --- a/oap-server/server-core/src/main/resources/META-INF/defines/indicator.def +++ /dev/null @@ -1,19 +0,0 @@ -# -# Licensed to the Apache Software Foundation (ASF) under one or more -# contributor license agreements. See the NOTICE file distributed with -# this work for additional information regarding copyright ownership. -# The ASF licenses this file to You under the Apache License, Version 2.0 -# (the "License"); you may not use this file except in compliance with -# the License. You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, software -# distributed under the License is distributed on an "AS IS" BASIS, -# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -# See the License for the specific language governing permissions and -# limitations under the License. -# -# - -org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgIndicator \ No newline at end of file diff --git a/oap-server/server-core/src/main/resources/META-INF/defines/worker.def b/oap-server/server-core/src/main/resources/META-INF/defines/worker.def deleted file mode 100644 index 5afc8db52..000000000 --- a/oap-server/server-core/src/main/resources/META-INF/defines/worker.def +++ /dev/null @@ -1,21 +0,0 @@ -# -# Licensed to the Apache Software Foundation (ASF) under one or more -# contributor license agreements. See the NOTICE file distributed with -# this work for additional information regarding copyright ownership. -# The ASF licenses this file to You under the Apache License, Version 2.0 -# (the "License"); you may not use this file except in compliance with -# the License. You may obtain a copy of the License at -# -# http://www.apache.org/licenses/LICENSE-2.0 -# -# Unless required by applicable law or agreed to in writing, software -# distributed under the License is distributed on an "AS IS" BASIS, -# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -# See the License for the specific language governing permissions and -# limitations under the License. -# -# - -org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgAggregateWorker -org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgRemoteWorker -org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgPersistentWorker \ No newline at end of file diff --git a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/TestAvgIndicator.java b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/TestAvgIndicator.java index 9865ff3ec..fae88930f 100644 --- a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/TestAvgIndicator.java +++ b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/define/TestAvgIndicator.java @@ -34,6 +34,10 @@ public class TestAvgIndicator extends AvgIndicator { return null; } + @Override public String name() { + return null; + } + @Override public void deserialize(RemoteData remoteData) { } @@ -41,10 +45,6 @@ public class TestAvgIndicator extends AvgIndicator { return null; } - @Override public String name() { - return null; - } - @Override public Map toMap() { return null; } diff --git a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/StorageInstallerTestCase.java b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/StorageInstallerTestCase.java index 533c1be07..aa4d3a00b 100644 --- a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/StorageInstallerTestCase.java +++ b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/storage/StorageInstallerTestCase.java @@ -20,8 +20,8 @@ package org.apache.skywalking.oap.server.core.storage; import java.util.LinkedList; import org.apache.skywalking.oap.server.core.*; -import org.apache.skywalking.oap.server.core.analysis.indicator.define.*; -import org.apache.skywalking.oap.server.core.storage.define.TableDefine; +import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer; +import org.apache.skywalking.oap.server.core.storage.model.*; import org.apache.skywalking.oap.server.library.client.Client; import org.apache.skywalking.oap.server.library.module.*; import org.junit.Test; @@ -34,8 +34,8 @@ import org.powermock.reflect.Whitebox; public class StorageInstallerTestCase { @Test - public void testInstall() throws StorageException, DuplicateProviderException, ServiceNotProvidedException, IndicatorDefineLoadException { - IndicatorMapper indicatorMapper = new IndicatorMapper(); + public void testInstall() throws StorageException, ServiceNotProvidedException { + StreamDataAnnotationContainer streamDataAnnotationContainer = new StreamDataAnnotationContainer(); CoreModuleProvider moduleProvider = Mockito.mock(CoreModuleProvider.class); CoreModule moduleDefine = Mockito.spy(CoreModule.class); ModuleManager moduleManager = Mockito.mock(ModuleManager.class); @@ -44,33 +44,33 @@ public class StorageInstallerTestCase { moduleProviders.add(moduleProvider); Mockito.when(moduleManager.find(CoreModule.NAME)).thenReturn(moduleDefine); - Mockito.when(moduleProvider.getService(IndicatorMapper.class)).thenReturn(indicatorMapper); + Mockito.when(moduleProvider.getService(StreamDataAnnotationContainer.class)).thenReturn(streamDataAnnotationContainer); - indicatorMapper.load(); +// streamDataAnnotationContainer.generate(); - TestStorageInstaller installer = new TestStorageInstaller(moduleManager); - installer.install(null); +// TestStorageInstaller installer = new TestStorageInstaller(moduleManager); +// installer.install(null); } - class TestStorageInstaller extends StorageInstaller { + class TestStorageInstaller extends ModelInstaller { public TestStorageInstaller(ModuleManager moduleManager) { super(moduleManager); } - @Override protected boolean isExists(Client client, TableDefine tableDefine) throws StorageException { + @Override protected boolean isExists(Client client, Model tableDefine) throws StorageException { return false; } - @Override protected void columnCheck(Client client, TableDefine tableDefine) throws StorageException { + @Override protected void columnCheck(Client client, Model tableDefine) throws StorageException { } - @Override protected void deleteTable(Client client, TableDefine tableDefine) throws StorageException { + @Override protected void deleteTable(Client client, Model tableDefine) throws StorageException { } - @Override protected void createTable(Client client, TableDefine tableDefine) throws StorageException { + @Override protected void createTable(Client client, Model tableDefine) throws StorageException { } } diff --git a/oap-server/server-library/library-util/src/main/java/org/apache/skywalking/oap/server/library/util/CollectionUtils.java b/oap-server/server-library/library-util/src/main/java/org/apache/skywalking/oap/server/library/util/CollectionUtils.java index 58e68a668..8faaffad5 100644 --- a/oap-server/server-library/library-util/src/main/java/org/apache/skywalking/oap/server/library/util/CollectionUtils.java +++ b/oap-server/server-library/library-util/src/main/java/org/apache/skywalking/oap/server/library/util/CollectionUtils.java @@ -16,7 +16,6 @@ * */ - package org.apache.skywalking.oap.server.library.util; import java.util.*; diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java index 12565a1be..db1db6374 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java @@ -64,7 +64,7 @@ public class StorageModuleElasticsearchProvider extends ModuleProvider { this.registerServiceImplementation(IBatchDAO.class, new BatchProcessEsDAO(elasticSearchClient, config.getBulkActions(), config.getBulkSize(), config.getFlushInterval(), config.getConcurrentRequests())); this.registerServiceImplementation(IPersistenceDAO.class, new PersistenceEsDAO(elasticSearchClient, nameSpace)); - this.registerServiceImplementation(IRegisterLockDAO.class, new RegisterLockDAOImpl(elasticSearchClient)); + this.registerServiceImplementation(IRegisterLockDAO.class, new RegisterLockDAOImpl(elasticSearchClient, 1000)); } @Override diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/ColumnTypeEsMapping.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/ColumnTypeEsMapping.java index 8e268c1ce..41bcb7804 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/ColumnTypeEsMapping.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/ColumnTypeEsMapping.java @@ -18,12 +18,12 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base; -import org.apache.skywalking.oap.server.core.storage.define.ColumnTypeMapping; +import org.apache.skywalking.oap.server.core.storage.model.DataTypeMapping; /** * @author peng-yongsheng */ -public class ColumnTypeEsMapping implements ColumnTypeMapping { +public class ColumnTypeEsMapping implements DataTypeMapping { @Override public String transform(Class type) { if (Integer.class.equals(type) || int.class.equals(type)) { diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java index 07868d808..95fb0949d 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsInstaller.java @@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base; import java.io.IOException; import org.apache.skywalking.oap.server.core.storage.*; -import org.apache.skywalking.oap.server.core.storage.define.*; +import org.apache.skywalking.oap.server.core.storage.model.*; import org.apache.skywalking.oap.server.library.client.Client; import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; import org.apache.skywalking.oap.server.library.module.ModuleManager; @@ -31,7 +31,7 @@ import org.slf4j.*; /** * @author peng-yongsheng */ -public class StorageEsInstaller extends StorageInstaller { +public class StorageEsInstaller extends ModelInstaller { private static final Logger logger = LoggerFactory.getLogger(StorageEsInstaller.class); @@ -46,7 +46,7 @@ public class StorageEsInstaller extends StorageInstaller { this.mapping = new ColumnTypeEsMapping(); } - @Override protected boolean isExists(Client client, TableDefine tableDefine) throws StorageException { + @Override protected boolean isExists(Client client, Model tableDefine) throws StorageException { ElasticSearchClient esClient = (ElasticSearchClient)client; try { return esClient.isExistsIndex(tableDefine.getName()); @@ -55,11 +55,11 @@ public class StorageEsInstaller extends StorageInstaller { } } - @Override protected void columnCheck(Client client, TableDefine tableDefine) throws StorageException { + @Override protected void columnCheck(Client client, Model tableDefine) throws StorageException { } - @Override protected void deleteTable(Client client, TableDefine tableDefine) throws StorageException { + @Override protected void deleteTable(Client client, Model tableDefine) throws StorageException { ElasticSearchClient esClient = (ElasticSearchClient)client; try { @@ -71,7 +71,7 @@ public class StorageEsInstaller extends StorageInstaller { } } - @Override protected void createTable(Client client, TableDefine tableDefine) throws StorageException { + @Override protected void createTable(Client client, Model tableDefine) throws StorageException { ElasticSearchClient esClient = (ElasticSearchClient)client; // mapping @@ -107,7 +107,7 @@ public class StorageEsInstaller extends StorageInstaller { .build(); } - private XContentBuilder createMappingBuilder(TableDefine tableDefine) throws IOException { + private XContentBuilder createMappingBuilder(Model tableDefine) throws IOException { XContentBuilder mappingBuilder = XContentFactory.jsonBuilder() .startObject() .startObject("_all") @@ -115,7 +115,7 @@ public class StorageEsInstaller extends StorageInstaller { .endObject() .startObject("properties"); - for (ColumnDefine columnDefine : tableDefine.getColumnDefines()) { + for (ModelColumn columnDefine : tableDefine.getColumns()) { mappingBuilder .startObject(columnDefine.getColumnName().getName()) .field("type", mapping.transform(columnDefine.getType())) diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/lock/RegisterLockDAOImpl.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/lock/RegisterLockDAOImpl.java index 9904cd4fe..676adb1ed 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/lock/RegisterLockDAOImpl.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/lock/RegisterLockDAOImpl.java @@ -35,11 +35,14 @@ public class RegisterLockDAOImpl extends EsDAO implements IRegisterLockDAO { private static final Logger logger = LoggerFactory.getLogger(RegisterLockDAOImpl.class); - public RegisterLockDAOImpl(ElasticSearchClient client) { + private final int timeout; + + public RegisterLockDAOImpl(ElasticSearchClient client, int timeout) { super(client); + this.timeout = timeout; } - @Override public boolean tryLock(Scope scope, int timeout) { + @Override public boolean tryLock(Scope scope) { String id = String.valueOf(scope.ordinal()); try { GetResponse response = getClient().get(RegisterLockIndex.NAME, id);