Series overview
Part 8 of 1844% complete
2026-05-16•16 min read

Shards, settings, and capacity planning for 300 million profiles

How many shards does a 300-million-profile index need? This chapter does not answer with a number. It gives you a method that produces one for your data and traffic: measure a representative sample, apply a worksheet, and validate the result with a load test. You measure the lab index field by field, work through the worksheet with inputs that are explicitly hypothetical, and run a real Rally benchmark against the lab so the validation step is something you have done, not only read about.

First, the settings that the worksheet reasons about: primary shards, replicas, refresh, the indexing buffer, and the translog. Each follows the pattern used throughout this series: what it is, why it matters at 300 million profiles, an example, and the common mistake.

You need the lab with user-profile-v1 loaded (chapter 05) and the export in data/profiles.ndjson. The Rally section also needs Docker and about 1 GB of free disk space. The chapter takes about 60 minutes.

Primary shards: fixed at creation

What it is. The number of primary shards splits an index’s documents into independent Lucene indices. Chapter 02 showed that it is fixed when the index is created.

Why it matters at 300 million profiles. Primaries set the size of each shard and the maximum parallelism of indexing. Every search runs on one copy of every shard. Too few primaries give you oversized shards that are slow to recover and rebalance. Too many give you per-shard overhead on every search and in the cluster state. Changing the count later means building a new index (chapter 14).

Example. The lab index has one primary holding one million profiles in 157.6 MB. Nothing about that number transfers to production except the measurement method, which Stage 1 uses.

Common mistake. Choosing a primary count by folklore: “one shard per node”, “always five”, “as many as the CPUs”. None of these knows your document size, growth, or query mix.

Replicas: adjustable at any time

What it is. Each replica is a complete copy of every primary, on a different node. The replica count can be changed at any time with a settings update.

Why it matters at 300 million profiles. Replicas provide the failure tolerance: with one replica, any single node can fail without losing data or availability. They also add search capacity, because any copy of a shard can answer a search. Every replica also multiplies disk use and indexing work, since each copy indexes every document itself.

Example. Chapter 02 changed number_of_replicas on a live index and watched the health go from yellow to green. The same request is how a production cluster adds search capacity for a seasonal peak, or drops replicas to zero during a bulk load into a new index version, which chapter 14 does.

Common mistake. Adding replicas to fix slow searches when the bottleneck is per-query work, such as an expensive query or too many shards. A replica adds capacity for more concurrent searches; it does not make one search faster.

Refresh, the indexing buffer, and the translog

These three settings govern the write path from chapter 02’s diagram: a write goes to an in-memory buffer and to the translog; a refresh turns the buffer into a searchable segment; a flush commits segments to disk. The defaults below were read from the lab cluster with include_defaults=true.

SettingLab defaultWhat it controls
index.refresh_interval1s (set explicitly in chapter 05)How stale search results may be, and how many small segments are created
indices.memory.index_buffer_size10% of heap, shared by the node’s active shardsHow much can be buffered in memory before segments must be written
index.translog.durabilityREQUESTWhether each write is fsynced to the translog before it is acknowledged
index.translog.sync_interval5sHow often the translog is fsynced when durability is async
index.translog.flush_threshold_size10gbTranslog size that triggers a flush

Why they matter at 300 million profiles. During a backfill, a one-second refresh makes a new segment every second on every shard, and the merges that follow compete with indexing. Turning refresh off for the load and restoring it afterwards, as chapter 05’s loader did, is the single largest indexing-speed setting you control. The indexing buffer settings rarely need changing; the default scales with the heap.

Trade-off. Setting index.translog.durability to async acknowledges writes before they are fsynced, so a crash can lose up to sync_interval of acknowledged writes. For a projection that can be rebuilt from PostgreSQL, and whose indexer would re-deliver unacknowledged batches anyway, that can be an acceptable exchange for indexing throughput. It is still a decision to make deliberately and write down, not a default to copy. The translog settings document every option.

Common mistake. Leaving refresh_interval: -1 on an index after a load. Nothing written afterwards ever becomes searchable, and the failure is silent. Chapter 05’s loader restores the interval in the same script that disables it.

Too many small shards, or too few large ones

