diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java index 62a329ec8..109bb7323 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/CoreModuleConfig.java @@ -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 downsampling; @Setter private int recordDataTTL; @Setter private int minuteMetricsDataTTL; 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 698ec7313..b116fec9e 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 @@ -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()); diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMeta.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMeta.java index a09198215..82cc939f4 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMeta.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/alarm/AlarmMeta.java @@ -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) { diff --git a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicator.java b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicator.java index 1a049a333..2ef52e4da 100644 --- a/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicator.java +++ b/oap-server/server-core/src/main/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicator.java @@ -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() { diff --git a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicatorTest.java b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicatorTest.java index 4ec330ea2..0a1c2c4d4 100644 --- a/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicatorTest.java +++ b/oap-server/server-core/src/test/java/org/apache/skywalking/oap/server/core/analysis/indicator/PercentIndicatorTest.java @@ -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 { diff --git a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java index 0bd15e600..831482914 100644 --- a/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java +++ b/oap-server/server-library/library-server/src/main/java/org/apache/skywalking/oap/server/library/server/grpc/GRPCServer.java @@ -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); } diff --git a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java index 70b755a9f..c42fce1af 100644 --- a/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java +++ b/oap-server/server-receiver-plugin/skywalking-mesh-receiver-plugin/src/main/java/org/apache/skywalking/aop/server/receiver/mesh/TelemetryDataDispatcher.java @@ -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); } diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java index c334431f2..ab8c72772 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/query/MetadataQueryEsDAO.java @@ -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);