# About scalability in data volume

**URL:** <https://discuss.elastic.co/t/about-scalability-in-data-volume/9469>\
**Category:** Elasticsearch\
**Created:** [October 25, 2012, 7:05am UTC](https://discuss.elastic.co/t/about-scalability-in-data-volume/9469 "2012-10-25T07:05:51Z")\
**Posts on this page:** 5\
**Page:** 1

<div class="post-metadata">

**Author:** ![Liang\_Yu\_Chou](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/liang_yu_chou/32/2666_2.png) [@Liang\_Yu\_Chou](https://discuss.elastic.co/u/Liang_Yu_Chou)\
**Post date:** [October 25, 2012, 7:05am UTC](https://discuss.elastic.co/t/about-scalability-in-data-volume/9469/1 "2012-10-25T07:05:51Z")

</div>

For each cluster, I know that I can scale out query capacity by adding  
replica nodes.  
But is it possible that I can scale out in data volume?

As I know all data are written into one of the primary shards first, then  
copied into replicas.  
However the number of primary shards has to be defined at the beginning.  
Doesn't that  
pose a limitation on the max number of instances(for primary shards) in the  
cluster?

The default number of primary shard is 5. For future scalability, is there  
any drawback if  
I set it to a big number?

--

---

<div class="post-metadata">

**Author:** ![radu\_gheorghe](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/radu_gheorghe/32/556_2.png) [@radu\_gheorghe](https://discuss.elastic.co/u/radu_gheorghe)\
**Post date:** [October 25, 2012, 9:46am UTC](https://discuss.elastic.co/t/about-scalability-in-data-volume/9469/2 "2012-10-25T09:46:29Z")

</div>

Hello Jerry,

The problem with having many shards is that each shard has an  
overhead. So while it's OK to over-shard in order to scale out, if you  
exaggerate things will get slow, especially on the query side.

Luckily, there are other solutions than starting off with a huge  
number of shards. It depends on your data to choose what fits best.  
For example, you can add more indices as you go along.

I suggest you take a look at this video, as it's exactly about this topic:

> **[Elastic — The Search AI Company](https://www.elastic.co)**
>
> Power insights and outcomes with The Elastic Search AI Platform. See into your data and find answers that matter with enterprise solutions designed to help you accelerate time to insight. Try Elastic ...

## Best regards, Radu

[http://sematext.com/](http://sematext.com/) -- Elasticsearch -- Solr -- Lucene

On Thu, Oct 25, 2012 at 10:05 AM, Jerry Chou [fishlet0528@gmail.com](mailto:fishlet0528@gmail.com) wrote:

> For each cluster, I know that I can scale out query capacity by adding  
> replica nodes.  
> But is it possible that I can scale out in data volume?
> 
> As I know all data are written into one of the primary shards first, then  
> copied into replicas.  
> However the number of primary shards has to be defined at the beginning.  
> Doesn't that  
> pose a limitation on the max number of instances(for primary shards) in the  
> cluster?
> 
> The default number of primary shard is 5. For future scalability, is there  
> any drawback if  
> I set it to a big number?

--

---

<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:** [October 28, 2012, 5:56pm UTC](https://discuss.elastic.co/t/about-scalability-in-data-volume/9469/3 "2012-10-28T17:56:46Z")

</div>

I add to Radu's answer that if you have many shards on a single node, you can hit the "too many open files" issue. Each shard is a full Lucene instance.

--  
David 😉  
Twitter : @dadoonet / @elasticsearchfr / @scrutmydocs

Le 25 oct. 2012 à 11:46, Radu Gheorghe [radu.gheorghe@sematext.com](mailto:radu.gheorghe@sematext.com) a écrit :

Hello Jerry,

The problem with having many shards is that each shard has an  
overhead. So while it's OK to over-shard in order to scale out, if you  
exaggerate things will get slow, especially on the query side.

Luckily, there are other solutions than starting off with a huge  
number of shards. It depends on your data to choose what fits best.  
For example, you can add more indices as you go along.

I suggest you take a look at this video, as it's exactly about this topic:

> **[Elastic — The Search AI Company](https://www.elastic.co)**
>
> Power insights and outcomes with The Elastic Search AI Platform. See into your data and find answers that matter with enterprise solutions designed to help you accelerate time to insight. Try Elastic ...

## Best regards, Radu

[http://sematext.com/](http://sematext.com/) -- Elasticsearch -- Solr -- Lucene

On Thu, Oct 25, 2012 at 10:05 AM, Jerry Chou [fishlet0528@gmail.com](mailto:fishlet0528@gmail.com) wrote:

> For each cluster, I know that I can scale out query capacity by adding  
> replica nodes.  
> But is it possible that I can scale out in data volume?
> 
> As I know all data are written into one of the primary shards first, then  
> copied into replicas.  
> However the number of primary shards has to be defined at the beginning.  
> Doesn't that  
> pose a limitation on the max number of instances(for primary shards) in the  
> cluster?
> 
> The default number of primary shard is 5. For future scalability, is there  
> any drawback if  
> I set it to a big number?

--

--

---

<div class="post-metadata">

**Author:** ![BillyEm](https://avatars.discourse-cdn.com/v4/letter/b/c4cdca/32.png) [@BillyEm](https://discuss.elastic.co/u/BillyEm)\
**Post date:** [October 29, 2012, 2:46am UTC](https://discuss.elastic.co/t/about-scalability-in-data-volume/9469/4 "2012-10-29T02:46:26Z")

</div>

Radu: It'd be nice if "overhead" were defined. Since resources are almost  
always in a cost-benefit relationship in a datacenter.

Jerry: as to sharding out before you need it. Take a look at routing, and  
installing your own hash algorithm. Remember Computer 101, about why there  
isn''t just 1 implementation of hash. You could use the two to distribute  
your data onto only the resources benefiting you. The PLAN to scale out via  
the same algorithmic relationships would have to be done as part of phase

1. Complexity changes orders of mag ... right if you don't. . As always,  
knowing your content, if you can, (sometimes not possible in the real world  
folks) will make a world of of difference.

g'luck.

On Thursday, October 25, 2012 3:05:51 AM UTC-4, Jerry Chou wrote:

> For each cluster, I know that I can scale out query capacity by adding  
> replica nodes.  
> But is it possible that I can scale out in data volume?
> 
> As I know all data are written into one of the primary shards first, then  
> copied into replicas.  
> However the number of primary shards has to be defined at the beginning.  
> Doesn't that  
> pose a limitation on the max number of instances(for primary shards) in  
> the cluster?
> 
> The default number of primary shard is 5. For future scalability, is there  
> any drawback if  
> I set it to a big number?

--

---

<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, 3:06am UTC](https://discuss.elastic.co/t/about-scalability-in-data-volume/9469/5 "2017-07-06T03:06:57Z")

</div>


