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.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Product / content search | Relevance, p99 < 200 ms, freshness ~seconds | Derived index via CDC, BM25 + business signals, replicas for QPS | Stale or missing documents; relevance regressions | Every live item findable within ~5 s; no deleted item shown |
| Log / observability analytics | Ingest 100K–1M+ docs/s, retention weeks–months | Data streams + ILM hot/warm/cold, few replicas, rollover by size | Ingest backpressure, disk watermarks, cost | Tolerates sampling; lose minutes, never days |
| Autocomplete / typeahead | p99 < 50 ms per keystroke | Small dedicated index, edge n-grams or completion suggester | Latency under burst; index bloat | Approximate is fine |
| Security / audit search | Retention 1+ years, rare queries | Frozen tier on object storage, searchable snapshots | Slow cold queries | Completeness 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#
| Role | Job | Sizing Rule |
|---|---|---|
| Master-eligible | Cluster state: mappings, shard allocation, node membership | Exactly 3 dedicated small nodes (e.g., 4–8 GB heap) for any cluster > ~10 data nodes |
| Data (hot) | Indexing + recent queries on fast NVMe | Heap ≤ ~30 GB (compressed object pointers), ~50% of RAM left for the OS page cache |
| Data (warm / cold) | Older, read-mostly indices on dense disks | Higher disk:RAM ratios (e.g., 1:100–1:200) |
| Data (frozen) | Searchable snapshots on object storage | Local cache only; queries take seconds |
| Coordinating-only | Fan out searches, reduce results | Add when heavy aggregations stress data nodes |
| Ingest | Pipeline 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.
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#
Every search touches one copy of every shard in the target indices unless routing narrows it. Two consequences drive design:
- 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.
- Deep pagination is quadratic-ish. Page 1,000 at size 10 makes each shard return 10,010 candidates. The default
index.max_result_windowof 10,000 exists for this reason — usesearch_afterwith 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_exceptionis a sizing signal, not a bug. doc_values(on-disk columnar, default for keyword/numeric) keep sorting and aggregations off-heap. Enablingfielddataontextfields 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 Type | Indexed As | Use For | Cost / Caveat |
|---|---|---|---|
text | Analyzed tokens in inverted index | Full-text match, relevance | Not sortable/aggregatable without fielddata |
keyword | Exact term + doc_values | Filters, facets, sorting, exact IDs | Max useful length ~256 chars (ignore_above) |
Multi-field text + .keyword | Both | Search and facet on the same value | ~2× index size for that field |
Numeric / date | BKD tree + doc_values | Range filters, sorting, histograms | Use keyword for IDs you never range over — faster term lookups |
nested | Hidden sub-documents | Arrays of objects needing per-object matching | Each nested object is a Lucene doc; 10K+ per parent hurts |
flattened | Single field of keywords | Arbitrary user-defined labels | No per-key typing, limited queries |
dense_vector | HNSW graph | kNN semantic search | RAM-hungry; ~4 bytes × dims per vector (less with quantization) |
enabled: false | Stored in _source only | Display-only blobs | Not 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
| Mistake | Symptom | Fix |
|---|---|---|
| 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 shard | Rollover 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 huge | Split with _split or reindex into more primaries |
| Uneven routing | One shard 5× others | Better 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#
Three rules make the pipeline correct:
- Key the topic by document ID so all changes to one document are ordered in one partition.
- 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. - 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_interval | Visible After | Indexing Throughput | Segment Churn | Pick When |
|---|---|---|---|---|
1s (default) | ~1 s | Baseline | High — many tiny segments | User-facing search that must feel instant |
5s–30s | 5–30 s | Often 20–50%+ higher | Lower | Logs, analytics, catalogs where seconds don't matter |
-1 (disabled) | On explicit refresh | Maximum | Minimal | Bulk loads and reindexes |
?refresh=wait_for per request | Before response returns | Per-request latency ≈ up to 1 s | None extra | Read-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#
| Setting | Loss on Node Crash | Throughput | Pick When |
|---|---|---|---|
translog.durability: request (default) | 0 acked writes | Baseline | Anything not trivially re-derivable |
translog.durability: async, sync_interval: 5s | Up to ~5 s | Higher | Logs with an upstream buffer (Kafka) that can replay |
wait_for_active_shards: 2 | Write fails if a replica is unavailable | Lower availability | Rarely worth it — the source of truth is elsewhere |
Dial 3: Query Cost vs Relevance Quality#
| Technique | Relevance Gain | Cost | Where It Runs |
|---|---|---|---|
| BM25 on analyzed fields | Baseline lexical relevance | Cheap | Query |
Field boosts, multi_match best_fields / cross_fields | Large for structured catalogs | Cheap | Query |
function_score / rank_feature (popularity, recency) | Business-aligned ranking | Moderate | Query |
rescore window (top 100–500) | Expensive signals only on candidates | Bounded | Query |
| Hybrid lexical + kNN (RRF) | Handles synonyms and intent | 2× work + vector RAM | Query |
| Learning-to-rank / external re-ranker | Highest | Model serving, feature pipeline, owners | Outside 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:
| Aggregation | Exact? | Error Source | Knob | Staff Default |
|---|---|---|---|---|
terms (top-N buckets) | No, across shards | Each shard returns only its local top shard_size | shard_size (default size × 1.5 + 10) | Raise shard_size for long-tail facets; report doc_count_error_upper_bound |
cardinality | No | HyperLogLog++ | precision_threshold (exact below it, max 40,000) | Fine for dashboards, never for billing |
percentiles | No | TDigest | compression | Fine for latency dashboards |
sum / min / max / avg | Yes (within the searched data) | — | — | — |
composite (paginated) | Yes | — | Page size | Use 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_atin 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-Pattern | Detection Signal | Blast Radius | Who Pays |
|---|---|---|---|
| ES as source of truth | No rebuild runbook; reindex "from ES" | Permanent data loss | Customers, product |
| Mapping explosion | Field count > 1,000; slow cluster.pending_tasks | Whole cluster via master | Every index owner |
| Oversharding | Shards per node > 600; shard sizes < 1 GB | Recovery time, heap | Platform on-call |
| Deep pagination / huge aggs | Coordinating heap spikes, circuit_breaking_exception | Concurrent queries | Other tenants |
| Leading wildcards | Slowlog entries with * prefix | CPU on all shards | Search latency SLO |
| High-frequency updates | docs.deleted > 30% of docs; merge throttling | Indexing throughput | Freshness SLO |
| Mixed search + logs | Search p99 tracks log ingest rate | Product search | Customers during incidents |
The Technology Landscape / Head-to-Head Comparison#
| Dimension | Elasticsearch | OpenSearch | Apache Solr | Postgres FTS (tsvector) | Vespa | ClickHouse / Loki (logs) |
|---|---|---|---|---|---|---|
| Core strength | Search + aggregations + large ecosystem | ES 7.10 fork, AWS-managed | Mature Lucene search, strong faceting | Search inside your OLTP DB | Large-scale ranking + ML serving | Cheap log/analytics storage |
| Scale model | Shards across nodes, auto-rebalance | Same | SolrCloud + ZooKeeper | Single primary | Content nodes, auto-distribution | Columnar / object storage |
| Relevance tooling | BM25, function_score, LTR plugin, kNN, RRF | Similar, neural plugins | BM25, LTR | Basic ranking functions | First-class multi-phase ranking | Not a relevance engine |
| Ops burden | Medium–high (shards, heap, ILM) | Same, or managed | High | None extra | High | Low–medium |
| License (2026) | Elastic License / SSPL / AGPLv3 option (added 2024) | Apache 2.0 | Apache 2.0 | PostgreSQL license | Apache 2.0 | Apache 2.0 / AGPL |
| Pick when | Default for product search and flexible analytics | AWS-native, want Apache license | Existing Solr expertise | < ~10M docs, simple relevance, want transactional freshness | Personalization and ranking are the product | Logs 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.
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.
Pattern 4: Hybrid Lexical + Vector Search#
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#
| Model | Isolation | Scale Ceiling | Pick When |
|---|---|---|---|
Shared index, tenant_id filter + routing | Logical | Millions of tenants | Default SaaS |
| Index per tenant | Per-index settings/mappings | Hundreds–low thousands (shard count explodes) | Few large tenants needing custom mappings |
| Cluster per large tenant / cell | Physical | Unlimited via more clusters | Enterprise 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.
| Pattern | Contract | Primary Risk | Owner |
|---|---|---|---|
| Alias + reindex | Zero-downtime schema change | Dual-write drift during catch-up | Search team |
| Data streams + ILM | Cost-tiered retention | Rollover misconfig → giant or tiny shards | Observability platform |
| CQRS read model | Freshness SLO | Missing deletes, event contract drift | Search team + source teams |
| Hybrid vector | Better recall | RAM cost, relevance regressions | Search relevance team |
| Tenant routing | Isolation + speed | Hot shard from a giant tenant | Search platform |
| Autocomplete index | p99 < 50 ms | Stale suggestions | Search team |
Scaling#
Numbers to Plan With#
| Quantity | Typical Figure | Note |
|---|---|---|
| Shard size | 10–50 GB | Recovery 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 p99 | Aggregations and deep queries far more |
| Heap | ≤ ~30 GB per node | Rest to page cache |
| Shards per node | Hundreds, cap 1,000 default | Fewer, bigger shards beat many small ones |
| Disk watermarks | 85% low / 90% high / 95% flood stage | Flood 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
filtercontext (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#
| Option | Mechanism | Pick When |
|---|---|---|
| Independent regional clusters fed by the same event log | Each 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 leader | Read-only replicas of one writer region |
| Cross-cluster search (CCS) | Query fans out to remote clusters | Occasional 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#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Red cluster | unassigned_shards > 0 on primaries | Affected indices | Restore or rebuild from source | Search platform |
| Flood stage | Node disk > 95% | Writes on affected indices | Free space, add nodes | Platform |
| Heap / GC spiral | Heap > 85%, GC > 10% | Whole cluster via node churn | Cancel tasks, stop relocations | Platform + query owners |
| Silent staleness | Freshness probe > 30 s, DLQ depth | User trust | Replay with versioning | Search team |
| Hot shard | Thread-pool rejections on one node | Tail latency for all | Reroute, isolate tenant | Platform |
| Mapping explosion | Field count, pending tasks | Master, all indices | Strict mappings, flattened | Index owners |
| Snapshot failures | Last successful snapshot age > 24 h | Recovery options | Fix repository, re-run | Platform |
When to Use vs. Alternatives#
| Requirement | Pick | Why | Not Elasticsearch Because |
|---|---|---|---|
| Relevance-ranked full-text over 10M+ docs | Elasticsearch / OpenSearch | Inverted index, BM25, analyzers | — |
| Faceted navigation + filters at high QPS | Elasticsearch | doc_values aggregations, filter cache | — |
| Ad-hoc log exploration, moderate volume | Elasticsearch | Flexible queries, Kibana | — |
| Search inside a small app (< 10M rows) | PostgreSQL FTS | No pipeline, transactional freshness | Pipeline + cluster cost unjustified |
| System of record | PostgreSQL / DynamoDB | Durability, transactions | Not a database |
| Logs at 10+ TB/day, mostly time+service filters | ClickHouse / Loki / object storage | 5–10× cheaper | Indexes every field |
| Metrics | Time Series DBs | Compression, downsampling | Per-document overhead |
| Key lookups by ID at 100K QPS | Redis / DynamoDB | O(1), cheap | Scatter-gather cost |
| Pure semantic search at 1B vectors | Dedicated vector DB or ANN service | Memory/compute specialized | RAM 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#
| Metric | Healthy | Page |
|---|---|---|
| Cluster status | green | red (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 rejections | 0 | Sustained > 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#
_cluster/healthand_cat/shards?h=index,shard,prirep,state,unassigned.reason— what's unassigned and why._cat/nodes?v&h=name,heap.percent,cpu,disk.used_percent— find the hot or full node._tasks?actions=*search*&detailed— find and cancel the runaway query.- Check the indexer: consumer lag, DLQ, canary freshness — a green cluster can still serve wrong answers.
- 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 Study | How Elasticsearch Is Used | Key Pattern |
|---|---|---|
| Search Indexing | The search tier itself | CDC → Kafka → indexer, alias-based reindex, relevance ownership |
| Metrics & Monitoring | Log search and exploration | Data streams, ILM tiers, cost per GB retained |
| Maps & Geospatial | Geo-filtered search (geo_point, geo_distance) | Combine geo filter with text relevance |
| News Feed | Post search and hashtag discovery | Time-bounded indices, recency boosting |
| Chat Messaging | Message search per user/conversation | Route by user or conversation, index a projection |
| Web Crawler | Searchable index of crawled documents | Bulk indexing, dedupe by content hash |
Every System Design Question Has an Elasticsearch Moment#
- E-commerce: facets (brand, price range, size) are
keyword/numeric aggregations infiltercontext — 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_idso 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#
| Scenario | Senior (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#
- Name the source of truth: "Postgres owns products; the index is derived and rebuildable in about 3 hours."
- Design the pipeline: "Outbox → Kafka keyed by product ID → idempotent indexer with external versioning; deletes are events."
- Design the mapping: "
dynamic: strict; title as text + keyword; facets as keyword with doc_values; volatile counters kept out." - Size the shards: "720 GB per copy at 30 GB per shard → 24 primaries, 1 replica across 3 AZs."
- State the freshness budget: "5 s end to end, measured by a canary document, alert at 60 s."
- 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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Each team runs its own clusters | Autonomy, custom tuning | 30 clusters, inconsistent ILM, oversharding everywhere, no shared rebuild tooling | Every team's on-call; duplicated expertise |
| One mega-cluster for everything | One team, one bill | Log spikes degrade product search; mapping explosion from 400 services | Search customers during incidents |
| Platform with workload-separated clusters (search cells + observability stack) | Isolation by SLO class, shared tooling, chargeback per TB | Platform team capacity; teams must accept standard templates | Platform 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.
| Scale | Footprint | Infra $/month | People | On-call Load | Dominant Risk |
|---|---|---|---|---|---|
| Startup — 5M docs product search, 50 QPS | 3 small nodes (combined roles), 1 replica | ~$1–1.5K | ~0.2 FTE | Rare | ES as source of truth; no rebuild path |
| Growth — 200M-doc catalog (1.4 TB with replica) + 2 TB/day logs, 14-day hot retention | Search: 8 data + 3 masters; Logs: ~20 hot + 10 warm nodes | Search ~$6K; logs ~$25–35K | 1–2 search + 1–2 observability engineers | 2–5 pages/month | Log 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 storage | Search ~$30K; logs ~$150–250K | Search platform 3–4, relevance team 3–5, observability 4–6 | 10+ pages/month fleet-wide | Correlated 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#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversal Cost | Why |
|---|---|---|---|
| Refresh interval, replicas, query DSL | Two-way | Minutes | Dynamic settings |
| Mapping of an existing field | One-way per index | Reindex (hours) — two-way if a rebuild path exists | Lucene segments are immutable |
| Primary shard count | One-way per index | _split/_shrink constraints or reindex | Routing depends on it |
| ES as the only copy of data | One-way | Build a source of truth after the fact; data may be lost | Can't rebuild from nothing |
| Vendor features (licensed CCR, ML, security features) / proprietary plugins | Mostly one-way | Re-implement on the other fork | ES and OpenSearch APIs keep diverging |
| Logs in a search engine | Two-way but expensive | Re-platform ingestion + dashboards: quarters | Every 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: strictfor 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#
| Signal | What 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 — From Elasticsearch to a Purpose-Built Code Search#
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