Feature/oap/annotation (#1531)

* Use annotation to instead of definition file.
This commit is contained in:
彭勇升 pengys 2018-08-08 20:09:02 +08:00 committed by 吴晟 Wu Sheng
parent 9ac0f07793
commit 428b2cd78a
60 changed files with 740 additions and 493 deletions

View File

@ -19,13 +19,9 @@
package org.apache.skywalking.apm.collector.analysis.segment.parser.provider.service;
import com.google.protobuf.InvalidProtocolBufferException;
import java.util.Base64;
import java.util.List;
import org.apache.skywalking.apm.network.language.agent.SpanObject;
import org.apache.skywalking.apm.network.language.agent.TraceSegmentObject;
import org.apache.skywalking.apm.network.language.agent.UniqueId;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.*;
import org.apache.skywalking.apm.network.language.agent.*;
import org.slf4j.*;
/**
* @author peng-yongsheng
@ -35,8 +31,11 @@ public class SegmentBase64Printer {
private static final Logger LOGGER = LoggerFactory.getLogger(SegmentBase64Printer.class);
public static void main(String[] args) throws InvalidProtocolBufferException {
String segmentBase64 = "CgoKCJbf2NPCLBAQEiAQ////////////ARiV39jTwiwg2+7Y08IsMNQPWANgARIlCAEYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAIQARif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIAxACGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgEEAMYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAUQBBif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIBhAFGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgHEAYYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAgQBxif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicICRAIGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgKEAkYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCAsQChif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEicIDBALGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAMSJwgNEAwYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgAxInCA4QDRif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADEvMCCA8QDhif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADggHIAhKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXIS8wIIEBAPGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAOCAcgCEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchLzAggREBAYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgA4IByAISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEvMCCBIQERif39jTwiwguezY08IsMJTIAkD///////////8BUAFYAmADggHIAhKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXIS8wIIExASGJ/f2NPCLCC57NjTwiwwlMgCQP///////////wFQAVgCYAOCAcgCEqEBCg1lcnJvciBtZXNzYWdlEo8BW0lORk9dIEJ1aWxkaW5nIGphcjogL1VzZXJzL3Blbmd5czUvY29kZS9za3ktd2Fsa2luZy9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC90YXJnZXQvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QtMS4wLWphci13aXRoLWRlcGVuZGVuY2llcy5qYXISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchLzAggUEBMYn9/Y08IsILns2NPCLDCUyAJA////////////AVABWAJgA4IByAISoQEKDWVycm9yIG1lc3NhZ2USjwFbSU5GT10gQnVpbGRpbmcgamFyOiAvVXNlcnMvcGVuZ3lzNS9jb2RlL3NreS13YWxraW5nL2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0L3RhcmdldC9jb2xsZWN0b3ItcGVyZm9ybWFuY2UtdGVzdC0xLjAtamFyLXdpdGgtZGVwZW5kZW5jaWVzLmphchKhAQoNZXJyb3IgbWVzc2FnZRKPAVtJTkZPXSBCdWlsZGluZyBqYXI6IC9Vc2Vycy9wZW5neXM1L2NvZGUvc2t5LXdhbGtpbmcvY29sbGVjdG9yLXBlcmZvcm1hbmNlLXRlc3QvdGFyZ2V0L2NvbGxlY3Rvci1wZXJmb3JtYW5jZS10ZXN0LTEuMC1qYXItd2l0aC1kZXBlbmRlbmNpZXMuamFyGP7//////////wEgCA==";
byte[] binarySegment = Base64.getDecoder().decode(segmentBase64);
if (args.length == 0) {
return;
}
byte[] binarySegment = Base64.getDecoder().decode(args[0]);
TraceSegmentObject segmentObject = TraceSegmentObject.parseFrom(binarySegment);
UniqueId segmentId = segmentObject.getTraceSegmentId();

View File

@ -19,12 +19,13 @@
package org.apache.skywalking.oap.server.core;
import java.util.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper;
import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper;
import org.apache.skywalking.oap.server.core.source.SourceReceiver;
import org.apache.skywalking.oap.server.core.remote.RemoteSenderService;
import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataClassGetter;
import org.apache.skywalking.oap.server.core.remote.client.RemoteClientManager;
import org.apache.skywalking.oap.server.core.server.*;
import org.apache.skywalking.oap.server.core.source.SourceReceiver;
import org.apache.skywalking.oap.server.core.storage.model.IModelGetter;
import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer;
import org.apache.skywalking.oap.server.library.module.ModuleDefine;
/**
@ -53,8 +54,9 @@ public class CoreModule extends ModuleDefine {
}
private void addInsideService(List<Class> classes) {
classes.add(IndicatorMapper.class);
classes.add(WorkerMapper.class);
classes.add(IModelGetter.class);
classes.add(StreamDataClassGetter.class);
classes.add(WorkerAnnotationContainer.class);
classes.add(RemoteClientManager.class);
classes.add(RemoteSenderService.class);
}

View File

@ -18,13 +18,17 @@
package org.apache.skywalking.oap.server.core;
import org.apache.skywalking.oap.server.core.analysis.indicator.define.*;
import org.apache.skywalking.oap.server.core.analysis.worker.define.*;
import java.io.IOException;
import org.apache.skywalking.oap.server.core.annotation.AnnotationScan;
import org.apache.skywalking.oap.server.core.cluster.*;
import org.apache.skywalking.oap.server.core.remote.*;
import org.apache.skywalking.oap.server.core.remote.annotation.*;
import org.apache.skywalking.oap.server.core.remote.client.RemoteClientManager;
import org.apache.skywalking.oap.server.core.server.*;
import org.apache.skywalking.oap.server.core.source.*;
import org.apache.skywalking.oap.server.core.storage.annotation.StorageAnnotationListener;
import org.apache.skywalking.oap.server.core.storage.model.IModelGetter;
import org.apache.skywalking.oap.server.core.worker.annotation.*;
import org.apache.skywalking.oap.server.library.module.*;
import org.apache.skywalking.oap.server.library.server.ServerException;
import org.apache.skywalking.oap.server.library.server.grpc.GRPCServer;
@ -41,14 +45,22 @@ public class CoreModuleProvider extends ModuleProvider {
private final CoreModuleConfig moduleConfig;
private GRPCServer grpcServer;
private JettyServer jettyServer;
private final IndicatorMapper indicatorMapper;
private final WorkerMapper workerMapper;
private final AnnotationScan annotationScan;
private final StorageAnnotationListener storageAnnotationListener;
private final StreamAnnotationListener streamAnnotationListener;
private final WorkerAnnotationListener workerAnnotationListener;
private final StreamDataAnnotationContainer streamDataAnnotationContainer;
private final WorkerAnnotationContainer workerAnnotationContainer;
public CoreModuleProvider() {
super();
this.moduleConfig = new CoreModuleConfig();
this.indicatorMapper = new IndicatorMapper();
this.workerMapper = new WorkerMapper();
this.annotationScan = new AnnotationScan();
this.storageAnnotationListener = new StorageAnnotationListener();
this.streamAnnotationListener = new StreamAnnotationListener();
this.workerAnnotationListener = new WorkerAnnotationListener();
this.streamDataAnnotationContainer = new StreamDataAnnotationContainer();
this.workerAnnotationContainer = new WorkerAnnotationContainer();
}
@Override public String name() {
@ -75,20 +87,27 @@ public class CoreModuleProvider extends ModuleProvider {
this.registerServiceImplementation(SourceReceiver.class, new SourceReceiverImpl(getManager()));
this.registerServiceImplementation(IndicatorMapper.class, indicatorMapper);
this.registerServiceImplementation(WorkerMapper.class, workerMapper);
this.registerServiceImplementation(StreamDataClassGetter.class, streamDataAnnotationContainer);
this.registerServiceImplementation(WorkerAnnotationContainer.class, workerAnnotationContainer);
this.registerServiceImplementation(RemoteClientManager.class, new RemoteClientManager(getManager()));
this.registerServiceImplementation(RemoteSenderService.class, new RemoteSenderService(getManager()));
this.registerServiceImplementation(IModelGetter.class, storageAnnotationListener);
annotationScan.registerListener(storageAnnotationListener);
annotationScan.registerListener(streamAnnotationListener);
annotationScan.registerListener(workerAnnotationListener);
}
@Override public void start() throws ModuleStartException {
grpcServer.addHandler(new RemoteServiceHandler(getManager()));
try {
indicatorMapper.load();
workerMapper.load(getManager());
} catch (IndicatorDefineLoadException | WorkerDefineLoadException e) {
annotationScan.scan(() -> {
streamDataAnnotationContainer.generate(streamAnnotationListener.getStreamClasses());
workerAnnotationContainer.load(getManager(), workerAnnotationListener.getWorkerClasses());
});
} catch (WorkerDefineLoadException | IOException e) {
throw new ModuleStartException(e.getMessage(), e);
}
}

View File

@ -19,6 +19,7 @@
package org.apache.skywalking.oap.server.core.analysis.data;
import java.util.*;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
/**
* @author peng-yongsheng

View File

@ -18,6 +18,8 @@
package org.apache.skywalking.oap.server.core.analysis.data;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
/**
* @author peng-yongsheng
*/

View File

@ -19,6 +19,7 @@
package org.apache.skywalking.oap.server.core.analysis.data;
import java.util.*;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
/**
* @author peng-yongsheng

View File

@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.SourceDispatcher;
import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper;
import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer;
import org.apache.skywalking.oap.server.core.source.Endpoint;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
@ -42,7 +42,7 @@ public class EndpointDispatcher implements SourceDispatcher<Endpoint> {
private void avg(Endpoint source) {
if (avgAggregator == null) {
WorkerMapper workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class);
WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class);
avgAggregator = (EndpointLatencyAvgAggregateWorker)workerMapper.findInstanceByClass(EndpointLatencyAvgAggregateWorker.class);
}

View File

@ -19,11 +19,13 @@
package org.apache.skywalking.oap.server.core.analysis.endpoint;
import org.apache.skywalking.oap.server.core.analysis.worker.AbstractAggregatorWorker;
import org.apache.skywalking.oap.server.core.worker.annotation.Worker;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
/**
* @author peng-yongsheng
*/
@Worker
public class EndpointLatencyAvgAggregateWorker extends AbstractAggregatorWorker<EndpointLatencyAvgIndicator> {
public EndpointLatencyAvgAggregateWorker(ModuleManager moduleManager) {

View File

@ -21,15 +21,19 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint;
import java.util.*;
import lombok.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.*;
import org.apache.skywalking.oap.server.core.remote.annotation.StreamData;
import org.apache.skywalking.oap.server.core.remote.grpc.proto.RemoteData;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
import org.apache.skywalking.oap.server.core.storage.annotation.*;
/**
* @author peng-yongsheng
*/
@StreamData
@StorageEntity(name = EndpointLatencyAvgIndicator.NAME)
public class EndpointLatencyAvgIndicator extends AvgIndicator {
private static final String NAME = "endpoint_latency_avg";
public static final String NAME = "endpoint_latency_avg";
private static final String ID = "id";
private static final String SERVICE_ID = "service_id";
private static final String SERVICE_INSTANCE_ID = "service_instance_id";

View File

@ -19,11 +19,13 @@
package org.apache.skywalking.oap.server.core.analysis.endpoint;
import org.apache.skywalking.oap.server.core.analysis.worker.AbstractPersistentWorker;
import org.apache.skywalking.oap.server.core.worker.annotation.Worker;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
/**
* @author peng-yongsheng
*/
@Worker
public class EndpointLatencyAvgPersistentWorker extends AbstractPersistentWorker<EndpointLatencyAvgIndicator> {
public EndpointLatencyAvgPersistentWorker(ModuleManager moduleManager) {

View File

@ -20,11 +20,13 @@ package org.apache.skywalking.oap.server.core.analysis.endpoint;
import org.apache.skywalking.oap.server.core.analysis.worker.AbstractRemoteWorker;
import org.apache.skywalking.oap.server.core.remote.selector.Selector;
import org.apache.skywalking.oap.server.core.worker.annotation.Worker;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
/**
* @author peng-yongsheng
*/
@Worker
public class EndpointLatencyAvgRemoteWorker extends AbstractRemoteWorker<EndpointLatencyAvgIndicator> {
public EndpointLatencyAvgRemoteWorker(ModuleManager moduleManager) {

View File

@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.core.analysis.indicator;
import java.util.Map;
import lombok.*;
import org.apache.skywalking.oap.server.core.analysis.data.StreamData;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.storage.annotation.Column;
/**

View File

@ -1,86 +0,0 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.analysis.indicator.define;
import java.io.*;
import java.net.URL;
import java.util.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.library.module.Service;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class IndicatorMapper implements Service {
private static final Logger logger = LoggerFactory.getLogger(IndicatorMapper.class);
private int id = 0;
private final Map<Class<Indicator>, Integer> classKeyMapping;
private final Map<Integer, Class<Indicator>> idKeyMapping;
public IndicatorMapper() {
this.classKeyMapping = new HashMap<>();
this.idKeyMapping = new HashMap<>();
}
@SuppressWarnings(value = "unchecked")
public void load() throws IndicatorDefineLoadException {
try {
List<String> indicatorClasses = new LinkedList<>();
Enumeration<URL> urlEnumeration = this.getClass().getClassLoader().getResources("META-INF/defines/indicator.def");
while (urlEnumeration.hasMoreElements()) {
URL definitionFileURL = urlEnumeration.nextElement();
logger.info("Load indicator definition file url: {}", definitionFileURL.getPath());
BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(definitionFileURL.openStream()));
Properties properties = new Properties();
properties.load(bufferedReader);
Enumeration defineItem = properties.propertyNames();
while (defineItem.hasMoreElements()) {
String fullNameClass = (String)defineItem.nextElement();
indicatorClasses.add(fullNameClass);
}
}
for (String indicatorClassName : indicatorClasses) {
Class<Indicator> indicatorClass = (Class<Indicator>)Class.forName(indicatorClassName);
id++;
classKeyMapping.put(indicatorClass, id);
idKeyMapping.put(id, indicatorClass);
}
} catch (IOException | ClassNotFoundException e) {
throw new IndicatorDefineLoadException(e.getMessage(), e);
}
}
public int findIdByClass(Class indicatorClass) {
return classKeyMapping.get(indicatorClass);
}
public Class<Indicator> findClassById(int id) {
return idKeyMapping.get(id);
}
public Collection<Class<Indicator>> indicatorClasses() {
return idKeyMapping.values();
}
}

View File

@ -24,18 +24,19 @@ import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.data.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper;
import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer;
import org.apache.skywalking.oap.server.core.worker.AbstractWorker;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public abstract class AbstractAggregatorWorker<INPUT extends Indicator> extends Worker<INPUT> {
public abstract class AbstractAggregatorWorker<INPUT extends Indicator> extends AbstractWorker<INPUT> {
private static final Logger logger = LoggerFactory.getLogger(AbstractAggregatorWorker.class);
private Worker worker;
private AbstractWorker worker;
private final ModuleManager moduleManager;
private final DataCarrier<INPUT> dataCarrier;
private final MergeDataCache<INPUT> mergeDataCache;
@ -85,7 +86,7 @@ public abstract class AbstractAggregatorWorker<INPUT extends Indicator> extends
private void onNext(INPUT data) {
if (worker == null) {
WorkerMapper workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class);
WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class);
worker = workerMapper.findInstanceByClass(nextWorkerClass());
}
worker.in(data);

View File

@ -22,6 +22,7 @@ import java.util.*;
import org.apache.skywalking.oap.server.core.analysis.data.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.storage.*;
import org.apache.skywalking.oap.server.core.worker.AbstractWorker;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.slf4j.*;
@ -30,7 +31,7 @@ import static java.util.Objects.nonNull;
/**
* @author peng-yongsheng
*/
public abstract class AbstractPersistentWorker<INPUT extends Indicator> extends Worker<INPUT> {
public abstract class AbstractPersistentWorker<INPUT extends Indicator> extends AbstractWorker<INPUT> {
private static final Logger logger = LoggerFactory.getLogger(AbstractPersistentWorker.class);

View File

@ -20,22 +20,23 @@ package org.apache.skywalking.oap.server.core.analysis.worker;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper;
import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer;
import org.apache.skywalking.oap.server.core.remote.RemoteSenderService;
import org.apache.skywalking.oap.server.core.remote.selector.Selector;
import org.apache.skywalking.oap.server.core.worker.AbstractWorker;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public abstract class AbstractRemoteWorker<INPUT extends Indicator> extends Worker<INPUT> {
public abstract class AbstractRemoteWorker<INPUT extends Indicator> extends AbstractWorker<INPUT> {
private static final Logger logger = LoggerFactory.getLogger(AbstractRemoteWorker.class);
private final ModuleManager moduleManager;
private RemoteSenderService remoteSender;
private WorkerMapper workerMapper;
private WorkerAnnotationContainer workerMapper;
public AbstractRemoteWorker(ModuleManager moduleManager) {
this.moduleManager = moduleManager;
@ -46,7 +47,7 @@ public abstract class AbstractRemoteWorker<INPUT extends Indicator> extends Work
remoteSender = moduleManager.find(CoreModule.NAME).getService(RemoteSenderService.class);
}
if (workerMapper == null) {
workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class);
workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class);
}
try {

View File

@ -1,100 +0,0 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.analysis.worker.define;
import java.io.*;
import java.lang.reflect.Constructor;
import java.net.URL;
import java.util.*;
import org.apache.skywalking.oap.server.core.analysis.worker.Worker;
import org.apache.skywalking.oap.server.library.module.*;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class WorkerMapper implements Service {
private static final Logger logger = LoggerFactory.getLogger(WorkerMapper.class);
private int id = 0;
private final Map<Class<Worker>, Integer> classKeyMapping;
private final Map<Integer, Class<Worker>> idKeyMapping;
private final Map<Class<Worker>, Worker> classKeyInstanceMapping;
private final Map<Integer, Worker> idKeyInstanceMapping;
public WorkerMapper() {
this.classKeyMapping = new HashMap<>();
this.idKeyMapping = new HashMap<>();
this.classKeyInstanceMapping = new HashMap<>();
this.idKeyInstanceMapping = new HashMap<>();
}
@SuppressWarnings(value = "unchecked")
public void load(ModuleManager moduleManager) throws WorkerDefineLoadException {
try {
List<String> workerClasses = new LinkedList<>();
Enumeration<URL> urlEnumeration = this.getClass().getClassLoader().getResources("META-INF/defines/worker.def");
while (urlEnumeration.hasMoreElements()) {
URL definitionFileURL = urlEnumeration.nextElement();
logger.info("Load worker definition file url: {}", definitionFileURL.getPath());
BufferedReader bufferedReader = new BufferedReader(new InputStreamReader(definitionFileURL.openStream()));
Properties properties = new Properties();
properties.load(bufferedReader);
Enumeration defineItem = properties.propertyNames();
while (defineItem.hasMoreElements()) {
String fullNameClass = (String)defineItem.nextElement();
workerClasses.add(fullNameClass);
}
}
for (String workerClassName : workerClasses) {
Class<Worker> workerClass = (Class<Worker>)Class.forName(workerClassName);
id++;
classKeyMapping.put(workerClass, id);
idKeyMapping.put(id, workerClass);
Constructor<Worker> constructor = workerClass.getDeclaredConstructor(ModuleManager.class);
Worker worker = constructor.newInstance(moduleManager);
classKeyInstanceMapping.put(workerClass, worker);
idKeyInstanceMapping.put(id, worker);
}
} catch (Exception e) {
throw new WorkerDefineLoadException(e.getMessage(), e);
}
}
public int findIdByClass(Class workerClass) {
return classKeyMapping.get(workerClass);
}
public Class<Worker> findClassById(int id) {
return idKeyMapping.get(id);
}
public Worker findInstanceByClass(Class workerClass) {
return classKeyInstanceMapping.get(workerClass);
}
public Worker findInstanceById(int id) {
return idKeyInstanceMapping.get(id);
}
}

View File

@ -16,14 +16,16 @@
*
*/
package org.apache.skywalking.oap.server.core.analysis.indicator.define;
package org.apache.skywalking.oap.server.core.annotation;
import java.lang.annotation.Annotation;
/**
* @author peng-yongsheng
*/
public class IndicatorDefineLoadException extends Exception {
public interface AnnotationListener {
public IndicatorDefineLoadException(String message, Throwable cause) {
super(message, cause);
}
Class<? extends Annotation> annotation();
void notify(Class aClass);
}

View File

@ -0,0 +1,56 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.annotation;
import com.google.common.collect.ImmutableSet;
import com.google.common.reflect.ClassPath;
import java.io.IOException;
import java.util.*;
/**
* @author peng-yongsheng
*/
public class AnnotationScan {
private final List<AnnotationListener> listeners;
public AnnotationScan() {
this.listeners = new LinkedList<>();
}
public void registerListener(AnnotationListener listener) {
listeners.add(listener);
}
public void scan(Runnable callBack) throws IOException {
ClassPath classpath = ClassPath.from(this.getClass().getClassLoader());
ImmutableSet<ClassPath.ClassInfo> classes = classpath.getTopLevelClassesRecursive("org.apache.skywalking");
for (ClassPath.ClassInfo classInfo : classes) {
Class<?> aClass = classInfo.load();
for (AnnotationListener listener : listeners) {
if (aClass.isAnnotationPresent(listener.annotation())) {
listener.notify(aClass);
}
}
}
callBack.run();
}
}

View File

@ -19,7 +19,7 @@
package org.apache.skywalking.oap.server.core.remote;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.remote.client.*;
import org.apache.skywalking.oap.server.core.remote.selector.*;
import org.apache.skywalking.oap.server.library.module.*;
@ -41,20 +41,20 @@ public class RemoteSenderService implements Service {
this.rollingSelector = new RollingSelector();
}
public void send(int nextWorkId, Indicator indicator, Selector selector) {
public void send(int nextWorkId, StreamData streamData, Selector selector) {
RemoteClientManager clientManager = moduleManager.find(CoreModule.NAME).getService(RemoteClientManager.class);
RemoteClient remoteClient;
switch (selector) {
case HashCode:
remoteClient = hashCodeSelector.select(clientManager.getRemoteClient(), indicator);
remoteClient.push(nextWorkId, indicator);
remoteClient = hashCodeSelector.select(clientManager.getRemoteClient(), streamData);
remoteClient.push(nextWorkId, streamData);
case Rolling:
remoteClient = rollingSelector.select(clientManager.getRemoteClient(), indicator);
remoteClient.push(nextWorkId, indicator);
remoteClient = rollingSelector.select(clientManager.getRemoteClient(), streamData);
remoteClient.push(nextWorkId, streamData);
case ForeverFirst:
remoteClient = foreverFirstSelector.select(clientManager.getRemoteClient(), indicator);
remoteClient.push(nextWorkId, indicator);
remoteClient = foreverFirstSelector.select(clientManager.getRemoteClient(), streamData);
remoteClient.push(nextWorkId, streamData);
}
}
}

View File

@ -19,11 +19,12 @@
package org.apache.skywalking.oap.server.core.remote;
import io.grpc.stub.StreamObserver;
import java.util.Objects;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper;
import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper;
import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataClassGetter;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.remote.grpc.proto.*;
import org.apache.skywalking.oap.server.core.worker.annotation.WorkerClassGetter;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.apache.skywalking.oap.server.library.server.grpc.GRPCHandler;
import org.slf4j.*;
@ -35,24 +36,33 @@ public class RemoteServiceHandler extends RemoteServiceGrpc.RemoteServiceImplBas
private static final Logger logger = LoggerFactory.getLogger(RemoteServiceHandler.class);
private final IndicatorMapper indicatorMapper;
private final WorkerMapper workerMapper;
private final ModuleManager moduleManager;
private StreamDataClassGetter streamDataClassGetter;
private WorkerClassGetter workerClassGetter;
public RemoteServiceHandler(ModuleManager moduleManager) {
this.indicatorMapper = moduleManager.find(CoreModule.NAME).getService(IndicatorMapper.class);
this.workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class);
this.moduleManager = moduleManager;
}
@Override public StreamObserver<RemoteMessage> call(StreamObserver<Empty> responseObserver) {
if (Objects.isNull(streamDataClassGetter)) {
streamDataClassGetter = moduleManager.find(CoreModule.NAME).getService(StreamDataClassGetter.class);
}
if (Objects.isNull(streamDataClassGetter)) {
workerClassGetter = moduleManager.find(CoreModule.NAME).getService(WorkerClassGetter.class);
}
return new StreamObserver<RemoteMessage>() {
@Override public void onNext(RemoteMessage message) {
int indicatorId = message.getIndicatorId();
int streamDataId = message.getStreamDataId();
int nextWorkerId = message.getNextWorkerId();
RemoteData remoteData = message.getRemoteData();
Class<Indicator> indicatorClass = indicatorMapper.findClassById(indicatorId);
Class<StreamData> streamDataClass = streamDataClassGetter.findClassById(streamDataId);
try {
indicatorClass.newInstance().deserialize(remoteData);
StreamData streamData = streamDataClass.newInstance();
streamData.deserialize(remoteData);
workerClassGetter.getClassById(nextWorkerId).newInstance().in(streamData);
} catch (InstantiationException | IllegalAccessException e) {
logger.warn(e.getMessage());
}

View File

@ -0,0 +1,49 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.remote.annotation;
import java.lang.annotation.Annotation;
import java.util.*;
import lombok.Getter;
import org.apache.skywalking.oap.server.core.annotation.AnnotationListener;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class StreamAnnotationListener implements AnnotationListener {
private static final Logger logger = LoggerFactory.getLogger(StreamAnnotationListener.class);
@Getter private final List<Class> streamClasses;
public StreamAnnotationListener() {
this.streamClasses = new LinkedList<>();
}
@Override public Class<? extends Annotation> annotation() {
return StreamData.class;
}
@Override public void notify(Class aClass) {
logger.info("The owner class of stream data annotation, class name: {}", aClass.getName());
streamClasses.add(aClass);
}
}

View File

@ -0,0 +1,29 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.remote.annotation;
import java.lang.annotation.*;
/**
* @author peng-yongsheng
*/
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
public @interface StreamData {
}

View File

@ -0,0 +1,59 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.remote.annotation;
import java.util.*;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class StreamDataAnnotationContainer implements StreamDataClassGetter {
private static final Logger logger = LoggerFactory.getLogger(StreamDataAnnotationContainer.class);
private int id = 0;
private final Map<Class<StreamData>, Integer> classMap;
private final Map<Integer, Class<StreamData>> idMap;
public StreamDataAnnotationContainer() {
this.classMap = new HashMap<>();
this.idMap = new HashMap<>();
}
@SuppressWarnings(value = "unchecked")
public synchronized void generate(List<Class> streamDataClasses) {
streamDataClasses.sort(Comparator.comparing(Class::getName));
for (Class streamDataClass : streamDataClasses) {
id++;
classMap.put(streamDataClass, id);
idMap.put(id, streamDataClass);
}
}
public int findIdByClass(Class streamDataClass) {
return classMap.get(streamDataClass);
}
@Override public Class<StreamData> findClassById(int id) {
return idMap.get(id);
}
}

View File

@ -16,21 +16,15 @@
*
*/
package org.apache.skywalking.oap.server.core.analysis.indicator.define;
package org.apache.skywalking.oap.server.core.remote.annotation;
import org.junit.*;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.library.module.Service;
/**
* @author peng-yongsheng
*/
public class IndicatorMapperTestCase {
public interface StreamDataClassGetter extends Service {
@Test
public void test() throws IndicatorDefineLoadException {
IndicatorMapper mapper = new IndicatorMapper();
mapper.load();
Assert.assertEquals(1, mapper.findIdByClass(TestAvgIndicator.class));
Assert.assertEquals(TestAvgIndicator.class, mapper.findClassById(1));
}
Class<StreamData> findClassById(int id);
}

View File

@ -23,9 +23,9 @@ import java.util.List;
import org.apache.skywalking.apm.commons.datacarrier.DataCarrier;
import org.apache.skywalking.apm.commons.datacarrier.buffer.BufferStrategy;
import org.apache.skywalking.apm.commons.datacarrier.consumer.IConsumer;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.cluster.RemoteInstance;
import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer;
import org.apache.skywalking.oap.server.core.remote.grpc.proto.*;
import org.apache.skywalking.oap.server.library.client.grpc.GRPCClient;
import org.slf4j.*;
@ -39,23 +39,23 @@ public class GRPCRemoteClient implements RemoteClient, Comparable<GRPCRemoteClie
private final GRPCClient client;
private final DataCarrier<RemoteMessage> carrier;
private final IndicatorMapper indicatorMapper;
private final StreamDataAnnotationContainer streamDataMapper;
public GRPCRemoteClient(IndicatorMapper indicatorMapper, RemoteInstance remoteInstance, int channelSize,
public GRPCRemoteClient(StreamDataAnnotationContainer streamDataMapper, RemoteInstance remoteInstance, int channelSize,
int bufferSize) {
this.indicatorMapper = indicatorMapper;
this.streamDataMapper = streamDataMapper;
this.client = new GRPCClient(remoteInstance.getHost(), remoteInstance.getPort());
this.carrier = new DataCarrier<>(channelSize, bufferSize);
this.carrier.setBufferStrategy(BufferStrategy.BLOCKING);
this.carrier.consume(new RemoteMessageConsumer(), 1);
}
@Override public void push(int nextWorkerId, Indicator indicator) {
int indicatorId = indicatorMapper.findIdByClass(indicator.getClass());
@Override public void push(int nextWorkerId, StreamData streamData) {
int streamDataId = streamDataMapper.findIdByClass(streamData.getClass());
RemoteMessage.Builder builder = RemoteMessage.newBuilder();
builder.setNextWorkerId(nextWorkerId);
builder.setIndicatorId(indicatorId);
builder.setRemoteData(indicator.serialize());
builder.setStreamDataId(streamDataId);
builder.setRemoteData(streamData.serialize());
this.carrier.produce(builder.build());
}

View File

@ -18,7 +18,7 @@
package org.apache.skywalking.oap.server.core.remote.client;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
/**
* @author peng-yongsheng
@ -29,5 +29,5 @@ public interface RemoteClient {
int getPort();
void push(int nextWorkerId, Indicator indicator);
void push(int nextWorkerId, StreamData streamData);
}

View File

@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.core.remote.client;
import java.util.*;
import java.util.concurrent.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper;
import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer;
import org.apache.skywalking.oap.server.core.cluster.*;
import org.apache.skywalking.oap.server.library.module.*;
import org.slf4j.*;
@ -33,7 +33,7 @@ public class RemoteClientManager implements Service {
private static final Logger logger = LoggerFactory.getLogger(RemoteClientManager.class);
private final ModuleManager moduleManager;
private IndicatorMapper indicatorMapper;
private StreamDataAnnotationContainer indicatorMapper;
private ClusterNodesQuery clusterNodesQuery;
private final List<RemoteClient> clientsA;
private final List<RemoteClient> clientsB;
@ -48,7 +48,7 @@ public class RemoteClientManager implements Service {
public void start() {
this.clusterNodesQuery = moduleManager.find(ClusterModule.NAME).getService(ClusterNodesQuery.class);
this.indicatorMapper = moduleManager.find(ClusterModule.NAME).getService(IndicatorMapper.class);
this.indicatorMapper = moduleManager.find(ClusterModule.NAME).getService(StreamDataAnnotationContainer.class);
Executors.newSingleThreadScheduledExecutor().scheduleAtFixedRate(this::refresh, 1, 2, TimeUnit.SECONDS);
}

View File

@ -19,8 +19,8 @@
package org.apache.skywalking.oap.server.core.remote.client;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.analysis.worker.define.WorkerMapper;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.worker.annotation.WorkerAnnotationContainer;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
/**
@ -46,8 +46,8 @@ public class SelfRemoteClient implements RemoteClient {
return port;
}
@Override public void push(int nextWorkerId, Indicator indicator) {
WorkerMapper workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerMapper.class);
workerMapper.findInstanceById(nextWorkerId).in(indicator);
@Override public void push(int nextWorkerId, StreamData streamData) {
WorkerAnnotationContainer workerMapper = moduleManager.find(CoreModule.NAME).getService(WorkerAnnotationContainer.class);
workerMapper.findInstanceById(nextWorkerId).in(streamData);
}
}

View File

@ -16,8 +16,9 @@
*
*/
package org.apache.skywalking.oap.server.core.analysis.data;
package org.apache.skywalking.oap.server.core.remote.data;
import org.apache.skywalking.oap.server.core.analysis.data.*;
import org.apache.skywalking.oap.server.core.remote.*;
/**

View File

@ -19,7 +19,7 @@
package org.apache.skywalking.oap.server.core.remote.selector;
import java.util.List;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.remote.client.RemoteClient;
import org.slf4j.*;
@ -30,7 +30,7 @@ public class ForeverFirstSelector implements RemoteClientSelector {
private static final Logger logger = LoggerFactory.getLogger(ForeverFirstSelector.class);
@Override public RemoteClient select(List<RemoteClient> clients, Indicator indicator) {
@Override public RemoteClient select(List<RemoteClient> clients, StreamData streamData) {
if (logger.isDebugEnabled()) {
logger.debug("clients size: {}", clients.size());
}

View File

@ -19,7 +19,7 @@
package org.apache.skywalking.oap.server.core.remote.selector;
import java.util.List;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.remote.client.RemoteClient;
/**
@ -27,9 +27,9 @@ import org.apache.skywalking.oap.server.core.remote.client.RemoteClient;
*/
public class HashCodeSelector implements RemoteClientSelector {
@Override public RemoteClient select(List<RemoteClient> clients, Indicator indicator) {
@Override public RemoteClient select(List<RemoteClient> clients, StreamData streamData) {
int size = clients.size();
int selectIndex = Math.abs(indicator.hashCode()) % size;
int selectIndex = Math.abs(streamData.hashCode()) % size;
return clients.get(selectIndex);
}
}

View File

@ -19,12 +19,12 @@
package org.apache.skywalking.oap.server.core.remote.selector;
import java.util.List;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.remote.client.RemoteClient;
/**
* @author peng-yongsheng
*/
public interface RemoteClientSelector {
RemoteClient select(List<RemoteClient> clients, Indicator indicator);
RemoteClient select(List<RemoteClient> clients, StreamData streamData);
}

View File

@ -19,7 +19,7 @@
package org.apache.skywalking.oap.server.core.remote.selector;
import java.util.List;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.remote.data.StreamData;
import org.apache.skywalking.oap.server.core.remote.client.RemoteClient;
/**
@ -29,7 +29,7 @@ public class RollingSelector implements RemoteClientSelector {
private int index = 0;
@Override public RemoteClient select(List<RemoteClient> clients, Indicator indicator) {
@Override public RemoteClient select(List<RemoteClient> clients, StreamData streamData) {
int size = clients.size();
index++;
int selectIndex = Math.abs(index) % size;

View File

@ -24,7 +24,7 @@ import org.apache.skywalking.oap.server.core.source.Scope;
* @author peng-yongsheng
*/
public interface IRegisterLockDAO extends DAO {
boolean tryLock(Scope scope, int timeout);
boolean tryLock(Scope scope);
void releaseLock(Scope scope);
}

View File

@ -1,81 +0,0 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.storage;
import java.util.*;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.analysis.indicator.define.IndicatorMapper;
import org.apache.skywalking.oap.server.core.storage.annotation.ColumnAnnotationRetrieval;
import org.apache.skywalking.oap.server.core.storage.define.*;
import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public abstract class StorageInstaller {
private static final Logger logger = LoggerFactory.getLogger(StorageInstaller.class);
private final ModuleManager moduleManager;
private final ColumnAnnotationRetrieval annotationRetrieval;
public StorageInstaller(ModuleManager moduleManager) {
this.moduleManager = moduleManager;
this.annotationRetrieval = new ColumnAnnotationRetrieval();
}
public final void install(Client client) throws StorageException {
IndicatorMapper indicatorMapper = moduleManager.find(CoreModule.NAME).getService(IndicatorMapper.class);
Collection<Class<Indicator>> indicatorClasses = indicatorMapper.indicatorClasses();
Boolean debug = System.getProperty("debug") != null;
for (Class<Indicator> indicatorClass : indicatorClasses) {
List<ColumnDefine> columnDefines = annotationRetrieval.retrieval(indicatorClass);
String tableName;
try {
tableName = indicatorClass.newInstance().name();
} catch (InstantiationException | IllegalAccessException e) {
throw new StorageException(e.getMessage());
}
TableDefine tableDefine = new TableDefine(tableName, columnDefines);
if (!isExists(client, tableDefine)) {
logger.info("table: {} not exists", tableDefine.getName());
createTable(client, tableDefine);
} else if (debug) {
logger.info("table: {} exists", tableDefine.getName());
deleteTable(client, tableDefine);
createTable(client, tableDefine);
}
columnCheck(client, tableDefine);
}
}
protected abstract boolean isExists(Client client, TableDefine tableDefine) throws StorageException;
protected abstract void columnCheck(Client client, TableDefine tableDefine) throws StorageException;
protected abstract void deleteTable(Client client, TableDefine tableDefine) throws StorageException;
protected abstract void createTable(Client client, TableDefine tableDefine) throws StorageException;
}

View File

@ -18,35 +18,48 @@
package org.apache.skywalking.oap.server.core.storage.annotation;
import java.lang.annotation.Annotation;
import java.lang.reflect.Field;
import java.util.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import org.apache.skywalking.oap.server.core.storage.define.*;
import lombok.Getter;
import org.apache.skywalking.oap.server.core.annotation.AnnotationListener;
import org.apache.skywalking.oap.server.core.storage.model.*;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class ColumnAnnotationRetrieval {
public class StorageAnnotationListener implements AnnotationListener, IModelGetter {
private static final Logger logger = LoggerFactory.getLogger(ColumnAnnotationRetrieval.class);
private static final Logger logger = LoggerFactory.getLogger(StorageAnnotationListener.class);
public List<ColumnDefine> retrieval(Class<Indicator> indicatorClass) {
if (logger.isDebugEnabled()) {
logger.debug("Retrieval column annotation from class {}", indicatorClass.getName());
}
List<ColumnDefine> columnDefines = new LinkedList<>();
retrieval(indicatorClass, columnDefines);
return columnDefines;
@Getter private final List<Model> models;
public StorageAnnotationListener() {
this.models = new LinkedList<>();
}
private void retrieval(Class clazz, List<ColumnDefine> columnDefines) {
@Override public Class<? extends Annotation> annotation() {
return StorageEntity.class;
}
@Override public void notify(Class aClass) {
logger.info("The owner class of storage annotation, class name: {}", aClass.getName());
List<ModelColumn> modelColumns = new LinkedList<>();
retrieval(aClass, modelColumns);
StorageEntity annotation = (StorageEntity)aClass.getAnnotation(StorageEntity.class);
models.add(new Model(annotation.name(), modelColumns));
}
private void retrieval(Class clazz, List<ModelColumn> modelColumns) {
Field[] fields = clazz.getDeclaredFields();
for (Field field : fields) {
if (field.isAnnotationPresent(Column.class)) {
Column column = field.getAnnotation(Column.class);
columnDefines.add(new ColumnDefine(new ColumnName(column.columnName(), column.columnName()), field.getType()));
modelColumns.add(new ModelColumn(new ColumnName(column.columnName(), column.columnName()), field.getType()));
if (logger.isDebugEnabled()) {
logger.debug("The field named {} with the {} type", column.columnName(), field.getType());
}
@ -54,7 +67,7 @@ public class ColumnAnnotationRetrieval {
}
if (Objects.nonNull(clazz.getSuperclass())) {
retrieval(clazz.getSuperclass(), columnDefines);
retrieval(clazz.getSuperclass(), modelColumns);
}
}
}

View File

@ -0,0 +1,30 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.storage.annotation;
import java.lang.annotation.*;
/**
* @author peng-yongsheng
*/
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
public @interface StorageEntity {
String name();
}

View File

@ -16,7 +16,7 @@
*
*/
package org.apache.skywalking.oap.server.core.storage.define;
package org.apache.skywalking.oap.server.core.storage.model;
/**
* @author peng-yongsheng

View File

@ -16,12 +16,12 @@
*
*/
package org.apache.skywalking.oap.server.core.storage.define;
package org.apache.skywalking.oap.server.core.storage.model;
/**
* @author peng-yongsheng
*/
public interface ColumnTypeMapping {
public interface DataTypeMapping {
String transform(Class<?> type);
}

View File

@ -0,0 +1,29 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.storage.model;
import java.util.List;
import org.apache.skywalking.oap.server.library.module.Service;
/**
* @author peng-yongsheng
*/
public interface IModelGetter extends Service {
List<Model> getModels();
}

View File

@ -16,27 +16,20 @@
*
*/
package org.apache.skywalking.oap.server.core.storage.define;
package org.apache.skywalking.oap.server.core.storage.model;
import java.util.List;
import lombok.Getter;
/**
* @author peng-yongsheng
*/
public class TableDefine {
private final String name;
private final List<ColumnDefine> columnDefines;
public class Model {
@Getter private final String name;
@Getter private final List<ModelColumn> columns;
public TableDefine(String name, List<ColumnDefine> columnDefines) {
public Model(String name, List<ModelColumn> columns) {
this.name = name;
this.columnDefines = columnDefines;
}
public final String getName() {
return name;
}
public final List<ColumnDefine> getColumnDefines() {
return columnDefines;
this.columns = columns;
}
}

View File

@ -16,25 +16,19 @@
*
*/
package org.apache.skywalking.oap.server.core.storage.define;
package org.apache.skywalking.oap.server.core.storage.model;
import lombok.Getter;
/**
* @author peng-yongsheng
*/
public class ColumnDefine {
private final ColumnName columnName;
private final Class<?> type;
public class ModelColumn {
@Getter private final ColumnName columnName;
@Getter private final Class<?> type;
public ColumnDefine(ColumnName columnName, Class<?> type) {
public ModelColumn(ColumnName columnName, Class<?> type) {
this.columnName = columnName;
this.type = type;
}
public final ColumnName getColumnName() {
return columnName;
}
public final Class<?> getType() {
return type;
}
}

View File

@ -0,0 +1,67 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.storage.model;
import java.util.List;
import org.apache.skywalking.oap.server.core.CoreModule;
import org.apache.skywalking.oap.server.core.storage.StorageException;
import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public abstract class ModelInstaller {
private static final Logger logger = LoggerFactory.getLogger(ModelInstaller.class);
private final ModuleManager moduleManager;
public ModelInstaller(ModuleManager moduleManager) {
this.moduleManager = moduleManager;
}
public final void install(Client client) throws StorageException {
IModelGetter modelGetter = moduleManager.find(CoreModule.NAME).getService(IModelGetter.class);
List<Model> models = modelGetter.getModels();
Boolean debug = System.getProperty("debug") != null;
for (Model model : models) {
if (!isExists(client, model)) {
logger.info("table: {} not exists", model.getName());
createTable(client, model);
} else if (debug) {
logger.info("table: {} exists", model.getName());
deleteTable(client, model);
createTable(client, model);
}
columnCheck(client, model);
}
}
protected abstract boolean isExists(Client client, Model model) throws StorageException;
protected abstract void columnCheck(Client client, Model model) throws StorageException;
protected abstract void deleteTable(Client client, Model model) throws StorageException;
protected abstract void createTable(Client client, Model model) throws StorageException;
}

View File

@ -0,0 +1,27 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.worker;
/**
* @author peng-yongsheng
*/
public abstract class AbstractWorker<INPUT> {
public abstract void in(INPUT input);
}

View File

@ -16,14 +16,14 @@
*
*/
package org.apache.skywalking.oap.server.core.analysis.worker;
package org.apache.skywalking.oap.server.core.worker.annotation;
import org.apache.skywalking.oap.server.core.analysis.indicator.Indicator;
import java.lang.annotation.*;
/**
* @author peng-yongsheng
*/
public abstract class Worker<INPUT extends Indicator> {
public abstract void in(INPUT input);
@Target(ElementType.TYPE)
@Retention(RetentionPolicy.RUNTIME)
public @interface Worker {
}

View File

@ -0,0 +1,84 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.worker.annotation;
import java.lang.reflect.Constructor;
import java.util.*;
import org.apache.skywalking.oap.server.core.worker.AbstractWorker;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class WorkerAnnotationContainer implements WorkerClassGetter {
private static final Logger logger = LoggerFactory.getLogger(WorkerAnnotationContainer.class);
private int id = 0;
private final Map<Class<AbstractWorker>, Integer> classKeyMapping;
private final Map<Integer, Class<AbstractWorker>> idKeyMapping;
private final Map<Class<AbstractWorker>, AbstractWorker> classKeyInstanceMapping;
private final Map<Integer, AbstractWorker> 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<Class> workerClasses) throws WorkerDefineLoadException {
if (Objects.isNull(workerClasses)) {
return;
}
try {
for (Class workerClass : workerClasses) {
id++;
classKeyMapping.put(workerClass, id);
idKeyMapping.put(id, workerClass);
Constructor<AbstractWorker> 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<AbstractWorker> 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);
}
}

View File

@ -0,0 +1,49 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.worker.annotation;
import java.lang.annotation.Annotation;
import java.util.*;
import lombok.Getter;
import org.apache.skywalking.oap.server.core.annotation.AnnotationListener;
import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class WorkerAnnotationListener implements AnnotationListener {
private static final Logger logger = LoggerFactory.getLogger(WorkerAnnotationListener.class);
@Getter private final List<Class> workerClasses;
public WorkerAnnotationListener() {
this.workerClasses = new LinkedList<>();
}
@Override public Class<? extends Annotation> annotation() {
return Worker.class;
}
@Override public void notify(Class aClass) {
logger.info("The owner class of worker annotation, class name: {}", aClass.getName());
workerClasses.add(aClass);
}
}

View File

@ -0,0 +1,29 @@
/*
* Licensed to the Apache Software Foundation (ASF) under one or more
* contributor license agreements. See the NOTICE file distributed with
* this work for additional information regarding copyright ownership.
* The ASF licenses this file to You under the Apache License, Version 2.0
* (the "License"); you may not use this file except in compliance with
* the License. You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*
*/
package org.apache.skywalking.oap.server.core.worker.annotation;
import org.apache.skywalking.oap.server.core.worker.AbstractWorker;
import org.apache.skywalking.oap.server.library.module.Service;
/**
* @author peng-yongsheng
*/
public interface WorkerClassGetter extends Service {
Class<AbstractWorker> getClassById(int workerId);
}

View File

@ -16,12 +16,12 @@
*
*/
package org.apache.skywalking.oap.server.core.analysis.worker.define;
package org.apache.skywalking.oap.server.core.worker.annotation;
/**
* @author peng-yongsheng
*/
public class WorkerDefineLoadException extends Exception {
public class WorkerDefineLoadException extends RuntimeException {
public WorkerDefineLoadException(String message, Throwable cause) {
super(message, cause);

View File

@ -28,7 +28,7 @@ service RemoteService {
message RemoteMessage {
int32 nextWorkerId = 1;
int32 indicatorId = 2;
int32 streamDataId = 2;
RemoteData remoteData = 3;
}

View File

@ -1,19 +0,0 @@
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
#
org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgIndicator

View File

@ -1,21 +0,0 @@
#
# Licensed to the Apache Software Foundation (ASF) under one or more
# contributor license agreements. See the NOTICE file distributed with
# this work for additional information regarding copyright ownership.
# The ASF licenses this file to You under the Apache License, Version 2.0
# (the "License"); you may not use this file except in compliance with
# the License. You may obtain a copy of the License at
#
# http://www.apache.org/licenses/LICENSE-2.0
#
# Unless required by applicable law or agreed to in writing, software
# distributed under the License is distributed on an "AS IS" BASIS,
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
# See the License for the specific language governing permissions and
# limitations under the License.
#
#
org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgAggregateWorker
org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgRemoteWorker
org.apache.skywalking.oap.server.core.analysis.endpoint.EndpointLatencyAvgPersistentWorker

View File

@ -34,6 +34,10 @@ public class TestAvgIndicator extends AvgIndicator {
return null;
}
@Override public String name() {
return null;
}
@Override public void deserialize(RemoteData remoteData) {
}
@ -41,10 +45,6 @@ public class TestAvgIndicator extends AvgIndicator {
return null;
}
@Override public String name() {
return null;
}
@Override public Map<String, Object> toMap() {
return null;
}

View File

@ -20,8 +20,8 @@ package org.apache.skywalking.oap.server.core.storage;
import java.util.LinkedList;
import org.apache.skywalking.oap.server.core.*;
import org.apache.skywalking.oap.server.core.analysis.indicator.define.*;
import org.apache.skywalking.oap.server.core.storage.define.TableDefine;
import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer;
import org.apache.skywalking.oap.server.core.storage.model.*;
import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.module.*;
import org.junit.Test;
@ -34,8 +34,8 @@ import org.powermock.reflect.Whitebox;
public class StorageInstallerTestCase {
@Test
public void testInstall() throws StorageException, DuplicateProviderException, ServiceNotProvidedException, IndicatorDefineLoadException {
IndicatorMapper indicatorMapper = new IndicatorMapper();
public void testInstall() throws StorageException, ServiceNotProvidedException {
StreamDataAnnotationContainer streamDataAnnotationContainer = new StreamDataAnnotationContainer();
CoreModuleProvider moduleProvider = Mockito.mock(CoreModuleProvider.class);
CoreModule moduleDefine = Mockito.spy(CoreModule.class);
ModuleManager moduleManager = Mockito.mock(ModuleManager.class);
@ -44,33 +44,33 @@ public class StorageInstallerTestCase {
moduleProviders.add(moduleProvider);
Mockito.when(moduleManager.find(CoreModule.NAME)).thenReturn(moduleDefine);
Mockito.when(moduleProvider.getService(IndicatorMapper.class)).thenReturn(indicatorMapper);
Mockito.when(moduleProvider.getService(StreamDataAnnotationContainer.class)).thenReturn(streamDataAnnotationContainer);
indicatorMapper.load();
// streamDataAnnotationContainer.generate();
TestStorageInstaller installer = new TestStorageInstaller(moduleManager);
installer.install(null);
// TestStorageInstaller installer = new TestStorageInstaller(moduleManager);
// installer.install(null);
}
class TestStorageInstaller extends StorageInstaller {
class TestStorageInstaller extends ModelInstaller {
public TestStorageInstaller(ModuleManager moduleManager) {
super(moduleManager);
}
@Override protected boolean isExists(Client client, TableDefine tableDefine) throws StorageException {
@Override protected boolean isExists(Client client, Model tableDefine) throws StorageException {
return false;
}
@Override protected void columnCheck(Client client, TableDefine tableDefine) throws StorageException {
@Override protected void columnCheck(Client client, Model tableDefine) throws StorageException {
}
@Override protected void deleteTable(Client client, TableDefine tableDefine) throws StorageException {
@Override protected void deleteTable(Client client, Model tableDefine) throws StorageException {
}
@Override protected void createTable(Client client, TableDefine tableDefine) throws StorageException {
@Override protected void createTable(Client client, Model tableDefine) throws StorageException {
}
}

View File

@ -16,7 +16,6 @@
*
*/
package org.apache.skywalking.oap.server.library.util;
import java.util.*;

View File

@ -64,7 +64,7 @@ public class StorageModuleElasticsearchProvider extends ModuleProvider {
this.registerServiceImplementation(IBatchDAO.class, new BatchProcessEsDAO(elasticSearchClient, config.getBulkActions(), config.getBulkSize(), config.getFlushInterval(), config.getConcurrentRequests()));
this.registerServiceImplementation(IPersistenceDAO.class, new PersistenceEsDAO(elasticSearchClient, nameSpace));
this.registerServiceImplementation(IRegisterLockDAO.class, new RegisterLockDAOImpl(elasticSearchClient));
this.registerServiceImplementation(IRegisterLockDAO.class, new RegisterLockDAOImpl(elasticSearchClient, 1000));
}
@Override

View File

@ -18,12 +18,12 @@
package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base;
import org.apache.skywalking.oap.server.core.storage.define.ColumnTypeMapping;
import org.apache.skywalking.oap.server.core.storage.model.DataTypeMapping;
/**
* @author peng-yongsheng
*/
public class ColumnTypeEsMapping implements ColumnTypeMapping {
public class ColumnTypeEsMapping implements DataTypeMapping {
@Override public String transform(Class<?> type) {
if (Integer.class.equals(type) || int.class.equals(type)) {

View File

@ -20,7 +20,7 @@ package org.apache.skywalking.oap.server.storage.plugin.elasticsearch.base;
import java.io.IOException;
import org.apache.skywalking.oap.server.core.storage.*;
import org.apache.skywalking.oap.server.core.storage.define.*;
import org.apache.skywalking.oap.server.core.storage.model.*;
import org.apache.skywalking.oap.server.library.client.Client;
import org.apache.skywalking.oap.server.library.client.elasticsearch.ElasticSearchClient;
import org.apache.skywalking.oap.server.library.module.ModuleManager;
@ -31,7 +31,7 @@ import org.slf4j.*;
/**
* @author peng-yongsheng
*/
public class StorageEsInstaller extends StorageInstaller {
public class StorageEsInstaller extends ModelInstaller {
private static final Logger logger = LoggerFactory.getLogger(StorageEsInstaller.class);
@ -46,7 +46,7 @@ public class StorageEsInstaller extends StorageInstaller {
this.mapping = new ColumnTypeEsMapping();
}
@Override protected boolean isExists(Client client, TableDefine tableDefine) throws StorageException {
@Override protected boolean isExists(Client client, Model tableDefine) throws StorageException {
ElasticSearchClient esClient = (ElasticSearchClient)client;
try {
return esClient.isExistsIndex(tableDefine.getName());
@ -55,11 +55,11 @@ public class StorageEsInstaller extends StorageInstaller {
}
}
@Override protected void columnCheck(Client client, TableDefine tableDefine) throws StorageException {
@Override protected void columnCheck(Client client, Model tableDefine) throws StorageException {
}
@Override protected void deleteTable(Client client, TableDefine tableDefine) throws StorageException {
@Override protected void deleteTable(Client client, Model tableDefine) throws StorageException {
ElasticSearchClient esClient = (ElasticSearchClient)client;
try {
@ -71,7 +71,7 @@ public class StorageEsInstaller extends StorageInstaller {
}
}
@Override protected void createTable(Client client, TableDefine tableDefine) throws StorageException {
@Override protected void createTable(Client client, Model tableDefine) throws StorageException {
ElasticSearchClient esClient = (ElasticSearchClient)client;
// mapping
@ -107,7 +107,7 @@ public class StorageEsInstaller extends StorageInstaller {
.build();
}
private XContentBuilder createMappingBuilder(TableDefine tableDefine) throws IOException {
private XContentBuilder createMappingBuilder(Model tableDefine) throws IOException {
XContentBuilder mappingBuilder = XContentFactory.jsonBuilder()
.startObject()
.startObject("_all")
@ -115,7 +115,7 @@ public class StorageEsInstaller extends StorageInstaller {
.endObject()
.startObject("properties");
for (ColumnDefine columnDefine : tableDefine.getColumnDefines()) {
for (ModelColumn columnDefine : tableDefine.getColumns()) {
mappingBuilder
.startObject(columnDefine.getColumnName().getName())
.field("type", mapping.transform(columnDefine.getType()))

View File

@ -35,11 +35,14 @@ public class RegisterLockDAOImpl extends EsDAO implements IRegisterLockDAO {
private static final Logger logger = LoggerFactory.getLogger(RegisterLockDAOImpl.class);
public RegisterLockDAOImpl(ElasticSearchClient client) {
private final int timeout;
public RegisterLockDAOImpl(ElasticSearchClient client, int timeout) {
super(client);
this.timeout = timeout;
}
@Override public boolean tryLock(Scope scope, int timeout) {
@Override public boolean tryLock(Scope scope) {
String id = String.valueOf(scope.ordinal());
try {
GetResponse response = getClient().get(RegisterLockIndex.NAME, id);