A search prototype is easy: put documents in Elasticsearch, OpenSearch or Solr, send a query, get ranked hits. The hard part starts when the corpus grows to billions of documents, traffic grows to thousands of queries per second, and the business wants a price change visible in results within seconds. Then every design choice that did not matter on one node, how many shards, how many copies, how writes arrive, how you change a mapping, turns into latency, cost or an outage.

This article is about that scaling and operating layer. It assumes you know what an inverted index, BM25 and scatter-gather are; if not, read Search System Architecture in Depth first, which covers the internals. Here we size a cluster from measurements, work through why fan-out makes the tail latency worse as you grow, build an indexing path that stays correct under replays, change the schema without downtime, and list the failures you should expect. The numbers in the worked example are illustrative, chosen to show the method; your own load tests replace them.

Advertisement

The three axes that scale independently

Search load grows along three axes, and each one is served by a different mechanism. Mixing them up is the most common sizing mistake.

Corpus size decides how many primary shards you need. A shard is a self-contained index over a slice of the documents; its size sets how long recovery takes when a node dies. More data means more shards, not more copies.

Query rate decides how many copies (replicas) of each shard you run. Every query visits one copy of each shard it needs, so throughput scales with complete copies. More shards add work, not throughput.

Freshness decides the write path. A daily rebuild is a batch job; a seconds-level target needs change data capture, a streaming indexer, and careful ordering and retries.

Source of truthOLTP databaseChange logCDC / outbox topicIndexer fleetenrich, bulk, versionBackfill jobfull rebuild, reindexAlias: productspoints at products_v7bulk writesClientsapps, APIsQuery serviceparse, route, hedgeShard 0N copiesShard 1N copiesShard NN copiesscatterMerge and reranktop-k per shardgatherZone A | Zone B | Zone Ccopies of every shard spread across zonesThree axes scale independently: corpus size sets shards, QPS sets copies, freshness sets the ingest path.
Search at scale: a change log feeds an idempotent indexer that writes through an alias; the query service scatters to one copy of every shard, hedges slow copies, and merges top-k results. Every zone holds at least one complete copy.

Sizing shards from measured bytes, not guesses

Start from a measurement. Index a sample of one million real documents with the production mapping and analyzers, force-merge it to one segment, and divide the on-disk size by the document count. That number, index bytes per document, often differs from the raw JSON size by a factor of two in either direction, because analyzers, doc values, stored fields and n-gram fields add up differently for every schema.

Next choose a target shard size. Elastic's guidance has long pointed at tens of gigabytes; the real constraints are recovery time (while a large shard re-copies you run one copy short), merge headroom, and per-shard overhead from thousands of tiny shards. Project growth over the horizon you want before resharding, because in Elasticsearch and OpenSearch the primary count is fixed at index creation; split and shrink exist, but they are operations, not settings.

import math

def plan_cluster(docs, bytes_per_doc, growth_per_year, years,
                 target_shard_gb, peak_qps, qps_per_copy, zones=3):
    """Shards from data volume, copies from load; both inputs measured."""
    future_docs = docs * (1 + growth_per_year) ** years
    index_gb = future_docs * bytes_per_doc / 1e9
    primaries = math.ceil(index_gb / target_shard_gb)

    # Lose a whole zone and still serve peak: survivors carry everything.
    copies_needed = peak_qps / qps_per_copy
    copies = math.ceil(copies_needed * zones / (zones - 1))
    copies = max(copies, zones)          # at least one copy per zone
    copies = math.ceil(copies / zones) * zones  # keep zones symmetric
    return dict(index_gb=round(index_gb), primaries=primaries,
                copies=copies, shard_copies=primaries * copies)

print(plan_cluster(docs=2e9, bytes_per_doc=1500, growth_per_year=0.4,
                   years=2, target_shard_gb=40, peak_qps=12_000,
                   qps_per_copy=3_000))
# {'index_gb': 5880, 'primaries': 147, 'copies': 6, 'shard_copies': 882}

In the worked example, two billion documents at 1,500 index bytes each, growing 40 percent a year for two years, give about 5.9 TB of primaries and 147 shards at 40 GB. One full copy serves 3,000 queries per second at the p99 target in the load test. Peak is 12,000, so four copies carry the load, but losing one of three zones must not overload the survivors, which raises it to six, two per zone: 882 shard copies and about 35 TB of disk before merge headroom. Doubling traffic changes copies, not shards.

Advertisement

Why fan-out makes the tail worse as you grow

