# Cluster questions

**URL:** <https://discuss.elastic.co/t/cluster-questions/3020>\
**Category:** Elasticsearch\
**Created:** [June 15, 2010, 6:46am UTC](https://discuss.elastic.co/t/cluster-questions/3020 "2010-06-15T06:46:57Z")\
**Posts on this page:** 8\
**Page:** 1

<div class="post-metadata">

**Author:** ![otisg](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/otisg/32/492_2.png) [@otisg](https://discuss.elastic.co/u/otisg)\
**Post date:** [June 15, 2010, 6:46am UTC](https://discuss.elastic.co/t/cluster-questions/3020/1 "2010-06-15T06:46:57Z")

</div>

Hi,

In ES, what controls which node/shard a doc will get indexed on?

What happens (or what does one need to do) when the search cluster is  
expanded? Is there a notion of (automatic) rebalancing?

Similarly, what happens or should be done when a node in a cluster  
goes down? Is there something that automatically replicates data that  
disappeared when that node went down?

Thanks,  
Otis

---

<div class="post-metadata">

**Author:** ![Lukas\_Vlcek1](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/lukas_vlcek1/32/819_2.png) [@Lukas\_Vlcek1](https://discuss.elastic.co/u/Lukas_Vlcek1)\
**Post date:** [June 15, 2010, 7:33am UTC](https://discuss.elastic.co/t/cluster-questions/3020/2 "2010-06-15T07:33:39Z")

</div>

Otis,

I am sure most of these questions are addressed in Elasticsearch  
presentation from Berlin:

> **[ElasticSearch at berlinbuzzwords 2010](https://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwords-2010)**
>
> ElasticSearch at berlinbuzzwords 2010 - Download as a PDF or view online for free

[http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwords-2010](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwords-2010)Automatic  
shard allocation - pages 90-108.

Also it is very easy to verify yourself, just change logging level to DEBUG  
for rootLogger in logging.yml. Then start one node and index some data. Then  
start second node and see log files. Once the second node is available some  
index shards are allocated to the new node. By default ES uses 5 shards with  
1 replica for each shard. If node goes down then you can use health API to  
see if you still have all data available for search (  
[http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluster/health/](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluster/health/)).  
It then depends on the speed in which particular nodes go down (crash, not  
regular shutdown) because if there is only one shard left and no replica is  
available then it should replicate itself to some other node (providing  
replica is set to 1 or more).

Regards,  
Lukas

On Tue, Jun 15, 2010 at 8:46 AM, Otis [otis.gospodnetic@gmail.com](mailto:otis.gospodnetic@gmail.com) wrote:

> Hi,
> 
> In ES, what controls which node/shard a doc will get indexed on?
> 
> What happens (or what does one need to do) when the search cluster is  
> expanded? Is there a notion of (automatic) rebalancing?
> 
> Similarly, what happens or should be done when a node in a cluster  
> goes down? Is there something that automatically replicates data that  
> disappeared when that node went down?
> 
> Thanks,  
> Otis

---

<div class="post-metadata">

**Author:** ![otisg](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/otisg/32/492_2.png) [@otisg](https://discuss.elastic.co/u/otisg)\
**Post date:** [June 18, 2010, 9:31pm UTC](https://discuss.elastic.co/t/cluster-questions/3020/3 "2010-06-18T21:31:20Z")

</div>

Thanks Lukas,

I looked at this the other day, but I don't think that answers my Qs,  
or at least I can't tell without hearing Shay's accompanying talk. 🙂

So, questions:

- Slide 94: cluster with settings: replicas = 1, shards = 2, and a  
node with 2 shards, OK

- Slide 95: 2nd node appears and end sup with the same shards as on  
node 1.  
-- Q: doesn't this mean that shards were replicated? Why, if  
number\_of\_replicas=1 ?

- Slide 96: 3rd and 4th nodes appear and after that each node in the  
cluster ends up with just 1 shard. This makes sense.  
-- Q: This happens automagically?  
-- Q: Does ES simply _copy_ the index/shard from one node to the  
other when it detects more nodes joined the cluster?

- Slide 107: a new index with 2 shards is added (shards=1,  
replicas=1), but 2 nodes get that new index, even though replicas=1.  
Huh?

So the above questions are really about understanding why an index/  
shard is being replicated even when replicas=1.

I think the above may also answer what happens when the cluster is  
expanded: ES detects new nodes joining and figures out that they can  
handle/host some of the indices or index replicas and somehow send  
them to the new nodes. I assume there are mechanisms in ES to let it  
spread the data evenly, though, I am guessing, it doesn't take into  
account query rates, so it is not yet able to migrate hot indices  
around automatically?

I think the above doesn't cover what happens when 1 or more nodes go  
down: I know ES detects that, but does it know which indices were on  
those nodes and does it automatically create more replicas for those  
indices in order to satisfy the replicas=X setting?

Thanks,  
Otis

On Jun 15, 3:33 am, Lukáš Vlček [lukas.vl...@gmail.com](mailto:lukas.vl...@gmail.com) wrote:

> Otis,
> 
> I am sure most of these questions are addressed in Elasticsearch  
> presentation from Berlin:[http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)...  
> [http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo...](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo...)Automatic  
> shard allocation - pages 90-108.
> 
> Also it is very easy to verify yourself, just change logging level to DEBUG  
> for rootLogger in logging.yml. Then start one node and index some data. Then  
> start second node and see log files. Once the second node is available some  
> index shards are allocated to the new node. By default ES uses 5 shards with  
> 1 replica for each shard. If node goes down then you can use health API to  
> see if you still have all data available for search ([http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...).  
> It then depends on the speed in which particular nodes go down (crash, not  
> regular shutdown) because if there is only one shard left and no replica is  
> available then it should replicate itself to some other node (providing  
> replica is set to 1 or more).
> 
> Regards,  
> Lukas
> 
> On Tue, Jun 15, 2010 at 8:46 AM, Otis [otis.gospodne...@gmail.com](mailto:otis.gospodne...@gmail.com) wrote:
> 
> > Hi,
> 
> > In ES, what controls which node/shard a doc will get indexed on?
> 
> > What happens (or what does one need to do) when the search cluster is  
> > expanded? Is there a notion of (automatic) rebalancing?
> 
> > Similarly, what happens or should be done when a node in a cluster  
> > goes down? Is there something that automatically replicates data that  
> > disappeared when that node went down?
> 
> > Thanks,  
> > Otis

---

<div class="post-metadata">

**Author:** ![Berkay\_Mollamustafao](https://avatars.discourse-cdn.com/v4/letter/b/22d042/32.png) [@Berkay\_Mollamustafao](https://discuss.elastic.co/u/Berkay_Mollamustafao)\
**Post date:** [June 18, 2010, 9:39pm UTC](https://discuss.elastic.co/t/cluster-questions/3020/4 "2010-06-18T21:39:25Z")

</div>

Otis,

number\_of\_replicas=1 means each shard has 1 replica, meaning there are 2  
copies of each shard.  
You seem to take number of replicas as number of copies. not exactly the  
same thing.

Regards,  
Berkay Mollamustafaoglu  
mberkay on yahoo, google and skype

On Fri, Jun 18, 2010 at 5:31 PM, Otis [otis.gospodnetic@gmail.com](mailto:otis.gospodnetic@gmail.com) wrote:

> Thanks Lukas,
> 
> I looked at this the other day, but I don't think that answers my Qs,  
> or at least I can't tell without hearing Shay's accompanying talk. 🙂
> 
> So, questions:
> 
> - Slide 94: cluster with settings: replicas = 1, shards = 2, and a  
> node with 2 shards, OK
> 
> - Slide 95: 2nd node appears and end sup with the same shards as on  
> node 1.  
> -- Q: doesn't this mean that shards were replicated? Why, if  
> number\_of\_replicas=1 ?
> 
> - Slide 96: 3rd and 4th nodes appear and after that each node in the  
> cluster ends up with just 1 shard. This makes sense.  
> -- Q: This happens automagically?  
> -- Q: Does ES simply _copy_ the index/shard from one node to the  
> other when it detects more nodes joined the cluster?
> 
> - Slide 107: a new index with 2 shards is added (shards=1,  
> replicas=1), but 2 nodes get that new index, even though replicas=1.  
> Huh?
> 
> So the above questions are really about understanding why an index/  
> shard is being replicated even when replicas=1.
> 
> I think the above may also answer what happens when the cluster is  
> expanded: ES detects new nodes joining and figures out that they can  
> handle/host some of the indices or index replicas and somehow send  
> them to the new nodes. I assume there are mechanisms in ES to let it  
> spread the data evenly, though, I am guessing, it doesn't take into  
> account query rates, so it is not yet able to migrate hot indices  
> around automatically?
> 
> I think the above doesn't cover what happens when 1 or more nodes go  
> down: I know ES detects that, but does it know which indices were on  
> those nodes and does it automatically create more replicas for those  
> indices in order to satisfy the replicas=X setting?
> 
> Thanks,  
> Otis
> 
> On Jun 15, 3:33 am, Lukáš Vlček [lukas.vl...@gmail.com](mailto:lukas.vl...@gmail.com) wrote:
> 
> > Otis,
> > 
> > I am sure most of these questions are addressed in Elasticsearch  
> > presentation from Berlin:  
> > [http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)...  
> > \<[http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)..  
> > .\>Automatic  
> > shard allocation - pages 90-108.
> > 
> > Also it is very easy to verify yourself, just change logging level to  
> > DEBUG  
> > for rootLogger in logging.yml. Then start one node and index some data.  
> > Then  
> > start second node and see log files. Once the second node is available  
> > some  
> > index shards are allocated to the new node. By default ES uses 5 shards  
> > with  
> > 1 replica for each shard. If node goes down then you can use health API  
> > to  
> > see if you still have all data available for search (  
> > [http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...).  
> > It then depends on the speed in which particular nodes go down (crash,  
> > not  
> > regular shutdown) because if there is only one shard left and no replica  
> > is  
> > available then it should replicate itself to some other node (providing  
> > replica is set to 1 or more).
> > 
> > Regards,  
> > Lukas
> > 
> > On Tue, Jun 15, 2010 at 8:46 AM, Otis [otis.gospodne...@gmail.com](mailto:otis.gospodne...@gmail.com)  
> > wrote:
> > 
> > > Hi,
> > 
> > > In ES, what controls which node/shard a doc will get indexed on?
> > 
> > > What happens (or what does one need to do) when the search cluster is  
> > > expanded? Is there a notion of (automatic) rebalancing?
> > 
> > > Similarly, what happens or should be done when a node in a cluster  
> > > goes down? Is there something that automatically replicates data that  
> > > disappeared when that node went down?
> > 
> > > Thanks,  
> > > Otis

---

<div class="post-metadata">

**Author:** ![Lukas\_Vlcek1](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/lukas_vlcek1/32/819_2.png) [@Lukas\_Vlcek1](https://discuss.elastic.co/u/Lukas_Vlcek1)\
**Post date:** [June 19, 2010, 6:08am UTC](https://discuss.elastic.co/t/cluster-questions/3020/5 "2010-06-19T06:08:24Z")

</div>

Hi Otis,

I know that Shay would be able to shed more light on this... but let us do  
our homework now (since many of your questions can be answered by running a  
test).

1. 

One node A is started and a new index is created (2 shards, 1 replica)  
curl -XPUT '[http://localhost:9200/twitter/](http://localhost:9200/twitter/)' -d '  
index :  
number\_of\_shards : 2  
number\_of\_replicas : 1  
'  
Now we can see node A has two shards allocated on it.  
The cluster status is YELLOW.  
(Cluster is not healthy but still all data is available for search.)

1. 

Second node B is started.  
Node B has index replicas on it. (because number\_of\_replicas=1)  
Cluster status is GREEN.  
(At tis point your cluster is healthy)

1. 

Third node C is started.  
One primary shard from node A is moved to C.  
So now index primary shards are located on A and C. Both replicas on B.  
Cluster status is GREEN.

1. 

Fourth node D is started.  
The node D got one of replicas from B.  
Cluster status is GREEN.

1. 

Now you can shutdown one node at a time until two nodes are left.  
Cluster status is GREEN.

1. Shutdown one node. Only one is left.  
Two shards and no replicas are on the last remaining node.  
Cluster status is YELLOW.

So as you can see if replicas=X setting can not be satisfied then you can  
learn this from cluster health status (  
[http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluster/health/](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluster/health/)  
).

Also see below for my other comments.

On Fri, Jun 18, 2010 at 11:31 PM, Otis [otis.gospodnetic@gmail.com](mailto:otis.gospodnetic@gmail.com) wrote:

> Thanks Lukas,
> 
> I looked at this the other day, but I don't think that answers my Qs,  
> or at least I can't tell without hearing Shay's accompanying talk. 🙂
> 
> So, questions:
> 
> - Slide 94: cluster with settings: replicas = 1, shards = 2, and a  
> node with 2 shards, OK
> 
> - Slide 95: 2nd node appears and end sup with the same shards as on  
> node 1.  
> -- Q: doesn't this mean that shards were replicated? Why, if  
> number\_of\_replicas=1 ?

This is obvious. number\_or\_replicas = 1 means there is one (non-primary)  
shard for each (primary) shard.  
If you had created index with 0 replicas and two shards then you would have  
got GREEN status in the first step (meaning replicas=X setting is satisfied  
).

> - Slide 96: 3rd and 4th nodes appear and after that each node in the  
> cluster ends up with just 1 shard. This makes sense.  
> -- Q: This happens automagically?

As you can see from the test above this is happening automagically 🙂

> -- Q: Does ES simply _copy_ the index/shard from one node to the  
> other when it detects more nodes joined the cluster?

I think it can take it from gateway but I will let Shay to provide more  
details.

> - Slide 107: a new index with 2 shards is added (shards=1,  
> replicas=1), but 2 nodes get that new index, even though replicas=1.  
> Huh?

This is correct. Isn't it? One primary shard on one node and the second node  
gets replica of that shard (if it wouldn't be possible to locate second node  
for replica then you would be able to learn this from health status).

> 

> So the above questions are really about understanding why an index/  
> shard is being replicated even when replicas=1.
> 
> I think the above may also answer what happens when the cluster is  
> expanded: ES detects new nodes joining and figures out that they can  
> handle/host some of the indices or index replicas and somehow send  
> them to the new nodes. I assume there are mechanisms in ES to let it  
> spread the data evenly, though, I am guessing, it doesn't take into  
> account query rates, so it is not yet able to migrate hot indices  
> around automatically?

I think this is not implemented yet. However, it seems that there is a lot  
of data available and implementing specific traffic load handlers would be  
possible (I would be surprised it that is not planed).

> I think the above doesn't cover what happens when 1 or more nodes go  
> down: I know ES detects that, but does it know which indices were on  
> those nodes and does it automatically create more replicas for those  
> indices in order to satisfy the replicas=X setting?

It is known which shards and its replicas are located on which nodes.  
See  
[http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluster/state/](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluster/state/)  
the  
"routing\_table.indices" part. (Note with current master you can expect to  
get more info from REST API)

> Thanks,  
> Otis
> 
> On Jun 15, 3:33 am, Lukáš Vlček [lukas.vl...@gmail.com](mailto:lukas.vl...@gmail.com) wrote:
> 
> > Otis,
> > 
> > I am sure most of these questions are addressed in Elasticsearch  
> > presentation from Berlin:  
> > [http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)...  
> > \<[http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)..  
> > .\>Automatic  
> > shard allocation - pages 90-108.
> > 
> > Also it is very easy to verify yourself, just change logging level to  
> > DEBUG  
> > for rootLogger in logging.yml. Then start one node and index some data.  
> > Then  
> > start second node and see log files. Once the second node is available  
> > some  
> > index shards are allocated to the new node. By default ES uses 5 shards  
> > with  
> > 1 replica for each shard. If node goes down then you can use health API  
> > to  
> > see if you still have all data available for search (  
> > [http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...).  
> > It then depends on the speed in which particular nodes go down (crash,  
> > not  
> > regular shutdown) because if there is only one shard left and no replica  
> > is  
> > available then it should replicate itself to some other node (providing  
> > replica is set to 1 or more).
> > 
> > Regards,  
> > Lukas
> > 
> > On Tue, Jun 15, 2010 at 8:46 AM, Otis [otis.gospodne...@gmail.com](mailto:otis.gospodne...@gmail.com)  
> > wrote:
> > 
> > > Hi,
> > 
> > > In ES, what controls which node/shard a doc will get indexed on?
> > 
> > > What happens (or what does one need to do) when the search cluster is  
> > > expanded? Is there a notion of (automatic) rebalancing?
> > 
> > > Similarly, what happens or should be done when a node in a cluster  
> > > goes down? Is there something that automatically replicates data that  
> > > disappeared when that node went down?
> > 
> > > Thanks,  
> > > Otis

---

<div class="post-metadata">

**Author:** ![otisg](https://sea2.discourse-cdn.com/elastic/user_avatar/discuss.elastic.co/otisg/32/492_2.png) [@otisg](https://discuss.elastic.co/u/otisg)\
**Post date:** [June 22, 2010, 4:48am UTC](https://discuss.elastic.co/t/cluster-questions/3020/6 "2010-06-22T04:48:06Z")

</div>

Thank you for the detailed reply, Lukáš!

The reason why I didn't get the number\_of\_replicas=X meaning is  
because in my mind "number of replicas" really means "number of  
identical copies", while in ES it means "the number of additional  
copies". Also, in HDFS there is a notion of a replication factor that  
matches "my" thinking: if replication factor is 3, that means there  
are 3 copies of each data block in HDFS. Isn't this the more common  
way of thinking about/counting replicas?

Oh, and you mentioned that maybe ES copies replicas from the gateway,  
but is gateway not an optional thing?

Thanks,  
Otis

On Jun 19, 2:08 am, Lukáš Vlček [lukas.vl...@gmail.com](mailto:lukas.vl...@gmail.com) wrote:

> Hi Otis,
> 
> I know that Shay would be able to shed more light on this... but let us do  
> our homework now (since many of your questions can be answered by running a  
> test).
> 
> 1. 
> 
> One node A is started and a new index is created (2 shards, 1 replica)  
> curl -XPUT '[http://localhost:9200/twitter/'-d](http://localhost:9200/twitter/'-d) '  
> index :  
> number\_of\_shards : 2  
> number\_of\_replicas : 1  
> '  
> Now we can see node A has two shards allocated on it.  
> The cluster status is YELLOW.  
> (Cluster is not healthy but still all data is available for search.)
> 
> 1. 
> 
> Second node B is started.  
> Node B has index replicas on it. (because number\_of\_replicas=1)  
> Cluster status is GREEN.  
> (At tis point your cluster is healthy)
> 
> 1. 
> 
> Third node C is started.  
> One primary shard from node A is moved to C.  
> So now index primary shards are located on A and C. Both replicas on B.  
> Cluster status is GREEN.
> 
> 1. 
> 
> Fourth node D is started.  
> The node D got one of replicas from B.  
> Cluster status is GREEN.
> 
> 1. 
> 
> Now you can shutdown one node at a time until two nodes are left.  
> Cluster status is GREEN.
> 
> 1. Shutdown one node. Only one is left.  
> Two shards and no replicas are on the last remaining node.  
> Cluster status is YELLOW.
> 
> So as you can see if replicas=X setting can not be satisfied then you can  
> learn this from cluster health status ([http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...  
> ).
> 
> Also see below for my other comments.
> 
> On Fri, Jun 18, 2010 at 11:31 PM, Otis [otis.gospodne...@gmail.com](mailto:otis.gospodne...@gmail.com) wrote:
> 
> > Thanks Lukas,
> 
> > I looked at this the other day, but I don't think that answers my Qs,  
> > or at least I can't tell without hearing Shay's accompanying talk. 🙂
> 
> > So, questions:
> 
> > - Slide 94: cluster with settings: replicas = 1, shards = 2, and a  
> > node with 2 shards, OK
> 
> > - Slide 95: 2nd node appears and end sup with the same shards as on  
> > node 1.  
> > -- Q: doesn't this mean that shards were replicated? Why, if  
> > number\_of\_replicas=1 ?
> 
> This is obvious. number\_or\_replicas = 1 means there is one (non-primary)  
> shard for each (primary) shard.  
> If you had created index with 0 replicas and two shards then you would have  
> got GREEN status in the first step (meaning replicas=X setting is satisfied  
> ).
> 
> > - Slide 96: 3rd and 4th nodes appear and after that each node in the  
> > cluster ends up with just 1 shard. This makes sense.  
> > -- Q: This happens automagically?
> 
> As you can see from the test above this is happening automagically 🙂
> 
> > -- Q: Does ES simply _copy_ the index/shard from one node to the  
> > other when it detects more nodes joined the cluster?
> 
> I think it can take it from gateway but I will let Shay to provide more  
> details.
> 
> > - Slide 107: a new index with 2 shards is added (shards=1,  
> > replicas=1), but 2 nodes get that new index, even though replicas=1.  
> > Huh?
> 
> This is correct. Isn't it? One primary shard on one node and the second node  
> gets replica of that shard (if it wouldn't be possible to locate second node  
> for replica then you would be able to learn this from health status).
> 
> > So the above questions are really about understanding why an index/  
> > shard is being replicated even when replicas=1.
> 
> > I think the above may also answer what happens when the cluster is  
> > expanded: ES detects new nodes joining and figures out that they can  
> > handle/host some of the indices or index replicas and somehow send  
> > them to the new nodes. I assume there are mechanisms in ES to let it  
> > spread the data evenly, though, I am guessing, it doesn't take into  
> > account query rates, so it is not yet able to migrate hot indices  
> > around automatically?
> 
> I think this is not implemented yet. However, it seems that there is a lot  
> of data available and implementing specific traffic load handlers would be  
> possible (I would be surprised it that is not planed).
> 
> > I think the above doesn't cover what happens when 1 or more nodes go  
> > down: I know ES detects that, but does it know which indices were on  
> > those nodes and does it automatically create more replicas for those  
> > indices in order to satisfy the replicas=X setting?
> 
> It is known which shards and its replicas are located on which nodes.  
> Seehttp://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste...  
> the  
> "routing\_table.indices" part. (Note with current master you can expect to  
> get more info from REST API)
> 
> > Thanks,  
> > Otis
> 
> > On Jun 15, 3:33 am, Lukáš Vlček [lukas.vl...@gmail.com](mailto:lukas.vl...@gmail.com) wrote:
> > 
> > > Otis,
> 
> > > I am sure most of these questions are addressed in Elasticsearch  
> > > presentation from Berlin:  
> > > [http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)...  
> > > \<[http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)..  
> > > .\>Automatic  
> > > shard allocation - pages 90-108.
> 
> > > Also it is very easy to verify yourself, just change logging level to  
> > > DEBUG  
> > > for rootLogger in logging.yml. Then start one node and index some data.  
> > > Then  
> > > start second node and see log files. Once the second node is available  
> > > some  
> > > index shards are allocated to the new node. By default ES uses 5 shards  
> > > with  
> > > 1 replica for each shard. If node goes down then you can use health API  
> > > to  
> > > see if you still have all data available for search (  
> > > [http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...).  
> > > It then depends on the speed in which particular nodes go down (crash,  
> > > not  
> > > regular shutdown) because if there is only one shard left and no replica  
> > > is  
> > > available then it should replicate itself to some other node (providing  
> > > replica is set to 1 or more).
> 
> > > Regards,  
> > > Lukas
> 
> > > On Tue, Jun 15, 2010 at 8:46 AM, Otis [otis.gospodne...@gmail.com](mailto:otis.gospodne...@gmail.com)  
> > > wrote:
> > > 
> > > > Hi,
> 
> > > > In ES, what controls which node/shard a doc will get indexed on?
> 
> > > > What happens (or what does one need to do) when the search cluster is  
> > > > expanded? Is there a notion of (automatic) rebalancing?
> 
> > > > Similarly, what happens or should be done when a node in a cluster  
> > > > goes down? Is there something that automatically replicates data that  
> > > > disappeared when that node went down?
> 
> > > > Thanks,  
> > > > Otis

---

<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:** [June 22, 2010, 8:50am UTC](https://discuss.elastic.co/t/cluster-questions/3020/7 "2010-06-22T08:50:04Z")

</div>

Gateway is required to support full cluster shutdown. See more here:  
[http://www.elasticsearch.com/blog/2010/02/16/searchengine\_time\_machine.html](http://www.elasticsearch.com/blog/2010/02/16/searchengine_time_machine.html).

Recovery from the gateway happens only when the first ever shard gets  
created, once its done, shards do recovery from other shard in the same  
replication group when they are moved around or instantiated.

-shay.banon

On Tue, Jun 22, 2010 at 7:48 AM, Otis [otis.gospodnetic@gmail.com](mailto:otis.gospodnetic@gmail.com) wrote:

> Thank you for the detailed reply, Lukáš!
> 
> The reason why I didn't get the number\_of\_replicas=X meaning is  
> because in my mind "number of replicas" really means "number of  
> identical copies", while in ES it means "the number of additional  
> copies". Also, in HDFS there is a notion of a replication factor that  
> matches "my" thinking: if replication factor is 3, that means there  
> are 3 copies of each data block in HDFS. Isn't this the more common  
> way of thinking about/counting replicas?
> 
> Oh, and you mentioned that maybe ES copies replicas from the gateway,  
> but is gateway not an optional thing?
> 
> Thanks,  
> Otis
> 
> On Jun 19, 2:08 am, Lukáš Vlček [lukas.vl...@gmail.com](mailto:lukas.vl...@gmail.com) wrote:
> 
> > Hi Otis,
> > 
> > I know that Shay would be able to shed more light on this... but let us  
> > do  
> > our homework now (since many of your questions can be answered by running  
> > a  
> > test).
> > 
> > 1. 
> > 
> > One node A is started and a new index is created (2 shards, 1 replica)  
> > curl -XPUT '[http://localhost:9200/twitter/'-d](http://localhost:9200/twitter/'-d) '  
> > index :  
> > number\_of\_shards : 2  
> > number\_of\_replicas : 1  
> > '  
> > Now we can see node A has two shards allocated on it.  
> > The cluster status is YELLOW.  
> > (Cluster is not healthy but still all data is available for search.)
> > 
> > 1. 
> > 
> > Second node B is started.  
> > Node B has index replicas on it. (because number\_of\_replicas=1)  
> > Cluster status is GREEN.  
> > (At tis point your cluster is healthy)
> > 
> > 1. 
> > 
> > Third node C is started.  
> > One primary shard from node A is moved to C.  
> > So now index primary shards are located on A and C. Both replicas on B.  
> > Cluster status is GREEN.
> > 
> > 1. 
> > 
> > Fourth node D is started.  
> > The node D got one of replicas from B.  
> > Cluster status is GREEN.
> > 
> > 1. 
> > 
> > Now you can shutdown one node at a time until two nodes are left.  
> > Cluster status is GREEN.
> > 
> > 1. Shutdown one node. Only one is left.  
> > Two shards and no replicas are on the last remaining node.  
> > Cluster status is YELLOW.
> > 
> > So as you can see if replicas=X setting can not be satisfied then you can  
> > learn this from cluster health status (  
> > [http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...  
> > ).
> > 
> > Also see below for my other comments.
> > 
> > On Fri, Jun 18, 2010 at 11:31 PM, Otis [otis.gospodne...@gmail.com](mailto:otis.gospodne...@gmail.com)  
> > wrote:
> > 
> > > Thanks Lukas,
> > 
> > > I looked at this the other day, but I don't think that answers my Qs,  
> > > or at least I can't tell without hearing Shay's accompanying talk. 🙂
> > 
> > > So, questions:
> > 
> > > - Slide 94: cluster with settings: replicas = 1, shards = 2, and a  
> > > node with 2 shards, OK
> > 
> > > - Slide 95: 2nd node appears and end sup with the same shards as on  
> > > node 1.  
> > > -- Q: doesn't this mean that shards were replicated? Why, if  
> > > number\_of\_replicas=1 ?
> > 
> > This is obvious. number\_or\_replicas = 1 means there is one (non-primary)  
> > shard for each (primary) shard.  
> > If you had created index with 0 replicas and two shards then you would  
> > have  
> > got GREEN status in the first step (meaning replicas=X setting is  
> > satisfied  
> > ).
> > 
> > > - Slide 96: 3rd and 4th nodes appear and after that each node in the  
> > > cluster ends up with just 1 shard. This makes sense.  
> > > -- Q: This happens automagically?
> > 
> > As you can see from the test above this is happening automagically 🙂
> > 
> > > -- Q: Does ES simply _copy_ the index/shard from one node to the  
> > > other when it detects more nodes joined the cluster?
> > 
> > I think it can take it from gateway but I will let Shay to provide more  
> > details.
> > 
> > > - Slide 107: a new index with 2 shards is added (shards=1,  
> > > replicas=1), but 2 nodes get that new index, even though replicas=1.  
> > > Huh?
> > 
> > This is correct. Isn't it? One primary shard on one node and the second  
> > node  
> > gets replica of that shard (if it wouldn't be possible to locate second  
> > node  
> > for replica then you would be able to learn this from health status).
> > 
> > > So the above questions are really about understanding why an index/  
> > > shard is being replicated even when replicas=1.
> > 
> > > I think the above may also answer what happens when the cluster is  
> > > expanded: ES detects new nodes joining and figures out that they can  
> > > handle/host some of the indices or index replicas and somehow send  
> > > them to the new nodes. I assume there are mechanisms in ES to let it  
> > > spread the data evenly, though, I am guessing, it doesn't take into  
> > > account query rates, so it is not yet able to migrate hot indices  
> > > around automatically?
> > 
> > I think this is not implemented yet. However, it seems that there is a  
> > lot  
> > of data available and implementing specific traffic load handlers would  
> > be  
> > possible (I would be surprised it that is not planed).
> > 
> > > I think the above doesn't cover what happens when 1 or more nodes go  
> > > down: I know ES detects that, but does it know which indices were on  
> > > those nodes and does it automatically create more replicas for those  
> > > indices in order to satisfy the replicas=X setting?
> > 
> > It is known which shards and its replicas are located on which nodes.  
> > Seehttp://  
> > [www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...  
> > the  
> > "routing\_table.indices" part. (Note with current master you can expect to  
> > get more info from REST API)
> > 
> > > Thanks,  
> > > Otis
> > 
> > > On Jun 15, 3:33 am, Lukáš Vlček [lukas.vl...@gmail.com](mailto:lukas.vl...@gmail.com) wrote:
> > > 
> > > > Otis,
> > 
> > > > I am sure most of these questions are addressed in Elasticsearch  
> > > > presentation from Berlin:  
> > > > [http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo).  
> > > > ..  
> > > > \<  
> > > > [http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo](http://www.slideshare.net/elasticsearch/elasticsearch-at-berlinbuzzwo)..  
> > > > .\>Automatic  
> > > > shard allocation - pages 90-108.
> > 
> > > > Also it is very easy to verify yourself, just change logging level to  
> > > > DEBUG  
> > > > for rootLogger in logging.yml. Then start one node and index some  
> > > > data.  
> > > > Then  
> > > > start second node and see log files. Once the second node is  
> > > > available  
> > > > some  
> > > > index shards are allocated to the new node. By default ES uses 5  
> > > > shards  
> > > > with  
> > > > 1 replica for each shard. If node goes down then you can use health  
> > > > API  
> > > > to  
> > > > see if you still have all data available for search (
> 
> [http://www.elasticsearch.com/docs/elasticsearch/rest\_api/admin/cluste](http://www.elasticsearch.com/docs/elasticsearch/rest_api/admin/cluste)...).
> 
> > > > It then depends on the speed in which particular nodes go down  
> > > > (crash,  
> > > > not  
> > > > regular shutdown) because if there is only one shard left and no  
> > > > replica  
> > > > is  
> > > > available then it should replicate itself to some other node  
> > > > (providing  
> > > > replica is set to 1 or more).
> > 
> > > > Regards,  
> > > > Lukas
> > 
> > > > On Tue, Jun 15, 2010 at 8:46 AM, Otis [otis.gospodne...@gmail.com](mailto:otis.gospodne...@gmail.com)  
> > > > wrote:
> > > > 
> > > > > Hi,
> > 
> > > > > In ES, what controls which node/shard a doc will get indexed on?
> > 
> > > > > What happens (or what does one need to do) when the search cluster  
> > > > > is  
> > > > > expanded? Is there a notion of (automatic) rebalancing?
> > 
> > > > > Similarly, what happens or should be done when a node in a cluster  
> > > > > goes down? Is there something that automatically replicates data  
> > > > > that  
> > > > > disappeared when that node went down?
> > 
> > > > > Thanks,  
> > > > > Otis

---

<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:23am UTC](https://discuss.elastic.co/t/cluster-questions/3020/8 "2017-07-06T04:23:06Z")

</div>


