From a5ad06ce46feca3585bae12e5b82be7095fe021f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?=E5=BD=AD=E5=8B=87=E5=8D=87=20pengys?= <8082209@qq.com> Date: Tue, 14 Aug 2018 14:39:03 +0800 Subject: [PATCH] Refactor register and analysis modules. (#1539) * Refactor register and analysis modules. * Fixed the startup error. --- ...encyAvgAggregateWorker.java => Const.java} | 31 ++--- .../oap/server/core/CoreModule.java | 2 - .../oap/server/core/CoreModuleProvider.java | 14 +- .../core/analysis/DispatcherManager.java | 5 +- .../analysis/endpoint/EndpointDispatcher.java | 18 +-- .../endpoint/EndpointLatencyAvgIndicator.java | 58 ++++---- .../EndpointLatencyAvgRemoteWorker.java | 43 ------ .../core/analysis/indicator/AvgIndicator.java | 3 +- .../core/analysis/indicator/Indicator.java | 10 +- .../annotation/IndicatorOperator.java} | 4 +- .../indicator/annotation/IndicatorType.java | 6 +- .../annotation/IndicatorTypeListener.java} | 23 ++-- ...ker.java => IndicatorAggregateWorker.java} | 71 ++++------ ...er.java => IndicatorPersistentWorker.java} | 57 ++++---- .../analysis/worker/IndicatorProcess.java | 64 +++++++++ ...Worker.java => IndicatorRemoteWorker.java} | 32 ++--- .../core/cache/EndpointCacheService.java | 95 ++++++++++++++ .../server/core/register/RegisterSource.java | 44 +++++++ .../core/register/endpoint/Endpoint.java | 124 ++++++++++++++++++ .../worker/RegisterDistinctWorker.java | 99 ++++++++++++++ .../worker/RegisterPersistentWorker.java | 81 ++++++++++++ .../register/worker/RegisterRemoteWorker.java | 52 ++++++++ .../core/remote/RemoteServiceHandler.java | 8 +- .../remote/client/RemoteClientManager.java | 4 +- .../core/remote/client/SelfRemoteClient.java | 11 +- .../core/source/SourceReceiverImpl.java | 5 +- ...PersistenceDAO.java => IIndicatorDAO.java} | 10 +- .../oap/server/core/storage/IRegisterDAO.java | 36 +++++ .../server/core/storage/StorageBuilder.java | 31 +++++ .../StorageDAO.java} | 9 +- .../oap/server/core/storage/StorageData.java | 26 ++++ .../server/core/storage/StorageModule.java | 2 +- .../core/storage/annotation/Column.java | 6 +- .../server/core/storage/annotation/Query.java | 26 ++++ .../annotation/StorageAnnotationListener.java | 4 +- .../storage/annotation/StorageEntity.java | 3 + .../StorageEntityAnnotationUtils.java | 46 +++++++ .../core/storage/cache/IEndpointCacheDAO.java | 32 +++++ .../server/core/worker/AbstractWorker.java | 8 ++ ...dException.java => WorkerIdGenerator.java} | 11 +- .../server/core/worker/WorkerInstances.java | 38 ++++++ .../annotation/WorkerAnnotationContainer.java | 84 ------------ .../indicator/define/TestAvgIndicator.java | 15 +-- .../resources/META-INF/defines/indicator.def | 19 --- .../resources/META-INF/defines/worker.def | 17 --- .../StorageModuleElasticsearchProvider.java | 2 +- .../plugin/elasticsearch/base/EsDAO.java | 39 ------ ...sistenceEsDAO.java => IndicatorEsDAO.java} | 35 +++-- .../elasticsearch/base/RegisterEsDAO.java | 99 ++++++++++++++ .../elasticsearch/base/StorageEsDAO.java} | 19 ++- .../cache/EndpointCacheEsDAO.java | 53 ++++++++ 51 files changed, 1151 insertions(+), 483 deletions(-) rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/{analysis/endpoint/EndpointLatencyAvgAggregateWorker.java => Const.java} (51%) delete mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgRemoteWorker.java rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/{worker/annotation/Worker.java => analysis/indicator/annotation/IndicatorOperator.java} (89%) rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/{worker/annotation/WorkerAnnotationListener.java => analysis/indicator/annotation/IndicatorTypeListener.java} (64%) rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/{AbstractAggregatorWorker.java => IndicatorAggregateWorker.java} (54%) rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/{AbstractPersistentWorker.java => IndicatorPersistentWorker.java} (68%) create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorProcess.java rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/{AbstractRemoteWorker.java => IndicatorRemoteWorker.java} (59%) create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/cache/EndpointCacheService.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/RegisterSource.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/endpoint/Endpoint.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterRemoteWorker.java rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/{IPersistenceDAO.java => IIndicatorDAO.java} (72%) create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterDAO.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageBuilder.java rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/{worker/annotation/WorkerClassGetter.java => storage/StorageDAO.java} (78%) create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageData.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Query.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageEntityAnnotationUtils.java create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/cache/IEndpointCacheDAO.java rename oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/{annotation/WorkerDefineLoadException.java => WorkerIdGenerator.java} (78%) create mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/WorkerInstances.java delete mode 100644 oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationContainer.java delete mode 100644 oap-server/server-core/src/test/resources/META-INF/defines/indicator.def delete mode 100644 oap-server/server-core/src/test/resources/META-INF/defines/worker.def rename oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/{PersistenceEsDAO.java => IndicatorEsDAO.java} (60%) create mode 100644 oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/RegisterEsDAO.java rename oap-server/{server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgPersistentWorker.java => server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsDAO.java} (58%) create mode 100644 oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/EndpointCacheEsDAO.java 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/Const.java similarity index 51% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgAggregateWorker.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/Const.java index d5388c3bf..55779f99a 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/Const.java @@ -16,23 +16,24 @@ * */ -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; +package org.apache.skywalking.oap.server.core; /** * @author peng-yongsheng */ -@Worker -public class EndpointLatencyAvgAggregateWorker extends AbstractAggregatorWorker { - - public EndpointLatencyAvgAggregateWorker(ModuleManager moduleManager) { - super(moduleManager); - } - - @Override public Class nextWorkerClass() { - return EndpointLatencyAvgRemoteWorker.class; - } +public class Const { + public static final int NONE = 0; + public static final String ID_SPLIT = "_"; + public static final int NONE_APPLICATION_ID = 1; + public static final int NONE_INSTANCE_ID = 1; + public static final int NONE_SERVICE_ID = 1; + public static final String NONE_SERVICE_NAME = "None"; + public static final String USER_CODE = "User"; + public static final String SEGMENT_SPAN_SPLIT = "S"; + public static final String UNKNOWN = "Unknown"; + public static final String EXCEPTION = "Exception"; + public static final String EMPTY_STRING = ""; + public static final String FILE_SUFFIX = "sw"; + public static final int SPAN_TYPE_VIRTUAL = 9; + public static final String DOMAIN_OPERATION_NAME = "{domain}"; } 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 2b0860a0d..27d1da6bf 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 @@ -25,7 +25,6 @@ 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; /** @@ -56,7 +55,6 @@ public class CoreModule extends ModuleDefine { private void addInsideService(List classes) { 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 9b428f89f..8e3d137b6 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 @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.core; import java.io.IOException; +import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.IndicatorTypeListener; 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.*; @@ -28,7 +29,6 @@ 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; @@ -48,9 +48,7 @@ public class CoreModuleProvider extends ModuleProvider { 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(); @@ -58,9 +56,7 @@ public class CoreModuleProvider extends ModuleProvider { 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() { @@ -85,10 +81,9 @@ public class CoreModuleProvider extends ModuleProvider { this.registerServiceImplementation(GRPCHandlerRegister.class, new GRPCHandlerRegisterImpl(grpcServer)); this.registerServiceImplementation(JettyHandlerRegister.class, new JettyHandlerRegisterImpl(jettyServer)); - this.registerServiceImplementation(SourceReceiver.class, new SourceReceiverImpl(getManager())); + this.registerServiceImplementation(SourceReceiver.class, new SourceReceiverImpl()); this.registerServiceImplementation(StreamDataClassGetter.class, streamDataAnnotationContainer); - this.registerServiceImplementation(WorkerAnnotationContainer.class, workerAnnotationContainer); this.registerServiceImplementation(RemoteClientManager.class, new RemoteClientManager(getManager())); this.registerServiceImplementation(RemoteSenderService.class, new RemoteSenderService(getManager())); @@ -96,7 +91,7 @@ public class CoreModuleProvider extends ModuleProvider { annotationScan.registerListener(storageAnnotationListener); annotationScan.registerListener(streamAnnotationListener); - annotationScan.registerListener(workerAnnotationListener); + annotationScan.registerListener(new IndicatorTypeListener(getManager())); } @Override public void start() throws ModuleStartException { @@ -105,9 +100,8 @@ public class CoreModuleProvider extends ModuleProvider { try { annotationScan.scan(() -> { streamDataAnnotationContainer.generate(streamAnnotationListener.getStreamClasses()); - workerAnnotationContainer.load(getManager(), workerAnnotationListener.getWorkerClasses()); }); - } catch (WorkerDefineLoadException | IOException e) { + } catch (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/DispatcherManager.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/DispatcherManager.java index 3dea269e6..4c0ab0f49 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/DispatcherManager.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/DispatcherManager.java @@ -21,7 +21,6 @@ package org.apache.skywalking.oap.server.core.analysis; import java.util.*; import org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointDispatcher; import org.apache.skywalking.oap.server.core.source.Scope; -import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.slf4j.*; /** @@ -33,9 +32,9 @@ public class DispatcherManager { private Map dispatcherMap; - public DispatcherManager(ModuleManager moduleManager) { + public DispatcherManager() { this.dispatcherMap = new HashMap<>(); - this.dispatcherMap.put(Scope.Endpoint, new EndpointDispatcher(moduleManager)); + this.dispatcherMap.put(Scope.Endpoint, new EndpointDispatcher()); } public SourceDispatcher getDispatcher(Scope scope) { 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 b7aef09ad..9bf8838aa 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 @@ -18,38 +18,24 @@ 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.worker.annotation.WorkerAnnotationContainer; +import org.apache.skywalking.oap.server.core.analysis.worker.IndicatorProcess; import org.apache.skywalking.oap.server.core.source.Endpoint; -import org.apache.skywalking.oap.server.library.module.ModuleManager; /** * @author peng-yongsheng */ public class EndpointDispatcher implements SourceDispatcher { - private final ModuleManager moduleManager; - private EndpointLatencyAvgAggregateWorker avgAggregator; - - public EndpointDispatcher(ModuleManager moduleManager) { - this.moduleManager = moduleManager; - } - @Override public void dispatch(Endpoint source) { avg(source); } private void avg(Endpoint source) { - if (avgAggregator == null) { - WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); - avgAggregator = (EndpointLatencyAvgAggregateWorker)workerMapper.findInstanceByClass(EndpointLatencyAvgAggregateWorker.class); - } - EndpointLatencyAvgIndicator indicator = new EndpointLatencyAvgIndicator(); indicator.setId(source.getId()); indicator.setTimeBucket(source.getTimeBucket()); indicator.combine(source.getLatency(), 1); - avgAggregator.in(indicator); + IndicatorProcess.INSTANCE.in(indicator); } } 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 6872acb55..718d357c1 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 @@ -20,20 +20,21 @@ 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.analysis.indicator.AvgIndicator; +import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.IndicatorType; 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.StorageBuilder; import org.apache.skywalking.oap.server.core.storage.annotation.*; /** * @author peng-yongsheng */ +@IndicatorType @StreamData -@StorageEntity(name = EndpointLatencyAvgIndicator.NAME) +@StorageEntity(name = "endpoint_latency_avg", builder = EndpointLatencyAvgIndicator.Builder.class) public class EndpointLatencyAvgIndicator extends AvgIndicator { - 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"; @@ -42,10 +43,6 @@ public class EndpointLatencyAvgIndicator extends AvgIndicator { @Setter @Getter @Column(columnName = SERVICE_ID) private int serviceId; @Setter @Getter @Column(columnName = SERVICE_INSTANCE_ID) private int serviceInstanceId; - @Override public String name() { - return NAME; - } - @Override public String id() { return String.valueOf(id); } @@ -99,27 +96,30 @@ public class EndpointLatencyAvgIndicator extends AvgIndicator { setValue(remoteData.getDataLongs(2)); } - @Override public Map toMap() { - Map map = new HashMap<>(); - map.put(ID, id); - map.put(SERVICE_ID, serviceId); - map.put(SERVICE_INSTANCE_ID, serviceInstanceId); - map.put(COUNT, getCount()); - map.put(SUMMATION, getSummation()); - map.put(VALUE, getValue()); - map.put(TIME_BUCKET, getTimeBucket()); - return map; - } + public static class Builder implements StorageBuilder { - @Override public Indicator newOne(Map dbMap) { - EndpointLatencyAvgIndicator indicator = new EndpointLatencyAvgIndicator(); - indicator.setId((Integer)dbMap.get(ID)); - indicator.setServiceId((Integer)dbMap.get(SERVICE_ID)); - indicator.setServiceInstanceId((Integer)dbMap.get(SERVICE_INSTANCE_ID)); - indicator.setCount((Integer)dbMap.get(COUNT)); - indicator.setSummation((Long)dbMap.get(SUMMATION)); - indicator.setValue((Long)dbMap.get(VALUE)); - indicator.setTimeBucket((Long)dbMap.get(TIME_BUCKET)); - return indicator; + @Override public EndpointLatencyAvgIndicator map2Data(Map dbMap) { + EndpointLatencyAvgIndicator indicator = new EndpointLatencyAvgIndicator(); + indicator.setId((Integer)dbMap.get(ID)); + indicator.setServiceId((Integer)dbMap.get(SERVICE_ID)); + indicator.setServiceInstanceId((Integer)dbMap.get(SERVICE_INSTANCE_ID)); + indicator.setCount((Integer)dbMap.get(COUNT)); + indicator.setSummation((Long)dbMap.get(SUMMATION)); + indicator.setValue((Long)dbMap.get(VALUE)); + indicator.setTimeBucket((Long)dbMap.get(TIME_BUCKET)); + return indicator; + } + + @Override public Map data2Map(EndpointLatencyAvgIndicator storageData) { + Map map = new HashMap<>(); + map.put(ID, storageData.getId()); + map.put(SERVICE_ID, storageData.getServiceId()); + map.put(SERVICE_INSTANCE_ID, storageData.getServiceInstanceId()); + map.put(COUNT, storageData.getCount()); + map.put(SUMMATION, storageData.getSummation()); + map.put(VALUE, storageData.getValue()); + map.put(TIME_BUCKET, storageData.getTimeBucket()); + return map; + } } } 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 deleted file mode 100644 index f36483af4..000000000 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgRemoteWorker.java +++ /dev/null @@ -1,43 +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.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) { - super(moduleManager); - } - - @Override public Selector selector() { - return Selector.HashCode; - } - - @Override public Class nextWorkerClass() { - return EndpointLatencyAvgPersistentWorker.class; - } -} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/AvgIndicator.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/AvgIndicator.java index 1bae8ec81..41f24c85a 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/AvgIndicator.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/AvgIndicator.java @@ -20,13 +20,12 @@ package org.apache.skywalking.oap.server.core.analysis.indicator; import lombok.*; import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.*; -import org.apache.skywalking.oap.server.core.remote.selector.Selector; import org.apache.skywalking.oap.server.core.storage.annotation.Column; /** * @author peng-yongsheng */ -@IndicatorType(selector = Selector.HashCode, needMerge = true) +@IndicatorOperator public abstract class AvgIndicator extends Indicator { protected static final String SUMMATION = "summation"; 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 439dee51d..fe29ed6ee 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 @@ -18,15 +18,15 @@ package org.apache.skywalking.oap.server.core.analysis.indicator; -import java.util.Map; import lombok.*; import org.apache.skywalking.oap.server.core.remote.data.StreamData; +import org.apache.skywalking.oap.server.core.storage.StorageData; import org.apache.skywalking.oap.server.core.storage.annotation.Column; /** * @author peng-yongsheng */ -public abstract class Indicator extends StreamData { +public abstract class Indicator extends StreamData implements StorageData { protected static final String TIME_BUCKET = "time_bucket"; @@ -35,10 +35,4 @@ public abstract class Indicator extends StreamData { public abstract String id(); public abstract void combine(Indicator indicator); - - public abstract String name(); - - public abstract Map toMap(); - - public abstract Indicator newOne(Map dbMap); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/Worker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorOperator.java similarity index 89% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/Worker.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorOperator.java index bc4c495fd..60ee11990 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/Worker.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorOperator.java @@ -16,7 +16,7 @@ * */ -package org.apache.skywalking.oap.server.core.worker.annotation; +package org.apache.skywalking.oap.server.core.analysis.indicator.annotation; import java.lang.annotation.*; @@ -25,5 +25,5 @@ import java.lang.annotation.*; */ @Target(ElementType.TYPE) @Retention(RetentionPolicy.RUNTIME) -public @interface Worker { +public @interface IndicatorOperator { } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorType.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorType.java index d1ad273e4..554325768 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorType.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorType.java @@ -19,15 +19,11 @@ package org.apache.skywalking.oap.server.core.analysis.indicator.annotation; import java.lang.annotation.*; -import org.apache.skywalking.oap.server.core.remote.selector.Selector; /** * @author peng-yongsheng */ @Target(ElementType.TYPE) -@Retention(RetentionPolicy.SOURCE) +@Retention(RetentionPolicy.RUNTIME) public @interface IndicatorType { - Selector selector(); - - boolean needMerge(); } 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/analysis/indicator/annotation/IndicatorTypeListener.java similarity index 64% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationListener.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/annotation/IndicatorTypeListener.java index 51f587517..bffcbbfae 100644 --- 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/analysis/indicator/annotation/IndicatorTypeListener.java @@ -16,34 +16,29 @@ * */ -package org.apache.skywalking.oap.server.core.worker.annotation; +package org.apache.skywalking.oap.server.core.analysis.indicator.annotation; import java.lang.annotation.Annotation; -import java.util.*; -import lombok.Getter; +import org.apache.skywalking.oap.server.core.analysis.worker.IndicatorProcess; import org.apache.skywalking.oap.server.core.annotation.AnnotationListener; -import org.slf4j.*; +import org.apache.skywalking.oap.server.library.module.ModuleManager; /** * @author peng-yongsheng */ -public class WorkerAnnotationListener implements AnnotationListener { +public class IndicatorTypeListener implements AnnotationListener { - private static final Logger logger = LoggerFactory.getLogger(WorkerAnnotationListener.class); + private final ModuleManager moduleManager; - @Getter private final List workerClasses; - - public WorkerAnnotationListener() { - this.workerClasses = new LinkedList<>(); + public IndicatorTypeListener(ModuleManager moduleManager) { + this.moduleManager = moduleManager; } @Override public Class annotation() { - return Worker.class; + return IndicatorType.class; } @Override public void notify(Class aClass) { - logger.info("The owner class of worker annotation, class name: {}", aClass.getName()); - - workerClasses.add(aClass); + IndicatorProcess.INSTANCE.create(moduleManager, aClass); } } 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/IndicatorAggregateWorker.java similarity index 54% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractAggregatorWorker.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorAggregateWorker.java index 3e420c9b2..8d0bf889c 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/IndicatorAggregateWorker.java @@ -21,44 +21,41 @@ package org.apache.skywalking.oap.server.core.analysis.worker; import java.util.*; import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; 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.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 AbstractWorker { +public class IndicatorAggregateWorker extends AbstractWorker { - private static final Logger logger = LoggerFactory.getLogger(AbstractAggregatorWorker.class); + private static final Logger logger = LoggerFactory.getLogger(IndicatorAggregateWorker.class); - private AbstractWorker worker; - private final ModuleManager moduleManager; - private final DataCarrier dataCarrier; - private final MergeDataCache mergeDataCache; + private AbstractWorker nextWorker; + private final DataCarrier dataCarrier; + private final MergeDataCache mergeDataCache; private int messageNum; - public AbstractAggregatorWorker(ModuleManager moduleManager) { - this.moduleManager = moduleManager; + IndicatorAggregateWorker(int workerId, AbstractWorker nextWorker) { + super(workerId); + this.nextWorker = nextWorker; this.mergeDataCache = new MergeDataCache<>(); this.dataCarrier = new DataCarrier<>(1, 10000); this.dataCarrier.consume(new AggregatorConsumer(this), 1); } - @Override public final void in(INPUT input) { - input.setEndOfBatchContext(new EndOfBatchContext(false)); - dataCarrier.produce(input); + @Override public final void in(Indicator indicator) { + indicator.setEndOfBatchContext(new EndOfBatchContext(false)); + dataCarrier.produce(indicator); } - private void onWork(INPUT message) { + private void onWork(Indicator indicator) { messageNum++; - aggregate(message); + aggregate(indicator); - if (messageNum >= 1000 || message.getEndOfBatchContext().isEndOfBatch()) { + if (messageNum >= 1000 || indicator.getEndOfBatchContext().isEndOfBatch()) { sendToNext(); messageNum = 0; } @@ -74,41 +71,31 @@ public abstract class AbstractAggregatorWorker extends } } - mergeDataCache.getLast().collection().forEach((INPUT key, INPUT data) -> { + mergeDataCache.getLast().collection().forEach((Indicator key, Indicator data) -> { if (logger.isDebugEnabled()) { logger.debug(data.toString()); } - onNext(data); + nextWorker.in(data); }); mergeDataCache.finishReadingLast(); } - private void onNext(INPUT data) { - if (worker == null) { - WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); - worker = workerMapper.findInstanceByClass(nextWorkerClass()); - } - worker.in(data); - } - - public abstract Class nextWorkerClass(); - - private void aggregate(INPUT message) { + private void aggregate(Indicator indicator) { mergeDataCache.writing(); - if (mergeDataCache.containsKey(message)) { - mergeDataCache.get(message).combine(message); + if (mergeDataCache.containsKey(indicator)) { + mergeDataCache.get(indicator).combine(indicator); } else { - mergeDataCache.put(message); + mergeDataCache.put(indicator); } mergeDataCache.finishWriting(); } - private class AggregatorConsumer implements IConsumer { + private class AggregatorConsumer implements IConsumer { - private final AbstractAggregatorWorker aggregator; + private final IndicatorAggregateWorker aggregator; - private AggregatorConsumer(AbstractAggregatorWorker aggregator) { + private AggregatorConsumer(IndicatorAggregateWorker aggregator) { this.aggregator = aggregator; } @@ -116,21 +103,21 @@ public abstract class AbstractAggregatorWorker extends } - @Override public void consume(List data) { - Iterator inputIterator = data.iterator(); + @Override public void consume(List data) { + Iterator inputIterator = data.iterator(); int i = 0; while (inputIterator.hasNext()) { - INPUT input = inputIterator.next(); + Indicator indicator = inputIterator.next(); i++; if (i == data.size()) { - input.getEndOfBatchContext().setEndOfBatch(true); + indicator.getEndOfBatchContext().setEndOfBatch(true); } - aggregator.onWork(input); + aggregator.onWork(indicator); } } - @Override public void onError(List data, Throwable t) { + @Override public void onError(List data, Throwable t) { logger.error(t.getMessage(), t); } 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/IndicatorPersistentWorker.java similarity index 68% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractPersistentWorker.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorPersistentWorker.java index a8bee9fbc..dd140c88b 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/IndicatorPersistentWorker.java @@ -31,26 +31,31 @@ import static java.util.Objects.nonNull; /** * @author peng-yongsheng */ -public abstract class AbstractPersistentWorker extends AbstractWorker { +public class IndicatorPersistentWorker extends AbstractWorker { - private static final Logger logger = LoggerFactory.getLogger(AbstractPersistentWorker.class); + private static final Logger logger = LoggerFactory.getLogger(IndicatorPersistentWorker.class); - private final MergeDataCache mergeDataCache; + private final String modelName; + private final MergeDataCache mergeDataCache; private final IBatchDAO batchDAO; - private final IPersistenceDAO persistenceDAO; - private final int blockBatchPersistenceSize = 1000; + private final IIndicatorDAO indicatorDAO; + private final int blockBatchPersistenceSize; - public AbstractPersistentWorker(ModuleManager moduleManager) { + IndicatorPersistentWorker(int workerId, String modelName, int batchSize, ModuleManager moduleManager, + IIndicatorDAO indicatorDAO) { + super(workerId); + this.modelName = modelName; + this.blockBatchPersistenceSize = batchSize; this.mergeDataCache = new MergeDataCache<>(); this.batchDAO = moduleManager.find(StorageModule.NAME).getService(IBatchDAO.class); - this.persistenceDAO = moduleManager.find(StorageModule.NAME).getService(IPersistenceDAO.class); + this.indicatorDAO = indicatorDAO; } - public final Window> getCache() { + public final Window> getCache() { return mergeDataCache; } - @Override public final void in(INPUT input) { + @Override public final void in(Indicator input) { if (getCache().currentCollectionSize() >= blockBatchPersistenceSize) { try { if (getCache().trySwitchPointer()) { @@ -86,33 +91,25 @@ public abstract class AbstractPersistentWorker extends return batchCollection; } - private List prepareBatch(MergeDataCollection collection) { + private List prepareBatch(MergeDataCollection collection) { List batchCollection = new LinkedList<>(); collection.collection().forEach((id, data) -> { - if (needMergeDBData()) { - INPUT dbData = null; + Indicator dbData = null; + try { + dbData = indicatorDAO.get(modelName, data); + } catch (Throwable t) { + logger.error(t.getMessage(), t); + } + if (nonNull(dbData)) { + dbData.combine(data); try { - dbData = persistenceDAO.get(data); + batchCollection.add(indicatorDAO.prepareBatchUpdate(modelName, dbData)); } catch (Throwable t) { logger.error(t.getMessage(), t); } - if (nonNull(dbData)) { - dbData.combine(data); - try { - batchCollection.add(persistenceDAO.prepareBatchUpdate(dbData)); - } catch (Throwable t) { - logger.error(t.getMessage(), t); - } - } else { - try { - batchCollection.add(persistenceDAO.prepareBatchInsert(data)); - } catch (Throwable t) { - logger.error(t.getMessage(), t); - } - } } else { try { - batchCollection.add(persistenceDAO.prepareBatchInsert(data)); + batchCollection.add(indicatorDAO.prepareBatchInsert(modelName, data)); } catch (Throwable t) { logger.error(t.getMessage(), t); } @@ -122,7 +119,7 @@ public abstract class AbstractPersistentWorker extends return batchCollection; } - private void cacheData(INPUT input) { + private void cacheData(Indicator input) { mergeDataCache.writing(); if (mergeDataCache.containsKey(input)) { mergeDataCache.get(input).combine(input); @@ -132,6 +129,4 @@ public abstract class AbstractPersistentWorker extends mergeDataCache.finishWriting(); } - - protected abstract boolean needMergeDBData(); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorProcess.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorProcess.java new file mode 100644 index 000000000..398d61d18 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorProcess.java @@ -0,0 +1,64 @@ +/* + * 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; + +import java.util.*; +import org.apache.skywalking.oap.server.core.UnexpectedException; +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.storage.annotation.StorageEntityAnnotationUtils; +import org.apache.skywalking.oap.server.core.worker.*; +import org.apache.skywalking.oap.server.library.module.ModuleManager; + +/** + * @author peng-yongsheng + */ +public enum IndicatorProcess { + INSTANCE; + + private Map, IndicatorAggregateWorker> entryWorkers = new HashMap<>(); + + public void in(Indicator indicator) { + entryWorkers.get(indicator.getClass()).in(indicator); + } + + public void create(ModuleManager moduleManager, Class indicatorClass) { + String modelName = StorageEntityAnnotationUtils.getModelName(indicatorClass); + Class builderClass = StorageEntityAnnotationUtils.getBuilder(indicatorClass); + + StorageDAO storageDAO = moduleManager.find(StorageModule.NAME).getService(StorageDAO.class); + IIndicatorDAO indicatorDAO; + try { + indicatorDAO = storageDAO.newIndicatorDao(builderClass.newInstance()); + } catch (InstantiationException | IllegalAccessException e) { + throw new UnexpectedException(""); + } + + IndicatorPersistentWorker persistentWorker = new IndicatorPersistentWorker(WorkerIdGenerator.INSTANCES.generate(), modelName, 1000, moduleManager, indicatorDAO); + WorkerInstances.INSTANCES.put(persistentWorker.getWorkerId(), persistentWorker); + + IndicatorRemoteWorker remoteWorker = new IndicatorRemoteWorker(WorkerIdGenerator.INSTANCES.generate(), moduleManager, persistentWorker); + WorkerInstances.INSTANCES.put(remoteWorker.getWorkerId(), remoteWorker); + + IndicatorAggregateWorker aggregateWorker = new IndicatorAggregateWorker(WorkerIdGenerator.INSTANCES.generate(), remoteWorker); + WorkerInstances.INSTANCES.put(aggregateWorker.getWorkerId(), aggregateWorker); + + entryWorkers.put(indicatorClass, aggregateWorker); + } +} 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/IndicatorRemoteWorker.java similarity index 59% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/AbstractRemoteWorker.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/worker/IndicatorRemoteWorker.java index fc9b48de6..0295864f9 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/IndicatorRemoteWorker.java @@ -20,7 +20,6 @@ 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.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; @@ -30,35 +29,24 @@ import org.slf4j.*; /** * @author peng-yongsheng */ -public abstract class AbstractRemoteWorker extends AbstractWorker { +public class IndicatorRemoteWorker extends AbstractWorker { - private static final Logger logger = LoggerFactory.getLogger(AbstractRemoteWorker.class); + private static final Logger logger = LoggerFactory.getLogger(IndicatorRemoteWorker.class); - private final ModuleManager moduleManager; - private RemoteSenderService remoteSender; - private WorkerAnnotationContainer workerMapper; + private final AbstractWorker nextWorker; + private final RemoteSenderService remoteSender; - public AbstractRemoteWorker(ModuleManager moduleManager) { - this.moduleManager = moduleManager; + IndicatorRemoteWorker(int workerId, ModuleManager moduleManager, AbstractWorker nextWorker) { + super(workerId); + this.remoteSender = moduleManager.find(CoreModule.NAME).getService(RemoteSenderService.class); + this.nextWorker = nextWorker; } - @Override public final void in(INPUT input) { - if (remoteSender == null) { - remoteSender = moduleManager.find(CoreModule.NAME).getService(RemoteSenderService.class); - } - if (workerMapper == null) { - workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); - } - + @Override public final void in(Indicator indicator) { try { - int nextWorkerId = workerMapper.findIdByClass(nextWorkerClass()); - remoteSender.send(nextWorkerId, input, selector()); + remoteSender.send(nextWorker.getWorkerId(), indicator, Selector.HashCode); } catch (Throwable e) { logger.error(e.getMessage(), e); } } - - public abstract Class nextWorkerClass(); - - public abstract Selector selector(); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/cache/EndpointCacheService.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/cache/EndpointCacheService.java new file mode 100644 index 000000000..71e581439 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/cache/EndpointCacheService.java @@ -0,0 +1,95 @@ +/* + * 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.cache; + +import com.google.common.cache.*; +import org.apache.skywalking.oap.server.core.Const; +import org.apache.skywalking.oap.server.core.register.endpoint.Endpoint; +import org.apache.skywalking.oap.server.core.storage.StorageModule; +import org.apache.skywalking.oap.server.core.storage.cache.IEndpointCacheDAO; +import org.apache.skywalking.oap.server.library.module.*; +import org.slf4j.*; + +import static java.util.Objects.*; + +/** + * @author peng-yongsheng + */ +public class EndpointCacheService implements Service { + + private static final Logger logger = LoggerFactory.getLogger(EndpointCacheService.class); + + private final ModuleManager moduleManager; + private IEndpointCacheDAO cacheDAO; + + public EndpointCacheService(ModuleManager moduleManager) { + this.moduleManager = moduleManager; + } + + private final Cache idCache = CacheBuilder.newBuilder().initialCapacity(1000).maximumSize(1000000).build(); + + private final Cache sequenceCache = CacheBuilder.newBuilder().initialCapacity(1000).maximumSize(1000000).build(); + + public int get(int serviceId, String serviceName, int srcSpanType) { + String id = serviceId + Const.ID_SPLIT + serviceName + Const.ID_SPLIT + srcSpanType; + + int endpointId = 0; + + try { + endpointId = idCache.get(id, () -> getCacheDAO().get(id)); + } catch (Throwable e) { + logger.error(e.getMessage(), e); + } + + if (serviceId == 0) { + endpointId = getCacheDAO().get(id); + if (endpointId != 0) { + idCache.put(id, endpointId); + } + } + return endpointId; + } + + public Endpoint get(int endpointId) { + Endpoint endpoint = null; + try { + endpoint = sequenceCache.get(endpointId, () -> getCacheDAO().get(endpointId)); + } catch (Throwable e) { + logger.error(e.getMessage(), e); + } + + if (isNull(endpoint)) { + endpoint = getCacheDAO().get(endpointId); + if (nonNull(endpoint)) { + sequenceCache.put(endpointId, endpoint); + } else { + logger.warn("Endpoint id {} is not in cache and persistent storage.", endpointId); + } + } + + return endpoint; + } + + private IEndpointCacheDAO getCacheDAO() { + if (isNull(cacheDAO)) { + cacheDAO = moduleManager.find(StorageModule.NAME).getService(IEndpointCacheDAO.class); + } + return cacheDAO; + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/RegisterSource.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/RegisterSource.java new file mode 100644 index 000000000..abb6c9edf --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/RegisterSource.java @@ -0,0 +1,44 @@ +/* + * 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.register; + +import lombok.*; +import org.apache.skywalking.oap.server.core.remote.data.StreamData; +import org.apache.skywalking.oap.server.core.storage.StorageData; +import org.apache.skywalking.oap.server.core.storage.annotation.Column; + +/** + * @author peng-yongsheng + */ +public abstract class RegisterSource extends StreamData implements StorageData { + + public static final String SEQUENCE = "sequence"; + protected static final String REGISTER_TIME = "register_time"; + protected static final String HEARTBEAT_TIME = "heartbeat_time"; + + @Getter @Setter @Column(columnName = SEQUENCE) private int sequence; + @Getter @Setter @Column(columnName = REGISTER_TIME) private long registerTime; + @Getter @Setter @Column(columnName = HEARTBEAT_TIME) private long heartbeatTime; + + public final void combine(RegisterSource registerSource) { + if (heartbeatTime < registerSource.getHeartbeatTime()) { + heartbeatTime = registerSource.getHeartbeatTime(); + } + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/endpoint/Endpoint.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/endpoint/Endpoint.java new file mode 100644 index 000000000..853695559 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/endpoint/Endpoint.java @@ -0,0 +1,124 @@ +/* + * 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.register.endpoint; + +import java.util.*; +import lombok.*; +import org.apache.skywalking.oap.server.core.Const; +import org.apache.skywalking.oap.server.core.register.RegisterSource; +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.StorageBuilder; +import org.apache.skywalking.oap.server.core.storage.annotation.*; + +/** + * @author peng-yongsheng + */ +@StreamData +@StorageEntity(name = "endpoint", builder = Endpoint.Builder.class) +public class Endpoint extends RegisterSource { + + private static final String SERVICE_ID = "service_id"; + private static final String NAME = "name"; + private static final String SRC_SPAN_TYPE = "src_span_type"; + + @Setter @Getter @Column(columnName = SERVICE_ID) private int serviceId; + @Setter @Getter @Column(columnName = NAME, matchQuery = true) private String name; + @Setter @Getter @Column(columnName = SRC_SPAN_TYPE) private int srcSpanType; + + @Override public String id() { + return String.valueOf(serviceId) + Const.ID_SPLIT + name + Const.ID_SPLIT + String.valueOf(srcSpanType); + } + + @Override public int hashCode() { + int result = 17; + result = 31 * result + serviceId; + result = 31 * result + name.hashCode(); + result = 31 * result + srcSpanType; + return result; + } + + @Override public boolean equals(Object obj) { + if (this == obj) + return true; + if (obj == null) + return false; + if (getClass() != obj.getClass()) + return false; + + Endpoint source = (Endpoint)obj; + if (serviceId != source.getServiceId()) + return false; + if (name.equals(source.getName())) + return false; + if (srcSpanType != source.getSrcSpanType()) + return false; + + return true; + } + + @Override public RemoteData.Builder serialize() { + RemoteData.Builder remoteBuilder = RemoteData.newBuilder(); + remoteBuilder.setDataIntegers(0, getSequence()); + remoteBuilder.setDataIntegers(1, serviceId); + remoteBuilder.setDataIntegers(2, srcSpanType); + + remoteBuilder.setDataLongs(0, getRegisterTime()); + remoteBuilder.setDataLongs(1, getHeartbeatTime()); + + remoteBuilder.setDataStrings(0, name); + return remoteBuilder; + } + + @Override public void deserialize(RemoteData remoteData) { + setSequence(remoteData.getDataIntegers(0)); + setServiceId(remoteData.getDataIntegers(1)); + setSrcSpanType(remoteData.getDataIntegers(2)); + + setRegisterTime(remoteData.getDataLongs(0)); + setHeartbeatTime(remoteData.getDataLongs(1)); + + setName(remoteData.getDataStrings(1)); + } + + public static class Builder implements StorageBuilder { + + @Override public Endpoint map2Data(Map dbMap) { + Endpoint endpoint = new Endpoint(); + endpoint.setSequence((Integer)dbMap.get(SEQUENCE)); + endpoint.setServiceId((Integer)dbMap.get(SERVICE_ID)); + endpoint.setName((String)dbMap.get(NAME)); + endpoint.setSrcSpanType((Integer)dbMap.get(SRC_SPAN_TYPE)); + endpoint.setRegisterTime((Long)dbMap.get(REGISTER_TIME)); + endpoint.setHeartbeatTime((Long)dbMap.get(HEARTBEAT_TIME)); + return endpoint; + } + + @Override public Map data2Map(Endpoint storageData) { + Map map = new HashMap<>(); + map.put(SEQUENCE, storageData.getSequence()); + map.put(SERVICE_ID, storageData.getServiceId()); + map.put(NAME, storageData.getName()); + map.put(SRC_SPAN_TYPE, storageData.getSrcSpanType()); + map.put(REGISTER_TIME, storageData.getRegisterTime()); + map.put(HEARTBEAT_TIME, storageData.getHeartbeatTime()); + return map; + } + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java new file mode 100644 index 000000000..640e1599e --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterDistinctWorker.java @@ -0,0 +1,99 @@ +/* + * 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.register.worker; + +import java.util.*; +import org.apache.skywalking.apm.commons.datacarrier.DataCarrier; +import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer; +import org.apache.skywalking.oap.server.core.analysis.data.EndOfBatchContext; +import org.apache.skywalking.oap.server.core.register.RegisterSource; +import org.apache.skywalking.oap.server.core.worker.AbstractWorker; +import org.slf4j.*; + +/** + * @author peng-yongsheng + */ +public class RegisterDistinctWorker extends AbstractWorker { + + private static final Logger logger = LoggerFactory.getLogger(RegisterDistinctWorker.class); + + private final AbstractWorker nextWorker; + private final DataCarrier dataCarrier; + private final Map sources; + private int messageNum; + + public RegisterDistinctWorker(int workerId, AbstractWorker nextWorker) { + super(workerId); + this.nextWorker = nextWorker; + this.sources = new HashMap<>(); + this.dataCarrier = new DataCarrier<>(1, 10000); + this.dataCarrier.consume(new AggregatorConsumer(this), 1); + } + + @Override public final void in(RegisterSource source) { + source.setEndOfBatchContext(new EndOfBatchContext(false)); + dataCarrier.produce(source); + } + + private void onWork(RegisterSource source) { + messageNum++; + + if (!sources.containsKey(source)) { + sources.get(source).combine(source); + } + + if (messageNum >= 1000 || source.getEndOfBatchContext().isEndOfBatch()) { + sources.values().forEach(nextWorker::in); + messageNum = 0; + } + } + + private class AggregatorConsumer implements IConsumer { + + private final RegisterDistinctWorker aggregator; + + private AggregatorConsumer(RegisterDistinctWorker aggregator) { + this.aggregator = aggregator; + } + + @Override public void init() { + } + + @Override public void consume(List sources) { + Iterator sourceIterator = sources.iterator(); + + int i = 0; + while (sourceIterator.hasNext()) { + RegisterSource source = sourceIterator.next(); + i++; + if (i == sources.size()) { + source.getEndOfBatchContext().setEndOfBatch(true); + } + aggregator.onWork(source); + } + } + + @Override public void onError(List sources, Throwable t) { + logger.error(t.getMessage(), t); + } + + @Override public void onExit() { + } + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java new file mode 100644 index 000000000..0e353737f --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterPersistentWorker.java @@ -0,0 +1,81 @@ +/* + * 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.register.worker; + +import java.util.*; +import org.apache.skywalking.oap.server.core.register.RegisterSource; +import org.apache.skywalking.oap.server.core.source.Scope; +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.*; + +/** + * @author peng-yongsheng + */ +public class RegisterPersistentWorker extends AbstractWorker { + + private static final Logger logger = LoggerFactory.getLogger(RegisterPersistentWorker.class); + + private final Scope scope; + private final String modelName; + private final Map sources; + private final IRegisterLockDAO registerLockDAO; + private final IRegisterDAO registerDAO; + + public RegisterPersistentWorker(int workerId, String modelName, ModuleManager moduleManager, + IRegisterDAO registerDAO, Scope scope) { + super(workerId); + this.modelName = modelName; + this.sources = new HashMap<>(); + this.registerDAO = registerDAO; + this.registerLockDAO = moduleManager.find(StorageModule.NAME).getService(IRegisterLockDAO.class); + this.scope = scope; + } + + @Override public final void in(RegisterSource registerSource) { + if (!sources.containsKey(registerSource)) { + sources.put(registerSource, registerSource); + } + if (registerSource.getEndOfBatchContext().isEndOfBatch()) { + + if (registerLockDAO.tryLock(scope)) { + try { + sources.values().forEach(source -> { + try { + RegisterSource newSource = registerDAO.get(modelName, registerSource.id()); + if (Objects.nonNull(newSource)) { + newSource.combine(newSource); + int sequence = registerDAO.max(modelName); + newSource.setSequence(sequence); + registerDAO.forceInsert(modelName, newSource); + } else { + registerDAO.forceUpdate(modelName, newSource); + } + } catch (Throwable t) { + logger.error(t.getMessage()); + } + }); + } finally { + registerLockDAO.releaseLock(scope); + } + } + } + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterRemoteWorker.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterRemoteWorker.java new file mode 100644 index 000000000..f8d14dd4e --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/register/worker/RegisterRemoteWorker.java @@ -0,0 +1,52 @@ +/* + * 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.register.worker; + +import org.apache.skywalking.oap.server.core.CoreModule; +import org.apache.skywalking.oap.server.core.register.RegisterSource; +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 class RegisterRemoteWorker extends AbstractWorker { + + private static final Logger logger = LoggerFactory.getLogger(RegisterRemoteWorker.class); + + private final AbstractWorker nextWorker; + private final RemoteSenderService remoteSender; + + RegisterRemoteWorker(int workerId, ModuleManager moduleManager, AbstractWorker nextWorker) { + super(workerId); + this.remoteSender = moduleManager.find(CoreModule.NAME).getService(RemoteSenderService.class); + this.nextWorker = nextWorker; + } + + @Override public final void in(RegisterSource indicator) { + try { + remoteSender.send(nextWorker.getWorkerId(), indicator, Selector.ForeverFirst); + } catch (Throwable e) { + logger.error(e.getMessage(), e); + } + } +} 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 9b94be896..16bcccf54 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 @@ -24,7 +24,7 @@ import org.apache.skywalking.oap.server.core.CoreModule; 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.core.worker.WorkerInstances; import org.apache.skywalking.oap.server.library.module.ModuleManager; import org.apache.skywalking.oap.server.library.server.grpc.GRPCHandler; import org.slf4j.*; @@ -38,7 +38,6 @@ public class RemoteServiceHandler extends RemoteServiceGrpc.RemoteServiceImplBas private final ModuleManager moduleManager; private StreamDataClassGetter streamDataClassGetter; - private WorkerClassGetter workerClassGetter; public RemoteServiceHandler(ModuleManager moduleManager) { this.moduleManager = moduleManager; @@ -48,9 +47,6 @@ public class RemoteServiceHandler extends RemoteServiceGrpc.RemoteServiceImplBas 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) { @@ -62,7 +58,7 @@ public class RemoteServiceHandler extends RemoteServiceGrpc.RemoteServiceImplBas try { StreamData streamData = streamDataClass.newInstance(); streamData.deserialize(remoteData); - workerClassGetter.getClassById(nextWorkerId).newInstance().in(streamData); + WorkerInstances.INSTANCES.get(nextWorkerId).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/client/RemoteClientManager.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/remote/client/RemoteClientManager.java index 361c98257..9de03ca09 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,8 +20,8 @@ package org.apache.skywalking.oap.server.core.remote.client; import java.util.*; import java.util.concurrent.*; -import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer; import org.apache.skywalking.oap.server.core.cluster.*; +import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer; import org.apache.skywalking.oap.server.library.module.*; import org.slf4j.*; @@ -96,7 +96,7 @@ public class RemoteClientManager implements Service { client = currentClientsMap.get(address); } else { if (remoteInstance.isSelf()) { - client = new SelfRemoteClient(moduleManager, remoteInstance.getHost(), remoteInstance.getPort()); + client = new SelfRemoteClient(remoteInstance.getHost(), remoteInstance.getPort()); } else { client = new GRPCRemoteClient(indicatorMapper, remoteInstance, 1, 3000); } 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 9e0ba7320..9b4f6bd1b 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 @@ -18,22 +18,18 @@ package org.apache.skywalking.oap.server.core.remote.client; -import org.apache.skywalking.oap.server.core.CoreModule; 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; +import org.apache.skywalking.oap.server.core.worker.WorkerInstances; /** * @author peng-yongsheng */ public class SelfRemoteClient implements RemoteClient { - private final ModuleManager moduleManager; private final String host; private final int port; - public SelfRemoteClient(ModuleManager moduleManager, String host, int port) { - this.moduleManager = moduleManager; + public SelfRemoteClient(String host, int port) { this.host = host; this.port = port; } @@ -47,7 +43,6 @@ public class SelfRemoteClient implements RemoteClient { } @Override public void push(int nextWorkerId, StreamData streamData) { - WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class); - workerMapper.findInstanceById(nextWorkerId).in(streamData); + WorkerInstances.INSTANCES.get(nextWorkerId).in(streamData); } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java index e744a480a..3b3028564 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/source/SourceReceiverImpl.java @@ -19,7 +19,6 @@ package org.apache.skywalking.oap.server.core.source; import org.apache.skywalking.oap.server.core.analysis.DispatcherManager; -import org.apache.skywalking.oap.server.library.module.ModuleManager; /** * @author peng-yongsheng @@ -28,8 +27,8 @@ public class SourceReceiverImpl implements SourceReceiver { private final DispatcherManager dispatcherManager; - public SourceReceiverImpl(ModuleManager moduleManager) { - this.dispatcherManager = new DispatcherManager(moduleManager); + public SourceReceiverImpl() { + this.dispatcherManager = new DispatcherManager(); } @Override public void receive(Source source) { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IPersistenceDAO.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IIndicatorDAO.java similarity index 72% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IPersistenceDAO.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IIndicatorDAO.java index 2a9d1f9a4..3dbf34806 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IPersistenceDAO.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IIndicatorDAO.java @@ -24,13 +24,13 @@ import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; /** * @author peng-yongsheng */ -public interface IPersistenceDAO extends DAO { +public interface IIndicatorDAO extends DAO { - INPUT get(INPUT input) throws IOException; + Indicator get(String modelName, Indicator indicator) throws IOException; - INSERT prepareBatchInsert(INPUT input) throws IOException; + INSERT prepareBatchInsert(String modelName, Indicator indicator) throws IOException; - UPDATE prepareBatchUpdate(INPUT input) throws IOException; + UPDATE prepareBatchUpdate(String modelName, Indicator indicator) throws IOException; - void deleteHistory(Long timeBucketBefore); + void deleteHistory(String modelName, Long timeBucketBefore); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterDAO.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterDAO.java new file mode 100644 index 000000000..18350a990 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/IRegisterDAO.java @@ -0,0 +1,36 @@ +/* + * 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.io.IOException; +import org.apache.skywalking.oap.server.core.register.RegisterSource; + +/** + * @author peng-yongsheng + */ +public interface IRegisterDAO extends DAO { + + int max(String modelName) throws IOException; + + RegisterSource get(String modelName, String id) throws IOException; + + void forceInsert(String modelName, RegisterSource source) throws IOException; + + void forceUpdate(String modelName, RegisterSource source) throws IOException; +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageBuilder.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageBuilder.java new file mode 100644 index 000000000..faf41dd67 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageBuilder.java @@ -0,0 +1,31 @@ +/* + * 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.Map; + +/** + * @author peng-yongsheng + */ +public interface StorageBuilder { + + T map2Data(Map dbMap); + + Map data2Map(T storageData); +} 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/storage/StorageDAO.java similarity index 78% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerClassGetter.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageDAO.java index 83b5b8d4a..b6cf4d643 100644 --- 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/storage/StorageDAO.java @@ -16,14 +16,15 @@ * */ -package org.apache.skywalking.oap.server.core.worker.annotation; +package org.apache.skywalking.oap.server.core.storage; -import org.apache.skywalking.oap.server.core.worker.AbstractWorker; +import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; import org.apache.skywalking.oap.server.library.module.Service; /** * @author peng-yongsheng */ -public interface WorkerClassGetter extends Service { - Class getClassById(int workerId); +public interface StorageDAO extends Service { + + IIndicatorDAO newIndicatorDao(StorageBuilder storageBuilder); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageData.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageData.java new file mode 100644 index 000000000..d7de8b828 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageData.java @@ -0,0 +1,26 @@ +/* + * 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; + +/** + * @author peng-yongsheng + */ +public interface StorageData { + String id(); +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageModule.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageModule.java index 62eb72d63..2e4beeb38 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageModule.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/StorageModule.java @@ -32,6 +32,6 @@ public class StorageModule extends ModuleDefine { } @Override public Class[] services() { - return new Class[] {IBatchDAO.class, IPersistenceDAO.class, IRegisterLockDAO.class}; + return new Class[] {IBatchDAO.class, StorageDAO.class, IRegisterLockDAO.class}; } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Column.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Column.java index aa6828b6e..adbf45e20 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Column.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Column.java @@ -23,8 +23,12 @@ import java.lang.annotation.*; /** * @author peng-yongsheng */ -@Target(ElementType.FIELD) +@Target({ElementType.FIELD}) @Retention(RetentionPolicy.RUNTIME) public @interface Column { String columnName(); + + boolean matchQuery() default false; + + boolean termQuery() default true; } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Query.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Query.java new file mode 100644 index 000000000..5f5de015c --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/Query.java @@ -0,0 +1,26 @@ +/* + * 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; + +/** + * @author peng-yongsheng + */ +public enum Query { + Term, Match +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageAnnotationListener.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageAnnotationListener.java index 894326cab..fe134aab3 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageAnnotationListener.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageAnnotationListener.java @@ -49,8 +49,8 @@ public class StorageAnnotationListener implements AnnotationListener, IModelGett List modelColumns = new LinkedList<>(); retrieval(aClass, modelColumns); - StorageEntity annotation = (StorageEntity)aClass.getAnnotation(StorageEntity.class); - models.add(new Model(annotation.name(), modelColumns)); + String modelName = StorageEntityAnnotationUtils.getModelName(aClass); + models.add(new Model(modelName, modelColumns)); } private void retrieval(Class clazz, List 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 index cfc7ce98f..7bf8cab97 100644 --- 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 @@ -19,6 +19,7 @@ package org.apache.skywalking.oap.server.core.storage.annotation; import java.lang.annotation.*; +import org.apache.skywalking.oap.server.core.storage.StorageBuilder; /** * @author peng-yongsheng @@ -27,4 +28,6 @@ import java.lang.annotation.*; @Retention(RetentionPolicy.RUNTIME) public @interface StorageEntity { String name(); + + Class builder(); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageEntityAnnotationUtils.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageEntityAnnotationUtils.java new file mode 100644 index 000000000..b754039de --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/annotation/StorageEntityAnnotationUtils.java @@ -0,0 +1,46 @@ +/* + * 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 org.apache.skywalking.oap.server.core.UnexpectedException; +import org.apache.skywalking.oap.server.core.storage.StorageBuilder; + +/** + * @author peng-yongsheng + */ +public class StorageEntityAnnotationUtils { + + public static String getModelName(Class aClass) { + if (aClass.isAnnotationPresent(StorageEntity.class)) { + StorageEntity annotation = (StorageEntity)aClass.getAnnotation(StorageEntity.class); + return annotation.name(); + } else { + throw new UnexpectedException(""); + } + } + + public static Class getBuilder(Class aClass) { + if (aClass.isAnnotationPresent(StorageEntity.class)) { + StorageEntity annotation = (StorageEntity)aClass.getAnnotation(StorageEntity.class); + return annotation.builder(); + } else { + throw new UnexpectedException(""); + } + } +} diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/cache/IEndpointCacheDAO.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/cache/IEndpointCacheDAO.java new file mode 100644 index 000000000..e010e7712 --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/storage/cache/IEndpointCacheDAO.java @@ -0,0 +1,32 @@ +/* + * 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.cache; + +import org.apache.skywalking.oap.server.core.register.endpoint.Endpoint; +import org.apache.skywalking.oap.server.core.storage.DAO; + +/** + * @author peng-yongsheng + */ +public interface IEndpointCacheDAO extends DAO { + + int get(String id); + + Endpoint get(int sequence); +} 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 index fd73d4c2a..c079a1111 100644 --- 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 @@ -18,10 +18,18 @@ package org.apache.skywalking.oap.server.core.worker; +import lombok.Getter; + /** * @author peng-yongsheng */ public abstract class AbstractWorker { + @Getter private final int workerId; + + public AbstractWorker(int workerId) { + this.workerId = workerId; + } + public abstract void in(INPUT input); } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerDefineLoadException.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/WorkerIdGenerator.java similarity index 78% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerDefineLoadException.java rename to oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/WorkerIdGenerator.java index 1b24afa55..fc7f5af92 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerDefineLoadException.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/WorkerIdGenerator.java @@ -16,14 +16,17 @@ * */ -package org.apache.skywalking.oap.server.core.worker.annotation; +package org.apache.skywalking.oap.server.core.worker; /** * @author peng-yongsheng */ -public class WorkerDefineLoadException extends RuntimeException { +public enum WorkerIdGenerator { + INSTANCES; - public WorkerDefineLoadException(String message, Throwable cause) { - super(message, cause); + private int workerId = 0; + + public synchronized int generate() { + return workerId++; } } diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/WorkerInstances.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/WorkerInstances.java new file mode 100644 index 000000000..a80dbd54d --- /dev/null +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/WorkerInstances.java @@ -0,0 +1,38 @@ +/* + * 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; + +import java.util.*; + +/** + * @author peng-yongsheng + */ +public enum WorkerInstances { + INSTANCES; + + private Map instances = new HashMap<>(); + + public void put(int workerId, AbstractWorker instance) { + instances.put(workerId, instance); + } + + public AbstractWorker get(int workerId) { + return instances.get(workerId); + } +} 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 deleted file mode 100644 index c06cabcf5..000000000 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/worker/annotation/WorkerAnnotationContainer.java +++ /dev/null @@ -1,84 +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.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/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 fae88930f..757215de2 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 @@ -18,9 +18,8 @@ package org.apache.skywalking.oap.server.core.analysis.indicator.define; -import java.util.Map; import lombok.*; -import org.apache.skywalking.oap.server.core.analysis.indicator.*; +import org.apache.skywalking.oap.server.core.analysis.indicator.AvgIndicator; import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData; /** @@ -34,22 +33,10 @@ public class TestAvgIndicator extends AvgIndicator { return null; } - @Override public String name() { - return null; - } - @Override public void deserialize(RemoteData remoteData) { } @Override public String id() { return null; } - - @Override public Map toMap() { - return null; - } - - @Override public Indicator newOne(Map dbMap) { - return null; - } } diff --git a/oap-server/server-core/src/test/resources/META-INF/defines/indicator.def b/oap-server/server-core/src/test/resources/META-INF/defines/indicator.def deleted file mode 100644 index 97491fffb..000000000 --- a/oap-server/server-core/src/test/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.indicator.define.TestAvgIndicator \ No newline at end of file diff --git a/oap-server/server-core/src/test/resources/META-INF/defines/worker.def b/oap-server/server-core/src/test/resources/META-INF/defines/worker.def deleted file mode 100644 index 33ebbb1f3..000000000 --- a/oap-server/server-core/src/test/resources/META-INF/defines/worker.def +++ /dev/null @@ -1,17 +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. -# -# \ No newline at end of file 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 db1db6374..e11c4a6f3 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 @@ -63,7 +63,7 @@ public class StorageModuleElasticsearchProvider extends ModuleProvider { elasticSearchClient = new ElasticSearchClient(config.getClusterNodes(), nameSpace); 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(StorageDAO.class, new StorageEsDAO(elasticSearchClient)); this.registerServiceImplementation(IRegisterLockDAO.class, new RegisterLockDAOImpl(elasticSearchClient, 1000)); } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/EsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/EsDAO.java index dd0a70bd8..e7c813873 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/EsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/EsDAO.java @@ -18,54 +18,15 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base; -import java.io.IOException; import org.apache.skywalking.oap.server.core.storage.AbstractDAO; import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; -import org.elasticsearch.action.search.SearchResponse; -import org.elasticsearch.search.aggregations.AggregationBuilders; -import org.elasticsearch.search.aggregations.metrics.max.Max; -import org.elasticsearch.search.builder.SearchSourceBuilder; -import org.slf4j.*; /** * @author peng-yongsheng */ public abstract class EsDAO extends AbstractDAO { - private static final Logger logger = LoggerFactory.getLogger(EsDAO.class); - public EsDAO(ElasticSearchClient client) { super(client); } - - protected final int getMaxId(String indexName, String columnName) { - SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder(); - searchSourceBuilder.aggregation(AggregationBuilders.max("agg").field(columnName)); - searchSourceBuilder.size(0); - return getResponse(indexName, searchSourceBuilder); - } - - protected final int getMinId(String indexName, String columnName) { - SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder(); - searchSourceBuilder.aggregation(AggregationBuilders.min("agg").field(columnName)); - searchSourceBuilder.size(0); - return getResponse(indexName, searchSourceBuilder); - } - - private int getResponse(String indexName, SearchSourceBuilder searchSourceBuilder) { - try { - SearchResponse searchResponse = getClient().search(indexName, searchSourceBuilder); - Max agg = searchResponse.getAggregations().get("agg"); - - int id = (int)agg.getValue(); - if (id == Integer.MAX_VALUE || id == Integer.MIN_VALUE) { - return 0; - } else { - return id; - } - } catch (IOException e) { - logger.error(e.getMessage(), e); - } - return 0; - } } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/PersistenceEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/IndicatorEsDAO.java similarity index 60% rename from oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/PersistenceEsDAO.java rename to oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/IndicatorEsDAO.java index 1f83b633a..37b4d4791 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/PersistenceEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/IndicatorEsDAO.java @@ -21,8 +21,7 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base; import java.io.IOException; import java.util.Map; import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; -import org.apache.skywalking.oap.server.core.storage.IPersistenceDAO; -import org.apache.skywalking.oap.server.library.client.NameSpace; +import org.apache.skywalking.oap.server.core.storage.*; import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; import org.elasticsearch.action.get.GetResponse; import org.elasticsearch.action.index.IndexRequest; @@ -32,48 +31,46 @@ import org.elasticsearch.common.xcontent.*; /** * @author peng-yongsheng */ -public class PersistenceEsDAO implements IPersistenceDAO { +public class IndicatorEsDAO extends EsDAO implements IIndicatorDAO { - private final ElasticSearchClient client; - private final NameSpace nameSpace; + private final StorageBuilder storageBuilder; - public PersistenceEsDAO(ElasticSearchClient client, NameSpace nameSpace) { - this.client = client; - this.nameSpace = nameSpace; + public IndicatorEsDAO(ElasticSearchClient client, StorageBuilder storageBuilder) { + super(client); + this.storageBuilder = storageBuilder; } - @Override public Indicator get(Indicator input) throws IOException { - GetResponse response = client.get(nameSpace.getNameSpace() + "_" + input.name(), input.id()); + @Override public Indicator get(String modelName, Indicator indicator) throws IOException { + GetResponse response = getClient().get(modelName, indicator.id()); if (response.isExists()) { - return input.newOne(response.getSource()); + return storageBuilder.map2Data(response.getSource()); } else { return null; } } - @Override public IndexRequest prepareBatchInsert(Indicator input) throws IOException { - Map objectMap = input.toMap(); + @Override public IndexRequest prepareBatchInsert(String modelName, Indicator indicator) throws IOException { + Map objectMap = storageBuilder.data2Map(indicator); XContentBuilder builder = XContentFactory.jsonBuilder().startObject(); for (String key : objectMap.keySet()) { builder.field(key, objectMap.get(key)); } builder.endObject(); - return client.prepareInsert(nameSpace.getNameSpace() + "_" + input.name(), input.id(), builder); + return getClient().prepareInsert(modelName, indicator.id(), builder); } - @Override public UpdateRequest prepareBatchUpdate(Indicator input) throws IOException { - Map objectMap = input.toMap(); + @Override public UpdateRequest prepareBatchUpdate(String modelName, Indicator indicator) throws IOException { + Map objectMap = storageBuilder.data2Map(indicator); XContentBuilder builder = XContentFactory.jsonBuilder().startObject(); for (String key : objectMap.keySet()) { builder.field(key, objectMap.get(key)); } builder.endObject(); - return client.prepareUpdate(nameSpace.getNameSpace() + "_" + input.name(), input.id(), builder); + return getClient().prepareUpdate(modelName, indicator.id(), builder); } - @Override public void deleteHistory(Long timeBucketBefore) { - + @Override public void deleteHistory(String modelName, Long timeBucketBefore) { } } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/RegisterEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/RegisterEsDAO.java new file mode 100644 index 000000000..13c6711dc --- /dev/null +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/RegisterEsDAO.java @@ -0,0 +1,99 @@ +/* + * 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.storage.plugin.elasticsearch.base; + +import java.io.IOException; +import java.util.Map; +import org.apache.skywalking.oap.server.core.register.RegisterSource; +import org.apache.skywalking.oap.server.core.storage.*; +import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; +import org.elasticsearch.action.get.GetResponse; +import org.elasticsearch.action.search.SearchResponse; +import org.elasticsearch.common.xcontent.*; +import org.elasticsearch.search.aggregations.AggregationBuilders; +import org.elasticsearch.search.aggregations.metrics.max.Max; +import org.elasticsearch.search.builder.SearchSourceBuilder; +import org.slf4j.*; + +/** + * @author peng-yongsheng + */ +public class RegisterEsDAO extends EsDAO implements IRegisterDAO { + + private static final Logger logger = LoggerFactory.getLogger(RegisterEsDAO.class); + + private final StorageBuilder storageBuilder; + + public RegisterEsDAO(ElasticSearchClient client, StorageBuilder storageBuilder) { + super(client); + this.storageBuilder = storageBuilder; + } + + @Override public RegisterSource get(String modelName, String id) throws IOException { + GetResponse response = getClient().get(modelName, id); + if (response.isExists()) { + return storageBuilder.map2Data(response.getSource()); + } else { + return null; + } + } + + @Override public void forceInsert(String modelName, RegisterSource source) throws IOException { + Map objectMap = storageBuilder.data2Map(source); + + XContentBuilder builder = XContentFactory.jsonBuilder().startObject(); + for (String key : objectMap.keySet()) { + builder.field(key, objectMap.get(key)); + } + builder.endObject(); + + getClient().forceInsert(modelName, source.id(), builder); + } + + @Override public void forceUpdate(String modelName, RegisterSource source) throws IOException { + Map objectMap = storageBuilder.data2Map(source); + + XContentBuilder builder = XContentFactory.jsonBuilder().startObject(); + for (String key : objectMap.keySet()) { + builder.field(key, objectMap.get(key)); + } + builder.endObject(); + + getClient().forceUpdate(modelName, source.id(), builder); + } + + @Override public int max(String modelName) throws IOException { + SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder(); + searchSourceBuilder.aggregation(AggregationBuilders.max(RegisterSource.SEQUENCE).field(RegisterSource.SEQUENCE)); + searchSourceBuilder.size(0); + return getResponse(modelName, searchSourceBuilder); + } + + private int getResponse(String modelName, SearchSourceBuilder searchSourceBuilder) throws IOException { + SearchResponse searchResponse = getClient().search(modelName, searchSourceBuilder); + Max agg = searchResponse.getAggregations().get(RegisterSource.SEQUENCE); + + int id = (int)agg.getValue(); + if (id == Integer.MAX_VALUE || id == Integer.MIN_VALUE) { + return 0; + } else { + return 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-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsDAO.java similarity index 58% rename from oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgPersistentWorker.java rename to oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsDAO.java index 99a823986..641582f7e 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/endpoint/EndpointLatencyAvgPersistentWorker.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/base/StorageEsDAO.java @@ -16,23 +16,22 @@ * */ -package org.apache.skywalking.oap.server.core.analysis.endpoint; +package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base; -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; +import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator; +import org.apache.skywalking.oap.server.core.storage.*; +import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; /** * @author peng-yongsheng */ -@Worker -public class EndpointLatencyAvgPersistentWorker extends AbstractPersistentWorker { +public class StorageEsDAO extends EsDAO implements StorageDAO { - public EndpointLatencyAvgPersistentWorker(ModuleManager moduleManager) { - super(moduleManager); + public StorageEsDAO(ElasticSearchClient client) { + super(client); } - @Override protected boolean needMergeDBData() { - return true; + @Override public IIndicatorDAO newIndicatorDao(StorageBuilder storageBuilder) { + return new IndicatorEsDAO(getClient(), storageBuilder); } } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/EndpointCacheEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/EndpointCacheEsDAO.java new file mode 100644 index 000000000..8645bc8fa --- /dev/null +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/EndpointCacheEsDAO.java @@ -0,0 +1,53 @@ +/* + * 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.storage.plugin.elasticsearch.cache; + +import org.apache.skywalking.oap.server.core.register.RegisterSource; +import org.apache.skywalking.oap.server.core.register.endpoint.Endpoint; +import org.apache.skywalking.oap.server.core.storage.cache.IEndpointCacheDAO; +import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient; +import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base.EsDAO; +import org.elasticsearch.action.get.GetResponse; + +/** + * @author peng-yongsheng + */ +public class EndpointCacheEsDAO extends EsDAO implements IEndpointCacheDAO { + + public EndpointCacheEsDAO(ElasticSearchClient client) { + super(client); + } + + @Override public int get(String id) { + try { + GetResponse response = getClient().get("", id); + if (response.isExists()) { + return response.getField(RegisterSource.SEQUENCE).getValue(); + } else { + return 0; + } + } catch (Throwable e) { + return 0; + } + } + + @Override public Endpoint get(int sequence) { + return null; + } +}