Hiring BarSupport

Design with Elasticsearch — Staff-Level Technology Guide

Technology guide38 min read6 diagrams

Why This Matters#

Elasticsearch is not a database. It is a derived, distributed inverted index that you must be able to throw away and rebuild — and the Staff question is never "can Elasticsearch store this?" but "what is the source of truth, how does data get from it into the index, and how long does a full rebuild take?" Teams that answer those three questions run Elasticsearch calmly. Teams that don't end up with a 40 TB cluster that is simultaneously the search engine, the analytics store, and the only copy of the data — and a red cluster status nobody knows how to recover.

It shows up in interviews because every product eventually needs search, and every platform eventually needs logs. Interviewers use it to test whether you understand the gap between finding data and storing data: near-real-time visibility (1 s refresh), no multi-document transactions, shard counts fixed at creation, scatter-gather reads, heap-bound aggregations, and relevance tuning as a product discipline rather than an infrastructure one.

The L5 answer says "index the products in Elasticsearch." The L6 answer says "Postgres is the source of truth; an outbox feeds Kafka; an indexer consumer writes to an alias products over a versioned index products_v7 with 12 primaries at ~30 GB each; refresh is 1 s so new listings appear within ~2 s; a full reindex from Kafka compaction takes 3 hours, which is our schema-change and disaster-recovery path." The L7 answer asks how many teams run their own clusters, what the org pays per TB of logs retained, and whether logs belong in a search engine at all.

The 60-Second Pitch#

"Elasticsearch gives us full-text search with relevance ranking, faceting, and fuzzy matching across hundreds of millions of documents in 10–100 ms. It scales horizontally by splitting indexes into shards across nodes, and each shard is a Lucene index of immutable segments. I'd treat it strictly as a derived read model: the source of truth stays in Postgres or DynamoDB, changes flow through CDC and Kafka, and we can rebuild the index from scratch in hours. That lets us change mappings freely with blue-green indexes behind an alias, and it means a corrupted cluster is an inconvenience, not a data-loss incident."

