# Shard balancing questions

**URL:** <https://discuss.elastic.co/t/shard-balancing-questions/168604>\
**Category:** Elasticsearch\
**Created:** [February 15, 2019, 2:29pm UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604 "2019-02-15T14:29:02Z")\
**Posts on this page:** 9\
**Page:** 1

<div class="post-metadata">

**Author:** ![Raimon\_Bosch](https://avatars.discourse-cdn.com/v4/letter/r/a8b319/32.png) [@Raimon\_Bosch](https://discuss.elastic.co/u/Raimon_Bosch)\
**Post date:** [February 15, 2019, 2:29pm UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/1 "2019-02-15T14:29:02Z")

</div>

Hi,

I am wondering what is the advantage of having primary shards spread amongst several nodes. Since you can use multithreading, it is not more effective to handle all the details of the same search query from the same node?

[https://www.elastic.co/guide/en/elasticsearch/reference/current/shards-allocation.html#\_shard\_balancing\_heuristics](https://www.elastic.co/guide/en/elasticsearch/reference/current/shards-allocation.html#_shard_balancing_heuristics)

Which would be the consequences on designing a policy that concentrates all the primary shards on the same node? Bear in mind that each index will go to a different node and you will keep you replicas away in another node.

Thanks in advance,

---

<div class="post-metadata">

**Author:** ![dadoonet](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/dadoonet/32/137187_2.png) [@dadoonet](https://discuss.elastic.co/u/dadoonet)\
**Post date:** [February 15, 2019, 2:41pm UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/2 "2019-02-15T14:41:15Z")

</div>

Instead you can use `preference=_local`. See [https://www.elastic.co/guide/en/elasticsearch/reference/current/search-request-preference.html](https://www.elastic.co/guide/en/elasticsearch/reference/current/search-request-preference.html)

---

<div class="post-metadata">

**Author:** ![Raimon\_Bosch](https://avatars.discourse-cdn.com/v4/letter/r/a8b319/32.png) [@Raimon\_Bosch](https://discuss.elastic.co/u/Raimon_Bosch)\
**Post date:** [February 15, 2019, 2:52pm UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/3 "2019-02-15T14:52:18Z")

</div>

I'm not sure if that will work because in this use case it is needed to throw the query in all the shards to have an accurate result. I was more looking for a sharding policy that puts all the primary nodes on the same node when possible. Maybe it's best to not use shards at all.

---

<div class="post-metadata">

**Author:** ![dadoonet](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/dadoonet/32/137187_2.png) [@dadoonet](https://discuss.elastic.co/u/dadoonet)\
**Post date:** [February 16, 2019, 8:23am UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/4 "2019-02-16T08:23:58Z")

</div>

It will query all shards even if some of them are not available locally. It will just have a preference for the local ones.

But I'm wondering if it's a just question you have or a real problem?  
I mean that you should not really try to change the default behavior unless you have a real problem. Do you?

---

<div class="post-metadata">

**Author:** ![Raimon\_Bosch](https://avatars.discourse-cdn.com/v4/letter/r/a8b319/32.png) [@Raimon\_Bosch](https://discuss.elastic.co/u/Raimon_Bosch)\
**Post date:** [February 16, 2019, 8:50am UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/5 "2019-02-16T08:50:42Z")

</div>

Ok, I'll try the param. The idea is to benchmark several configurations before going live.

Cheers,

---

<div class="post-metadata">

**Author:** ![DavidTurner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/davidturner/32/22453_2.png) [@DavidTurner](https://discuss.elastic.co/u/DavidTurner)\
**Post date:** [February 16, 2019, 10:25am UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/6 "2019-02-16T10:25:05Z")

</div>

> [@Raimon\_Bosch](#):
>
> I was more looking for a sharding policy that puts all the primary nodes on the same node when possible.

I think you are perhaps confused about the role of primary shards. From the point of view of a search they're just another shard copy, no different from any other replica. The main difference between primaries and replicas is that primaries perform a small amount of extra coordination when indexing a document. Using `_preference=local` will try and keep a search on the local node regardless of whether the shards on the local node are primaries or replicas.

---

<div class="post-metadata">

**Author:** ![Raimon\_Bosch](https://avatars.discourse-cdn.com/v4/letter/r/a8b319/32.png) [@Raimon\_Bosch](https://discuss.elastic.co/u/Raimon_Bosch)\
**Post date:** [February 16, 2019, 10:42am UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/7 "2019-02-16T10:42:44Z")

</div>

An algorithm that keeps the N shards of an index in the same node works for us. Even if all are primary, replica or mixed.

I guess that for recovery purposes I would keep all the primaries together in one node, and its replicas together in another. So the probability of losing the data is less high.

So from that point of view, my interest has nothing to do with primaries or replicas, just with the fact of performing the search in multithread in the same node instead of multithread across several nodes. But anyway, maybe Elastic is designed so the latencies are negligible for this situation.

---

<div class="post-metadata">

**Author:** ![DavidTurner](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/davidturner/32/22453_2.png) [@DavidTurner](https://discuss.elastic.co/u/DavidTurner)\
**Post date:** [February 16, 2019, 12:22pm UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/8 "2019-02-16T12:22:27Z")

</div>

> [@Raimon\_Bosch](#):
>
> I guess that for recovery purposes I would keep all the primaries together in one node, and its replicas together in another. So the probability of losing the data is less high.

The primary/replica distinction doesn't make any difference in terms of data durability either.

> [@Raimon\_Bosch](#):
>
> So from that point of view, my interest has nothing to do with primaries or replicas, just with the fact of performing the search in multithread in the same node instead of multithread across several nodes. But anyway, maybe Elastic is designed so the latencies are negligible for this situation.

If CPU were the limiting factor then it probably would be better to stay on a single node using lots of threads, but there are other things like I/O bandwidth that don't scale with the number of processors and this is simplistically why it can be better to use multiple nodes. Network latency isn't normally a major concern, but if it is then you can use [shard allocation awareness](https://www.elastic.co/guide/en/elasticsearch/reference/current/allocation-awareness.html) to try and keep searches within a single zone (e.g. node, or rack) where the latency is best.

---

<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:** [March 16, 2019, 12:22pm UTC](https://discuss.elastic.co/t/shard-balancing-questions/168604/9 "2019-03-16T12:22:32Z")

</div>

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