Elastic’s shard sizing guidance is the primary source for this section. Its main recommendations: aim for shards between 10 GB and 50 GB, keep each shard below 200 million documents, and benchmark with production data, queries, and hardware. It also states the hard limit: a single shard can hold at most 2,147,483,519 documents.

Too many small shards cost you in three ways. Each shard has fixed memory and CPU overhead. A search touches one copy of every shard, and each shard’s part of the query runs on one thread, so a search over hundreds of tiny shards can exhaust the search thread pool, which has 10 threads on the lab’s 6-core machine. And every shard adds to the cluster state the master must maintain. Elasticsearch caps a cluster at 1,000 non-frozen shards per node by default, as cluster.max_shards_per_node. Elasticsearch 9.5 batches the query phase for shards on the same data node into one round trip, which lowers the network overhead of many shards; it does not remove the per-shard work.

Too few large shards cost you in recovery and flexibility. When a node fails, its shards are rebuilt on other nodes by copying whole shards. A 200 GB shard takes far longer to recover and to rebalance than a 25 GB one, and the cluster runs without full redundancy for that whole time. Very large shards also slow some searches and merges.

Principle. For a single search-heavy index of this size, a moderate number of shards in the middle of the recommended range is usually the starting point. The worksheet in Stage 2 turns that into a number, and Stage 3 tests it.

Stage 1 — Measure the sample

Every worksheet input that can be measured, should be. Three measurements come from the lab, and you would take the same three from an anonymised sample of production data.

Average source document size. Compute the average size of a document as compact JSON from the export:

Terminal window
python3 - <<'EOF'
import json
count = total = 0
with open("data/profiles.ndjson", encoding="utf-8") as f:
for i, line in enumerate(f):
if i % 2 == 1: # every second line is a document
doc = json.dumps(json.loads(line), separators=(",", ":"), ensure_ascii=False)
total += len(doc.encode("utf-8"))
count += 1
print(count, total, round(total / count, 1))
EOF
1000000 373633780 373.6

Index size on disk. GET _cat/indices/user-profile-v1?h=store.size,docs.count reports 157.6 MB for 1,000,000 documents. The ratio of index size to raw source size is 157.6 MB ÷ 373.6 MB, about 0.42.

Where the bytes go. The analyse index disk usage API breaks the size down per field and per data structure. It reads every segment, so it must be allowed explicitly:

POST user-profile-v1/_disk_usage?run_expensive_tasks=true
FieldTotalInverted indexStored fieldsDoc valuesPoints
_source87.8 MB87.8 MB
email13.2 MB7.3 MB5.9 MB
fullName.prefix8.3 MB8.3 MB
_id7 MB3.7 MB3.2 MB
userId6.5 MB3.6 MB2.8 MB
mobileNumber6.3 MB3.4 MB2.9 MB
updatedAt4.9 MB2.2 MB2.6 MB
createdAt4.2 MB1.6 MB2.6 MB
fullName.keyword3.3 MB1.9 MB1.4 MB
fullName1.9 MB1.9 MB

The table shows the ten largest fields from the response. Three findings are useful. The stored _source is more than half the index. email is the largest indexed field, because every value is unique and it carries doc values it may not need: no query shape sorts or aggregates by email, so "doc_values": false on email and mobileNumber is a candidate saving for the next index version. And the autocomplete sub-field costs more than the main fullName field, as chapter 07 measured.

Warning — The lab’s ratio of 0.42 is not a planning number. The seed data repeats 30 first names, 20 surnames, and 10 cities, and repetitive data compresses unusually well. Real profiles have millions of distinct names and emails. Measure the ratio on a real, anonymised sample of at least a few million documents, merged to a single segment, before trusting it.

Measure merged size, not live size. The same million documents measured 0.33 GB in the Rally run in Stage 3, straight after bulk indexing with 16 segments, and 157.6 MB in user-profile-v1 once merges had run. Take the size from an index that has settled, or after a _forcemerge to one segment on a copy, or the worksheet will overstate storage.

Stage 2 — Work through the sizing worksheet

The worksheet turns measurements and requirements into shard counts, node counts, and disk. Every input in this example is hypothetical. The column labelled “Source” says where the real value must come from.

