Cluster operationshard5-8 years

A team wants to add primary shards to an existing, growing index without reindexing, expecting it to rebalance data across the new shards the way adding a database read replica doesn't require touching existing data. Why doesn't this work, and what's the actual mechanism that makes the primary shard count permanent?

Which shard a document lives on is decided by hash(document_id) % number_of_primary_shards — pure arithmetic, computed fresh every time, not a stored assignment that could be moved. If the primary count changes, the modulus in that formula changes, so nearly every existing document's hash(id) % new_count lands on a different shard than the one it's actually stored on — there's no way to shift data to match a new formula except to read every document out and write it back in through the new routing, which is exactly what a reindex is. Adding replicas doesn't have this problem because replicas are just copies of an already-decided primary; adding primaries changes the decision itself.

The lesson behind it →
More on Cluster operations