Elasticsearch / OpenSearch concept page: distributed Lucene cluster (3 master + 3 data + coord), inverted index, sharding, query/fetch phases, aggregations, split-brain, mapping explosion
Реляционная БД отлично хранит и достаёт по primary key. Но как только пользователь хочет искать по любому слову в любом поле — LIKE '%iphone%' ставит на колени даже отлично индексированный Postgres. Полнотекстовый поиск, ранжирование по релевантности, faceted-фильтры, агрегации по миллиардам логов — это другой класс задачи, нужен inverted index + sharding + distributed merge.
Elasticsearch (и его OSS-форк OpenSearch) — это distributed Lucene: каждый shard — полноценный Lucene-индекс с inverted-структурой, поверх — cluster state machine, query coordinator и аналитический движок с агрегациями. Используется как secondary search index в продуктах (GitHub code search, Stack Overflow), как основной backend для ELK-стека (logs/metrics/observability), для e-commerce-фасетов (Amazon-style фильтры), для security analytics (SIEM).
«Elasticsearch — это distributed Lucene с координатором: документ инвертируется на одном shard, поиск делается fan-out + merge, durability через WAL (translog), видимость в search появляется через 1 секунду после индексации (refresh). Это secondary index, не source of truth.»
Три ключевых уровня:
hash(_id) % primary_shards. У каждого shard есть primary + N replicas (для HA + read throughput).Два жизненных цикла внутри shard:
node.roles: [coordinating_only]).minimum_master_nodes). Их три, потому что нужно нечётное число для кворума и переживание отказа одного.Edges:
client -> coord — единственная точка входа для приложения.coord -> shard-X-primary — write path (после hash-routing).coord -> shard-X-{p|r} — read path (coord может выбрать любую копию).shard-X-p -> shard-X-r — synchronous replication primary → replicas (wait_for_active_shards).master-N <-> master-M — consensus для cluster state.master-1 -> coord — push нового cluster state (когда меняется mapping, allocation, node up/down).Запись документа, end-to-end. Coord хеширует _id, форвардит на primary shard. Primary пишет в translog (WAL — для durability), параллельно анализирует текст (tokenize + lowercase + stemming), кладёт в in-memory buffer, форвардит на replicas. После ack'а от replicas — отвечает клиенту 201 Created. Но документ ещё не видим в search — он сидит в buffer'е до следующего refresh'а (1s default). Через 1s buffer становится новым Lucene segment, и документ появляется в выдаче.
Поиск + агрегация — два фазы. Query phase: coord фанаутит запрос на одну реплику каждого shard (round-robin), каждый shard локально находит top-K документов через Lucene + считает partial term counters для агрегации. Fetch phase: coord мерджит top-K от всех shards в global top-K, потом фанаутит запросы на конкретные shards за полным _source. Возвращает клиенту hits + aggregations. Это scatter-gather, и его стоимость = max(shard latency), плюс merge overhead.
Failure recovery. data-2 умирает (hardware failure). Master через 30s heartbeat timeout помечает node как gone, через кворум согласует новый cluster state: replica shard-B на data-1 промотируется в primary, для shard-B нужен новый replica — master выбирает свободный slot на data-3, новый primary стримит туда segments + translog. ~5 min — кластер снова GREEN. Writes во время recovery не блокируются (после промоушна), но redundancy временно деградирована (YELLOW).
Split-brain до 7.0. Сетевой раздел отрезает master-1 от master-2/3. master-2/3 выбирают нового лидера (себя), а master-1 не знает про partition и продолжает считать себя лидером. Клиенты на обеих сторонах пишут — divergent cluster states. После heal'а — данные несовместимы, merge без потери данных невозможен. Fix до 7.0: discovery.zen.minimum_master_nodes = N/2+1 — меньшая сторона не может elect leader. Fix 7.0+: vote-based quorum через cluster.initial_master_nodes. Урок: всегда 3 dedicated master nodes (нечётное), правильный quorum, и отдельные master/data роли.
Mapping explosion. Dynamic mapping без strict mode + миллионы уникальных field names (типичный случай — логи с request_param.search_query.iphone15, где значение запроса попадает в имя поля). Каждое новое поле — изменение cluster state, broadcast на все nodes. После 250K полей cluster state — 1.5 GB, master heap exhausted, full GC pauses 30s, cluster RED, индексирование отвергает запросы. Fix: index.mapping.total_fields.limit, dynamic: strict (или runtime), или тип flattened для подобъектов с произвольными ключами.
Elasticsearch vs Postgres FTS vs Algolia. Три точки на спектре ops cost / scale / features:
tsvector + GIN + опционально pg_trgm для typo-tolerance): нулевая операционка (одна БД для всего), хорош до ~10M документов и QPS < 100. Минусы: нет distributed sharding, нет BM25 (только ts_rank), aggregations через GROUP BY дорого, faceted search катастрофически медленный на nested.Решение: до 10M docs и QPS < 100 — Postgres FTS. От 10M документов или сложные facets/aggregations — Elasticsearch/OpenSearch (предпочтительно OpenSearch на AWS из-за лицензии Elastic SSPL). Algolia — только если consumer UX latency < 50ms критичен и бюджет позволяет. ES/Algolia всегда secondary index — primary store Postgres/MySQL, индексируем через CDC (Debezium / outbox pattern), а не double-write.
Refresh interval: 1s default vs увеличенный для bulk. refresh_interval определяет, через сколько после индексации документ становится виден в search. Default 1s — комфортно для NRT UX (поиск через секунду после INSERT уже находит), но создаёт ~1 segment в секунду на shard, что давит на segment merging и тратит CPU/IO.
Решение:
refresh_interval=1s (NRT обязателен для UX).refresh_interval=-1 + replicas=0, после загрузки → forced refresh + восстановить replicas → 5–10x быстрее.refresh_interval=30s — меньше segments, меньше merge, latency search +30s допустим для логов.request_param.X.Y.Z) → mapping explosion → master OOM → cluster RED.GET by id, ES не нужен.terms aggregation на поле с миллионами уникальных значений может выжрать heap координатора. Используйте composite aggregation с пагинацией.tsvector + GIN + опционально pg_trgm) дешевле в эксплуатации и достаточно функционален.durability=request vs async.completion suggester или edge ngram.