#InputHypothetical valueSource of the real value
AProfiles today300,000,000PostgreSQL
BGrowth over the planning horizon (24 months)30%Product forecast
CAverage source document size700 bytesStage 1 on a real sample
DIndex-to-source ratio1.1Stage 1 on a real sample, merged
ETarget shard size25 GBElastic’s 10–50 GB range, low end for latency-sensitive search
FPeak search rate1,500 requests/sTraffic forecast
Gp99 latency target150 msProduct requirement
HSearch rate one node sustains within G400 requests/sStage 3 on production-like hardware
IPeak change rate from the outbox2,000 profiles/sPostgreSQL change statistics
JBackfill window for a full rebuild8 hoursOperational requirement
KNode failures to tolerate1Availability requirement

And the calculations, each rounded up:

StepFormulaHypothetical result
Documents at horizonA × (1 + B)390 million
Primary data sizeA × (1 + B) × C × D390M × 700 B × 1.1 ≈ 300 GB
Primary shardsprimary size ÷ E300 ÷ 25 = 12
Documents per shard390M ÷ 1232.5 million, below the 200 million guidance
Replicasat least K; more if F needs itStart with 1; the search calculation below may raise it
Nodes for searchF ÷ H, plus K1,500 ÷ 400 = 3.75, so 4, plus 1 = 5
Shard copies12 × (1 + replicas)24 with one replica
Nodes, adjusted for even placementa count that divides 246, with 4 shard copies each
Data per node300 GB × 2 copies ÷ 6100 GB
Data per node during a reindexdoubled while v1 and v2 coexist (chapter 14)200 GB
Minimum disk per nodereindex size ÷ 0.85 (low watermark)about 235 GB
Backfill indexing rate needed390M ÷ Jabout 13,500 documents/s, to be measured in Stage 3

Read the result as a hypothesis to test, not a design: 12 primaries, 1 replica, 6 data nodes with at least 235 GB of disk each. Three things in it are easy to miss.

The reindex doubles storage. Zero-downtime reindexing (chapter 14) keeps the old and new index versions side by side, so disk must fit both. A cluster sized only for steady state cannot reindex without adding capacity first.

Disk watermarks cap usable disk. By default, Elasticsearch stops allocating new shards to a node above 85% disk use, moves shards away above 90%, and blocks writes above 95%, with absolute headroom caps on very large disks; the allocation settings reference lists them. Plan to stay below the low watermark, including during reindexing and after a node failure.

Failure tolerance is a capacity requirement. With one node down, the remaining five must hold every shard copy the cluster can still allocate and serve the peak search rate. That is why the search calculation adds K nodes, rather than assuming every node is always available.

Memory is sized per node, not per index. Give each Elasticsearch process a heap of no more than 50% of the node’s memory and below the compressed-pointer threshold, which the JVM settings reference puts at about 26 GB on most systems. The rest of the memory is the operating system’s file cache, and search speed depends on how much of the index fits in it.

Stage 3 — Validate with Rally

The worksheet’s most important input, H, can only be measured. Rally is Elastic’s benchmarking tool. It indexes a corpus, runs a schedule of operations at a controlled rate, and reports throughput and latency percentiles. This stage builds a small custom track from the lab data and runs it against the lab, using the official Docker image so there is nothing to install.

