# More shards = more throughput?

**URL:** <https://discuss.elastic.co/t/more-shards-more-throughput/252223>\
**Category:** Elasticsearch\
**Created:** [October 15, 2020, 3:13pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223 "2020-10-15T15:13:21Z")\
**Posts on this page:** 10\
**Page:** 1

<div class="post-metadata">

**Author:** ![David\_Hopkins](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/david_hopkins/32/77272_2.png) [@David\_Hopkins](https://discuss.elastic.co/u/David_Hopkins)\
**Post date:** [October 15, 2020, 3:13pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/1 "2020-10-15T15:13:21Z")

</div>

I've been trying to better understand the effect sharing has on performance but the theory doesn't seem to match the reality so clearly one of those things must be wrong 😉

After watching this video on [Quantitative Cluster Sizing](https://www.elastic.co/elasticon/conf/2016/sf/quantitative-cluster-sizing) where an increase to the number of shards seems increase throughput -- which in my limited understanding seems to make sense.

My setup is about as vanilla/reliable as I can go: 3 node cluster running on [elastic.co](http://elastic.co) (compute optimised) and using a k6 loadtest client running a simple loop to post new documents as fast as it can.

I start the test with a single partition index and push the script to the point where response times start to degrade. I then delete the index and recreate it with two partitions.

What seems to happen is that a 2-partition index consistently performs _no_ better (and sometimes worse) than the single partition. This seems counter intuitive to me but perhaps I am missing something?

 ![Screenshot 2020-10-15 at 15.58.35](https://us1.discourse-cdn.com/elastic/original/3X/f/a/fa1dd7c579ec2d014bade56f61f73d6169d099fc.png)

Thanks!

---

<div class="post-metadata">

**Author:** ![Christian\_Dahlqvist](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/christian_dahlqvist/32/4617_2.png) [@Christian\_Dahlqvist](https://discuss.elastic.co/u/Christian_Dahlqvist)\
**Post date:** [October 15, 2020, 3:20pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/2 "2020-10-15T15:20:03Z")

</div>

How many concurrent connections are you using? What is your bulk size? What is the size of your documents? Do you have 3 data nodes in your cluster or 2 data nodes plus an arbiter?

---

<div class="post-metadata">

**Author:** ![David\_Hopkins](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/david_hopkins/32/77272_2.png) [@David\_Hopkins](https://discuss.elastic.co/u/David_Hopkins)\
**Post date:** [October 15, 2020, 3:41pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/3 "2020-10-15T15:41:42Z")

</div>

There's no bulking as such, I'm just making a single index request over http as fast as the k6 loadtest tool will allow me to do. I realise this isn't the optimal way but was curious to see how far this could get me.

Is an arbiter equivalent to master? I have a three nodes across three AWS zones where one of them is a master (and one is called a Tiebreaker.)

Thanks

---

<div class="post-metadata">

**Author:** ![Christian\_Dahlqvist](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/christian_dahlqvist/32/4617_2.png) [@Christian\_Dahlqvist](https://discuss.elastic.co/u/Christian_Dahlqvist)\
**Post date:** [October 15, 2020, 4:39pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/4 "2020-10-15T16:39:57Z")

</div>

If you are sending single documents using a single thread you are likely limited by the load testing tool rather than Elasticsearch, which is probably why you will get the same results no matter how many shards you have. I would recommend setting up a more realistic test and suspect [this webinar](https://www.elastic.co/elasticon/conf/2018/sf/the-seven-deadly-sins-of-elasticsearch-benchmarking) might be useful.

When it comes to optimal number of shards the results will vary depending on use case, but I would expect one or two active shards per node would be sufficient to saturate throughput and would expect throughput to potentially start decreasing after that. More active shards does in general not result in better throughput.

---

<div class="post-metadata">

**Author:** ![David\_Hopkins](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/david_hopkins/32/77272_2.png) [@David\_Hopkins](https://discuss.elastic.co/u/David_Hopkins)\
**Post date:** [October 15, 2020, 5:41pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/5 "2020-10-15T17:41:50Z")

</div>

Agreed that the tool is probably part of the reason but noticed that when I increased the shards I did see a proportional increase in search throughput so it seems the cluster is at least partially responsible.

In the below image I ran the same test with 1 and 2 shards respectively and you can see how searches goes up but indexing actually goes down. I thought it might be my local disk that was the bottleneck which I why I wanted to try the same test on the official Elastic cluster but I found I got the same results.

 ![Screenshot 2020-10-15 at 18.36.39](https://us1.discourse-cdn.com/elastic/original/3X/7/2/723d429645fef4102949df328b9059b5f1fd1906.png)

Do you have any thoughts on why this might be?

Thank you for the link I will definitely watch the webinar a bit later!

---

<div class="post-metadata">

**Author:** ![Christian\_Dahlqvist](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/christian_dahlqvist/32/4617_2.png) [@Christian\_Dahlqvist](https://discuss.elastic.co/u/Christian_Dahlqvist)\
**Post date:** [October 15, 2020, 6:07pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/6 "2020-10-15T18:07:51Z")

</div>

Search scales differently compared to indexing. Queries against multiple shards can be executed in parallel which means having 2 shards is likely to give a performance gain compared to a single primary shard, at least for low query concurrency.

---

<div class="post-metadata">

**Author:** ![astanton1978](https://avatars.discourse-cdn.com/v4/letter/a/e480ec/32.png) [@astanton1978](https://discuss.elastic.co/u/astanton1978)\
**Post date:** [October 15, 2020, 6:57pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/7 "2020-10-15T18:57:54Z")

</div>

If that testing tool does not do multiple concurrent requests, its not testing a realistic scenario.

Once you try concurrent single document writes/indexing, you are going to hit a limit with the indexing buffer far sooner than you would think. At that point, if you cant tune the index buffer, the only other options are to bulk index everything possible, and if that or throwing more hardware at it isnt enough. stick a load limiting/regulating queue or service in front of ES to receive and batch up the requests to ES

---

<div class="post-metadata">

**Author:** ![David\_Hopkins](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/david_hopkins/32/77272_2.png) [@David\_Hopkins](https://discuss.elastic.co/u/David_Hopkins)\
**Post date:** [October 15, 2020, 7:58pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/8 "2020-10-15T19:58:32Z")

</div>

> If that testing tool does not do multiple concurrent requests

Yes k6 performs concurrent http requests.

> you are going to hit a limit with the indexing buffer far sooner than you would think

Perhaps I am hitting this buffer but isn't that a per node limitation? My query is around comparing the throughput of a single node/single shard with two.

---

<div class="post-metadata">

**Author:** ![Christian\_Dahlqvist](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/christian_dahlqvist/32/4617_2.png) [@Christian\_Dahlqvist](https://discuss.elastic.co/u/Christian_Dahlqvist)\
**Post date:** [October 15, 2020, 8:15pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/9 "2020-10-15T20:15:16Z")

</div>

Indexing single documents is quite inefficient as data is synced to disk per request which results in a lot of disk I/O, which is often what limits Indexing throughput in Elasticsearch.

---

<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:** [November 12, 2020, 8:15pm UTC](https://discuss.elastic.co/t/more-shards-more-throughput/252223/10 "2020-11-12T20:15:17Z")

</div>

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