From fbabcb1782762a8a14afaa6136e3c82cb5e5666a Mon Sep 17 00:00:00 2001 From: Jared Tan Date: Fri, 1 Nov 2019 18:32:06 +0800 Subject: [PATCH] make query max window size configurable. (#3765) * fix elasticsearch query data window size too large error. * make query max window size configurable. --- docker/oap/docker-entrypoint.sh | 1 + .../server-starter/src/main/assembly/application.yml | 1 + .../server-starter/src/main/resources/application.yml | 1 + .../elasticsearch/StorageModuleElasticsearchConfig.java | 3 ++- .../elasticsearch/StorageModuleElasticsearchProvider.java | 6 +++--- .../cache/NetworkAddressInventoryCacheEsDAO.java | 8 +++++--- .../elasticsearch/cache/ServiceInventoryCacheEsDAO.java | 8 +++++--- 7 files changed, 18 insertions(+), 10 deletions(-) diff --git a/docker/oap/docker-entrypoint.sh b/docker/oap/docker-entrypoint.sh index b75e6387e..eaca66fc2 100755 --- a/docker/oap/docker-entrypoint.sh +++ b/docker/oap/docker-entrypoint.sh @@ -116,6 +116,7 @@ cat <> ${var_application_file} bulkSize: \${SW_STORAGE_ES_BULK_SIZE:20} # flush the bulk every 20mb flushInterval: \${SW_STORAGE_ES_FLUSH_INTERVAL:10} # flush the bulk every 10 seconds whatever the number of requests concurrentRequests: \${SW_STORAGE_ES_CONCURRENT_REQUESTS:2} # the number of concurrent requests + resultWindowMaxSize: \${SW_STORAGE_ES_QUERY_MAX_WINDOW_SIZE:10000} metadataQueryMaxSize: \${SW_STORAGE_ES_QUERY_MAX_SIZE:5000} segmentQueryMaxSize: \${SW_STORAGE_ES_QUERY_SEGMENT_SIZE:200} EOT diff --git a/oap-server/server-starter/src/main/assembly/application.yml b/oap-server/server-starter/src/main/assembly/application.yml index c18191dba..83a9926b5 100644 --- a/oap-server/server-starter/src/main/assembly/application.yml +++ b/oap-server/server-starter/src/main/assembly/application.yml @@ -91,6 +91,7 @@ storage: # bulkActions: ${SW_STORAGE_ES_BULK_ACTIONS:1000} # Execute the bulk every 1000 requests # flushInterval: ${SW_STORAGE_ES_FLUSH_INTERVAL:10} # flush the bulk every 10 seconds whatever the number of requests # concurrentRequests: ${SW_STORAGE_ES_CONCURRENT_REQUESTS:2} # the number of concurrent requests +# resultWindowMaxSize: ${SW_STORAGE_ES_QUERY_MAX_WINDOW_SIZE:10000} # metadataQueryMaxSize: ${SW_STORAGE_ES_QUERY_MAX_SIZE:5000} # segmentQueryMaxSize: ${SW_STORAGE_ES_QUERY_SEGMENT_SIZE:200} h2: diff --git a/oap-server/server-starter/src/main/resources/application.yml b/oap-server/server-starter/src/main/resources/application.yml index 8b0d0fe20..e23b68666 100755 --- a/oap-server/server-starter/src/main/resources/application.yml +++ b/oap-server/server-starter/src/main/resources/application.yml @@ -90,6 +90,7 @@ storage: bulkActions: ${SW_STORAGE_ES_BULK_ACTIONS:1000} # Execute the bulk every 1000 requests flushInterval: ${SW_STORAGE_ES_FLUSH_INTERVAL:10} # flush the bulk every 10 seconds whatever the number of requests concurrentRequests: ${SW_STORAGE_ES_CONCURRENT_REQUESTS:2} # the number of concurrent requests + resultWindowMaxSize: ${SW_STORAGE_ES_QUERY_MAX_WINDOW_SIZE:10000} metadataQueryMaxSize: ${SW_STORAGE_ES_QUERY_MAX_SIZE:5000} segmentQueryMaxSize: ${SW_STORAGE_ES_QUERY_SEGMENT_SIZE:200} # h2: diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchConfig.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchConfig.java index 80e2c99c4..d8cea5e5b 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchConfig.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchConfig.java @@ -23,7 +23,7 @@ import lombok.Setter; import org.apache.skywalking.oap.server.library.module.ModuleConfig; /** - * @author peng-yongsheng + * @author peng-yongsheng, jian.tan */ @Getter public class StorageModuleElasticsearchConfig extends ModuleConfig { @@ -41,6 +41,7 @@ public class StorageModuleElasticsearchConfig extends ModuleConfig { @Setter private String password; @Getter @Setter String trustStorePath; @Getter @Setter String trustStorePass; + @Setter private int resultWindowMaxSize = 10000; @Setter private int metadataQueryMaxSize = 5000; @Setter private int segmentQueryMaxSize = 200; @Setter private int recordDataTTL = 7; diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java index d0766b58c..a025900e9 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/StorageModuleElasticsearchProvider.java @@ -71,7 +71,7 @@ import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.query.Trace import org.apache.skywalking.oap.server.storage.plugin.elasticsearch.ttl.ElasticsearchStorageTTL; /** - * @author peng-yongsheng + * @author peng-yongsheng, jian.tan */ public class StorageModuleElasticsearchProvider extends ModuleProvider { @@ -110,10 +110,10 @@ public class StorageModuleElasticsearchProvider extends ModuleProvider { this.registerServiceImplementation(IRegisterLockDAO.class, new RegisterLockDAOImpl(elasticSearchClient)); this.registerServiceImplementation(IHistoryDeleteDAO.class, new HistoryDeleteEsDAO(getManager(), elasticSearchClient, new ElasticsearchStorageTTL())); - this.registerServiceImplementation(IServiceInventoryCacheDAO.class, new ServiceInventoryCacheEsDAO(elasticSearchClient)); + this.registerServiceImplementation(IServiceInventoryCacheDAO.class, new ServiceInventoryCacheEsDAO(elasticSearchClient, config.getResultWindowMaxSize())); this.registerServiceImplementation(IServiceInstanceInventoryCacheDAO.class, new ServiceInstanceInventoryCacheDAO(elasticSearchClient)); this.registerServiceImplementation(IEndpointInventoryCacheDAO.class, new EndpointInventoryCacheEsDAO(elasticSearchClient)); - this.registerServiceImplementation(INetworkAddressInventoryCacheDAO.class, new NetworkAddressInventoryCacheEsDAO(elasticSearchClient)); + this.registerServiceImplementation(INetworkAddressInventoryCacheDAO.class, new NetworkAddressInventoryCacheEsDAO(elasticSearchClient, config.getResultWindowMaxSize())); this.registerServiceImplementation(ITopologyQueryDAO.class, new TopologyQueryEsDAO(elasticSearchClient)); this.registerServiceImplementation(IMetricsQueryDAO.class, new MetricsQueryEsDAO(elasticSearchClient)); diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/NetworkAddressInventoryCacheEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/NetworkAddressInventoryCacheEsDAO.java index 58154d095..57338e228 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/NetworkAddressInventoryCacheEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/NetworkAddressInventoryCacheEsDAO.java @@ -32,16 +32,18 @@ import org.elasticsearch.search.builder.SearchSourceBuilder; import org.slf4j.*; /** - * @author peng-yongsheng + * @author peng-yongsheng, jian.tan */ public class NetworkAddressInventoryCacheEsDAO extends EsDAO implements INetworkAddressInventoryCacheDAO { private static final Logger logger = LoggerFactory.getLogger(NetworkAddressInventoryCacheEsDAO.class); private final NetworkAddressInventory.Builder builder = new NetworkAddressInventory.Builder(); + private final int resultWindowMaxSize; - public NetworkAddressInventoryCacheEsDAO(ElasticSearchClient client) { + public NetworkAddressInventoryCacheEsDAO(ElasticSearchClient client, int resultWindowMaxSize) { super(client); + this.resultWindowMaxSize = resultWindowMaxSize; } @Override public int getAddressId(String networkAddress) { @@ -84,7 +86,7 @@ public class NetworkAddressInventoryCacheEsDAO extends EsDAO implements INetwork try { SearchSourceBuilder searchSourceBuilder = new SearchSourceBuilder(); searchSourceBuilder.query(QueryBuilders.rangeQuery(NetworkAddressInventory.LAST_UPDATE_TIME).gte(lastUpdateTime)); - searchSourceBuilder.size(Integer.MAX_VALUE); + searchSourceBuilder.size(resultWindowMaxSize); SearchResponse response = getClient().search(NetworkAddressInventory.INDEX_NAME, searchSourceBuilder); diff --git a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/ServiceInventoryCacheEsDAO.java b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/ServiceInventoryCacheEsDAO.java index f20d9d1d8..23fc8f32e 100644 --- a/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/ServiceInventoryCacheEsDAO.java +++ b/oap-server/server-storage-plugin/storage-elasticsearch-plugin/src/main/java/org/apache/skywalking/oap/server/storage/plugin/elasticsearch/cache/ServiceInventoryCacheEsDAO.java @@ -33,16 +33,18 @@ import org.elasticsearch.search.builder.SearchSourceBuilder; import org.slf4j.*; /** - * @author peng-yongsheng + * @author peng-yongsheng, jian.tan */ public class ServiceInventoryCacheEsDAO extends EsDAO implements IServiceInventoryCacheDAO { private static final Logger logger = LoggerFactory.getLogger(ServiceInventoryCacheEsDAO.class); private final ServiceInventory.Builder builder = new ServiceInventory.Builder(); + private final int resultWindowMaxSize; - public ServiceInventoryCacheEsDAO(ElasticSearchClient client) { + public ServiceInventoryCacheEsDAO(ElasticSearchClient client, int resultWindowMaxSize) { super(client); + this.resultWindowMaxSize = resultWindowMaxSize; } @Override public int getServiceId(String serviceName) { @@ -99,7 +101,7 @@ public class ServiceInventoryCacheEsDAO extends EsDAO implements IServiceInvento boolQuery.must().add(QueryBuilders.rangeQuery(ServiceInventory.LAST_UPDATE_TIME).gte(lastUpdateTime)); searchSourceBuilder.query(boolQuery); - searchSourceBuilder.size(Integer.MAX_VALUE); + searchSourceBuilder.size(resultWindowMaxSize); SearchResponse response = getClient().search(ServiceInventory.INDEX_NAME, searchSourceBuilder);