Security note — Benchmark a dedicated index, never the production alias. This track creates and deletes its own user-profile-bench index. Point it at a real cluster only with a benchmark-specific user whose privileges are limited to that index.

  1. Generate the track’s index body and corpus from files you already have. Create rally/make-track-files.py:

    rally/make-track-files.py
    """Build the Rally track's index body and corpus from the lab's files."""
    import json
    from pathlib import Path
    track = Path("rally/user-profile-track")
    track.mkdir(parents=True, exist_ok=True)
    # Index body: the v1 definition without aliases, with shard and replica
    # counts turned into track parameters.
    body = json.loads(Path("es/user-profile-v1.json").read_text())
    body.pop("aliases")
    index_settings = body["settings"]["index"]
    index_settings["number_of_shards"] = "@SHARDS@"
    index_settings["number_of_replicas"] = "@REPLICAS@"
    text = json.dumps(body, indent=2)
    text = text.replace('"@SHARDS@"', "{{ number_of_shards | default(1) }}")
    text = text.replace('"@REPLICAS@"', "{{ number_of_replicas | default(0) }}")
    (track / "index.json").write_text(text + "\n")
    # Corpus: every second line of the bulk export is a document.
    lines = Path("data/profiles.ndjson").read_text(encoding="utf-8").splitlines()
    docs = lines[1::2]
    (track / "documents.json").write_text("\n".join(docs) + "\n", encoding="utf-8")
    (track / "documents-1k.json").write_text("\n".join(docs[:1000]) + "\n", encoding="utf-8")
    print(f"documents: {len(docs)}, bytes: {(track / 'documents.json').stat().st_size}")
    Terminal window
    python3 rally/make-track-files.py
    documents: 1000000, bytes: 418633780

    The {{ ... }} placeholders are Rally track parameters, so you can rerun the same track with different shard and replica counts. documents-1k.json is the small corpus Rally uses in test mode.

  2. Define the track. It recreates the index, bulk-indexes the corpus with four clients, and then runs three searches that match the Search API’s query shapes, each at a fixed rate:

    rally/user-profile-track/track.json
    {
    "version": 2,
    "description": "User-profile search: bulk indexing plus the Search API's three query shapes",
    "indices": [
    { "name": "user-profile-bench", "body": "index.json" }
    ],
    "corpora": [
    {
    "name": "user-profiles",
    "documents": [
    {
    "source-file": "documents.json",
    "document-count": 1000000,
    "uncompressed-bytes": 418633780,
    "target-index": "user-profile-bench"
    }
    ]
    }
    ],
    "schedule": [
    { "operation": { "operation-type": "delete-index" } },
    { "operation": { "operation-type": "create-index" } },
    {
    "operation": { "operation-type": "cluster-health", "request-params": { "wait_for_status": "green" }, "retry-until-success": true }
    },
    {
    "operation": { "name": "bulk-index", "operation-type": "bulk", "bulk-size": 5000 },
    "warmup-time-period": 10,
    "clients": 4
    },
    { "operation": { "name": "refresh", "operation-type": "refresh" } },
    {
    "operation": {
    "name": "exact-lookup-email",
    "operation-type": "search",
    "index": "user-profile-bench",
    "body": { "query": { "term": { "email": "priya.kumar.1@example.com" } } }
    },
    "warmup-iterations": 100, "iterations": 500, "target-throughput": 50, "clients": 2
    },
    {
    "operation": {
    "name": "autocomplete-filtered",
    "operation-type": "search",
    "index": "user-profile-bench",
    "body": {
    "size": 10,
    "_source": ["userId", "fullName", "city"],
    "query": { "bool": {
    "must": { "match": { "fullName.prefix": { "query": "prashant ku", "operator": "and" } } },
    "filter": [ { "term": { "state": "bihar" } }, { "term": { "accountStatus": "ACTIVE" } } ]
    } }
    }
    },
    "warmup-iterations": 100, "iterations": 500, "target-throughput": 50, "clients": 2
    },
    {
    "operation": {
    "name": "name-search-sorted",
    "operation-type": "search",
    "index": "user-profile-bench",
    "body": {
    "size": 20,
    "_source": ["userId", "fullName", "city", "updatedAt"],
    "query": { "bool": {
    "must": { "match": { "fullName": { "query": "prashant", "fuzziness": "AUTO" } } },
    "filter": [ { "term": { "accountStatus": "ACTIVE" } }, { "range": { "updatedAt": { "gte": "now-3y" } } } ]
    } },
    "sort": [ { "_score": "desc" }, { "updatedAt": "desc" }, { "userId": "asc" } ]
    }
    },
    "warmup-iterations": 100, "iterations": 500, "target-throughput": 50, "clients": 2
    }
    ]
    }
  3. Check the track in test mode, which uses the 1,000-document corpus and a handful of iterations. The index uses the profile-name-synonyms set, which already exists in the lab. From the lab directory:

    Terminal window
    set -a; source .env; set +a
    docker run --rm --network host \
    -v "$PWD/rally/user-profile-track:/rally/track" \
    elastic/rally:latest race \
    --track-path=/rally/track --target-hosts=localhost:9200 --pipeline=benchmark-only \
    --client-options="basic_auth_user:'elastic',basic_auth_password:'${ELASTIC_PASSWORD}'" \
    --test-mode
    ...
    | error rate | bulk-index | 0 | % |
    | error rate | exact-lookup-email | 0 | % |
    | error rate | autocomplete-filtered | 0 | % |
    | error rate | name-search-sorted | 0 | % |
    [INFO] SUCCESS (took 11 seconds)

    --pipeline=benchmark-only tells Rally to benchmark an existing cluster instead of provisioning one. --network host lets the container reach localhost:9200; on Docker Desktop for macOS or Windows, use --target-hosts=host.docker.internal:9200 instead. The image version used here reports esrally 2.13.0.

  4. Run the full race by removing --test-mode. In the lab it took about two and a half minutes. Condensed from the summary report:

    | Metric | Task | Value | Unit |
    | Mean Throughput | bulk-index | 65399.6 | docs/s |
    | 50th percentile service time | exact-lookup-email | 1.54554 | ms |
    | 99th percentile service time | exact-lookup-email | 3.40097 | ms |
    | 50th percentile service time | autocomplete-filtered | 2.22327 | ms |
    | 99th percentile service time | autocomplete-filtered | 3.58815 | ms |
    | 50th percentile service time | name-search-sorted | 5.01858 | ms |
    | 99th percentile service time | name-search-sorted | 7.92158 | ms |
    | error rate | (all tasks) | 0 | % |

    Delete the benchmark index afterwards with DELETE user-profile-bench.