The Staff-level insight: the three hard parts of Elasticsearch are not in Elasticsearch. They are the indexing pipeline (freshness, ordering, deletes), the relevance model (what "good results" means, who owns it, how it's measured), and the lifecycle (shard sizing, retention, reindex). The query DSL is the easy part.

IntentConstraintStrategyFailure ModeCorrectness Bar
Product / content searchRelevance, p99 < 200 ms, freshness ~secondsDerived index via CDC, BM25 + business signals, replicas for QPSStale or missing documents; relevance regressionsEvery live item findable within ~5 s; no deleted item shown
Log / observability analyticsIngest 100K–1M+ docs/s, retention weeks–monthsData streams + ILM hot/warm/cold, few replicas, rollover by sizeIngest backpressure, disk watermarks, costTolerates sampling; lose minutes, never days
Autocomplete / typeaheadp99 < 50 ms per keystrokeSmall dedicated index, edge n-grams or completion suggesterLatency under burst; index bloatApproximate is fine
Security / audit searchRetention 1+ years, rare queriesFrozen tier on object storage, searchable snapshotsSlow cold queriesCompleteness matters — every event retained

🎯 Staff Move: "I'll treat the index as a disposable read model. Postgres owns the truth, an outbox publishes every change, and the indexer is idempotent on document version. That gives us one sentence for schema changes and disaster recovery alike: rebuild from the log into a new index and flip the alias."


Architecture & Internals#

Only the internals that change design decisions.

Cluster, Nodes, Roles#

RoleJobSizing Rule
Master-eligibleCluster state: mappings, shard allocation, node membershipExactly 3 dedicated small nodes (e.g., 4–8 GB heap) for any cluster > ~10 data nodes
Data (hot)Indexing + recent queries on fast NVMeHeap ≤ ~30 GB (compressed object pointers), ~50% of RAM left for the OS page cache
Data (warm / cold)Older, read-mostly indices on dense disksHigher disk:RAM ratios (e.g., 1:100–1:200)
Data (frozen)Searchable snapshots on object storageLocal cache only; queries take seconds
Coordinating-onlyFan out searches, reduce resultsAdd when heavy aggregations stress data nodes
IngestPipeline processors (grok, enrich)Or do transformation upstream in the pipeline

The cluster state is managed by an elected master (Raft-like voting configuration since 7.x — the old minimum_master_nodes split-brain footgun is gone). Cluster state size is a scaling limit: every index, shard, and mapping field lives in it and is published to every node. Tens of thousands of shards or hundreds of thousands of fields make every mapping update and node join slow.

Index → Shards → Segments#

An index is split into primary shards (fixed at creation — changing requires _split, _shrink, or a reindex) and replica shards (changeable anytime). Each shard is a full Lucene index made of immutable segments.

  • Refresh (default every 1 s, only on indices searched recently) writes the in-memory buffer to a new searchable segment. This is why ES is near-real-time: a write acked at t=0 is searchable at up to t≈1 s.
  • Translog — every write is appended to a per-shard transaction log, fsynced per request by default (index.translog.durability: request), so acked writes survive a crash.
  • Flush commits segments to disk and trims the translog.
  • Merge combines small segments into larger ones and physically drops deleted documents. Updates are delete + reindex, so update-heavy indices carry 10–30%+ deleted docs until merges catch up.
Diagram: Index → Shards → Segments

Write Path: Primary–Replica#

A write goes to a coordinating node → routed by hash(_routing or _id) % num_primary_shards → the primary indexes it → forwards to all in-sync replicas in parallel → acks when they succeed. Sequence numbers and primary terms order operations per shard and let a promoted replica resync. Replicas that fail are removed from the in-sync set by the master rather than blocking writes.

Consequence: replicas multiply indexing cost. number_of_replicas: 1 means every document is analyzed and indexed twice. For bulk initial loads, set replicas to 0 and refresh_interval: -1, then restore both — commonly 2–3× faster ingest.

Read Path: Query Then Fetch#

Diagram: Read Path: Query Then Fetch

Every search touches one copy of every shard in the target indices unless routing narrows it. Two consequences drive design:

  1. Tail latency is the slowest shard. 60 shards means 60 chances for a GC pause, a merge, or a hot node to set your p99.
  2. Deep pagination is quadratic-ish. Page 1,000 at size 10 makes each shard return 10,010 candidates. The default index.max_result_window of 10,000 exists for this reason — use search_after with a point-in-time (PIT) for deep scrolling.

Memory: Heap vs Page Cache#

Elasticsearch performance is mostly filesystem cache performance: Lucene segments are memory-mapped. The heap holds cluster state, request buffers, field data, and aggregation state. Rules:

  • Heap ≤ ~30–31 GB (compressed oops cutoff) and ≤ 50% of RAM; a 64 GB node gets ~30 GB heap and ~34 GB of page cache.
  • Circuit breakers (parent breaker ~95% of heap with real-memory accounting) reject requests that would OOM the node — a 429/circuit_breaking_exception is a sizing signal, not a bug.
  • doc_values (on-disk columnar, default for keyword/numeric) keep sorting and aggregations off-heap. Enabling fielddata on text fields pulls them onto the heap — a classic node killer.

Data Modeling / Core Usage — "The Entire Game"#

In Elasticsearch, the entire game is the mapping and the shard plan. Queries are cheap to change; mappings and primary shard counts require a reindex. Design both for the queries you'll run and the volume you'll have in 12 months.

Mapping: Decide Every Field's Job#

Field TypeIndexed AsUse ForCost / Caveat
textAnalyzed tokens in inverted indexFull-text match, relevanceNot sortable/aggregatable without fielddata
keywordExact term + doc_valuesFilters, facets, sorting, exact IDsMax useful length ~256 chars (ignore_above)
Multi-field text + .keywordBothSearch and facet on the same value~2× index size for that field
Numeric / dateBKD tree + doc_valuesRange filters, sorting, histogramsUse keyword for IDs you never range over — faster term lookups
nestedHidden sub-documentsArrays of objects needing per-object matchingEach nested object is a Lucene doc; 10K+ per parent hurts
flattenedSingle field of keywordsArbitrary user-defined labelsNo per-key typing, limited queries
dense_vectorHNSW graphkNN semantic searchRAM-hungry; ~4 bytes × dims per vector (less with quantization)
enabled: falseStored in _source onlyDisplay-only blobsNot searchable
PUT products_v7
{
  "settings": {
    "number_of_shards": 12,
    "number_of_replicas": 1,
    "refresh_interval": "1s",
    "analysis": {
      "analyzer": {
        "autocomplete": { "tokenizer": "standard", "filter": ["lowercase", "edge_ngram_2_15"] }
      },
      "filter": { "edge_ngram_2_15": { "type": "edge_ngram", "min_gram": 2, "max_gram": 15 } }
    }
  },
  "mappings": {
    "dynamic": "strict",
    "properties": {
      "product_id":  { "type": "keyword" },
      "tenant_id":   { "type": "keyword" },
      "title":       { "type": "text", "fields": { "raw": { "type": "keyword" }, "ac": { "type": "text", "analyzer": "autocomplete", "search_analyzer": "standard" } } },
      "description": { "type": "text" },
      "brand":       { "type": "keyword" },
      "price_cents": { "type": "long" },
      "in_stock":    { "type": "boolean" },
      "popularity":  { "type": "rank_feature" },
      "updated_at":  { "type": "date" },
      "source_version": { "type": "long" }
    }
  }
}

"dynamic": "strict" is the single most important production setting for application indices: an unknown field fails the write instead of silently growing the mapping. Dynamic mapping on user-controlled JSON is how clusters reach 50,000 fields and multi-second cluster-state updates (the default index.mapping.total_fields.limit is 1,000 for a reason).

Denormalize — There Are No Joins#

Elasticsearch has no real joins. nested and join (parent-child) fields exist but cost query latency and memory. The Staff default: denormalize at index time. The product document carries brand name, category path, seller rating — the indexer assembles it from the source of truth. The price you pay is update fan-out: renaming a brand touches every product document with that brand. Size it: 1 brand × 200K products at ~5K docs/s bulk = 40 s of reindex — fine. A category rename touching 50M docs is a scheduled job, not a live update.

Shard Sizing#

primary_shards = ceil(projected_index_size_in_12_months / target_shard_size)
target_shard_size: 10–50 GB (search: ~10–30 GB; logs: ~30–50 GB); keep docs per shard under ~200M
total_shards_per_node: stay well under the default cap of 1,000; aim for < ~20 shards per GB of heap
replicas: 1 for HA; add more only for read QPS — each replica adds a full indexing cost
MistakeSymptomFix
Oversharding (e.g., 5 primaries per daily index × 365 days × 3 envs)Tens of thousands of tiny shards, huge cluster state, slow recovery, heap overhead per shardRollover by size, not by day; shrink old indices; fewer indices
Undersharding (1 primary, 500 GB)Slow recovery (hours per shard), can't spread load, merges hugeSplit with _split or reindex into more primaries
Uneven routingOne shard 5× othersBetter routing key or routing partition size

Custom Routing#

By default documents route by _id. Routing by tenant_id puts one tenant's documents on one shard, so tenant-scoped searches hit one shard instead of all — latency and throughput win by ~N× for N shards. Risk: a large tenant makes a hot, huge shard. index.routing_partition_size spreads each routing value over a subset of shards (e.g., 4 of 48) as a middle ground.

The Indexing Pipeline Is Part of the Data Model#

Diagram: The Indexing Pipeline Is Part of the Data Model

Three rules make the pipeline correct:

  1. Key the topic by document ID so all changes to one document are ordered in one partition.
  2. Index with external versioning (version_type=external, version = source row version or LSN). Out-of-order or replayed events become no-ops instead of overwriting newer data.
  3. Deletes are events too. Soft-delete in the source emits a tombstone; the indexer issues a delete. Missing deletes are the most common "search shows a product that no longer exists" bug.

The Tunable Tradeoff — Freshness × Throughput × Relevance Cost#

Dial 1: Refresh Interval (freshness vs indexing throughput)#

refresh_intervalVisible AfterIndexing ThroughputSegment ChurnPick When
1s (default)~1 sBaselineHigh — many tiny segmentsUser-facing search that must feel instant
5s–30s5–30 sOften 20–50%+ higherLowerLogs, analytics, catalogs where seconds don't matter
-1 (disabled)On explicit refreshMaximumMinimalBulk loads and reindexes
?refresh=wait_for per requestBefore response returnsPer-request latency ≈ up to 1 sNone extraRead-your-write for a specific user action

Never use ?refresh=true on the hot path — it forces a segment per request and can collapse indexing throughput by 10× under load.

Dial 2: Durability vs Ingest Rate#

SettingLoss on Node CrashThroughputPick When
translog.durability: request (default)0 acked writesBaselineAnything not trivially re-derivable
translog.durability: async, sync_interval: 5sUp to ~5 sHigherLogs with an upstream buffer (Kafka) that can replay
wait_for_active_shards: 2Write fails if a replica is unavailableLower availabilityRarely worth it — the source of truth is elsewhere

Dial 3: Query Cost vs Relevance Quality#

TechniqueRelevance GainCostWhere It Runs
BM25 on analyzed fieldsBaseline lexical relevanceCheapQuery
Field boosts, multi_match best_fields / cross_fieldsLarge for structured catalogsCheapQuery
function_score / rank_feature (popularity, recency)Business-aligned rankingModerateQuery
rescore window (top 100–500)Expensive signals only on candidatesBoundedQuery
Hybrid lexical + kNN (RRF)Handles synonyms and intent2× work + vector RAMQuery
Learning-to-rank / external re-rankerHighestModel serving, feature pipeline, ownersOutside ES or plugin
search_latency ≈ max_over_shards(query_phase) + merge(shards × size) + fetch(size) + network
indexing_cost  ≈ docs/s × (1 + replicas) × analysis_cost_per_doc
visibility     ≈ pipeline_lag (CDC + Kafka + indexer) + refresh_interval

Aggregation Accuracy Is Also a Dial#

Distributed aggregations are computed per shard and merged, so some are approximate by design:

AggregationExact?Error SourceKnobStaff Default
terms (top-N buckets)No, across shardsEach shard returns only its local top shard_sizeshard_size (default size × 1.5 + 10)Raise shard_size for long-tail facets; report doc_count_error_upper_bound
cardinalityNoHyperLogLog++precision_threshold (exact below it, max 40,000)Fine for dashboards, never for billing
percentilesNoTDigestcompressionFine for latency dashboards
sum / min / max / avgYes (within the searched data)———
composite (paginated)Yes—Page sizeUse for exports and "all buckets" requirements

"Top 10 sellers by revenue" from a terms aggregation can be wrong at the margin; if finance consumes it, compute it in the warehouse. Relevance tuning follows the same discipline: keep a labeled judgment set (a few hundred queries with graded results), measure nDCG or precision@10 offline before shipping a boost change, then A/B test click-through online. Relevance changes without an evaluation set are opinions.

🎯 Staff Move: "Freshness isn't the refresh interval — it's pipeline lag plus refresh. Our budget is 5 seconds end to end: ~2 for CDC and Kafka, ~2 for bulk batching, 1 for refresh. I'll alert on the end-to-end number by comparing updated_at in the source to the time it becomes searchable, not on any single stage."


Anti-Patterns — What Kills Elasticsearch Deployments#

1. Elasticsearch as the Source of Truth#

No multi-document transactions, eventual visibility, historical shard-level data-loss scenarios under partitions, and schema changes that require reindexing from something. If ES is the only copy, every mapping change and every red cluster is a data-loss risk. Fix: a durable source (Postgres, DynamoDB, S3 + Kafka) and a tested rebuild path.

2. Mapping Explosion#

Dynamic mapping on arbitrary JSON (user attributes, log fields from 400 services). Cluster state grows, every new field triggers a master update, heap climbs. Fix: dynamic: strict for app indices, flattened for arbitrary key-value bags, a field-count budget per log source.

3. Oversharding#

Daily indices × 5 primaries × 1 replica × 90 days × 30 services = 27,000 shards. Each shard has fixed heap and file-handle overhead; recovery after a node restart takes hours. Fix: data streams with rollover at ~50 GB or 30 days, ILM shrink/force-merge on warm, one index per data type rather than per service.

4. Deep Pagination and Huge Aggregations#

from: 50000 or a terms aggregation over a high-cardinality field with size: 100000. Coordinating-node heap spikes; circuit breakers trip; neighbors' queries fail. Fix: search_after + PIT; composite aggregations with pagination; cap size and bucket counts in an API layer.

5. Leading Wildcards and Regex on the Hot Path#

*phone or .*shoe.* scans the whole term dictionary per shard. Fix: n-gram or wildcard field type designed for it; reject leading wildcards at the API.

6. Updating Documents at High Frequency#

Updating a view counter on a product document 1,000 times/min: each update is a full delete + reindex of the whole document, creating deleted-doc bloat and merge pressure. Fix: keep volatile counters out of the search document; refresh popularity in a periodic batch (e.g., every 15 min) or use a separate small index.

7. One Cluster for Search and Logs#

A log ingest spike (10× during an incident — exactly when you need both) starves product search. Fix: separate clusters by workload; the search cluster has an SLO, the log cluster has a budget.

Anti-PatternDetection SignalBlast RadiusWho Pays
ES as source of truthNo rebuild runbook; reindex "from ES"Permanent data lossCustomers, product
Mapping explosionField count > 1,000; slow cluster.pending_tasksWhole cluster via masterEvery index owner
OvershardingShards per node > 600; shard sizes < 1 GBRecovery time, heapPlatform on-call
Deep pagination / huge aggsCoordinating heap spikes, circuit_breaking_exceptionConcurrent queriesOther tenants
Leading wildcardsSlowlog entries with * prefixCPU on all shardsSearch latency SLO
High-frequency updatesdocs.deleted > 30% of docs; merge throttlingIndexing throughputFreshness SLO
Mixed search + logsSearch p99 tracks log ingest rateProduct searchCustomers during incidents

The Technology Landscape / Head-to-Head Comparison#

DimensionElasticsearchOpenSearchApache SolrPostgres FTS (tsvector)VespaClickHouse / Loki (logs)
Core strengthSearch + aggregations + large ecosystemES 7.10 fork, AWS-managedMature Lucene search, strong facetingSearch inside your OLTP DBLarge-scale ranking + ML servingCheap log/analytics storage
Scale modelShards across nodes, auto-rebalanceSameSolrCloud + ZooKeeperSingle primaryContent nodes, auto-distributionColumnar / object storage
Relevance toolingBM25, function_score, LTR plugin, kNN, RRFSimilar, neural pluginsBM25, LTRBasic ranking functionsFirst-class multi-phase rankingNot a relevance engine
Ops burdenMedium–high (shards, heap, ILM)Same, or managedHighNone extraHighLow–medium
License (2026)Elastic License / SSPL / AGPLv3 option (added 2024)Apache 2.0Apache 2.0PostgreSQL licenseApache 2.0Apache 2.0 / AGPL
Pick whenDefault for product search and flexible analyticsAWS-native, want Apache licenseExisting Solr expertise< ~10M docs, simple relevance, want transactional freshnessPersonalization and ranking are the productLogs at TB/day where search is secondary

Elasticsearch vs Postgres full-text: Postgres FTS with a GIN index is excellent up to ~10M documents, gives transactional freshness (searchable on commit), and costs nothing extra. It lacks BM25-style relevance tuning, fuzzy matching at scale, and distributed aggregations. Start there if search is a secondary feature; move to ES when relevance becomes a product surface.

Elasticsearch for logs: it works and it's familiar, but indexing every field of every log line costs ~1.5–3× raw size on disk plus replicas and heap. Columnar stores (ClickHouse) or label-indexed object storage (Loki) often cut log cost 5–10× when most queries are "filter by service and time, then grep." See Metrics & Monitoring.

🎯 Staff Insight: The 2021 relicense split the ecosystem into Elasticsearch and OpenSearch; Elastic added an AGPL option in 2024. In an interview, "Elasticsearch or OpenSearch — I'd pick based on whether we want Elastic's managed features or an Apache-licensed AWS service" signals you treat the license as an input, not trivia.


Patterns#

Pattern 1: Blue-Green Reindex Behind an Alias#

Applications only ever reference an alias (products). A mapping change creates products_v8, a rebuild consumer replays the source (Kafka from offset 0, or a snapshot + catch-up), dual-writes live changes to v7 and v8, validates doc counts and a relevance test suite, then swaps the alias atomically (_aliases with remove + add in one call). Rollback = swap back. Rebuild time is an SLO: measure it quarterly.

Diagram: Pattern 1: Blue-Green Reindex Behind an Alias

Pattern 2: Data Streams + ILM for Time-Series#

Logs and events go into a data stream with an ILM policy: hot (rollover at 50 GB per primary or 1 day, NVMe, 1 replica) → warm after 7 days (force-merge to 1 segment, shrink, dense disks) → cold after 30 days (0–1 replica, searchable snapshot) → frozen after 90 days (object storage) → delete at 365 days. Tiering commonly cuts storage cost 3–10× versus keeping everything hot.

Pattern 3: Search Over Multiple Sources (CQRS Read Model)#

The search document is a projection built from several services (catalog, pricing, inventory, reviews). The indexer joins them at write time. Ownership is explicit: each source team owns its event contract; the search team owns the projection and its rebuild. See Data Pipeline Patterns and Search Indexing.

BM25 for exact terms (SKUs, brand names) plus dense_vector kNN for meaning, fused with Reciprocal Rank Fusion. Memory math: 50M docs × 768 dims × 4 bytes ≈ 150 GB of vectors before graph overhead — int8 quantization cuts ~4×. Vectors change the cluster's sizing from disk-bound to RAM-bound.

Pattern 5: Tenant-Routed Multi-Tenancy#

ModelIsolationScale CeilingPick When
Shared index, tenant_id filter + routingLogicalMillions of tenantsDefault SaaS
Index per tenantPer-index settings/mappingsHundreds–low thousands (shard count explodes)Few large tenants needing custom mappings
Cluster per large tenant / cellPhysicalUnlimited via more clustersEnterprise isolation contracts

Always enforce the tenant filter in a service layer (or document-level security), never trust the client to add it.

Pattern 6: Autocomplete as a Separate Index#

A small index of popular queries and entity names with edge n-grams or the completion suggester, 1–3 shards, many replicas, fully in page cache. Keeps keystroke traffic (5–10× search QPS) off the main index.

PatternContractPrimary RiskOwner
Alias + reindexZero-downtime schema changeDual-write drift during catch-upSearch team
Data streams + ILMCost-tiered retentionRollover misconfig → giant or tiny shardsObservability platform
CQRS read modelFreshness SLOMissing deletes, event contract driftSearch team + source teams
Hybrid vectorBetter recallRAM cost, relevance regressionsSearch relevance team
Tenant routingIsolation + speedHot shard from a giant tenantSearch platform
Autocomplete indexp99 < 50 msStale suggestionsSearch team

Scaling#

Numbers to Plan With#

QuantityTypical FigureNote
Shard size10–50 GBRecovery of a 50 GB shard over 1 Gbps ≈ 7 min; over 10 Gbps ≈ 1 min (throttled by default)
Indexing per data node~10K–50K docs/s (1 KB docs, 1 replica)Varies 5× with analyzers and mapping
Search latency (well-sized)10–50 ms p50, 100–300 ms p99Aggregations and deep queries far more
Heap≤ ~30 GB per nodeRest to page cache
Shards per nodeHundreds, cap 1,000 defaultFewer, bigger shards beat many small ones
Disk watermarks85% low / 90% high / 95% flood stageFlood stage makes indices read-only
Storage amplification~1.1–2× raw source per copy, × (1 + replicas)Depends on _source, doc_values, n-grams

Scaling Reads vs Writes#

  • Read QPS: add replicas (each replica is another copy serving queries) and nodes to hold them. Adaptive replica selection routes to the least-loaded copy.
  • Write throughput: add primaries (needs reindex/split) and nodes; use bulk requests of 5–15 MB; increase refresh_interval; reduce replicas during backfills.
  • Query latency: fewer shards per query (routing, time-based index selection), smaller result windows, filters in filter context (cached, no scoring), and pre-computed fields instead of scripts.

Capacity Sketch#

Catalog: 200M products × 3 KB source → ~600 GB raw
Index size per copy ≈ 1.2 × raw ≈ 720 GB → with 1 replica ≈ 1.44 TB
Target 30 GB shards → 24 primaries (720 / 30)
Data nodes: 64 GB RAM, 30 GB heap, 1.5 TB NVMe, keep disk < 70% → ~1 TB usable each → 2 nodes of data minimum,
  but page cache wants hot set in RAM: aim ~10–20% of index in page cache → 150–300 GB RAM → 6–9 data nodes
Search peak 5K QPS × 24 shards = 120K shard-queries/s → ~15–20K per node at 8 nodes → CPU check before launch

Multi-Region#

OptionMechanismPick When
Independent regional clusters fed by the same event logEach region's indexer consumes the global Kafka topic (or a replicated one)Default — each region rebuildable, no cross-region coupling
Cross-cluster replication (CCR, licensed)Follower indices pull from leaderRead-only replicas of one writer region
Cross-cluster search (CCS)Query fans out to remote clustersOccasional global queries over regional data (logs)

🎯 Staff Move: "I'd rather replicate the event log than the index. Each region builds its own index from the same Kafka topic, so a region is self-sufficient, and a corrupted index in one region never propagates to another."


Failure Modes & Recovery#

1. Red Cluster — Unassigned Primary Shards#

Symptom: Cluster health red; searches return partial results (_shards.failed > 0); writes to affected indices fail.

Root cause: A node with the only copy of a primary died (replicas 0), or disk hit the high watermark and allocation refused, or a corrupted shard.

Detection: cluster.health.status, unassigned_shards, _cluster/allocation/explain for the reason.

Fix: Bring the node back or restore the index from snapshot; if the index is derived, rebuild it from the source — the fastest path when a rebuild runbook exists. allocate_stale_primary accepts data loss and should be a last resort with sign-off.

Prevention: ≥ 1 replica on anything user-facing; shard allocation awareness across AZs (cluster.routing.allocation.awareness.attributes: zone) so a replica never shares an AZ with its primary. Owner: search platform.

2. Disk Flood Stage → Indices Go Read-Only#

Symptom: Writes fail with cluster_block_exception ... read-only-allow-delete; indexer lag climbs; freshness SLO breached.

Root cause: A node crossed 95% disk. Often a log spike, a failed ILM rollover producing one giant index, or a merge temporarily needing extra space.

Detection: Per-node disk %; indexer_lag_seconds; ILM step errors (_ilm/explain).

Fix: Free space (delete old indices, add nodes, move shards); since 7.4, the block releases automatically once usage falls below the high watermark.

Prevention: Alert at 75% per node; ILM delete phase always configured; capacity plan with 30% free headroom. Owner: observability platform for log clusters, search platform for app clusters.

3. Heap Pressure and GC Death Spiral#

Symptom: Long old-gen GC pauses (> 5 s), nodes leaving and rejoining the cluster, cascading shard relocations that add more load.

Root cause: fielddata on text fields, huge terms aggregations, too many shards per node, oversized bulk requests, or cluster state bloat.

Detection: jvm.mem.heap_used_percent > 85% sustained; GC time > 10% of wall-clock; breakers.*.tripped counts.

Fix: Kill offending tasks (_tasks/_cancel), clear fielddata cache, temporarily lower cluster.routing.allocation.node_concurrent_recoveries to stop relocation storms.

Prevention: Disallow fielddata on text; cap aggregation sizes at the API; shard budget per node. Owner: search platform; query owners for abusive aggregations.

4. Stale or Ghost Documents (Silent Pipeline Failure)#

Symptom: Users find products that were deleted hours ago, or new listings never appear. Cluster is green.

Root cause: Indexer consumer stuck on a poison message; mapping rejects (mapper_parsing_exception) silently dropped; deletes not emitted by the source; replayed events without external versioning overwrote newer docs.

Detection: Consumer lag on the change topic; DLQ depth; end-to-end freshness probe (write a canary document to the source every minute, alert if not searchable within 30 s); periodic count reconciliation (source vs index, per tenant).

Fix: Replay from the last good offset with external versioning; drain DLQ after mapping fix; targeted rebuild for affected tenants.

Prevention: DLQ with alerting, canary freshness probe, nightly reconciliation job. Owner: search team — the pipeline is theirs even though ES is healthy.

5. Hot Shard / Hot Node#

Symptom: One node at 100% CPU, search p99 dominated by that node; others idle.

Root cause: A giant tenant with custom routing, uneven shard placement, a new index whose shards all landed on new empty nodes, or heavy queries targeting one time range.

Detection: Per-node search/indexing thread pool queue and rejections (thread_pool.search.rejected); per-shard size skew.

Fix: Move shards (_cluster/reroute), increase routing_partition_size, set index.routing.allocation.total_shards_per_node to force spreading.

Prevention: Isolate top-N tenants to dedicated indices or cells; balance checks in CI for index templates. Owner: search platform.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Red clusterunassigned_shards > 0 on primariesAffected indicesRestore or rebuild from sourceSearch platform
Flood stageNode disk > 95%Writes on affected indicesFree space, add nodesPlatform
Heap / GC spiralHeap > 85%, GC > 10%Whole cluster via node churnCancel tasks, stop relocationsPlatform + query owners
Silent stalenessFreshness probe > 30 s, DLQ depthUser trustReplay with versioningSearch team
Hot shardThread-pool rejections on one nodeTail latency for allReroute, isolate tenantPlatform
Mapping explosionField count, pending tasksMaster, all indicesStrict mappings, flattenedIndex owners
Snapshot failuresLast successful snapshot age > 24 hRecovery optionsFix repository, re-runPlatform
Diagram: Operational Reality Matrix

When to Use vs. Alternatives#

RequirementPickWhyNot Elasticsearch Because
Relevance-ranked full-text over 10M+ docsElasticsearch / OpenSearchInverted index, BM25, analyzers—
Faceted navigation + filters at high QPSElasticsearchdoc_values aggregations, filter cache—
Ad-hoc log exploration, moderate volumeElasticsearchFlexible queries, Kibana—
Search inside a small app (< 10M rows)PostgreSQL FTSNo pipeline, transactional freshnessPipeline + cluster cost unjustified
System of recordPostgreSQL / DynamoDBDurability, transactionsNot a database
Logs at 10+ TB/day, mostly time+service filtersClickHouse / Loki / object storage5–10× cheaperIndexes every field
MetricsTime Series DBsCompression, downsamplingPer-document overhead
Key lookups by ID at 100K QPSRedis / DynamoDBO(1), cheapScatter-gather cost
Pure semantic search at 1B vectorsDedicated vector DB or ANN serviceMemory/compute specializedRAM cost of HNSW on ES nodes

When NOT to use Elasticsearch: as the only copy of any data; for strongly consistent reads-after-write across users; for joins-heavy relational queries; for metrics; for high-frequency counter updates.


Operational Concerns#

Settings Every Production Cluster Should Have#

3 dedicated master-eligible nodes; data nodes spread across 3 AZs with allocation awareness
heap = min(50% RAM, ~30 GB); bootstrap.memory_lock: true; swap off
index templates with dynamic: strict (app) or explicit field budget (logs)
action.destructive_requires_name: true      # no wildcard deletes
search.max_buckets and API-layer caps on size/from/aggregation cardinality
index.max_result_window left at 10,000 — use search_after + PIT for deep paging
ILM on every time-based index; snapshots to object storage every 30–60 min (SLM)
slowlog thresholds: query warn 1s / info 500ms; indexing warn 5s

Upgrades#

Rolling upgrades within a major version: disable replica allocation, stop one node, upgrade, restart, re-enable, wait for green — node by node. Major version upgrades can't read indices from two majors back, so old indices must be reindexed or they block the upgrade — another reason the rebuild path matters. For large clusters, a blue-green cluster upgrade (build new cluster, dual-index from Kafka, switch traffic) is safer than rolling.

The Dashboard the On-Call Actually Uses#

MetricHealthyPage
Cluster statusgreenred (any) / yellow > 30 min
Search p99 (per index)< 300 ms> 1 s for 10 min
End-to-end freshness (canary)< 5 s> 60 s
Indexer consumer lag< 10 s> 5 min
Heap used % (max node)< 75%> 85% sustained
Search / write thread-pool rejections0Sustained > 0
Disk used % (max node)< 70%> 85%
Last successful snapshot< 1 h> 24 h
Shards per node (max)< 500> 900

What the On-Call Actually Does#

  1. _cluster/health and _cat/shards?h=index,shard,prirep,state,unassigned.reason — what's unassigned and why.
  2. _cat/nodes?v&h=name,heap.percent,cpu,disk.used_percent — find the hot or full node.
  3. _tasks?actions=*search*&detailed — find and cancel the runaway query.
  4. Check the indexer: consumer lag, DLQ, canary freshness — a green cluster can still serve wrong answers.
  5. Prefer rebuilding a derived index over heroic shard surgery when a rebuild runbook exists.

Interview Application — Staff-Level Plays#

Which Case Studies Use Elasticsearch#

Case StudyHow Elasticsearch Is UsedKey Pattern
Search IndexingThe search tier itselfCDC → Kafka → indexer, alias-based reindex, relevance ownership
Metrics & MonitoringLog search and explorationData streams, ILM tiers, cost per GB retained
Maps & GeospatialGeo-filtered search (geo_point, geo_distance)Combine geo filter with text relevance
News FeedPost search and hashtag discoveryTime-bounded indices, recency boosting
Chat MessagingMessage search per user/conversationRoute by user or conversation, index a projection
Web CrawlerSearchable index of crawled documentsBulk indexing, dedupe by content hash

Every System Design Question Has an Elasticsearch Moment#

  • E-commerce: facets (brand, price range, size) are keyword/numeric aggregations in filter context — cached and unscored; the text query alone drives relevance.
  • Ticketing: event search is ES; seat availability is never read from ES — it's stale by design.
  • Chat: message search routes by conversation_id so a query hits one shard; permission checks are filters, not post-processing.
  • Observability: logs for 14 days hot in ES, 1 year in object storage, queried rarely through searchable snapshots.

L5 → L6 → L7 Responses#

ScenarioSenior (L5)Staff (L6)Principal (L7)
"Add search""Index documents in Elasticsearch and query with match.""Derived index fed by outbox → Kafka → idempotent indexer with external versioning; alias over versioned index; 5 s freshness SLO measured by a canary; rebuild in 3 h.""Is search a product surface or a feature? If it's a surface, fund a relevance owner and an offline evaluation set — infra alone won't make results good. If a feature, Postgres FTS may be enough."
"Change the mapping""Update the mapping.""Most changes need a reindex: build v8 from the log, dual-write, validate counts and relevance suite, swap alias, keep v7 7 days.""Rebuild time is an org SLO. I'd require every index to have a tested rebuild under 6 hours — it's also our DR and major-upgrade path."
"Cluster is slow""Add nodes.""Check shard count and size first — oversharding and heavy aggregations are the usual causes. Cap aggregation sizes at the API, move logs off the search cluster.""Why do search and logs share a cluster? Separate by workload and give each a budget owner; the log cluster's cost curve shouldn't threaten a revenue SLO."
"Store logs""Elasticsearch with daily indices.""Data streams, rollover at 50 GB, hot 7 d → warm → frozen, 1 replica hot only; ~$X per TB-month computed per tier.""At 20 TB/day, ES costs millions per year. I'd price a columnar/object-storage stack and keep ES only for the 7-day interactive window."
"ES or OpenSearch?""Elasticsearch — it's the standard.""Whichever our cloud manages well; features we need (kNN, CCR) decide.""License and vendor risk are architectural inputs: keep clients on the common API subset so switching stays a two-way door."
Why "Add search" separates levels

Indexing documents and running match queries works in a demo and is how most teams start. The Staff answer knows the failure lives in the pipeline — missed deletes, out-of-order updates, no rebuild — and builds versioning, freshness measurement, and a rebuild path in from day one. The Principal answer notices that "good search" is a relevance problem with a product owner and an evaluation loop, and that without one the best infrastructure still returns bad results.

Why "Store logs" separates levels

Elasticsearch for logs is a well-trodden, reasonable choice. The Staff answer makes it sustainable with lifecycle tiers and shard discipline. The Principal answer puts a dollar figure on retention and asks whether a search engine is the right storage tier for data that is written once and read almost never.

The Staff Elasticsearch Checklist#

  1. Name the source of truth: "Postgres owns products; the index is derived and rebuildable in about 3 hours."
  2. Design the pipeline: "Outbox → Kafka keyed by product ID → idempotent indexer with external versioning; deletes are events."
  3. Design the mapping: "dynamic: strict; title as text + keyword; facets as keyword with doc_values; volatile counters kept out."
  4. Size the shards: "720 GB per copy at 30 GB per shard → 24 primaries, 1 replica across 3 AZs."
  5. State the freshness budget: "5 s end to end, measured by a canary document, alert at 60 s."
  6. Protect the cluster: "API caps on page depth and aggregation size, no leading wildcards, logs on a separate cluster."

🎯 Staff Insight: What NOT to use Elasticsearch for: the only copy of anything, inventory or balance reads that must be current, or high-frequency counters. "Search finds the candidates; the source of truth confirms them" — for example, re-checking stock on the product page — is the sentence that shows you know its contract.


The Principal Lens#

Why L7 Sees This Problem Differently#

At Staff level, Elasticsearch is a cluster to size and a pipeline to make correct. At Principal level, it's usually two different businesses wearing the same logo: a revenue-critical product search with a latency SLO and a relevance owner, and a fast-growing log store whose cost scales with every service's verbosity. Treating them as one platform makes the log bill threaten the search SLO and hides a multi-million-dollar retention decision inside an infrastructure budget. The L7 work is separating those, pricing retention explicitly, and standardizing the rebuild path so that mapping changes, major upgrades, license moves, and disasters all use the same tested procedure.

The Org-Level Fault Line#

One central search/log platform vs team-owned clusters — and search vs logs on the same stack.

OptionWhat WorksWhat BreaksWho Pays
Each team runs its own clustersAutonomy, custom tuning30 clusters, inconsistent ILM, oversharding everywhere, no shared rebuild toolingEvery team's on-call; duplicated expertise
One mega-cluster for everythingOne team, one billLog spikes degrade product search; mapping explosion from 400 servicesSearch customers during incidents
Platform with workload-separated clusters (search cells + observability stack)Isolation by SLO class, shared tooling, chargeback per TBPlatform team capacity; teams must accept standard templatesPlatform headcount; teams pay per GB retained

Cost Model#

Assumptions: managed ES/OpenSearch on memory-optimized data nodes (~8 vCPU / 64 GB ≈ $600/month each, approximate), SSD storage ~$0.10–0.15/GB-month, object storage $0.02/GB-month for snapshots and frozen tier, 3 small dedicated masters ($300/month total at small scale). Engineer ~$25K/month loaded.

ScaleFootprintInfra $/monthPeopleOn-call LoadDominant Risk
Startup — 5M docs product search, 50 QPS3 small nodes (combined roles), 1 replica~$1–1.5K~0.2 FTERareES as source of truth; no rebuild path
Growth — 200M-doc catalog (1.4 TB with replica) + 2 TB/day logs, 14-day hot retentionSearch: 8 data + 3 masters; Logs: ~20 hot + 10 warm nodesSearch ~$6K; logs ~$25–35K1–2 search + 1–2 observability engineers2–5 pages/monthLog cost growth, mapping explosion, freshness bugs
Enterprise — 3 regional search cells + 20 TB/day logs, 30 d hot/warm, 1 y frozen~40 search nodes; ~200 log nodes + frozen tier on object storageSearch ~$30K; logs ~$150–250KSearch platform 3–4, relevance team 3–5, observability 4–610+ pages/month fleet-wideCorrelated failure: bad template fleet-wide, major-upgrade debt, license/vendor change

The biggest lever at enterprise scale is log retention policy: dropping DEBUG in production, sampling high-volume INFO, and moving days 8–30 to a frozen tier routinely cuts the log bill 50%+ — a product-and-policy decision, not an infrastructure one.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal CostWhy
Refresh interval, replicas, query DSLTwo-wayMinutesDynamic settings
Mapping of an existing fieldOne-way per indexReindex (hours) — two-way if a rebuild path existsLucene segments are immutable
Primary shard countOne-way per index_split/_shrink constraints or reindexRouting depends on it
ES as the only copy of dataOne-wayBuild a source of truth after the fact; data may be lostCan't rebuild from nothing
Vendor features (licensed CCR, ML, security features) / proprietary pluginsMostly one-wayRe-implement on the other forkES and OpenSearch APIs keep diverging
Logs in a search engineTwo-way but expensiveRe-platform ingestion + dashboards: quartersEvery team's dashboards depend on it

The Standard I'd Write#

RFC: Search and Log Indexing Standard (v1)

Scope: All Elasticsearch/OpenSearch clusters and indices in production.

MUST

  • Declare the source of truth for every application index and maintain a rebuild runbook, tested every 6 months, with measured duration.
  • Access indices through aliases only; mapping changes ship as new versioned indices.
  • Use dynamic: strict for application indices; log sources declare a field budget (default 200).
  • Attach an ILM policy with a delete phase to every time-based index.
  • Separate clusters for revenue-critical search and for logs/observability.
  • Enforce API-layer limits: page depth ≤ 10,000 (deep paging via search_after), aggregation bucket caps, no leading wildcards.

SHOULD

  • Measure end-to-end freshness with a canary document per index.
  • Use only the API subset common to Elasticsearch and OpenSearch unless an architecture review approves otherwise.

Exceptions: Approved by the search platform team; time-boxed to 2 quarters.

Success metrics: 100% of app indices with a tested rebuild < 6 h; log $ per GB-ingested down 40% in a year; zero search incidents caused by log workloads; shards per node < 500 fleet-wide.

What I'd Tell the VP#

Search is a revenue feature and logs are an operating cost, but today they run on the same infrastructure, so a spike in logging during an incident can slow down product search exactly when we need both. I want to separate them, which costs about one extra cluster and removes that risk. On logs, we're paying roughly $200K a month mostly to keep data nobody reads after the first week; a retention policy with cheaper storage tiers should cut that by about half. Finally, every search index will be rebuildable from our primary databases within a few hours, which turns disasters, upgrades, and schema changes into routine operations instead of emergencies.

Principal Interview Signals#

SignalWhat It Sounds Like
Separates workloads by business class"Product search has an SLO and a revenue owner; logs have a budget. They shouldn't share a failure domain."
Prices retention"Seven days hot, then frozen — at 20 TB/day that's the difference between ~$250K and ~$80K a month."
Rebuild as a universal procedure"The same rebuild path covers mapping changes, major upgrades, a fork migration, and DR. I test it twice a year and publish the duration."
Relevance as an org function"Without an evaluation set and an owner, relevance regresses silently. I'd fund that before adding vector search."
Vendor/license awareness"We stay on the common API subset so moving between Elasticsearch and OpenSearch stays a two-way door."

Staff answers that L7 interviewers find insufficient:

  • "Rollover at 50 GB with hot-warm-cold ILM." — Correct mechanics; doesn't ask whether the data should be retained at all or who pays for it.
  • "We'll reindex behind an alias." — Right for one index; silent on whether every team can, and how long it takes.
  • "Add vector search to improve relevance." — A technique, not a measured outcome; no evaluation loop or owner.

🧭 Principal Move: "Before sizing this cluster I'd ask two questions: can we rebuild it from somewhere else, and who decides how long we keep the data? If nobody owns either answer, that's the design problem — the shard count is arithmetic."


In the Wild#

Wikipedia — CirrusSearch#

Wikimedia publicly runs Wikipedia's on-site search on Elasticsearch through its CirrusSearch extension, with indices updated from page-edit events and clusters in multiple data centers. Indices are rebuilt from the wiki databases when mappings or analyzers change, and language-specific analysis chains handle hundreds of languages.

Staff insight: The database of record (the wiki's MediaWiki databases) never moves into Elasticsearch; the index is a projection that can be rebuilt. That's the contract to state first in any search interview.

GitHub publicly described how its earlier code search on Elasticsearch struggled with the scale and shape of code (exact substring and symbol search across huge repositories), and how it built a custom Rust engine ("Blackbird") with n-gram indices, launched in 2023. GitHub continued to use Elasticsearch for other search surfaces.

Staff insight: Knowing when a general-purpose engine stops fitting — when the query shape (substring over code) fights the inverted-index model — is a build-vs-buy judgment. See Build vs Buy.

Stack Overflow — Search on Elasticsearch#

Stack Overflow has publicly documented using Elasticsearch for site search alongside SQL Server as the system of record, on a deliberately small hardware footprint. Search indices are fed from the primary database and can be rebuilt from it.

Staff insight: A small, well-understood cluster behind a relational source of truth serves a very large site — scale in search comes from shard discipline and a clean pipeline, not node count.


Practice Drill#

Prompt: "Your marketplace's search shows sold-out and deleted listings for up to 40 minutes, and last week a mapping change required a 'full reindex' that took the team two days and caused a partial outage. Search and logs share one 60-node cluster. What's your plan?"

Staff Answer

Three problems, in order of customer impact. Freshness: 40-minute staleness is a pipeline problem, not a cluster problem. I'd trace a listing change end to end: likely the indexer batches by time with retries blocking the partition, deletes aren't emitted as events, or refresh=false with a long interval. Fix: outbox → Kafka keyed by listing ID, idempotent indexer with version_type=external, deletes as tombstone events, DLQ for mapping rejects, and a canary listing updated every minute with an alert if it isn't searchable in 30 s. Target: 5 s end to end. Stock is also re-checked at render time from the source of truth, so a stale index can't sell a sold-out item. Reindex: introduce aliases, build listings_v2 with replicas 0 and refresh −1 from a Kafka replay or a source snapshot plus catch-up, validate counts per category and a relevance test set, swap the alias, keep v1 for 7 days — target under 4 hours with no outage. Shared cluster: move logs to their own cluster with ILM (hot 7 d, then frozen) so incident log spikes stop hurting search. Metrics: canary freshness, consumer lag, DLQ depth, search p99, per-node heap. Owners: search team owns the pipeline and rebuild; observability platform owns the log cluster.

Why this is L6:

  • Diagnoses staleness in the pipeline, not the cluster, and makes freshness measurable end to end.
  • Turns reindexing into a routine, reversible alias swap with validation.
  • Separates workloads by SLO and names owners and metrics.

What L7 adds:

  • Makes "tested rebuild under 6 hours" a platform-wide standard, since the same gap likely exists in other teams' indices.
  • Prices log retention and introduces chargeback so log growth has an owner.
  • Recognizes that showing unavailable listings is a trust and revenue issue, and sets a freshness SLO that the business agrees to.

Quick Reference Card#

Contract:      derived read model — source of truth elsewhere, rebuildable in hours
Visibility:    near-real-time; refresh 1 s default; freshness = pipeline lag + refresh
Durability:    translog fsync per request (default); replicas = extra indexing cost
Shards:        primaries fixed at creation; 10–50 GB each; < ~200M docs; < 1,000 per node
Heap:          ≤ ~30 GB and ≤ 50% RAM; rest is page cache (that's where speed lives)
Masters:       3 dedicated master-eligible nodes; allocation awareness across 3 AZs
Reads:         query-then-fetch across one copy of every shard; p99 = slowest shard
Paging:        max_result_window 10,000; search_after + PIT for deep paging
Mappings:      dynamic: strict; text + keyword multi-fields; no fielddata on text
Pipeline:      outbox → Kafka keyed by doc ID → bulk 5–15 MB, external versioning, deletes as events
Schema change: new versioned index + alias swap; keep old 7 days
Logs:          data streams + ILM; rollover 50 GB; hot 7 d → warm → frozen → delete
Aggregations:  terms and cardinality are approximate across shards; composite for exact exports
Bulk loads:    replicas 0 + refresh_interval -1, then restore; 2–3× faster ingest
Routing:       route by tenant/conversation so a query hits 1 shard instead of N
Relevance:     judgment set + offline nDCG before any boost change; then A/B
Disk:          watermarks 85 / 90 / 95% (flood = read-only); alert at 75%
Red flags:     ES as only copy, refresh=true per request, leading wildcards, deep from/size,
               search and logs on one cluster, daily indices × 5 shards, counters in search docs
  1. Loading the index…