A query that must visit every shard is only as fast as the slowest shard that answers it. If each shard copy is slow (say, stalled by garbage collection or a merge) one percent of the time, a query that touches N shards is slow with probability 1 - 0.99N. At 10 shards that is about 10 percent; at 100 shards, about 63 percent. The median barely moves as you add shards, but the p99 collapses toward the slowest-component latency. This is the same effect Dean and Barroso described as the tail at scale, and search is its textbook case.

Three defences follow: touch fewer shards, ask another copy when one is slow, and answer without the slowest shard when the product allows it.

Cutting fan-out with routing and partitioning

Most large search systems are document-partitioned: each shard indexes a subset of documents, so every query must visit every shard. What you can do is route. If queries naturally carry a partition key, a tenant, a marketplace, a country, put each key's documents in one shard or a small group with a routing value, and send the query only there. A multi-tenant SaaS search with ten thousand customers should not scatter a customer's query across 147 shards when all of that customer's documents live in one. The risk is skew: one giant tenant makes one hot shard. The fix is to route large tenants to dedicated indices and share the rest, the same hot-key reasoning covered in the sharding guide.

Time is the other natural partition. Logs, events and news are queried mostly by recent windows, so one index per day or week, queried only where it overlaps the window, cuts fan-out sharply and turns retention into a cheap index drop.

Hedged scatter-gather and partial results

For queries that must touch every shard, the query service can hedge: send the request to the replica with the best recent latency, and if it has not answered within roughly the p95 latency, send the same request to a second replica and take whichever answers first. Only about five percent of shard requests get hedged, so the extra load is small, and a stalled copy stops setting your p99. Elasticsearch's adaptive replica selection already ranks copies by recent response time and queue size; the hedge you add in your own query tier.

import asyncio

async def query_shard_group(shard, replicas, request, hedge_after_ms, timeout_ms):
    """Best replica first; if slow, hedge to a second. First answer wins."""
    order = sorted(replicas, key=lambda r: r.ewma_latency_ms)   # adaptive selection
    first = asyncio.create_task(order[0].search(shard, request))
    done, _ = await asyncio.wait({first}, timeout=hedge_after_ms / 1000)
    if done:
        return first.result()
    second = asyncio.create_task(order[1].search(shard, request))
    done, pending = await asyncio.wait({first, second}, timeout=timeout_ms / 1000,
                                       return_when=asyncio.FIRST_COMPLETED)
    for t in pending:
        t.cancel()
    return next(iter(done)).result() if done else None   # None = shard missing

async def search(request, shard_map, k=20):
    tasks = [query_shard_group(s, reps, request, hedge_after_ms=40, timeout_ms=150)
             for s, reps in shard_map.items()]
    results = await asyncio.gather(*tasks)
    missing = sum(r is None for r in results)
    hits = sorted((h for r in results if r for h in r.top_k), key=lambda h: -h.score)[:k]
    return dict(hits=hits, partial=missing > 0, shards_missing=missing)

Partial results are the last line. With a hard per-shard timeout, the service returns the hits it has and flags the response as partial: usually right for a catalogue browse, wrong for a compliance search, so make it a per-endpoint policy, and graph the partial rate as an early warning of a degrading node or zone. Pair this with admission control that rejects or downgrades expensive requests (deep pagination, huge aggregations, leading wildcards) when data-node queues grow; load shedding covers the policies.

An indexing path that stays correct under replays

Dual writes from the application to the database and the search engine are the classic shortcut and the classic source of drift: one write fails and the index stays silently wrong. The robust design reads committed changes from the database, either with change data capture on its log or with a transactional outbox, into a durable topic keyed by document ID. An indexer fleet consumes the topic, enriches each document (joins, denormalised fields, computed ranking features), and bulk-writes it.

Two properties make it correct. Keying the topic by document ID keeps each document's changes in order within a partition. Writing with an external version, the source's commit sequence number, makes every write idempotent and order-safe: the engine rejects a write whose version is not newer than the stored one, so a replay after a crash, or a backfill racing the stream, cannot overwrite fresh data with stale data.

def handle_change_batch(events, es):
    """Apply CDC events idempotently and in the right order per document.

    Each event carries the row's commit sequence (e.g. an LSN or an
    updated_at version). version_type=external makes the engine reject any
    write whose version is not higher than the stored one, so replays and
    out-of-order delivery cannot resurrect stale data.
    """
    actions = []
    for ev in events:
        meta = {"_index": "products", "_id": ev.key,
                "version": ev.commit_seq, "version_type": "external"}
        if ev.op == "delete":
            actions.append({"delete": meta})
        else:
            actions.append({"index": meta})
            actions.append(enrich(ev.row))       # joins, denormalised fields
    resp = es.bulk(operations=actions)
    for item in resp["items"]:
        status = next(iter(item.values()))["status"]
        if status == 409:        # version conflict: we already have newer data
            continue
        if status >= 500 or status == 429:
            raise RetryableBatchError(item)       # do not commit the offset
    record_freshness_lag(max(ev.commit_ts for ev in events))

