# On scaling

**URL:** <https://discuss.elastic.co/t/on-scaling/6048>\
**Category:** Elasticsearch\
**Created:** [December 2, 2011, 6:43am UTC](https://discuss.elastic.co/t/on-scaling/6048 "2011-12-02T06:43:15Z")\
**Posts on this page:** 11\
**Page:** 1

<div class="post-metadata">

**Author:** ![BARNEY](https://avatars.discourse-cdn.com/v4/letter/b/48db29/32.png) [@BARNEY](https://discuss.elastic.co/u/BARNEY)\
**Post date:** [December 2, 2011, 6:43am UTC](https://discuss.elastic.co/t/on-scaling/6048/1 "2011-12-02T06:43:15Z")

</div>

Let us assume we create index with 2 shards and 1 replica on two nodes  
in a elasticsearch box1.After indexing some documents,disk space(say  
1000GB) in the box1 will get filled.  
To index more documents, I will have 2 more nodes(one more box say  
box2 has disk space 1000GB).  
What is the mechanism to use total 4 nodes(two boxes) with only one  
index having memory 2000GB?

What I am thinking is like If we create 2 more primary shards to the  
box1 dynamically,then created shards will get allocated on the second  
box.To do this I want to know how to create shards dynamically?

Apart from creating dynamic shards,any other mechanism to scale?

---

<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:** [December 2, 2011, 6:51am UTC](https://discuss.elastic.co/t/on-scaling/6048/2 "2011-12-02T06:51:18Z")

</div>

You can't modify shards after index creation.

HTH  
David 😉  
@dadoonet

Le 2 déc. 2011 à 07:43, BARNEY [kalyanc007@gmail.com](mailto:kalyanc007@gmail.com) a écrit :

> Let us assume we create index with 2 shards and 1 replica on two nodes  
> in a elasticsearch box1.After indexing some documents,disk space(say  
> 1000GB) in the box1 will get filled.  
> To index more documents, I will have 2 more nodes(one more box say  
> box2 has disk space 1000GB).  
> What is the mechanism to use total 4 nodes(two boxes) with only one  
> index having memory 2000GB?
> 
> What I am thinking is like If we create 2 more primary shards to the  
> box1 dynamically,then created shards will get allocated on the second  
> box.To do this I want to know how to create shards dynamically?
> 
> Apart from creating dynamic shards,any other mechanism to scale?

---

<div class="post-metadata">

**Author:** ![drewr](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/drewr/32/7803_2.png) [@drewr](https://discuss.elastic.co/u/drewr)\
**Post date:** [December 2, 2011, 4:32pm UTC](https://discuss.elastic.co/t/on-scaling/6048/3 "2011-12-02T16:32:28Z")

</div>

BARNEY wrote:

> Let us assume we create index with 2 shards and 1 replica on two  
> nodes in a elasticsearch box1.After indexing some documents,disk  
> space(say 1000GB) in the box1 will get filled. To index more  
> documents, I will have 2 more nodes(one more box say box2 has disk  
> space 1000GB). What is the mechanism to use total 4 nodes(two  
> boxes) with only one index having memory 2000GB?
> 
> What I am thinking is like If we create 2 more primary shards to  
> the box1 dynamically,then created shards will get allocated on the  
> second box.To do this I want to know how to create shards  
> dynamically?
> 
> Apart from creating dynamic shards,any other mechanism to scale?

As David mentioned, shards are not configurable after index  
creation. You have a few options, and you may want to do all of them  
in time.

The first is to create your index with more shards. This will give  
you more horizontal node growth. You can't store a 10TiB index on  
100 1TiB nodes if you only have 2 shards, because each shard would be  
5TiB. So, calculate a reasonable number given the size of the  
resources you have available.

The next tool at your disposal when you've exhausted the previous  
step's capacity are aliases. You'll want to roll over into a new  
index and create an alias that points to the old and new index. You  
don't even have to do this step if your older data doesn't need to be  
queried seamlessly. You could either disable it or leave it to be  
queried manually, but an alias allows you to query both indices with  
no client knowledge that they're separate. You can also set up alias  
filters to enforce logical partitions of data.

At some point you'll exhaust the limits of that cluster. The only  
option then will be to create another cluster and either migrate data  
to it or use both clusters together. We do the latter through load  
balancers with great success, managing 20 clusters of 1PiB of data.

-Drew

---

<div class="post-metadata">

**Author:** ![Michael\_Sick](https://avatars.discourse-cdn.com/v4/letter/m/22d042/32.png) [@Michael\_Sick](https://discuss.elastic.co/u/Michael_Sick)\
**Post date:** [December 2, 2011, 4:59pm UTC](https://discuss.elastic.co/t/on-scaling/6048/4 "2011-12-02T16:59:05Z")

</div>

Drew,

1PB! That's the biggest ES cluster I've seen referenced to date. Maybe you  
could post some experiences/lessons learned to the "Please, tell about the  
success story about ES usage on production" thread?

--Mike

On Fri, Dec 2, 2011 at 11:32 AM, Drew Raines [aaraines@gmail.com](mailto:aaraines@gmail.com) wrote:

> BARNEY wrote:
> 
> > Let us assume we create index with 2 shards and 1 replica on two  
> > nodes in a elasticsearch box1.After indexing some documents,disk  
> > space(say 1000GB) in the box1 will get filled. To index more  
> > documents, I will have 2 more nodes(one more box say box2 has disk  
> > space 1000GB). What is the mechanism to use total 4 nodes(two  
> > boxes) with only one index having memory 2000GB?
> > 
> > What I am thinking is like If we create 2 more primary shards to  
> > the box1 dynamically,then created shards will get allocated on the  
> > second box.To do this I want to know how to create shards  
> > dynamically?
> > 
> > Apart from creating dynamic shards,any other mechanism to scale?
> 
> As David mentioned, shards are not configurable after index  
> creation. You have a few options, and you may want to do all of them  
> in time.
> 
> The first is to create your index with more shards. This will give  
> you more horizontal node growth. You can't store a 10TiB index on  
> 100 1TiB nodes if you only have 2 shards, because each shard would be  
> 5TiB. So, calculate a reasonable number given the size of the  
> resources you have available.
> 
> The next tool at your disposal when you've exhausted the previous  
> step's capacity are aliases. You'll want to roll over into a new  
> index and create an alias that points to the old and new index. You  
> don't even have to do this step if your older data doesn't need to be  
> queried seamlessly. You could either disable it or leave it to be  
> queried manually, but an alias allows you to query both indices with  
> no client knowledge that they're separate. You can also set up alias  
> filters to enforce logical partitions of data.
> 
> At some point you'll exhaust the limits of that cluster. The only  
> option then will be to create another cluster and either migrate data  
> to it or use both clusters together. We do the latter through load  
> balancers with great success, managing 20 clusters of 1PiB of data.
> 
> -Drew

---

<div class="post-metadata">

**Author:** ![BARNEY](https://avatars.discourse-cdn.com/v4/letter/b/48db29/32.png) [@BARNEY](https://discuss.elastic.co/u/BARNEY)\
**Post date:** [December 5, 2011, 6:08am UTC](https://discuss.elastic.co/t/on-scaling/6048/5 "2011-12-05T06:08:48Z")

</div>

Thanks Drew.  
I have one more option.I will create one more index with 4 shards and  
1 replica in the same cluster.Then I migrate old one(2shards  
1replica) to newly created index(4shards 1replica).  
which one is more practical and give good performance?

On Dec 2, 9:32 pm, Drew Raines [aarai...@gmail.com](mailto:aarai...@gmail.com) wrote:

> BARNEYwrote:
> 
> > Let us assume we create index with 2 shards and 1 replica on two  
> > nodes in a elasticsearch box1.After indexing some documents,disk  
> > space(say 1000GB) in the box1 will get filled. To index more  
> > documents, I will have 2 more nodes(one more box say box2 has disk  
> > space 1000GB). What is the mechanism to use total 4 nodes(two  
> > boxes) with only one index having memory 2000GB?
> 
> > What I am thinking is like If we create 2 more primary shards to  
> > the box1 dynamically,then created shards will get allocated on the  
> > second box.To do this I want to know how to create shards  
> > dynamically?
> 
> > Apart from creating dynamic shards,any other mechanism to scale?
> 
> As David mentioned, shards are not configurable after index  
> creation. You have a few options, and you may want to do all of them  
> in time.
> 
> The first is to create your index with more shards. This will give  
> you more horizontal node growth. You can't store a 10TiB index on  
> 100 1TiB nodes if you only have 2 shards, because each shard would be  
> 5TiB. So, calculate a reasonable number given the size of the  
> resources you have available.
> 
> The next tool at your disposal when you've exhausted the previous  
> step's capacity are aliases. You'll want to roll over into a new  
> index and create an alias that points to the old and new index. You  
> don't even have to do this step if your older data doesn't need to be  
> queried seamlessly. You could either disable it or leave it to be  
> queried manually, but an alias allows you to query both indices with  
> no client knowledge that they're separate. You can also set up alias  
> filters to enforce logical partitions of data.
> 
> At some point you'll exhaust the limits of that cluster. The only  
> option then will be to create another cluster and either migrate data  
> to it or use both clusters together. We do the latter through load  
> balancers with great success, managing 20 clusters of 1PiB of data.
> 
> -Drew

---

<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:** [December 5, 2011, 7:25pm UTC](https://discuss.elastic.co/t/on-scaling/6048/6 "2011-12-05T19:25:14Z")

</div>

Which one compared to what? The simplest answer is over-allocate shards,  
for a cluster with 4 nodes, create an index with 6-10 shards with 1  
replica, this will allow you to grow to 6-10 nodes index size wise, and  
6-10 \* (number\_of\_replicas + 1) nodes. (just for that index, if you have  
more indices, obviously you can add more nodes before you get to a "limit"  
of 1 shard per node).

On Mon, Dec 5, 2011 at 8:08 AM, BARNEY [kalyanc007@gmail.com](mailto:kalyanc007@gmail.com) wrote:

> Thanks Drew.  
> I have one more option.I will create one more index with 4 shards and  
> 1 replica in the same cluster.Then I migrate old one(2shards  
> 1replica) to newly created index(4shards 1replica).  
> which one is more practical and give good performance?
> 
> On Dec 2, 9:32 pm, Drew Raines [aarai...@gmail.com](mailto:aarai...@gmail.com) wrote:
> 
> > BARNEYwrote:
> > 
> > > Let us assume we create index with 2 shards and 1 replica on two  
> > > nodes in a elasticsearch box1.After indexing some documents,disk  
> > > space(say 1000GB) in the box1 will get filled. To index more  
> > > documents, I will have 2 more nodes(one more box say box2 has disk  
> > > space 1000GB). What is the mechanism to use total 4 nodes(two  
> > > boxes) with only one index having memory 2000GB?
> > 
> > > What I am thinking is like If we create 2 more primary shards to  
> > > the box1 dynamically,then created shards will get allocated on the  
> > > second box.To do this I want to know how to create shards  
> > > dynamically?
> > 
> > > Apart from creating dynamic shards,any other mechanism to scale?
> > 
> > As David mentioned, shards are not configurable after index  
> > creation. You have a few options, and you may want to do all of them  
> > in time.
> > 
> > The first is to create your index with more shards. This will give  
> > you more horizontal node growth. You can't store a 10TiB index on  
> > 100 1TiB nodes if you only have 2 shards, because each shard would be  
> > 5TiB. So, calculate a reasonable number given the size of the  
> > resources you have available.
> > 
> > The next tool at your disposal when you've exhausted the previous  
> > step's capacity are aliases. You'll want to roll over into a new  
> > index and create an alias that points to the old and new index. You  
> > don't even have to do this step if your older data doesn't need to be  
> > queried seamlessly. You could either disable it or leave it to be  
> > queried manually, but an alias allows you to query both indices with  
> > no client knowledge that they're separate. You can also set up alias  
> > filters to enforce logical partitions of data.
> > 
> > At some point you'll exhaust the limits of that cluster. The only  
> > option then will be to create another cluster and either migrate data  
> > to it or use both clusters together. We do the latter through load  
> > balancers with great success, managing 20 clusters of 1PiB of data.
> > 
> > -Drew

---

<div class="post-metadata">

**Author:** ![Michael\_Sick](https://avatars.discourse-cdn.com/v4/letter/m/22d042/32.png) [@Michael\_Sick](https://discuss.elastic.co/u/Michael_Sick)\
**Post date:** [December 5, 2011, 7:46pm UTC](https://discuss.elastic.co/t/on-scaling/6048/7 "2011-12-05T19:46:54Z")

</div>

Shay,

Any advice on what the min/max # of shards/node would be?

On Mon, Dec 5, 2011 at 2:25 PM, Shay Banon [kimchy@gmail.com](mailto:kimchy@gmail.com) wrote:

> Which one compared to what? The simplest answer is over-allocate shards,  
> for a cluster with 4 nodes, create an index with 6-10 shards with 1  
> replica, this will allow you to grow to 6-10 nodes index size wise, and  
> 6-10 \* (number\_of\_replicas + 1) nodes. (just for that index, if you have  
> more indices, obviously you can add more nodes before you get to a "limit"  
> of 1 shard per node).
> 
> On Mon, Dec 5, 2011 at 8:08 AM, BARNEY [kalyanc007@gmail.com](mailto:kalyanc007@gmail.com) wrote:
> 
> > Thanks Drew.  
> > I have one more option.I will create one more index with 4 shards and  
> > 1 replica in the same cluster.Then I migrate old one(2shards  
> > 1replica) to newly created index(4shards 1replica).  
> > which one is more practical and give good performance?
> > 
> > On Dec 2, 9:32 pm, Drew Raines [aarai...@gmail.com](mailto:aarai...@gmail.com) wrote:
> > 
> > > BARNEYwrote:
> > > 
> > > > Let us assume we create index with 2 shards and 1 replica on two  
> > > > nodes in a elasticsearch box1.After indexing some documents,disk  
> > > > space(say 1000GB) in the box1 will get filled. To index more  
> > > > documents, I will have 2 more nodes(one more box say box2 has disk  
> > > > space 1000GB). What is the mechanism to use total 4 nodes(two  
> > > > boxes) with only one index having memory 2000GB?
> > > 
> > > > What I am thinking is like If we create 2 more primary shards to  
> > > > the box1 dynamically,then created shards will get allocated on the  
> > > > second box.To do this I want to know how to create shards  
> > > > dynamically?
> > > 
> > > > Apart from creating dynamic shards,any other mechanism to scale?
> > > 
> > > As David mentioned, shards are not configurable after index  
> > > creation. You have a few options, and you may want to do all of them  
> > > in time.
> > > 
> > > The first is to create your index with more shards. This will give  
> > > you more horizontal node growth. You can't store a 10TiB index on  
> > > 100 1TiB nodes if you only have 2 shards, because each shard would be  
> > > 5TiB. So, calculate a reasonable number given the size of the  
> > > resources you have available.
> > > 
> > > The next tool at your disposal when you've exhausted the previous  
> > > step's capacity are aliases. You'll want to roll over into a new  
> > > index and create an alias that points to the old and new index. You  
> > > don't even have to do this step if your older data doesn't need to be  
> > > queried seamlessly. You could either disable it or leave it to be  
> > > queried manually, but an alias allows you to query both indices with  
> > > no client knowledge that they're separate. You can also set up alias  
> > > filters to enforce logical partitions of data.
> > > 
> > > At some point you'll exhaust the limits of that cluster. The only  
> > > option then will be to create another cluster and either migrate data  
> > > to it or use both clusters together. We do the latter through load  
> > > balancers with great success, managing 20 clusters of 1PiB of data.
> > > 
> > > -Drew

---

<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:** [December 7, 2011, 3:34pm UTC](https://discuss.elastic.co/t/on-scaling/6048/8 "2011-12-07T15:34:00Z")

</div>

Hard to say, really depends on the environment and data you have. Note  
though, if you have a single index with 100 shards (for example), then it  
means a distributed search across 100 shards (unless you use routing).

On Mon, Dec 5, 2011 at 9:46 PM, Michael Sick \<  
[michael.sick@serenesoftware.com](mailto:michael.sick@serenesoftware.com)\> wrote:

> Shay,
> 
> Any advice on what the min/max # of shards/node would be?
> 
> On Mon, Dec 5, 2011 at 2:25 PM, Shay Banon [kimchy@gmail.com](mailto:kimchy@gmail.com) wrote:
> 
> > Which one compared to what? The simplest answer is over-allocate shards,  
> > for a cluster with 4 nodes, create an index with 6-10 shards with 1  
> > replica, this will allow you to grow to 6-10 nodes index size wise, and  
> > 6-10 \* (number\_of\_replicas + 1) nodes. (just for that index, if you have  
> > more indices, obviously you can add more nodes before you get to a "limit"  
> > of 1 shard per node).
> > 
> > On Mon, Dec 5, 2011 at 8:08 AM, BARNEY [kalyanc007@gmail.com](mailto:kalyanc007@gmail.com) wrote:
> > 
> > > Thanks Drew.  
> > > I have one more option.I will create one more index with 4 shards and  
> > > 1 replica in the same cluster.Then I migrate old one(2shards  
> > > 1replica) to newly created index(4shards 1replica).  
> > > which one is more practical and give good performance?
> > > 
> > > On Dec 2, 9:32 pm, Drew Raines [aarai...@gmail.com](mailto:aarai...@gmail.com) wrote:
> > > 
> > > > BARNEYwrote:
> > > > 
> > > > > Let us assume we create index with 2 shards and 1 replica on two  
> > > > > nodes in a elasticsearch box1.After indexing some documents,disk  
> > > > > space(say 1000GB) in the box1 will get filled. To index more  
> > > > > documents, I will have 2 more nodes(one more box say box2 has disk  
> > > > > space 1000GB). What is the mechanism to use total 4 nodes(two  
> > > > > boxes) with only one index having memory 2000GB?
> > > > 
> > > > > What I am thinking is like If we create 2 more primary shards to  
> > > > > the box1 dynamically,then created shards will get allocated on the  
> > > > > second box.To do this I want to know how to create shards  
> > > > > dynamically?
> > > > 
> > > > > Apart from creating dynamic shards,any other mechanism to scale?
> > > > 
> > > > As David mentioned, shards are not configurable after index  
> > > > creation. You have a few options, and you may want to do all of them  
> > > > in time.
> > > > 
> > > > The first is to create your index with more shards. This will give  
> > > > you more horizontal node growth. You can't store a 10TiB index on  
> > > > 100 1TiB nodes if you only have 2 shards, because each shard would be  
> > > > 5TiB. So, calculate a reasonable number given the size of the  
> > > > resources you have available.
> > > > 
> > > > The next tool at your disposal when you've exhausted the previous  
> > > > step's capacity are aliases. You'll want to roll over into a new  
> > > > index and create an alias that points to the old and new index. You  
> > > > don't even have to do this step if your older data doesn't need to be  
> > > > queried seamlessly. You could either disable it or leave it to be  
> > > > queried manually, but an alias allows you to query both indices with  
> > > > no client knowledge that they're separate. You can also set up alias  
> > > > filters to enforce logical partitions of data.
> > > > 
> > > > At some point you'll exhaust the limits of that cluster. The only  
> > > > option then will be to create another cluster and either migrate data  
> > > > to it or use both clusters together. We do the latter through load  
> > > > balancers with great success, managing 20 clusters of 1PiB of data.
> > > > 
> > > > -Drew

---

<div class="post-metadata">

**Author:** ![drewr](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/drewr/32/7803_2.png) [@drewr](https://discuss.elastic.co/u/drewr)\
**Post date:** [December 7, 2011, 8:22pm UTC](https://discuss.elastic.co/t/on-scaling/6048/9 "2011-12-07T20:22:09Z")

</div>

Michael Sick wrote:

> Any advice on what the min/max # of shards/node would be?

Depends on how much searching you're doing. On a non-busy cluster  
we've gotten beyond 250 shards/node on AWS m1.xlarge (15GiB RAM).  
Avg size of those shards is 100-200GiB.

-Drew

---

<div class="post-metadata">

**Author:** ![drewr](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/drewr/32/7803_2.png) [@drewr](https://discuss.elastic.co/u/drewr)\
**Post date:** [January 3, 2012, 5:47pm UTC](https://discuss.elastic.co/t/on-scaling/6048/10 "2012-01-03T17:47:37Z")

</div>

Michael Sick wrote:

> 1PB! That's the biggest ES cluster I've seen referenced to  
> date. Maybe you could post some experiences/lessons learned to the  
> "Please, tell about the success story about ES usage on production"  
> thread?

I haven't had the time to properly contribute to that thread, but  
just to clarify: we don't have a PB in a single cluster. It's made  
up of about 15 clusters ranging from 10-150TB.

I don't have any doubt that ES can scale to a PB though. It's  
extremely horizontal given enough shards. We just choose not to  
architect that way. You run into non-ES issues at that scale (and  
well before it). 🙂

-Drew

---

<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:43am UTC](https://discuss.elastic.co/t/on-scaling/6048/11 "2017-07-06T03:43:57Z")

</div>


