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

Elasticsearch в Go-сервисе почти всегда второе хранилище: товары, заказы и пользователи живут в PostgreSQL, а в индекс уезжает копия того, по чему ищут. Из этого следует всё остальное: клиент нужен для двух вещей, писать документы и выполнять запросы, а главная инженерная задача не в коде поиска, а в том, чтобы копия не отставала от правды и чтобы падение кластера не роняло сервис.

Официальный клиент один, github.com/elastic/go-elasticsearch/v8. Внутри него два клиента: низкоуровневый, где запрос это JSON в io.Reader и ответ это сырое тело, и типизированный, где запрос собирается из структур пакета typedapi. В статье основной инструмент типизированный, низкоуровневый появляется там, где он удобнее: в массовой индексации.

Обязательно

Подключение и версии

Клиент создаётся один раз на процесс и безопасен для горутин: внутри пул соединений net/http и балансировка по узлам из Addresses.

import "github.com/elastic/go-elasticsearch/v8"

es, err := elasticsearch.NewTypedClient(elasticsearch.Config{
    Addresses:     []string{"https://es-1:9200", "https://es-2:9200"},
    Username:      "app",
    Password:      "secret",
    CACert:        caPEM,
    RetryOnStatus: []int{502, 503, 504, 429},
    MaxRetries:    3,
})

Что здесь важно.

  • Мажорная версия клиента равна мажорной версии кластера. Клиент 8.x с кластером 8.x. Клиент 7.x умеет говорить с кластером 8.x в режиме совместимости (ELASTIC_CLIENT_APIVERSIONING=true), наоборот нет: при первом запросе клиент проверяет заголовок X-Elastic-Product и отказывается работать с чужим сервером.
  • Повторы по кодам. По умолчанию клиент повторяет запрос на 502, 503 и 504 до трёх раз, перебирая узлы. 429 (кластер отбивает нагрузку) в список стоит добавить руками: это именно та ситуация, когда через секунду станет легче.
  • Таймауты задаёт не клиент, а контекст. Каждый вызов заканчивается .Do(ctx); дедлайн контекста обрывает HTTP-запрос. Для поиска из обработчика HTTP это 300-500 мс, для фоновой переиндексации минуты.
  • Безопасность включена по умолчанию в образах 8.x: HTTPS, пользователь elastic, самоподписанный сертификат. В проде сертификат берут из секрета и кладут в CACert, в тестах его отдаёт контейнер.

Клиент умеет CloudID и APIKey для Elastic Cloud, а DiscoverNodesOnStart подхватывает список узлов у самого кластера. За балансировщиком (Kubernetes Service) достаточно одного адреса.

Индекс: маппинг раньше первого документа

Если первый документ прилетит в несуществующий индекс, Elasticsearch создаст его сам и угадает типы: строка станет text с подполем .keyword, целое число long, дробное float, строка в формате ISO-даты date. Угадывание не отменить: тип поля в существующем индексе менять нельзя, только пересоздавать индекс. Поэтому маппинг задаёт код сервиса до первой записи.

import (
    "github.com/elastic/go-elasticsearch/v8/typedapi/indices/create"
    "github.com/elastic/go-elasticsearch/v8/typedapi/types"
)

func createIndex(ctx context.Context, es *elasticsearch.TypedClient, name string) error {
    _, err := es.Indices.Create(name).Request(&create.Request{
        Settings: &types.IndexSettings{
            NumberOfShards:   strPtr("1"),
            NumberOfReplicas: strPtr("1"),
        },
        Mappings: &types.TypeMapping{
            Properties: map[string]types.Property{
                "title":       types.NewTextProperty(),
                "description": types.NewTextProperty(),
                "category":    types.NewKeywordProperty(),
                "price":       types.NewScaledFloatNumberProperty(),
                "updated_at":  types.NewDateProperty(),
            },
        },
    }).Do(ctx)
    return err
}