What these numbers are. They were measured on one 6-core desktop CPU, with Elasticsearch in a container capped at a 1 GB heap, Rally running on the same machine, one shard, and synthetic data. They show that the track works and how to read it. They say nothing about production capacity, and they must not be used as input H.

Reading the report. Rally reports two timings. Service time is how long Elasticsearch took to answer. Latency adds any time the request waited because the client was behind schedule. At a fixed target-throughput the two should be close; when latency climbs far above service time, the cluster has fallen behind the offered load, which is the saturation point the worksheet needs. The summary report reference defines every metric.

From a working track to a trustworthy measurement

The lab track proves the mechanics. A benchmark that can set worksheet input H also needs the following.

  • Production-like hardware and topology. The same instance types, disks, heap, node count, shard count, and replica count as the candidate design, with Rally on separate machines.
  • Realistic data. An anonymised production sample, or synthetic data generated with production’s distributions of name frequency, name length, and city skew, at a scale where each shard reaches its target size.
  • Varied queries. This track repeats the same three bodies, so caches make it optimistic. Real tracks draw query values from recorded or generated lists through Rally’s parameter sources.
  • Mixed load. Run searches while the outbox indexer writes at its peak rate, using parallel tasks, because merges and refreshes change search latency.
  • A throughput ramp. Repeat the search tasks at increasing target-throughput until p99 latency crosses the target G. The last rate that met the target is H.
  • Failure scenarios. Repeat the ramp with one node stopped, to confirm the cluster still meets G at peak with K nodes lost.

When the measured H differs from the hypothesis, change the worksheet and repeat. That loop, not a formula, is capacity planning.

Common mistakes with capacity planning

  • Copying a shard count from another system. Document size, growth, and query mix all differ.
  • Sizing for steady state only. Reindexing doubles storage; node failure removes capacity. Both must fit.
  • Benchmarking with the lab’s data. Synthetic, repetitive data compresses and caches far better than real data.
  • Reading service time as user latency. Add the Search API’s own time and the network, and measure at a fixed rate, not flat out.
  • Planning once. Revisit the worksheet when profile growth, query mix, or peak traffic changes materially.

What you built, and what comes next

You measured the lab index field by field, worked through a sizing worksheet from measurements and requirements to a shard, replica, node, and disk plan, and built a Rally track that runs against the lab and reports throughput and latency percentiles. The hypothetical plan, 12 primaries, 1 replica, and 6 nodes, is an example of the method’s output, not a recommendation.

What this chapter cannot provide is the real inputs. Those come from your data, your traffic forecast, and a benchmark on hardware you intend to run.

Chapter 09 moves to the application side: the Elasticsearch Java API Client, its low-level REST client, and Spring Data Elasticsearch, what each is good at from Kotlin, and where each one stops being the right tool.

ElasticsearchPerformance

Type to search the site.

↑↓ navigate⏎ openPowered by Pagefind