# Msearch performance decreases when increasing the number of shards

**URL:** <https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724>\
**Category:** Elasticsearch\
**Created:** [August 2, 2018, 9:11am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724 "2018-08-02T09:11:08Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![vincent.arnaud90](https://avatars.discourse-cdn.com/v4/letter/v/c89c15/32.png) [@vincent.arnaud90](https://discuss.elastic.co/u/vincent.arnaud90)\
**Post date:** [August 2, 2018, 9:11am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/1 "2018-08-02T09:11:08Z")

</div>

**Elasticsearch version** : 6.1.1

**Plugins installed** : [analysis-icu, analysis-phonetic]

**Description of the problem** :

I am testing the msearch capability on the rally's geonames set with quite a lot of basic match queries (1500) and even though each query is very fast I have noticed that the more shards I set for my index, the longer the whole msearch takes... no matter how many nodes/machines I set up in the cluster.

Thread pool might be a bottleneck on a single machine but I would have expected some performance gain with horizontal scaling.

**Steps to reproduce** :

Create three geonames indices with the mapping provided in rally's tracks :  
[https://github.com/elastic/rally-tracks/blob/master/geonames/index.json](https://github.com/elastic/rally-tracks/blob/master/geonames/index.json)  
Each one with a different number of shards (let's say 5, 10, 20). For example : geonames-5, geonames-10, geonames-20

Populate them with the rally's geonames data set (I had actually let rally create the first index and used the wonderful reindex API to populate the others)

Use the terms list ( [https://github.com/elastic/rally-tracks/blob/master/geonames/terms.txt](https://github.com/elastic/rally-tracks/blob/master/geonames/terms.txt) ) to generate the msearch "bodies" with some command like

```
shuf terms.txt | head -1500 > subterms.txt
awk '{ print "{\"index\": \"geonames-5\"}\n{\"query\": {\"match\": {\"name\": {\"query\": \"" $0 "\" }}}}"}' subterms.txt > msearch.geonames-5
awk '{ print "{\"index\": \"geonames-10\"}\n{\"query\": {\"match\": {\"name\": {\"query\": \"" $0 "\" }}}}"}' subterms.txt > msearch.geonames-10
awk '{ print "{\"index\": \"geonames-20\"}\n{\"query\": {\"match\": {\"name\": {\"query\": \"" $0 "\" }}}}"}' subterms.txt > msearch.geonames-20

```

run the msearches on each index.

Thanks for your help!

---

<div class="post-metadata">

**Author:** ![vincent.arnaud90](https://avatars.discourse-cdn.com/v4/letter/v/c89c15/32.png) [@vincent.arnaud90](https://discuss.elastic.co/u/vincent.arnaud90)\
**Post date:** [August 14, 2018, 2:06pm UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/2 "2018-08-14T14:06:28Z")

</div>

I am wondering whether it could be seen as an issue or not. Any advice?

---

<div class="post-metadata">

**Author:** ![danielmitterdorfer](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/danielmitterdorfer/32/110510_2.png) [@danielmitterdorfer](https://discuss.elastic.co/u/danielmitterdorfer)\
**Post date:** [August 31, 2018, 10:27am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/3 "2018-08-31T10:27:21Z")

</div>

Hi,

Elasticsearch has several limits in place that could explain this behavior. First of all, the number of concurrent searches is limited based on the number of nodes and the size of the search thread pool (see also `max_concurrent_searches` in the [msearch docs](https://www.elastic.co/guide/en/elasticsearch/reference/current/search-multi-search.html)). After a node has received an msearch requests, it sends each query individually to a node in the cluster (and respecting `max_concurrent_searches`). That node will coordinate the search request and will query individual shards. In order to avoid overwhelming the cluster there is another limit per query in place which allows at most 256 or number of nodes \* default number of shards (5) concurrent shard requests whichever of those two numbers is smaller (see [source code](https://github.com/elastic/elasticsearch/blob/1c105716f99450397d7a475aaccb640c57ead543/core/src/main/java/org/elasticsearch/action/search/TransportSearchAction.java#L344-L345)). To be clear: the default number of shards is the out-of-the-box default, not the configured default for that index. In the case of a three node cluster it would query 3 (nodes) \* 5 (shards) = 15 shards at most per search request.

As the geonames index is rather small (around 6GB), I guess the search time is dominated by:

- Wait time spent in the search queue of the individual nodes
- Wait time for results to arrive at the coordinating node

and finally also the limit of maximum number of concurrent shard requests.

I think it could make sense to retry this experiment with a much larger index.

Daniel

---

<div class="post-metadata">

**Author:** ![vakarami](https://avatars.discourse-cdn.com/v4/letter/v/47e85d/32.png) [@vakarami](https://discuss.elastic.co/u/vakarami)\
**Post date:** [September 2, 2018, 5:40am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/4 "2018-09-02T05:40:56Z")

</div>

Is there any config to limit `MaxConcurrentShardRequests`? because in our case 256 is still to much.  
Also we use ES 2.4 and It seems that this limit was added in ES 6.x. So there was no limit before that, yeah?

---

<div class="post-metadata">

**Author:** ![danielmitterdorfer](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/danielmitterdorfer/32/110510_2.png) [@danielmitterdorfer](https://discuss.elastic.co/u/danielmitterdorfer)\
**Post date:** [September 3, 2018, 5:18am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/5 "2018-09-03T05:18:59Z")

</div>

Hi,

> [@vakarami](#):
>
> Is there any config to limit `MaxConcurrentShardRequests` ? because in our case 256 is still to much.

See the [docs](https://www.elastic.co/guide/en/elasticsearch/reference/current/search.html#search-concurrency-and-parallelism) how to limit it. This has been introduced in 5.6.0 for the search API with [Limit the number of concurrent shard requests per search request by s1monw · Pull Request #25632 · elastic/elasticsearch · GitHub](https://github.com/elastic/elasticsearch/pull/25632) and [Expose `max_concurrent_shard_requests` in `_msearch` by s1monw · Pull Request #33016 · elastic/elasticsearch · GitHub](https://github.com/elastic/elasticsearch/pull/33016) will bring the same parameter for msearch.

Daniel

---

<div class="post-metadata">

**Author:** ![vincent.arnaud90](https://avatars.discourse-cdn.com/v4/letter/v/c89c15/32.png) [@vincent.arnaud90](https://discuss.elastic.co/u/vincent.arnaud90)\
**Post date:** [September 3, 2018, 7:09am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/6 "2018-09-03T07:09:38Z")

</div>

Hi Daniel and thanks a lot for the answer and details! It fits all my tests : I had first tried to play with the `max_concurrent_searches` setting (this [defaultMaxConcurrentSearches](https://github.com/elastic/elasticsearch/blob/237650e9c054149fd08213b38a81a3666c1868e5/server/src/main/java/org/elasticsearch/action/search/TransportMultiSearchAction.java#L105) looked pretty interesting) but without any success... which kind of determined the coordinating node as the bottleneck. Then splitting the msearch and sending each part to different client nodes helped.

Oh btw, is there any better and finer way to follow a threadpool than `watch -n 0.1` ?

---

<div class="post-metadata">

**Author:** ![danielmitterdorfer](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/danielmitterdorfer/32/110510_2.png) [@danielmitterdorfer](https://discuss.elastic.co/u/danielmitterdorfer)\
**Post date:** [September 3, 2018, 7:39am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/7 "2018-09-03T07:39:04Z")

</div>

Hi Vincent,

> [@vincent.arnaud90](#):
>
> Oh btw, is there any better and finer way to follow a threadpool than `watch -n 0.1` ?

I suppose you mean by that, that you call the [node stats API](https://www.elastic.co/guide/en/elasticsearch/reference/current/cluster-nodes-stats.html) periodically? That's the usual approach to monitor thread pool statistics from the command line. You can limit node stats to thread pool statistics e.g. by issuing `curl -X GET "es_host:9200/_nodes/stats/thread_pool?pretty"` and use a tool like [`jq`](https://stedolan.github.io/jq/) to only show the fields you're interested in. Alternatively, you can use [X-Pack monitoring](https://www.elastic.co/guide/en/elasticsearch/reference/current/es-monitoring.html) to view thread pool statistics over time.

Daniel

---

<div class="post-metadata">

**Author:** ![vincent.arnaud90](https://avatars.discourse-cdn.com/v4/letter/v/c89c15/32.png) [@vincent.arnaud90](https://discuss.elastic.co/u/vincent.arnaud90)\
**Post date:** [September 3, 2018, 11:47am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/8 "2018-09-03T11:47:23Z")

</div>

Yes, I'm calling the thread\_pool stats periodically. I am going to give another look at the X-pack monitoring, then.

Thank you again!

Vincent

---

<div class="post-metadata">

**Author:** ![system](https://us1.discourse-cdn.com/elastic/original/3X/1/a/1ac57faf039f6b580b3f104ef42a2a89e41014de.png) [@system](https://discuss.elastic.co/u/system)\
**Post date:** [October 1, 2018, 11:47am UTC](https://discuss.elastic.co/t/msearch-performance-decreases-when-increasing-the-number-of-shards/142724/9 "2018-10-01T11:47:34Z")

</div>

This topic was automatically closed 28 days after the last reply. New replies are no longer allowed.