Правила, которые экономят переиндексацию: по чему ищут словами, то text; по чему фильтруют, сортируют и агрегируют, то keyword; деньги это scaled_float с множителем или целое число копеек, но не float; в индексе лежит только то, что участвует в поиске и в выдаче, а не весь агрегат целиком. Пересоздание индекса в проде делается через псевдоним, об этом ниже.

Сам документ в Go это обычная структура с json-тегами. Имена полей в тегах и в маппинге обязаны совпадать, иначе поле попадёт в индекс как «угаданное».

type Product struct {
    ID          int64     `json:"id"`
    Title       string    `json:"title"`
    Description string    `json:"description"`
    Category    string    `json:"category"`
    Price       float64   `json:"price"`
    UpdatedAt   time.Time `json:"updated_at"`
}

Запись одного документа

Идентификатор документа это первичный ключ из базы, переведённый в строку. Тогда повторная индексация того же товара перезаписывает документ, а не плодит дубли, и удаление по id работает без поиска.

func indexOne(ctx context.Context, es *elasticsearch.TypedClient, p Product) error {
    _, err := es.Index("products").Id(strconv.FormatInt(p.ID, 10)).Request(p).Do(ctx)
    return err
}

func deleteOne(ctx context.Context, es *elasticsearch.TypedClient, id int64) error {
    _, err := es.Delete("products", strconv.FormatInt(id, 10)).Do(ctx)
    return err
}

Документ не виден в поиске сразу после записи. Индекс обновляется раз в refresh_interval, по умолчанию раз в секунду. Тест, который записал товар и тут же его ищет, получит пустую выдачу. Для таких мест есть параметр Refresh:

import "github.com/elastic/go-elasticsearch/v8/typedapi/types/enums/refresh"

_, err := es.Index("products").Id(id).Request(p).Refresh(refresh.Waitfor).Do(ctx)

refresh.Waitfor дожидается ближайшего планового обновления, refresh.True заставляет кластер обновить индекс немедленно. Второе дорого и в массовой записи запрещено: каждый принудительный refresh создаёт новый сегмент. В обработчике «пользователь сохранил и сразу увидел» уместен Waitfor, в фоновой синхронизации не нужен ни тот, ни другой.

Поиск: запрос, фильтры, страницы

Типичный запрос каталога: слова пользователя по заголовку и описанию, фильтр по категории, сортировка по релевантности, страница.

import (
    "github.com/elastic/go-elasticsearch/v8/typedapi/core/search"
    "github.com/elastic/go-elasticsearch/v8/typedapi/types/enums/sortorder"
)

func searchProducts(ctx context.Context, es *elasticsearch.TypedClient,
    q, category string, page, size int) ([]Hit, int64, error) {
    must := []types.Query{{
        MultiMatch: &types.MultiMatchQuery{
            Query:     q,
            Fields:    []string{"title^3", "description"},
            Fuzziness: "AUTO",
        },
    }}
    var filter []types.Query
    if category != "" {
        filter = append(filter, types.Query{
            Term: map[string]types.TermQuery{"category": {Value: category}},
        })
    }
    res, err := es.Search().Index("products").Request(&search.Request{
        Query: &types.Query{Bool: &types.BoolQuery{Must: must, Filter: filter}},
        From:  intPtr(page * size),
        Size:  intPtr(size),
        Sort: []types.SortCombinations{
            types.SortOptions{Score_: &types.ScoreSort{Order: &sortorder.Desc}},
            types.SortOptions{SortOptions: map[string]types.FieldSort{
                "updated_at": {Order: &sortorder.Desc},
            }},
        },
    }).Do(ctx)
    if err != nil {
        return nil, 0, err
    }
    hits := make([]Hit, 0, len(res.Hits.Hits))
    for _, h := range res.Hits.Hits {
        var p Product
        if err := json.Unmarshal(h.Source_, &p); err != nil {
            return nil, 0, err
        }
        score := 0.0
        if h.Score_ != nil {
            score = float64(*h.Score_)
        }
        hits = append(hits, Hit{Product: p, Score: score})
    }
    var total int64
    if res.Hits.Total != nil {
        total = res.Hits.Total.Value
    }
    return hits, total, nil
}

