Optimize collector settings and indicator (#1841)
* Support maxConcurrentCallsPerConnection and maxMessageSize in grpcServer * Make percent percent is 1/10000 * Fix test cases. * Make service instance indicator dispatch based on server side only. * Set init data right in AlarmMeta. * Make endpoint query based on server side only.
This commit is contained in:
parent
948ef6d1c4
commit
b6241b2ad2
|
|
@ -32,6 +32,8 @@ public class CoreModuleConfig extends ModuleConfig {
|
|||
@Setter private String restContextPath;
|
||||
@Setter private String gRPCHost;
|
||||
@Setter private int gRPCPort;
|
||||
@Setter private int maxConcurrentCallsPerConnection;
|
||||
@Setter private int maxMessageSize;
|
||||
private final List<String> downsampling;
|
||||
@Setter private int recordDataTTL;
|
||||
@Setter private int minuteMetricsDataTTL;
|
||||
|
|
|
|||
|
|
@ -22,26 +22,58 @@ import java.io.IOException;
|
|||
import org.apache.skywalking.oap.server.core.analysis.indicator.annotation.IndicatorTypeListener;
|
||||
import org.apache.skywalking.oap.server.core.analysis.record.annotation.RecordTypeListener;
|
||||
import org.apache.skywalking.oap.server.core.annotation.AnnotationScan;
|
||||
import org.apache.skywalking.oap.server.core.cache.*;
|
||||
import org.apache.skywalking.oap.server.core.cluster.*;
|
||||
import org.apache.skywalking.oap.server.core.config.*;
|
||||
import org.apache.skywalking.oap.server.core.query.*;
|
||||
import org.apache.skywalking.oap.server.core.cache.CacheUpdateTimer;
|
||||
import org.apache.skywalking.oap.server.core.cache.EndpointInventoryCache;
|
||||
import org.apache.skywalking.oap.server.core.cache.NetworkAddressInventoryCache;
|
||||
import org.apache.skywalking.oap.server.core.cache.ServiceInstanceInventoryCache;
|
||||
import org.apache.skywalking.oap.server.core.cache.ServiceInventoryCache;
|
||||
import org.apache.skywalking.oap.server.core.cluster.ClusterModule;
|
||||
import org.apache.skywalking.oap.server.core.cluster.ClusterRegister;
|
||||
import org.apache.skywalking.oap.server.core.cluster.RemoteInstance;
|
||||
import org.apache.skywalking.oap.server.core.config.ComponentLibraryCatalogService;
|
||||
import org.apache.skywalking.oap.server.core.config.DownsamplingConfigService;
|
||||
import org.apache.skywalking.oap.server.core.config.IComponentLibraryCatalogService;
|
||||
import org.apache.skywalking.oap.server.core.query.AggregationQueryService;
|
||||
import org.apache.skywalking.oap.server.core.query.AlarmQueryService;
|
||||
import org.apache.skywalking.oap.server.core.query.MetadataQueryService;
|
||||
import org.apache.skywalking.oap.server.core.query.MetricQueryService;
|
||||
import org.apache.skywalking.oap.server.core.query.TopologyQueryService;
|
||||
import org.apache.skywalking.oap.server.core.query.TraceQueryService;
|
||||
import org.apache.skywalking.oap.server.core.register.annotation.InventoryTypeListener;
|
||||
import org.apache.skywalking.oap.server.core.register.service.*;
|
||||
import org.apache.skywalking.oap.server.core.remote.*;
|
||||
import org.apache.skywalking.oap.server.core.remote.annotation.*;
|
||||
import org.apache.skywalking.oap.server.core.register.service.EndpointInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.register.service.IEndpointInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.register.service.INetworkAddressInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.register.service.IServiceInstanceInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.register.service.IServiceInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.register.service.NetworkAddressInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.register.service.ServiceInstanceInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.register.service.ServiceInventoryRegister;
|
||||
import org.apache.skywalking.oap.server.core.remote.RemoteSenderService;
|
||||
import org.apache.skywalking.oap.server.core.remote.RemoteServiceHandler;
|
||||
import org.apache.skywalking.oap.server.core.remote.annotation.StreamAnnotationListener;
|
||||
import org.apache.skywalking.oap.server.core.remote.annotation.StreamDataAnnotationContainer;
|
||||
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.*;
|
||||
import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegister;
|
||||
import org.apache.skywalking.oap.server.core.server.GRPCHandlerRegisterImpl;
|
||||
import org.apache.skywalking.oap.server.core.server.JettyHandlerRegister;
|
||||
import org.apache.skywalking.oap.server.core.server.JettyHandlerRegisterImpl;
|
||||
import org.apache.skywalking.oap.server.core.source.SourceReceiver;
|
||||
import org.apache.skywalking.oap.server.core.source.SourceReceiverImpl;
|
||||
import org.apache.skywalking.oap.server.core.storage.PersistenceTimer;
|
||||
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.storage.ttl.DataTTLKeeperTimer;
|
||||
import org.apache.skywalking.oap.server.library.module.*;
|
||||
import org.apache.skywalking.oap.server.library.module.ModuleConfig;
|
||||
import org.apache.skywalking.oap.server.library.module.ModuleDefine;
|
||||
import org.apache.skywalking.oap.server.library.module.ModuleProvider;
|
||||
import org.apache.skywalking.oap.server.library.module.ModuleStartException;
|
||||
import org.apache.skywalking.oap.server.library.module.ServiceNotProvidedException;
|
||||
import org.apache.skywalking.oap.server.library.server.ServerException;
|
||||
import org.apache.skywalking.oap.server.library.server.grpc.GRPCServer;
|
||||
import org.apache.skywalking.oap.server.library.server.jetty.JettyServer;
|
||||
import org.slf4j.*;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
/**
|
||||
* @author peng-yongsheng
|
||||
|
|
@ -82,6 +114,12 @@ public class CoreModuleProvider extends ModuleProvider {
|
|||
|
||||
@Override public void prepare() throws ServiceNotProvidedException {
|
||||
grpcServer = new GRPCServer(moduleConfig.getGRPCHost(), moduleConfig.getGRPCPort());
|
||||
if (moduleConfig.getMaxConcurrentCallsPerConnection() > 0) {
|
||||
grpcServer.setMaxConcurrentCallsPerConnection(moduleConfig.getMaxConcurrentCallsPerConnection());
|
||||
}
|
||||
if (moduleConfig.getMaxMessageSize() > 0) {
|
||||
grpcServer.setMaxMessageSize(moduleConfig.getMaxMessageSize());
|
||||
}
|
||||
grpcServer.initialize();
|
||||
|
||||
jettyServer = new JettyServer(moduleConfig.getRestHost(), moduleConfig.getRestPort(), moduleConfig.getRestContextPath());
|
||||
|
|
|
|||
|
|
@ -20,6 +20,7 @@ package org.apache.skywalking.oap.server.core.alarm;
|
|||
|
||||
import lombok.Getter;
|
||||
import lombok.Setter;
|
||||
import org.apache.skywalking.oap.server.core.Const;
|
||||
import org.apache.skywalking.oap.server.core.source.Scope;
|
||||
|
||||
/**
|
||||
|
|
@ -33,7 +34,7 @@ public class AlarmMeta {
|
|||
public AlarmMeta(String indicatorName, Scope scope) {
|
||||
this.indicatorName = indicatorName;
|
||||
this.scope = scope;
|
||||
this.id = id;
|
||||
this.id = Const.EMPTY_STRING;
|
||||
}
|
||||
|
||||
public AlarmMeta(String indicatorName, Scope scope, String id) {
|
||||
|
|
|
|||
|
|
@ -54,7 +54,7 @@ public abstract class PercentIndicator extends Indicator implements IntValueHold
|
|||
}
|
||||
|
||||
@Override public void calculate() {
|
||||
percentage = (int)(match * 100 / total);
|
||||
percentage = (int)(match * 10000 / total);
|
||||
}
|
||||
|
||||
@Override public int getValue() {
|
||||
|
|
|
|||
|
|
@ -36,7 +36,7 @@ public class PercentIndicatorTest {
|
|||
|
||||
impl.calculate();
|
||||
|
||||
Assert.assertEquals(33, impl.getValue());
|
||||
Assert.assertEquals(3333, impl.getValue());
|
||||
|
||||
impl = new PercentIndicatorImpl();
|
||||
impl.combine(new EqualMatch(), true, true);
|
||||
|
|
@ -45,7 +45,7 @@ public class PercentIndicatorTest {
|
|||
|
||||
impl.calculate();
|
||||
|
||||
Assert.assertEquals(66, impl.getValue());
|
||||
Assert.assertEquals(6666, impl.getValue());
|
||||
}
|
||||
|
||||
@Test
|
||||
|
|
@ -64,7 +64,7 @@ public class PercentIndicatorTest {
|
|||
|
||||
impl.calculate();
|
||||
|
||||
Assert.assertEquals(50, impl.getValue());
|
||||
Assert.assertEquals(5000, impl.getValue());
|
||||
}
|
||||
|
||||
public class PercentIndicatorImpl extends PercentIndicator {
|
||||
|
|
|
|||
|
|
@ -37,6 +37,8 @@ public class GRPCServer implements Server {
|
|||
|
||||
private final String host;
|
||||
private final int port;
|
||||
private int maxConcurrentCallsPerConnection;
|
||||
private int maxMessageSize;
|
||||
private io.grpc.Server server;
|
||||
private NettyServerBuilder nettyServerBuilder;
|
||||
private SslContextBuilder sslContextBuilder;
|
||||
|
|
@ -46,6 +48,16 @@ public class GRPCServer implements Server {
|
|||
public GRPCServer(String host, int port) {
|
||||
this.host = host;
|
||||
this.port = port;
|
||||
this.maxConcurrentCallsPerConnection = 4;
|
||||
this.maxMessageSize = Integer.MAX_VALUE;
|
||||
}
|
||||
|
||||
public void setMaxConcurrentCallsPerConnection(int maxConcurrentCallsPerConnection) {
|
||||
this.maxConcurrentCallsPerConnection = maxConcurrentCallsPerConnection;
|
||||
}
|
||||
|
||||
public void setMaxMessageSize(int maxMessageSize) {
|
||||
this.maxMessageSize = maxMessageSize;
|
||||
}
|
||||
|
||||
/**
|
||||
|
|
@ -79,6 +91,7 @@ public class GRPCServer implements Server {
|
|||
public void initialize() {
|
||||
InetSocketAddress address = new InetSocketAddress(host, port);
|
||||
nettyServerBuilder = NettyServerBuilder.forAddress(address);
|
||||
nettyServerBuilder = nettyServerBuilder.maxConcurrentCallsPerConnection(maxConcurrentCallsPerConnection).maxMessageSize(maxMessageSize);
|
||||
logger.info("Server started, host {} listening on {}", host, port);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -80,10 +80,10 @@ public class TelemetryDataDispatcher {
|
|||
if (org.apache.skywalking.apm.network.common.DetectPoint.server.equals(metric.getDetectPoint())) {
|
||||
toAll(decorator, minuteTimeBucket);
|
||||
toService(decorator, minuteTimeBucket);
|
||||
toServiceInstance(decorator, minuteTimeBucket);
|
||||
toEndpoint(decorator, minuteTimeBucket);
|
||||
}
|
||||
toServiceRelation(decorator, minuteTimeBucket);
|
||||
toServiceInstance(decorator, minuteTimeBucket);
|
||||
toServiceInstanceRelation(decorator, minuteTimeBucket);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -149,6 +149,8 @@ public class MetadataQueryEsDAO extends EsDAO implements IMetadataQueryDAO {
|
|||
boolQueryBuilder.must().add(QueryBuilders.matchQuery(matchCName, keyword));
|
||||
}
|
||||
|
||||
boolQueryBuilder.must().add(QueryBuilders.termQuery(EndpointInventory.DETECT_POINT, DetectPoint.SERVER.ordinal()));
|
||||
|
||||
sourceBuilder.query(boolQueryBuilder);
|
||||
sourceBuilder.size(limit);
|
||||
|
||||
|
|
|
|||
Loading…
Reference in New Issue