# Elasticsearch sharding

**URL:** <https://discuss.elastic.co/t/elasticsearch-sharding/3925>\
**Category:** Elasticsearch\
**Created:** [February 14, 2011, 2:04pm UTC](https://discuss.elastic.co/t/elasticsearch-sharding/3925 "2011-02-14T14:04:19Z")\
**Posts on this page:** 4\
**Page:** 1

<div class="post-metadata">

**Author:** ![mbx](https://avatars.discourse-cdn.com/v4/letter/m/2bfe46/32.png) [@mbx](https://discuss.elastic.co/u/mbx)\
**Post date:** [February 14, 2011, 2:04pm UTC](https://discuss.elastic.co/t/elasticsearch-sharding/3925/1 "2011-02-14T14:04:19Z")

</div>

Hi,  
for my application i've tested Lucene, but quickly had problems with

> 10 million documents in one index on one server.  
> Query latency growed up to several seconds and for some keywords  
> result sets had been very large - up to out of memory.  
> Now i'm looking for a better solution.  
> During my research it seems that sharding is state of the art (Solr,  
> Elasticsearch) to work with large indices ( \>20 Mio docs).

My Questions:  
Why do i need sharding?  
If i want to search over several shards (e.g. jan2010 upto dec2010)  
i've to merge the results. Isn't it the same work than searching in a  
large 2010-index?

How does sharding work for search? I think understood how the  
documents are hashed and distributed to different shards but where are  
they merged?

Sorry for the beginner questions.  
Tank you!  
-mbx

---

<div class="post-metadata">

**Author:** ![Karussell1](https://avatars.discourse-cdn.com/v4/letter/k/50afbb/32.png) [@Karussell1](https://discuss.elastic.co/u/Karussell1)\
**Post date:** [February 14, 2011, 6:51pm UTC](https://discuss.elastic.co/t/elasticsearch-sharding/3925/2 "2011-02-14T18:51:18Z")

</div>

Results are merged when querying. there are different query types.  
take a look into:

[http://www.marcsturlese.com/2010/02/12/elasticsearch/](http://www.marcsturlese.com/2010/02/12/elasticsearch/)

Hitting several indices instead of one allows you to use several CPUs  
or to put the shards on different servers.

More on that subject should be explained by the author of ES 🙂 (i'm  
not an expert ...)

On 14 Feb., 15:04, mbx [myze...@googlemail.com](mailto:myze...@googlemail.com) wrote:

> Hi,  
> for my application i've tested Lucene, but quickly had problems with\>10 million documents in one index on one server.
> 
> Query latency growed up to several seconds and for some keywords  
> result sets had been very large - up to out of memory.  
> Now i'm looking for a better solution.  
> During my research it seems that sharding is state of the art (Solr,  
> Elasticsearch) to work with large indices ( \>20 Mio docs).
> 
> My Questions:  
> Why do i need sharding?  
> If i want to search over several shards (e.g. jan2010 upto dec2010)  
> i've to merge the results. Isn't it the same work than searching in a  
> large 2010-index?
> 
> How does sharding work for search? I think understood how the  
> documents are hashed and distributed to different shards but where are  
> they merged?
> 
> Sorry for the beginner questions.  
> Tank you!  
> -mbx

---

<div class="post-metadata">

**Author:** ![kimchy](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/kimchy/32/44952_2.png) [@kimchy](https://discuss.elastic.co/u/kimchy)\
**Post date:** [February 16, 2011, 2:03am UTC](https://discuss.elastic.co/t/elasticsearch-sharding/3925/3 "2011-02-16T02:03:18Z")

</div>

Thats the gist of it. Spreading a the search over multiple shards (that exist on several machines) will mean searching over a smaller index (per shard), and then joining the results back.  
On Monday, February 14, 2011 at 8:51 PM, Karussell wrote:

> Results are merged when querying. there are different query types.  
> take a look into:
> 
> [http://www.marcsturlese.com/2010/02/12/elasticsearch/](http://www.marcsturlese.com/2010/02/12/elasticsearch/)
> 
> Hitting several indices instead of one allows you to use several CPUs  
> or to put the shards on different servers.
> 
> More on that subject should be explained by the author of ES 🙂 (i'm  
> not an expert ...)
> 
> On 14 Feb., 15:04, mbx [myze...@googlemail.com](mailto:myze...@googlemail.com) wrote:
> 
> > Hi,  
> > for my application i've tested Lucene, but quickly had problems with\>10 million documents in one index on one server.
> > 
> > Query latency growed up to several seconds and for some keywords  
> > result sets had been very large - up to out of memory.  
> > Now i'm looking for a better solution.  
> > During my research it seems that sharding is state of the art (Solr,  
> > Elasticsearch) to work with large indices ( \>20 Mio docs).
> > 
> > My Questions:  
> > Why do i need sharding?  
> > If i want to search over several shards (e.g. jan2010 upto dec2010)  
> > i've to merge the results. Isn't it the same work than searching in a  
> > large 2010-index?
> > 
> > How does sharding work for search? I think understood how the  
> > documents are hashed and distributed to different shards but where are  
> > they merged?
> > 
> > Sorry for the beginner questions.  
> > Tank you!  
> > -mbx

---

<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:** [July 6, 2017, 4:12am UTC](https://discuss.elastic.co/t/elasticsearch-sharding/3925/4 "2017-07-06T04:12:01Z")

</div>


