Replicas help reads, not writes
Each primary fans every write out to all in-sync replicas before acking the client, so adding replicas linearly adds read capacity and HA but does not increase write throughput.
Scene 06
Replicas help reads, not writes
- Watch
- Try it
- Predict
- Capture
Five primary shards, one replica each. Indexing clients send writes to the primary, which fans them out to every replica before acknowledging. Reads can hit any copy — primary or replica — picked by adaptive_replica_selection. Cluster status is GREEN.
Highlighted lines are the ones running in the diagram right now.
def index(doc):shard = hash(doc._id) % num_primary_shardsprimary = routing_table.primary_for(shard)primary.applyLocally(doc)acks = []for r in primary.in_sync_replicas:acks.append(r.send(doc)) # parallel fan-outwait_all(acks) # block on slowestreturn ok(client) # ack now
def search(query):for shard in all_primary_shards:copies = primary + in_sync_replicascopy = adaptive_replica_selection(copies)scatter(copy, query)# reads scale linearly with (1 + num_replicas)
Where this sits in Build a distributed search engine (Elasticsearch / OpenSearch style)
Scene 06 of 12. Replicas linearly add read capacity and HA but do not speed up writes — every primary fans every write out to all in-sync replicas before acking.
All 12 scenes in Build a distributed search engine (Elasticsearch / OpenSearch style) · Every curriculum