Три вещи, которые видно только на практике.

  • must считает релевантность, filter только отсекает. Категория, наличие на складе, диапазон цены идут в filter: он кэшируется и не портит порядок выдачи. Слова пользователя идут в must.
  • _source приходит сырым JSON (json.RawMessage), разбирать его в структуру задача кода. Второй вариант: не разбирать вовсе, а взять из попаданий только идентификаторы и достать карточки из PostgreSQL. Так выдача всегда свежая, а индекс может хранить минимум полей.
  • total врёт по умолчанию после 10 000: кластер перестаёт считать точно ради скорости. Для «найдено 12 345 товаров» этого хватает, для пагинации «последняя страница» нет: нужен TrackTotalHits: true.

Глубокие страницы ограничены index.max_result_window, по умолчанию 10 000 документов: from + size больше этого числа это ошибка запроса, а не пустая выдача. Листать дальше можно только search_after, об этом в разделе «Глубже».

Массовая индексация

Записывать документы по одному это один HTTP-запрос на документ. Для первичной загрузки каталога и для переиндексации нужен bulk: один запрос несёт сотни и тысячи операций. У типизированного клиента это Bulk():

func bulkIndex(ctx context.Context, es *elasticsearch.TypedClient, index string, products []Product) error {
    bulk := es.Bulk().Index(index)
    for _, p := range products {
        id := strconv.FormatInt(p.ID, 10)
        if err := bulk.IndexOp(types.IndexOperation{Id_: &id}, p); err != nil {
            return err
        }
    }
    res, err := bulk.Do(ctx)
    if err != nil {
        return err
    }
    if res.Errors {
        for _, item := range res.Items {
            for op, r := range item {
                if r.Error != nil {
                    slog.Warn("bulk item failed", "op", op, "id", r.Id_, "reason", r.Error.Reason)
                }
            }
        }
        return errors.New("bulk: part of items failed")
    }
    return nil
}

Bulk не атомарен. Запрос вернётся с кодом 200, а половина операций внутри может быть отклонена: res.Errors == true и у каждого элемента свой статус. Код обязан пройти по Items и решить, что делать с упавшими: повторить, записать в лог, остановить переиндексацию.

Для долгих загрузок удобнее esutil.BulkIndexer из низкоуровневого пакета: он сам режет поток на запросы по размеру, держит несколько воркеров и вызывает обработчик на каждую неудачу. Ему нужен низкоуровневый *elasticsearch.Client, который создаётся из того же Config через NewClient.

import "github.com/elastic/go-elasticsearch/v8/esutil"

bi, err := esutil.NewBulkIndexer(esutil.BulkIndexerConfig{
    Client:        es,
    Index:         "products_v2",
    NumWorkers:    4,
    FlushBytes:    5 << 20,
    FlushInterval: 5 * time.Second,
})
for _, p := range batch {
    body, _ := json.Marshal(p)
    err := bi.Add(ctx, esutil.BulkIndexerItem{
        Action:     "index",
        DocumentID: strconv.FormatInt(p.ID, 10),
        Body:       strings.NewReader(string(body)),
        OnFailure: func(_ context.Context, item esutil.BulkIndexerItem,
            res esutil.BulkIndexerResponseItem, err error) {
            slog.Warn("reindex item failed", "id", item.DocumentID, "reason", res.Error.Reason, "err", err)
        },
    })
}
if err := bi.Close(ctx); err != nil {
    return err
}
st := bi.Stats()
if st.NumFailed > 0 {
    return fmt.Errorf("reindex: %d of %d failed", st.NumFailed, st.NumAdded)
}

Размер одного bulk-запроса держат в 5-15 МБ: меньше это лишние запросы, больше это память координирующего узла и риск упереться в http.max_content_length (100 МБ). Перед полной переиндексацией реплики временно ставят в 0 и refresh_interval в -1, после возвращают.

