← назад к разделу

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, одним атомарным запросом перекинуть псевдоним, удалить старый индекс через день.

сейчас псевдоним products смотрит на products_v1 шаг 1 создать products_v2 с новым маппингом шаг 2 переиндексировать v1 в v2 (reindex) шаг 3 update_aliases: снять с v1, повесить на v2 одной операцией шаг 4 удалить products_v1

Приложение ходит только в псевдоним, поэтому смена маппинга это новый индекс и атомарное переключение, без простоя и правок кода.

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.

Что почитать дальше