Elasticsearch в Python-сервисе почти всегда второе хранилище: товары, заказы и пользователи живут в PostgreSQL, а в индекс уезжает копия того, по чему ищут. Из этого следует всё остальное: клиент нужен для двух вещей, писать документы и выполнять запросы, а главная инженерная задача не в коде поиска, а в том, чтобы копия не отставала от правды и чтобы падение кластера не роняло сервис.
Официальный клиент один, пакет elasticsearch. В нём два класса с одинаковым набором методов: синхронный Elasticsearch и AsyncElasticsearch для asyncio. Запрос это словари Python с той же структурой, что JSON в документации кластера, ответ это объект, который читается как словарь. В статье основной инструмент асинхронный клиент, а для массовой индексации — помощники из elasticsearch.helpers.
Подключение и версии
Клиент создаётся один раз на процесс и безопасен для одновременных вызовов: внутри пул соединений и балансировка по узлам из hosts. В FastAPI его открывают в lifespan и закрывают там же через await es.close().
from elasticsearch import AsyncElasticsearch
es = AsyncElasticsearch(
hosts=["https://es-1:9200", "https://es-2:9200"],
basic_auth=("app", settings.es_password),
ca_certs="/etc/es/ca.crt",
request_timeout=2,
max_retries=3,
retry_on_timeout=False,
)
Что здесь важно.
- Мажорная версия клиента равна мажорной версии кластера. Клиент 9.x с кластером 9.x. Клиент предыдущей мажорной версии работает с новым кластером в режиме совместимости: он сам шлёт заголовок
compatible-withсо своей версией, и кластер отвечает в старом формате. Наоборот нет: при первом запросе клиент проверяет заголовокX-Elastic-Productи версию и отказывается работать с чужим сервером (UnsupportedProductError). - Повторы по кодам. По умолчанию клиент повторяет запрос на 429, 502, 503 и 504 до трёх раз, перебирая узлы. Таймаут по умолчанию не повторяется (
retry_on_timeout=False), и это правильно для поиска из обработчика: второй такой же запрос только добавит нагрузку. - Таймаут задаёт клиент, а не контекст. По умолчанию
request_timeoutравен десяти секундам на запрос, что для поиска из обработчика HTTP непозволительно много. Общий таймаут ставят в конструкторе, а для отдельного вызова переопределяют черезes.options(request_timeout=0.5).search(...):optionsвозвращает копию клиента с другими настройками, соединения у неё общие. Для фоновой переиндексации тем же способом ставят минуты. - Безопасность включена по умолчанию в образах 8.x и 9.x: HTTPS, пользователь
elastic, самоподписанный сертификат. В проде сертификат берут из секрета и кладут вca_certs, в тестах контейнер безопасность выключает.
Клиент умеет cloud_id и api_key для Elastic Cloud, а sniff_on_start=True подхватывает список узлов у самого кластера. За балансировщиком (Kubernetes Service) достаточно одного адреса.
Индекс: маппинг раньше первого документа
Если первый документ прилетит в несуществующий индекс, Elasticsearch создаст его сам и угадает типы: строка станет text с подполем .keyword, целое число long, дробное float, строка в формате ISO-даты date. Угадывание не отменить: тип поля в существующем индексе менять нельзя, только пересоздавать индекс. Поэтому маппинг задаёт код сервиса до первой записи.
PRODUCT_MAPPINGS = {
"properties": {
"title": {"type": "text"},
"description": {"type": "text"},
"category": {"type": "keyword"},
"price": {"type": "scaled_float", "scaling_factor": 100},
"updated_at": {"type": "date"},
}
}
async def create_index(es: AsyncElasticsearch, name: str) -> None:
await es.indices.create(
index=name,
settings={"number_of_shards": 1, "number_of_replicas": 1},
mappings=PRODUCT_MAPPINGS,
)
Правила, которые экономят переиндексацию: по чему ищут словами, то text; по чему фильтруют, сортируют и агрегируют, то keyword; деньги это scaled_float с множителем или целое число копеек, но не float; в индексе лежит только то, что участвует в поиске и в выдаче, а не весь агрегат целиком. Пересоздание индекса в проде делается через псевдоним, об этом ниже.
Сам документ в Python это модель Pydantic, отдельная от модели ответа API. Имена её полей и имена в маппинге обязаны совпадать, иначе поле попадёт в индекс как «угаданное».
from datetime import datetime
from decimal import Decimal
from pydantic import BaseModel
class ProductDoc(BaseModel):
id: int
title: str
description: str
category: str
price: Decimal
updated_at: datetime
model_dump(mode="json") превращает её в словарь, который клиент сериализует сам: Decimal станет числом, datetime строкой ISO.
Запись одного документа
Идентификатор документа это первичный ключ из базы, переведённый в строку. Тогда повторная индексация того же товара перезаписывает документ, а не плодит дубли, и удаление по id работает без поиска.
async def index_one(es: AsyncElasticsearch, product: ProductDoc) -> None:
await es.index(index="products", id=str(product.id), document=product.model_dump(mode="json"))
async def delete_one(es: AsyncElasticsearch, product_id: int) -> None:
await es.options(ignore_status=404).delete(index="products", id=str(product_id))
ignore_status=404 нужен удалению: без него отсутствующий документ это исключение NotFoundError, а для синхронизации «товар уже удалён» — нормальный исход.
Документ не виден в поиске сразу после записи. Индекс обновляется раз в refresh_interval, по умолчанию раз в секунду. Тест, который записал товар и тут же его ищет, получит пустую выдачу. Для таких мест есть параметр refresh:
await es.index(index="products", id=str(product.id), document=doc, refresh="wait_for")
"wait_for" дожидается ближайшего планового обновления, "true" заставляет кластер обновить индекс немедленно. Второе дорого и в массовой записи запрещено: каждый принудительный refresh создаёт новый сегмент. В обработчике «пользователь сохранил и сразу увидел» уместен wait_for, в фоновой синхронизации не нужен ни тот, ни другой.
Поиск: запрос, фильтры, страницы
Типичный запрос каталога: слова пользователя по заголовку и описанию, фильтр по категории, сортировка по релевантности, страница.
class Hit(BaseModel):
product: ProductDoc
score: float
async def search_products(es: AsyncElasticsearch, q: str, category: str | None, page: int, size: int) -> tuple[list[Hit], int]:
must = [{"multi_match": {"query": q, "fields": ["title^3", "description"], "fuzziness": "AUTO"}}]
filters = [{"term": {"category": category}}] if category else []
resp = await es.options(request_timeout=0.5).search(
index="products",
query={"bool": {"must": must, "filter": filters}},
from_=page * size,
size=size,
sort=[{"_score": "desc"}, {"updated_at": "desc"}],
track_total_hits=True,
)
hits = [
Hit(product=ProductDoc.model_validate(hit["_source"]), score=hit["_score"] or 0.0)
for hit in resp["hits"]["hits"]
]
return hits, resp["hits"]["total"]["value"]
Три вещи, которые видно только на практике.
mustсчитает релевантность,filterтолько отсекает. Категория, наличие на складе, диапазон цены идут вfilter: он кэшируется и не портит порядок выдачи. Слова пользователя идут вmust._sourceприходит словарём, и превратить его в модель — задача кода:model_validateзаодно проверит, что индекс отдаёт то, что сервис ожидает. Второй вариант: не разбирать вовсе, а взять из попаданий только идентификаторы (source=False) и достать карточки из PostgreSQL. Так выдача всегда свежая, а индекс может хранить минимум полей.totalврёт по умолчанию после 10 000: кластер перестаёт считать точно ради скорости. Для «найдено 12 345 товаров» этого хватает, для пагинации «последняя страница» нет: нуженtrack_total_hits=True.
Параметр страницы называется from_ с подчёркиванием: from в Python зарезервировано. Глубокие страницы ограничены index.max_result_window, по умолчанию 10 000 документов: from_ + size больше этого числа это BadRequestError, а не пустая выдача. Листать дальше можно только search_after, об этом в разделе «Глубже».
Массовая индексация
Записывать документы по одному это один HTTP-запрос на документ. Для первичной загрузки каталога и для переиндексации нужен bulk: один запрос несёт сотни и тысячи операций. Руками его собирать не нужно, есть помощник:
from elasticsearch import helpers
def actions(index: str, products: list[ProductDoc]):
for p in products:
yield {"_op_type": "index", "_index": index, "_id": str(p.id), "_source": p.model_dump(mode="json")}
async def bulk_index(es: AsyncElasticsearch, index: str, products: list[ProductDoc]) -> None:
ok, failed = await helpers.async_bulk(es, actions(index, products), raise_on_error=False)
if failed:
for item in failed:
op, result = next(iter(item.items()))
log.warning("bulk item failed", op=op, id=result.get("_id"), reason=result.get("error", {}).get("reason"))
raise ReindexFailed(len(failed), ok)
Bulk не атомарен. Запрос вернётся с кодом 200, а половина операций внутри может быть отклонена: у каждого элемента свой статус. Помощник разбирает ответ сам: с raise_on_error=False возвращает число успешных и список упавших, с умолчанием raise_on_error=True бросает BulkIndexError после обработки пачки, и упавшие лежат в exc.errors. Код обязан решить, что делать с ними: повторить, записать в лог, остановить переиндексацию.
Пачки помощник режет сам: по 500 действий или по 100 МБ, что наступит раньше (chunk_size, max_chunk_bytes). Для долгих загрузок удобнее потоковый вариант: async_streaming_bulk принимает тот же генератор, отдаёт результат по каждому документу и умеет повторять пачки, отбитые кодом 429, с нарастающей паузой:
async for ok, item in helpers.async_streaming_bulk(
es, actions("products_v2", batch), chunk_size=1000, max_retries=3, initial_backoff=2, yield_ok=False,
):
op, result = next(iter(item.items()))
log.warning("reindex item failed", id=result.get("_id"), reason=result.get("error", {}).get("reason"))
yield_ok=False оставляет в цикле только неудачи, остальное помощник пропускает молча. Размер одного bulk-запроса держат в 5–15 МБ: меньше это лишние запросы, больше это память координирующего узла и риск упереться в http.max_content_length (100 МБ). Перед полной переиндексацией реплики временно ставят в 0 и refresh_interval в -1, после возвращают.
Как держать индекс в согласии с базой
Это главный вопрос интеграции, и у него четыре ответа, от простого к надёжному.
Двойная запись из обработчика. Сохранили в PostgreSQL, следом вызвали index_one. Просто и неверно: транзакция в базе уже закоммичена, а вызов в индекс упал по таймауту, и товар навсегда пропал из поиска. Повтор в цикле внутри запроса не спасает: процесс могут убить между двумя записями.
Transactional outbox. В той же транзакции, что и товар, в таблицу outbox пишется событие «товар изменён». Отдельная задача или воркер читает outbox и отправляет события в индекс (или в Kafka, откуда их читает индексатор). Гарантия «хотя бы раз» на стороне базы, повторы безопасны, потому что индексация по id идемпотентна. Рецепт целиком разобран в паттернах распределённых систем на Python.
CDC через Debezium. Коннектор читает WAL PostgreSQL и кладёт изменения строк в Kafka, потребитель на aiokafka собирает из них документы и шлёт bulk. Код сервиса про индекс не знает вообще. Цена: Kafka Connect в инфраструктуре и документ, который собирается из «сырых» строк таблиц, а не из доменной модели.
Полная переиндексация по расписанию. Ночью выгрузить всё из базы в новый индекс и переключить псевдоним. Сама по себе это не синхронизация, а страховка: любой из трёх способов выше со временем накапливает расхождения (упавшие повторы, ручные правки в базе), и регулярная пересборка их стирает.
Рабочая связка для каталога: outbox для изменений в течение дня плюс ночная пересборка. Удаление товара это тоже событие, иначе он останется в выдаче.
Псевдоним: смена индекса без простоя
Приложение никогда не ходит в индекс по имени products_v2, только по псевдониму products. Тогда переиндексация выглядит так: создать products_v3 с новым маппингом, залить его потоковым bulk, одним атомарным запросом перекинуть псевдоним, удалить старый индекс через день.
Приложение ходит только в псевдоним, поэтому смена маппинга это новый индекс и атомарное переключение, без простоя и правок кода.
async def switch_alias(es: AsyncElasticsearch, alias: str, old_index: str, new_index: str) -> None:
await es.indices.update_aliases(actions=[
{"remove": {"index": old_index, "alias": alias}},
{"add": {"index": new_index, "alias": alias}},
])
Оба действия в одном запросе применяются атомарно: нет момента, когда поиск не видит ни одного индекса. Пока идёт заливка, изменения из outbox пишут в оба индекса, иначе новый индекс к моменту переключения уже отстанет.
Когда кластер недоступен
Поиск редко бывает критичным для денег: заказ оформляется без него. Поэтому обработчик поиска получает короткий request_timeout и честный ответ 503 при ConnectionError или ConnectionTimeout, а не ожидание до бесконечности и не перекладывание запроса в PostgreSQL с ILIKE (это только положит базу вместе с поиском). У самого запроса есть серверный таймаут: timeout="500ms" в search вернёт то, что успели найти шарды, с флагом resp["timed_out"], и код решает, показывать ли частичную выдачу.
Запись в индекс через outbox недоступности не боится: события подождут в таблице, воркер повторит с паузой. Единственное, что нужно, это метрика отставания (возраст самого старого необработанного события) и алерт на неё.
Три ловушки, которые встречаются чаще других.
- Elasticsearch не транзакционен. Нет «записать два документа или ни одного», нет чтения своих незакоммиченных записей. Всё, что требует атомарности, живёт в PostgreSQL.
- Поле с неверно угаданным типом всплывает через месяц, когда понадобилась сортировка по цене, а
priceоказалсяtext. Лечится только пересозданием индекса через псевдоним. - Большие документы. Описание с HTML на мегабайт делает каждый запрос медленнее. В индексе лежит очищенный текст и поля для фильтров; картинки, историю цен и остатки по складам в него не кладут.
Тесты
Юнит-тесты сервиса поиска не ходят в кластер: репозиторий поиска это Protocol с методом search(query) -> list[Hit], и в тестах обработчика его подменяет заглушка. Сам клиент тестируется интеграционно, через контейнер:
import pytest
from testcontainers.community.elasticsearch import ElasticSearchContainer
@pytest.fixture(scope="session")
def es_url() -> str:
with ElasticSearchContainer("docker.elastic.co/elasticsearch/elasticsearch:8.17.0", mem_limit="2g") as container:
yield f"http://{container.get_container_host_ip()}:{container.get_exposed_port(9200)}"
@pytest.fixture
async def es(es_url: str, request):
client = AsyncElasticsearch(hosts=[es_url])
index = f"test-{request.node.name.lower()}"
await create_index(client, index)
yield client, index
await client.options(ignore_status=404).indices.delete(index=index)
await client.close()
Контейнер выключает безопасность для образов 8.x и 9.x (xpack.security.enabled=false), поэтому в тестах нет ни пароля, ни сертификата, в отличие от прода. Запускается он десятки секунд и просит около гигабайта памяти, поэтому такие тесты живут за маркером integration, поднимают один контейнер на сессию и используют свой индекс на каждый тест, а после записи зовут refresh="true": в тесте это допустимо.
Глубже: search_after вместо глубокой пагинациирасширенное
from_ + size заставляет каждый шард собрать from_ + size лучших документов и отдать координатору; на странице 500 это 5 000 документов с каждого шарда ради 10 нужных. search_after принимает значения сортировки последнего документа предыдущей страницы (hit["sort"]) и продолжает с них: нагрузка одинакова для первой и для тысячной страницы. Условие одно: сортировка должна быть однозначной, поэтому к _score или дате добавляют поле-разделитель (идентификатор). Чтобы выдача не «плыла» при записи между страницами, запрос открывают через point in time — await es.open_point_in_time(index="products", keep_alive="1m") — и передают pit={"id": pit_id, "keep_alive": "1m"} в каждый следующий search. Для каталога, где дальше десятой страницы никто не ходит, это избыточно; для выгрузок и обходов индекса обязательно.
Глубже: что кладут в индекс, а что оставляют в базерасширенное
Индекс это производные данные, их всегда можно пересобрать из базы, и это главный критерий: если поле нельзя восстановить из PostgreSQL, ему в индексе не место. Внутрь идёт то, по чему ищут (текст), фильтруют (категория, бренд, наличие), сортируют (цена, дата, рейтинг) и то, что показывают прямо в выдаче (название, картинка-превью, цена). Остатки по складам, которые меняются каждую минуту, в индекс не пишут: их подтягивают из базы по идентификаторам из выдачи. Тогда обновление товара в индексе происходит при редактировании карточки, а не при каждой продаже.
Коротко
- Один
AsyncElasticsearchна процесс вlifespan, версия клиента совпадает с мажорной версией кластера; таймаут по умолчанию 10 секунд, его опускают в конструкторе и переопределяют черезoptions; 429 уже в списке повторов, таймауты не повторяются. - Маппинг создаёт код до первого документа:
textдля поиска,keywordдля фильтров, деньги безfloat; документ — отдельная модель Pydantic черезmodel_dump(mode="json"). - Идентификатор документа равен первичному ключу, тогда запись идемпотентна, а удаление не требует поиска; отсутствие при удалении гасят
ignore_status=404. - Документ виден через
refresh_interval(1 с);refresh="wait_for"только там, где пользователь ждёт результат, в bulk никогда. mustсчитает релевантность,filterотсекает;_sourceразбирают черезmodel_validateили берут только id и читают карточки из базы; страница —from_с подчёркиванием.async_bulkвозвращает 200 и при частичных ошибках:raise_on_error=Falseи разбор упавших,BulkIndexErrorсerrorsиначе; для больших загрузокasync_streaming_bulkсmax_retriesиyield_ok=False.- Синхронизация с базой через outbox или CDC плюс ночная пересборка; двойная запись из обработчика теряет документы.
- Приложение ходит в индекс по псевдониму; смена индекса одним
update_aliases. - Падение кластера это 503 на поиске и очередь в outbox на записи, а не
ILIKEпо базе. - Глубокие страницы через
search_afterиopen_point_in_time,from_ + sizeограничен 10 000.
Что почитать дальше
- Как устроен Elasticsearch — инвертированный индекс, шарды, сегменты и почему refresh не бесплатен.
- Запросы и релевантность —
match,bool, анализаторы и как читать_score. - Elasticsearch в проде — ILM, снимки, размер шардов и мониторинг кластера.
- Поиск: PostgreSQL FTS или Elasticsearch — когда полнотекстового поиска в базе достаточно.
- Паттерны распределённых систем на Python — outbox и идемпотентный потребитель, на которых держится синхронизация.
- Интеграционные тесты на Python — контейнеры на сессию, маркеры и изоляция данных между тестами.