Как держать индекс в согласии с базой

Это главный вопрос интеграции, и у него четыре ответа, от простого к надёжному.

Двойная запись из обработчика. Сохранили в PostgreSQL, следом вызвали indexOne. Просто и неверно: транзакция в базе уже закоммичена, а вызов в индекс упал по таймауту, и товар навсегда пропал из поиска. Повтор в цикле внутри запроса не спасает: процесс могут убить между двумя записями.

Transactional outbox. В той же транзакции, что и товар, в таблицу outbox пишется событие «товар изменён». Отдельная горутина или воркер читает outbox и отправляет события в индекс (или в Kafka, откуда их читает индексатор). Гарантия «хотя бы раз» на стороне базы, повторы безопасны, потому что индексация по id идемпотентна. Рецепт целиком разобран в паттернах распределённых систем на Go.

CDC через Debezium. Коннектор читает WAL PostgreSQL и кладёт изменения строк в Kafka, Go-потребитель собирает из них документы и шлёт bulk. Код сервиса про индекс не знает вообще. Цена: Kafka Connect в инфраструктуре и документ, который собирается из «сырых» строк таблиц, а не из доменной модели.

Полная переиндексация по расписанию. Ночью выгрузить всё из базы в новый индекс и переключить псевдоним. Сама по себе это не синхронизация, а страховка: любой из трёх способов выше со временем накапливает расхождения (упавшие повторы, ручные правки в базе), и регулярная пересборка их стирает.

Рабочая связка для каталога: outbox для изменений в течение дня плюс ночная пересборка. Удаление товара это тоже событие, иначе он останется в выдаче.

Псевдоним: смена индекса без простоя

Приложение никогда не ходит в индекс по имени products_v2, только по псевдониму products. Тогда переиндексация выглядит так: создать products_v3 с новым маппингом, залить его bulk-индексатором, одним атомарным запросом перекинуть псевдоним, удалить старый индекс через день.

import "github.com/elastic/go-elasticsearch/v8/typedapi/indices/updatealiases"

func switchAlias(ctx context.Context, es *elasticsearch.TypedClient, alias, oldIndex, newIndex string) error {
    _, err := es.Indices.UpdateAliases().Request(&updatealiases.Request{
        Actions: []types.IndicesAction{
            {Remove: &types.RemoveAction{Index: &oldIndex, Alias: &alias}},
            {Add: &types.AddAction{Index: &newIndex, Alias: &alias}},
        },
    }).Do(ctx)
    return err
}

Оба действия в одном запросе применяются атомарно: нет момента, когда поиск не видит ни одного индекса. Пока идёт заливка, изменения из outbox пишут в оба индекса, иначе новый индекс к моменту переключения уже отстанет.

Когда кластер недоступен

Поиск редко бывает критичным для денег: заказ оформляется без него. Поэтому обработчик поиска получает жёсткий дедлайн в контексте и честный ответ 503 при ошибке, а не ожидание до бесконечности и не перекладывание запроса в PostgreSQL с ILIKE (это только положит базу вместе с поиском). У самого запроса есть серверный таймаут: Timeout: "500ms" в search.Request вернёт то, что успели найти шарды, с флагом res.TimedOut == true, и код решает, показывать ли частичную выдачу.

Запись в индекс через outbox недоступности не боится: события подождут в таблице, воркер повторит с паузой. Единственное, что нужно, это метрика отставания (возраст самого старого необработанного события) и алерт на неё.

Три ловушки, которые встречаются чаще других.

  • Elasticsearch не транзакционен. Нет «записать два документа или ни одного», нет чтения своих незакоммиченных записей. Всё, что требует атомарности, живёт в PostgreSQL.
  • Поле с неверно угаданным типом всплывает через месяц, когда понадобилась сортировка по цене, а price оказался text. Лечится только пересозданием индекса через псевдоним.
  • Большие документы. Описание с HTML на мегабайт делает каждый запрос медленнее. В индексе лежит очищенный текст и поля для фильтров; картинки, историю цен и остатки по складам в него не кладут.