Measure freshness end to end: the time from database commit to the document being visible in search, which includes the engine's refresh interval. Alert on its p99, not the average. Bigger bulk batches index more efficiently but add latency and costlier retries; start at a few megabytes and tune from throughput and rejections.

Changing the schema without downtime

Sooner or later you need a new analyzer, a new field type or a different shard count, and existing indices cannot change those in place. The pattern is to read and write through an alias and build a new index beside the old one.

  1. Create products_v8 with the new mapping and settings, replicas set to zero and refresh disabled for fast bulk loading.
  2. Start dual consumption: a second indexer group reads the change topic from the current offset and writes to v8, with external versions so ordering is safe.
  3. Backfill v8 from a snapshot of the source of truth. Because writes carry versions, backfill and stream can overlap without harm.
  4. Restore replicas and the refresh interval, wait for green, then verify: document counts per partition, a sample of IDs compared field by field, and replayed production queries compared on overlap of the top results and on latency.
  5. Swap the alias atomically, keep v7 and its consumers running for a rollback window, then delete them.
POST /_aliases
{
  "actions": [
    { "remove": { "index": "products_v7", "alias": "products" } },
    { "add":    { "index": "products_v8", "alias": "products" } }
  ]
}

Teams skip verification and regret it: an analyzer change can shift recall silently, and diffing the top ten for a day of replayed queries catches it.

Tiering, caching and cost control

At scale, cost is dominated by memory and fast disk holding rarely queried data. Time-partitioned data can move through tiers: recent indices on fast nodes with more replicas, older ones on dense nodes with fewer copies, the oldest in snapshots restored on demand. The rollover API starts a new write index at a size or age limit, keeping shard sizes even.

Caching helps less than expected, because the long tail of queries rarely repeats. Engine caches for filters and aggregations help a lot; a result cache mostly helps head queries such as popular category pages. The caching guide covers invalidation, the hard part when freshness targets are seconds.

Failure modes and what to watch

FailureSymptomMitigation
Node or zone lossCluster yellow, recoveries saturate the network, p99 climbsSize copies for zone loss; throttle recovery bandwidth; keep shards small enough to recover quickly
Mapping explosionHeap pressure, slow cluster state updatesStrict mappings, a field-count limit, flatten free-form attributes into one keyword field
Hot shardOne node's CPU and queue far above peersSplit large tenants into their own indices, check routing skew, rebalance
Indexing backlogFreshness lag grows, bulk rejections (HTTP 429)Scale the indexer, lengthen refresh during catch-up, shed optional enrichments
Deep paginationCoordinator memory spikes on page 500Cap from+size, use search_after for export
Stale or resurrected docsDeleted items appear in resultsExternal versions, delete events in the same ordered stream, periodic reconciliation

Dashboard per endpoint: query rate, p99 and partial rate; per node: queue depth, rejections, GC pauses, shard count; globally: freshness lag p99 and recovery activity. The p99, partial rate and freshness lag move before users complain.

What to do next

  1. Measure index bytes per document on one million real documents with your production mapping, and write it down next to your projected growth.
  2. Load-test one complete copy of the index with replayed production queries, and record the QPS it sustains at your p99 target. Size copies with the zone-loss formula above.
  3. Graph per-shard latency and compute how many shards your common queries touch. If a partition key exists, add routing.
  4. Add a hedge at roughly the per-shard p95 and a per-endpoint partial-result policy, and graph the hedge and partial rates.
  5. Replace any dual writes with a CDC or outbox stream keyed by document ID and written with external versions; alert on freshness lag p99.
  6. Put every index behind an alias today, and rehearse one full reindex with verification before you need it in an incident.
Key takeaway: Search scales along three separate axes: corpus size sets shard count, query rate sets the number of complete copies, and freshness sets the write path. Size both from measurements, plan for losing a zone, and remember that fan-out turns rare slow shards into a common slow query, so route to fewer shards, hedge the slow ones, and decide where partial results are acceptable. Feed the index from an ordered, versioned change stream and change schemas only through aliases with a verified rebuild.