Тесты

Юнит-тесты сервиса поиска не ходят в кластер: репозиторий поиска это интерфейс (Search(ctx, query) ([]Hit, error)), и в тестах обработчика его подменяет заглушка. Сам клиент тестируется интеграционно, через контейнер:

import tces "github.com/testcontainers/testcontainers-go/modules/elasticsearch"

c, err := tces.Run(ctx, "docker.elastic.co/elasticsearch/elasticsearch:8.17.0")
if err != nil {
    t.Fatal(err)
}
t.Cleanup(func() { c.Terminate(ctx) })
es, err := elasticsearch.NewTypedClient(elasticsearch.Config{
    Addresses: []string{c.Settings.Address},
    Username:  "elastic",
    Password:  c.Settings.Password,
    CACert:    c.Settings.CACert,
})

Контейнер отдаёт адрес, пароль и сертификат, потому что в образах 8.x безопасность включена. Запускается он десятки секунд и просит около гигабайта памяти, поэтому такие тесты живут за build-тегом integration, поднимают один контейнер на пакет в TestMain и используют свой индекс на каждый тест (имя с t.Name()), а после записи зовут Refresh(refresh.True): в тесте это допустимо.

Дополнительно: при первом чтении можно пропустить

Глубже: search_after вместо глубокой пагинациирасширенное

from + size заставляет каждый шард собрать from + size лучших документов и отдать координатору; на странице 500 это 5 000 документов с каждого шарда ради 10 нужных. search_after принимает значения сортировки последнего документа предыдущей страницы и продолжает с них: нагрузка одинакова для первой и для тысячной страницы. Условие одно: сортировка должна быть однозначной, поэтому к _score или дате добавляют поле-разделитель (идентификатор). Чтобы выдача не «плыла» при записи между страницами, запрос открывают через point in time (OpenPointInTime) и передают его id в каждом следующем запросе. Для каталога, где дальше десятой страницы никто не ходит, это избыточно; для выгрузок и обходов индекса обязательно.

Глубже: что кладут в индекс, а что оставляют в базерасширенное

Индекс это производные данные, их всегда можно пересобрать из базы, и это главный критерий: если поле нельзя восстановить из PostgreSQL, ему в индексе не место. Внутрь идёт то, по чему ищут (текст), фильтруют (категория, бренд, наличие), сортируют (цена, дата, рейтинг) и то, что показывают прямо в выдаче (название, картинка-превью, цена). Остатки по складам, которые меняются каждую минуту, в индекс не пишут: их подтягивают из базы по идентификаторам из выдачи. Тогда обновление товара в индексе происходит при редактировании карточки, а не при каждой продаже.

Коротко

  • Один TypedClient на процесс, версия клиента совпадает с мажорной версией кластера, таймауты через контекст, в RetryOnStatus добавлен 429.
  • Маппинг создаёт код до первого документа: text для поиска, keyword для фильтров, деньги без float.
  • Идентификатор документа равен первичному ключу, тогда запись идемпотентна, а удаление не требует поиска.
  • Документ виден через refresh_interval (1 с); Refresh(refresh.Waitfor) только там, где пользователь ждёт результат, в bulk никогда.
  • must считает релевантность, filter отсекает; _source разбирают из json.RawMessage или берут только id и читают карточки из базы.
  • Bulk возвращает 200 и при частичных ошибках: проверять res.Errors и Items; для больших загрузок esutil.BulkIndexer.
  • Синхронизация с базой через outbox или CDC плюс ночная пересборка; двойная запись из обработчика теряет документы.
  • Приложение ходит в индекс по псевдониму; смена индекса одним UpdateAliases.
  • Падение кластера это 503 на поиске и очередь в outbox на записи, а не ILIKE по базе.
  • Глубокие страницы через search_after и point in time, from + size ограничен 10 000.

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