Сбор данных в OpenSearch: API, клиенты и bulk-загрузка

Как отправлять данные в OpenSearch из приложений на Go, Java, Rust и Python: REST и bulk, официальные клиенты, обработка ошибок и живой бенчмарк single vs bulk на 3.5.0

OpenSearch принимает данные через REST API, и на первый взгляд всё просто: POST-запрос с JSON — и документ в индексе. Именно так и выглядит первый пример в любой документации, и для одного документа это действительно вся история.

Но в реальной системе документ редко приходит один. Как только их становится тысячи в секунду, простая картинка обрастает вопросами: отправлять по одному или пакетами, что делать, если часть документов в пакете не прошла валидацию, а часть прошла, и каким клиентом всё это писать из Go, Java, Rust или Python вместо голого curl. В этой статье — сначала голый REST и почему пакетная загрузка (_bulk) не опция, а необходимость, а дальше по разделам — официальные клиенты каждого языка и обработка ошибок.

Сбор данных в OpenSearch: клиенты и bulk

О версии. Примеры в статье проверены на OpenSearch 3.5.0 — том же стенде, что и в предыдущих статьях серии. Сниппеты используют канонический порт 9200; воспроизводимый стенд с клиентами на всех четырёх языках лежит в digital-cookbook, opensearch/ingestion/. Стенд из cookbook слушает 9214 (чтобы не конфликтовать с другими демо серии), поэтому при запуске его примеров используйте порт 9214, а не 9200.

В статье

REST напрямую

Самый простой способ положить документ в OpenSearch — прямой HTTP-запрос, без единой строчки клиентского кода. Одна операция — POST /_doc, где часть пути после индекса указывает, что OpenSearch сам сгенерирует _id документа:

curl -s -u admin:'<пароль>' -k \
  -X POST "https://localhost:9200/app-logs/_doc" \
  -H 'Content-Type: application/json' \
  -d '{"level":"error","message":"disk full","status":500}'

Вот что реально вернул стенд 3.5.0:

{
    "_index": "app-logs",
    "_id": "tXkBJ58Bsx2eYVE9pjVi",
    "_version": 1,
    "result": "created",
    "_shards": {
        "total": 2,
        "successful": 1,
        "failed": 0
    },
    "_seq_no": 0,
    "_primary_term": 1
}

result: created — документ записан. Блок _shards показывает, на скольких копиях шарда операция подтвердилась: total: 2 — сколько копий (primary плюс реплики) в принципе должны её выполнить по настройкам индекса, а successful: 1 — сколько подтвердили. На одноузловом стенде у индекса по умолчанию заявлена одна реплика, но разместить её негде, поэтому подтверждает только primary — отсюда 1 из 2; на кластере с живыми репликами successful дорастёт до числа доступных копий, а failed покажет, если какая-то из них не ответила. Для одного документа этого достаточно. Проблема в том, что каждый такой запрос — это отдельное TCP-соединение (или как минимум отдельный HTTP round-trip поверх keep-alive), отдельный разбор заголовков, отдельная точка синхронизации на стороне кластера. При тысяче документов в секунду тысяча отдельных POST /_doc создаёт overhead, который не имеет отношения к самой записи данных — об этом подробнее в следующем разделе.

Для пакетной записи в OpenSearch есть отдельный эндпоинт — POST /_bulk, и у него нетипичный формат тела запроса: не JSON-массив, а NDJSON (newline-delimited JSON) — построчный поток, где на каждый документ приходится две строки. Первая строка — action-объект: что делать (index, create, update, delete) и в каком индексе. Вторая строка — сам документ (для index/create) или тело операции (для update). Между строками — обычный перевод строки, а не запятая, как в обычном JSON-массиве, и это принципиально: _bulk разбирает тело построчно, а не как единое дерево, что и позволяет ему стримить произвольно большой пакет, не удерживая весь JSON в памяти целиком при парсинге.

curl -s -u admin:'<пароль>' -k \
  -X POST "https://localhost:9200/_bulk" \
  -H 'Content-Type: application/x-ndjson' \
  --data-binary $'{"index":{"_index":"service-a"}}\n{"level":"info","message":"started"}\n{"index":{"_index":"service-b"}}\n{"level":"warn","message":"slow query"}\n'

Обратите внимание на завершающий \n в конце последней строки — это не стилистическая деталь, а обязательное требование: без перевода строки после последнего документа OpenSearch не сможет корректно разобрать тело запроса и вернёт ошибку парсинга. Клиентские библиотеки (Go, Java, Rust, Python — дальше в статье) берут этот момент на себя, но при ручной сборке NDJSON-тела это самая частая причина, почему bulk-запрос падает без видимой причины в самих данных.

Ответ _bulk устроен иначе, чем ответ на одиночный _doc: вместо одного результата — массив items, по одному элементу на каждую операцию из запроса, и общий флаг errors, который сразу говорит, была ли в пакете хоть одна неудачная операция. Вот реальный ответ стенда 3.5.0 на пакет из двух документов:

{
  "took" : 624,
  "errors" : false,
  "items" : [
    {
      "index" : {
        "_index" : "service-a",
        "_id" : "dxdXJ58BQTHTQRaRsCLu",
        "_version" : 1,
        "result" : "created",
        "_shards" : { "total" : 2, "successful" : 1, "failed" : 0 },
        "_seq_no" : 0,
        "_primary_term" : 1,
        "status" : 201
      }
    },
    {
      "index" : {
        "_index" : "service-b",
        "_id" : "eBdXJ58BQTHTQRaRsCLu",
        "_version" : 1,
        "result" : "created",
        "_shards" : { "total" : 2, "successful" : 1, "failed" : 0 },
        "_seq_no" : 0,
        "_primary_term" : 1,
        "status" : 201
      }
    }
  ]
}

errors: false и оба статуса 201 (Created) — пакет прошёл целиком. Но errors: false — это статус пакета в целом, а не гарантия, что каждый документ действительно записан по отдельности: bulk-операция не атомарна, и в общем случае часть документов может пройти, а часть — нет. Что происходит в этом случае и как это разбирать в клиентском коде — тема раздела «Обработка ошибок» дальше в статье.

Почему bulk

Разница между «по документу за раз» и «пакетом» — не архитектурная тонкость, а кратный эффект, который легко измерить. У каждого HTTP-запроса к OpenSearch есть накладные расходы, не зависящие от размера самого документа: установка/переиспользование соединения, разбор заголовков, десериализация JSON, поиск нужного шарда, точка синхронизации _shards в ответе. При отправке по одному документу эти расходы платятся за каждую отдельную запись; при _bulk — один раз на весь пакет, а сама работа по записи документов на стороне кластера батчится внутри одного вызова.

flowchart LR A["Приложение"] -->|"N × POST /_doc"| S1["OpenSearch"] B["Приложение"] -->|"1 × POST /_bulk (N докум.)"| S2["OpenSearch"]

flowchart LR
    A["Приложение"] -->|"N × POST /_doc"| S1["OpenSearch"]
    B["Приложение"] -->|"1 × POST /_bulk (N докум.)"| S2["OpenSearch"]
Single: N запросов против bulk: один запрос на батч

Вот во что это выливается на практике — реальный замер на стенде 3.5.0, single-node, локально (сеть между клиентом и кластером не участвует, поэтому цифры иллюстративные: на проде с сетевым расстоянием между приложением и кластером разница будет ещё заметнее, а абсолютные docs/s — другими):

способ docs/s время на 3000 док
single (по одному) 131 22.95 s
bulk, батч 500 6016 0.50 s
bulk, батч 3000 6816 0.44 s

По одному документу — 131 docs/s, почти 23 секунды на 3000 записей. Пакетами по 500 — уже 6016 docs/s, то есть в 45-50 раз быстрее на том же объёме данных и на том же кластере: единственное, что изменилось, — как документы сгруппированы в HTTP-запросы. Дальнейшее укрупнение батча с 500 до всего пакета целиком (3000) добавляет ещё немного — 6816 против 6016 docs/s, разница уже в пределах 15%. Первый переход (single → bulk) убирает почти весь overhead одиночных запросов; второй (батч 500 → батч 3000) упирается в то, что сама запись документов на диск и в память уже не бесплатна, и дальше расти особо некуда — отдача от увеличения батча падает. Отсюда практический вывод, к которому в статье ещё не раз вернёмся: bulk обязателен при сколько-нибудь заметном потоке документов, а конкретный размер батча — это настройка, которую имеет смысл подбирать эмпирически под свою нагрузку и размер документа, а не бесконечно увеличивать в расчёте на линейный рост.

Чтобы подбор шёл не от нуля, разумная стартовая точка — батч в 500–1000 документов или ограничение по объёму тела запроса (условно 5–15 МБ, смотря по размеру документа), плюс flush по таймауту, чтобы редкий поток не застревал в полупустом буфере. Дальше — по метрикам: если растёт доля ответов 429 (кластер отвергает запись) или задержка bulk-запросов, размер батча и число параллельных bulk-запросов надо уменьшать, а не наращивать. Клиентские bulk-хелперы (Java BulkProcessor, Python helpers.bulk) как раз инкапсулируют логику «флашить по размеру ИЛИ по таймауту» — к ним вернёмся в разделах про клиентов.

Дальше — как то же самое (создание клиента, single index, bulk и разбор ответа) выглядит из кода, а не из curl. Начнём с Go.

Go: opensearch-go

Официальный Go-клиент — opensearch-go, в статье используется v4.6.0 (Go 1.24+). У v4 два слоя API: низкоуровневый opensearch.Client — тонкая обёртка над HTTP, и opensearchapi.Client поверх него — с типизированными запросами и ответами для конкретных эндпоинтов (Index, Bulk, Search и так далее). В демо ниже используется второй — меньше ручной сборки JSON, разбор ответа bulk уже структурирован.

Клиент создаётся один раз на процесс и переиспользуется — как http.Client, он потокобезопасен и держит пул соединений:

client, err := opensearchapi.NewClient(opensearchapi.Config{
	Client: opensearch.Config{
		Addresses: []string{"https://localhost:9200"},
		Username:  "admin",
		Password:  "<пароль>",
		Transport: &http.Transport{
			TLSClientConfig: &tls.Config{InsecureSkipVerify: true}, // demo only
		},
	},
})
if err != nil {
	log.Fatalf("client: %v", err)
}

InsecureSkipVerify: true в TLSClientConfig отключает проверку сертификата целиком — в этом демо-примере кластер поднят с самоподписанным сертификатом, и без этой опции клиент откажется соединяться. В проде так делать не нужно: либо валидный сертификат от доверенного CA, либо явный RootCAs со своим CA-пулом в том же tls.Config.

Одиночная индексация через opensearchapi.Client — метод Index, тело запроса — io.Reader:

body := `{"ts":"2026-07-03T10:00:00Z","level":"error","service":"service-a","message":"disk full","status":500}`
ir, err := client.Index(ctx, opensearchapi.IndexReq{
	Index: "app-logs-go",
	Body:  strings.NewReader(body),
})
if err != nil {
	log.Fatalf("index: %v", err)
}
fmt.Printf("single index -> result=%s id=%s\n", ir.Result, ir.ID)

Заметьте: err здесь — это ошибка транспорта (соединение не установлено, таймаут, обрыв). HTTP-код 4xx/5xx в bulk-подобных ответах err не породит — это касается и Index, но особенно важно для Bulk ниже, где успешный HTTP-ответ ничего не говорит о судьбе отдельных документов.

Для bulk opensearchapi.Client не строит NDJSON-тело за вас — принимает уже готовый io.Reader с тем же построчным форматом, что и в разделе «REST напрямую»: action-строка, документ-строка, и так на каждый документ:

var sb strings.Builder
docs := []LogDoc{
	{Ts: "2026-07-03T10:01:00Z", Level: "info", Service: "service-a", Message: "ok", Status: 200},
	{Ts: "2026-07-03T10:02:00Z", Level: "warn", Service: "service-b", Message: "retry", Status: 429},
}
for _, d := range docs {
	sb.WriteString(`{"index":{"_index":"app-logs-go"}}` + "\n")
	sb.WriteString(fmt.Sprintf(`{"ts":%q,"level":%q,"service":%q,"message":%q,"status":%d}`+"\n",
		d.Ts, d.Level, d.Service, d.Message, d.Status))
}
br, err := client.Bulk(ctx, opensearchapi.BulkReq{Body: strings.NewReader(sb.String())})
if err != nil {
	log.Fatalf("bulk: %v", err)
}

Собирать NDJSON строками через fmt.Sprintf — нормально для демо и для небольшого числа простых полей, но на реальных документах с вложенной структурой лучше encoding/json.Marshal на каждую строку — меньше риска сломать экранирование вручную.

А вот разбор ответа — тот самый момент, который решает, узнаете ли вы о частичном отказе пакета:

// ВАЖНО: HTTP-успех != все документы записаны. Проверяем каждый item.
fmt.Printf("bulk -> errors=%v, items=%d\n", br.Errors, len(br.Items))
for i, item := range br.Items {
	op := item["index"]
	if op.Status >= 300 {
		fmt.Printf("  item %d FAILED status=%d type=%s\n", i, op.Status, op.Error.Type)
	} else {
		fmt.Printf("  item %d ok status=%d\n", i, op.Status)
	}
}

br.Items[i] — это map[string]opensearchapi.BulkRespItem, ключ — тип операции ("index", "create", "update", "delete"), поэтому item["index"] берёт результат именно операции индексации; для смешанного пакета с разными типами операций ключ пришлось бы определять по тому, что реально было отправлено на этой позиции. Реальный вывод на стенде 3.5.0 для пакета из двух документов:

single index -> result=created id=...
bulk -> errors=false, items=2
  item 0 ok status=201
  item 1 ok status=201

Для потока документов (логи, метрики, события) паттерн из двух вызовов «накопил — отправил» вручную неудобен масштабировать: нужен буфер, который наполняется до нужного размера или таймаута, и пул воркеров, которые эти буферы отправляют параллельно, не блокируя производителя событий. Сам opensearch-go не навязывает конкретную реализацию такого пайплайна — в отличие от opensearch-java с его BulkProcessor (следующий раздел), здесь это на стороне приложения: канал с документами, N воркеров, каждый копит свой батч и делает client.Bulk по накоплении или по тикеру. Плюс подхода — полный контроль над backpressure и обработкой ошибок; минус — писать и тестировать эту часть самому.

Полный код демо — создание клиента, single index, bulk, разбор partial failures — лежит в digital-cookbook, opensearch/ingestion/clients/go/.

Java: opensearch-java

Официальный Java-клиент — opensearch-java, в статье используется 3.9.0 на JDK 21. В отличие от Go-клиента, здесь явно разделены транспорт и клиент: транспорт отвечает за HTTP/TLS/аутентификацию и создаётся через builder, а OpenSearchClient поверх него уже даёт типизированные методы вроде index() и bulk(), где документы — обычные Java-объекты (POJO или, как в демо, record), а не строки JSON.

Транспорт в демо собран на ApacheHttpClient5TransportBuilder — Basic Auth и TLS настраиваются на уровне HTTP-клиента, а не самого OpenSearchClient:

HttpHost host = new HttpHost("https", "localhost", 9200);
BasicCredentialsProvider cp = new BasicCredentialsProvider();
cp.setCredentials(new AuthScope(host), new UsernamePasswordCredentials("admin", "<пароль>".toCharArray()));
SSLContext ssl = SSLContextBuilder.create().loadTrustMaterial(null, (chain, authType) -> true).build(); // demo only

OpenSearchTransport transport = ApacheHttpClient5TransportBuilder.builder(host)
    .setMapper(new JacksonJsonpMapper())
    .setHttpClientConfigCallback(hc -> {
      var tls = ClientTlsStrategyBuilder.create().setSslContext(ssl)
          .setHostnameVerifier(NoopHostnameVerifier.INSTANCE)
          .setTlsDetailsFactory(sslEngine -> new TlsDetails(sslEngine.getSession(), sslEngine.getApplicationProtocol()))
          .build();
      var cm = PoolingAsyncClientConnectionManagerBuilder.create().setTlsStrategy(tls).build();
      return hc.setDefaultCredentialsProvider(cp).setConnectionManager(cm);
    }).build();
OpenSearchClient client = new OpenSearchClient(transport);

loadTrustMaterial(null, (chain, authType) -> true) принимает любой сертификат без проверки — как и InsecureSkipVerify в Go, это только для самоподписанного демо-стенда; в проде — доверенный CA. setMapper(new JacksonJsonpMapper()) подключает Jackson как JSON-биндер: именно он превращает Java-объекты в тело запроса и обратно, что и делает возможной типизацию ниже.

Single index — типизированный, документ передаётся как record LogDoc, а не строка JSON:

public record LogDoc(String ts, String level, String service, String message, int status) {}

IndexResponse ir = client.index(i -> i.index("app-logs-java")
    .document(new LogDoc("2026-07-03T10:00:00Z", "error", "service-a", "disk full", 500)));
System.out.println("single index -> result=" + ir.result() + " id=" + ir.id());

JacksonJsonpMapper сериализует record через стандартную рефлексию по геттерам-аксессорам record’а (ts(), level() и так далее) — отдельно описывать маппинг полей не нужно, если имена полей и JSON-ключи совпадают.

Bulk собирается из списка типизированных BulkOperation, каждая операция оборачивает тот же .document(...), что и single index:

List<LogDoc> docs = List.of(
    new LogDoc("2026-07-03T10:01:00Z", "info", "service-a", "ok", 200),
    new LogDoc("2026-07-03T10:02:00Z", "warn", "service-b", "retry", 429));
List<BulkOperation> ops = new ArrayList<>();
for (LogDoc d : docs) {
  ops.add(BulkOperation.of(op -> op.index(idx -> idx.index("app-logs-java").document(d))));
}
BulkResponse br = client.bulk(new BulkRequest.Builder().operations(ops).build());

Здесь клиент сам собирает NDJSON-тело из списка операций — в отличие от Go-варианта, руками строку \n-разделённого JSON собирать не нужно. Разбор ответа устроен симметрично Go-клиенту: общий флаг errors() и статус каждого элемента items():

// HTTP-успех != все записаны: проверяем errors() и каждый item.
System.out.println("bulk -> errors=" + br.errors() + ", items=" + br.items().size());
for (int i = 0; i < br.items().size(); i++) {
  var it = br.items().get(i);
  System.out.println("  item " + i + " status=" + it.status() + (it.error() != null ? " ERROR " + it.error().type() : " ok"));
}

it.error() возвращает null для успешных операций и объект с типом/причиной для отказавших — тот же принцип, что и op.Error в Go, только через null-проверку вместо кода состояния. Вывод на стенде 3.5.0:

single index -> result=Created id=...
bulk -> errors=false, items=2
  item 0 status=201 ok
  item 1 status=201 ok

В Spring Boot-приложении транспорт и OpenSearchClient из примера выше оформляются одним бином в конфигурации (@Bean OpenSearchClient openSearchClient(...)), собирающим тот же ApacheHttpClient5TransportBuilder из значений application.yml вместо хардкода — сам код клиента не меняется, меняется только то, откуда берутся адрес и креды. Для равномерного потока документов вместо ручного накопления списка BulkOperation есть BulkIngester (пакет org.opensearch.client.opensearch.helpers) — обёртка с авто-flush по размеру батча или таймауту и повторными попытками, аналог BulkProcessor из клиента для Elasticsearch; в демо-примере она не используется, но для продакшен-пайплайна с непрерывным потоком это более естественная отправная точка, чем ручной список операций.

Полный код демо лежит в digital-cookbook, opensearch/ingestion/clients/java/.

Rust: opensearch-rs

Официальный Rust-клиент — opensearch-rs, в статье используется crate opensearch версии 2.4.0. Как и в Java-клиенте, здесь отдельно собирается транспорт (Transport) и отдельно — сам OpenSearch-клиент поверх него, но с поправкой на экосистему Rust: транспорт строится через TransportBuilder с явным пулом соединений, а весь клиент по умолчанию асинхронный — методы index() и bulk() возвращают future, которые нужно .await-ить внутри tokio-рантайма.

use opensearch::auth::Credentials;
use opensearch::cert::CertificateValidation;
use opensearch::http::transport::{SingleNodeConnectionPool, TransportBuilder};
use opensearch::OpenSearch;
use url::Url;

let url = Url::parse("https://localhost:9200")?;
let pool = SingleNodeConnectionPool::new(url);
let transport = TransportBuilder::new(pool)
    .auth(Credentials::Basic("admin".into(), "<пароль>".into()))
    .cert_validation(CertificateValidation::None) // demo only
    .build()?;
let client = OpenSearch::new(transport);

SingleNodeConnectionPool — самый простой пул, для одного узла (для кластера с несколькими нодами есть CloudConnectionPool / MultiNodeConnectionPool со стратегией round-robin по списку адресов). CertificateValidation::None отключает проверку TLS-сертификата целиком — тот же смысл, что у InsecureSkipVerify в Go и loadTrustMaterial(..., true) в Java: только для самоподписанного демо-стенда, в проде — CertificateValidation::Full с доверенным CA.

Документ описывается обычной структурой с #[derive(Serialize)]opensearch-rs полагается на serde для сериализации, а не на собственный маппинг:

use serde::Serialize;

#[derive(Serialize)]
struct LogDoc {
    ts: &'static str,
    level: &'static str,
    service: &'static str,
    message: &'static str,
    status: u16,
}

Вход в приложение — асинхронная main, помеченная атрибутом #[tokio::main], который разворачивает её в обычную синхронную main, запускающую tokio-рантайм и внутри него — асинхронную функцию:

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let doc = LogDoc {
        ts: "2026-07-03T10:00:00Z", level: "error",
        service: "service-a", message: "disk full", status: 500,
    };
    let resp = client
        .index(IndexParts::Index("app-logs-rust"))
        .body(doc)
        .send()
        .await?;
    println!("single index -> status={}", resp.status_code());
    // ...
}

.body(doc) принимает саму структуру и сериализует её через serde под капотом — передавать готовый JSON-объект не обязательно, хотя и serde_json::Value тоже подойдёт как тело.

Bulk устроен иначе, чем в Go или Java, и здесь есть особенность, о которую легко споткнуться: тело bulk-запроса — это Vec<JsonBody<_>>, а не Vec<serde_json::Value>, потому что Value не реализует трейт Body, которого требует API client.bulk(...). Каждый Value нужно явно обернуть через .into() в JsonBody:

use opensearch::http::request::JsonBody;
use serde_json::{json, Value};

let docs = [
    LogDoc { ts: "2026-07-03T10:01:00Z", level: "info", service: "service-a", message: "ok", status: 200 },
    LogDoc { ts: "2026-07-03T10:02:00Z", level: "warn", service: "service-b", message: "retry", status: 429 },
];
let mut body: Vec<JsonBody<Value>> = Vec::new();
for d in &docs {
    body.push(json!({"index": {}}).into());
    body.push(serde_json::to_value(d)?.into());
}
let resp = client.bulk(BulkParts::Index("app-logs-rust")).body(body).send().await?;

Структура та же, что и в NDJSON из раздела «REST напрямую» — action-объект, затем документ, на каждую операцию по паре элементов вектора, — просто opensearch-rs собирает построчный NDJSON-формат из вектора значений сам, а не ждёт готовую строку. .into() на каждом json!(...) и на результате serde_json::to_value(d)? — это и есть та самая обёртка ValueJsonBody<Value>, без которой код не скомпилируется.

Ответ bulk() — не типизированная структура, как в Java, а сырой Response, из которого JSON нужно достать явно через .json::<Value>().await?, и дальше разбирать как обычное дерево serde_json::Value:

let out: Value = resp.json().await?;
// HTTP-успех != все записаны: проверяем errors и каждый item.
println!("bulk -> errors={}", out["errors"]);
for (i, item) in out["items"].as_array().unwrap().iter().enumerate() {
    println!("  item {} status={}", i, item["index"]["status"]);
}

Тот же принцип, что и в Go/Java-клиентах: общий флаг errors в теле ответа сигнализирует, был ли в пакете хоть один отказ, а статус конкретного документа смотрится по индексу в массиве items. Реальный вывод на стенде 3.5.0:

single index -> status=201 Created
bulk -> errors=false
  item 0 status=201
  item 1 status=201

opensearch-rs — самый «низкоуровневый» из четырёх клиентов статьи в том смысле, что он не прячет от вас JSON: там, где Go и Java дают структурированные BulkRespItem/BulkResponse.items(), здесь ответ — это serde_json::Value, и разбор конкретных полей — на вашей стороне. Для сервиса на Rust с высокой пропускной способностью (тот же профиль задач, что и у Go-клиента: логи, метрики, события) это осознанный компромисс экосистемы: меньше готовых абстракций вроде BulkProcessor/BulkIngester, зато полный контроль над тем, как батчи собираются и во что превращается ответ.

Полный код демо лежит в digital-cookbook, opensearch/ingestion/clients/rust/.

Python: opensearch-py

Официальный Python-клиент — opensearch-py, в статье используется версия 3.2.0. В отличие от Go/Java/Rust-клиентов, здесь нет отдельного шага «собрать транспорт, затем клиент» — конструктор OpenSearch(...) принимает адрес, креды и TLS-параметры одним вызовом, без builder’ов и промежуточных объектов:

from opensearchpy import OpenSearch, helpers

client = OpenSearch(
    hosts=[{"host": "localhost", "port": 9200}],
    http_auth=("admin", "<пароль>"),
    use_ssl=True,
    verify_certs=False,   # demo only
    ssl_show_warn=False,
)

verify_certs=False — тот же смысл, что и InsecureSkipVerify в Go, loadTrustMaterial в Java и CertificateValidation::None в Rust: отключает проверку сертификата для самоподписанного демо-стенда, в проде вместо этого — ca_certs с путём к доверенному CA-бандлу. ssl_show_warn=False дополнительно приглушает предупреждения urllib3 о непроверенном TLS, которые иначе будут сыпаться в лог на каждый запрос.

Single index — документ передаётся обычным dict, без промежуточной сериализации:

resp = client.index(index="app-logs-python", body={
    "ts": "2026-07-03T10:00:00Z", "level": "error",
    "service": "service-a", "message": "disk full", "status": 500,
})
print("single index ->", resp["result"])

resp — это уже распарсенный dict (клиент делает json.loads за вас), поэтому resp["result"] доступен сразу, без отдельного шага разбора ответа, как в Rust.

Для пакетной загрузки opensearch-py не заставляет вручную собирать NDJSON или следить за размером батча — за это отвечает помощник helpers.bulk, которому достаточно скормить список простых dict-описаний операций, а разбиение на чанки нужного размера он берёт на себя:

docs = [
    {"ts": "2026-07-03T10:01:00Z", "level": "info", "service": "service-a", "message": "ok", "status": 200},
    {"ts": "2026-07-03T10:02:00Z", "level": "warn", "service": "service-b", "message": "retry", "status": 429},
]
actions = [{"_index": "app-logs-python", "_source": d} for d in docs]
# raise_on_error=False -> вернёт список ошибок вместо исключения; проверяем!
success, errors = helpers.bulk(client, actions, raise_on_error=False)
print("bulk -> success:", success, "errors:", errors)

Каждый элемент actions — словарь с _index и _source (по умолчанию операция index; для update/delete есть отдельные ключи _op_type и так далее). helpers.bulk сам решает, сколько документов уложить в один _bulk-запрос — по умолчанию чанками по 500 (настраивается через chunk_size), и для двух документов из примера это один запрос, а для тысяч — уже несколько последовательных вызовов _bulk под капотом, без явного цикла в вашем коде.

Ключевой момент — параметр raise_on_error. По умолчанию он True, и тогда helpers.bulk бросает исключение BulkIndexError при первой же неудачной операции в пакете, обрывая обработку. С raise_on_error=False функция вместо исключения возвращает кортеж (success, errors): число успешно записанных документов и список описаний отказавших операций — то есть тот же принцип «HTTP/вызов прошёл, но не все документы записаны», что и errors/items в остальных клиентах статьи, только выраженный через возвращаемое значение, а не через поле ответа. Не проверить errors здесь так же опасно, как не проверить br.Errors в Go или out["errors"] в Rust — только вместо тихо проигнорированного поля ответа это будет тихо проигнорированный элемент кортежа. Реальный вывод на стенде 3.5.0:

single index -> created
bulk -> success: 2 errors: []

Из четырёх клиентов статьи opensearch-py — самый быстрый способ написать рабочий скрипт: без транспорта, без builder’ов, dict вместо типизированных структур. Это делает его естественным выбором для разовых миграций, ETL-скриптов и админских утилит, где важнее скорость написания, чем строгая типизация или максимальная пропускная способность в проде — для устойчивого высоконагруженного сервиса на Python те же helpers.bulk с chunk_size и raise_on_error=False вполне годятся и в проде, но конкурировать по throughput с ручным контролем батчей на Go или Rust не станут.

Полный код демо лежит в digital-cookbook, opensearch/ingestion/clients/python/.

Обработка ошибок

Во всех четырёх клиентах статьи один и тот же рефрен уже прозвучал по разу: HTTP-успех bulk-запроса ничего не говорит о судьбе отдельных документов внутри пакета. Пора показать, что это значит на практике, а не только предупреждать в комментариях к коду.

Partial failures. Bulk-операция не атомарна: OpenSearch пытается выполнить каждую операцию из пакета независимо, и пакет из десяти документов вполне может вернуться с девятью 201 Created и одним 400. При этом HTTP-статус самого ответа на /_bulk — почти всегда 200 OK, даже если половина пакета отвалилась: 200 здесь означает «сервер принял и обработал запрос», а не «все операции внутри него успешны». Единственный надёжный сигнал — поле errors в теле ответа и статус каждого элемента items.

Вот что реально вернул стенд 3.5.0 на пакет из двух документов в metrics-strict, где первый прошёл валидацию, а второй — нет: в поле value типа integer пришла строка "oops" вместо числа.

{
    "took": 12,
    "errors": true,
    "items": [
        {
            "index": {
                "_index": "metrics-strict",
                "_id": "vnkBJ58Bsx2eYVE9qzUg",
                "_version": 1,
                "result": "created",
                "_shards": {
                    "total": 2,
                    "successful": 1,
                    "failed": 0
                },
                "_seq_no": 0,
                "_primary_term": 1,
                "status": 201
            }
        },
        {
            "index": {
                "_index": "metrics-strict",
                "_id": "v3kBJ58Bsx2eYVE9qzUg",
                "status": 400,
                "error": {
                    "type": "mapper_parsing_exception",
                    "reason": "failed to parse field [value] of type [integer] in document with id 'v3kBJ58Bsx2eYVE9qzUg'. Preview of field's value: 'oops'",
                    "caused_by": {
                        "type": "number_format_exception",
                        "reason": "For input string: \"oops\""
                    }
                }
            }
        }
    ]
}

errors: true на верхнем уровне — единственный по-настоящему дешёвый сигнал «в пакете что-то пошло не так», но он не говорит, что именно и с каким документом: за деталями всё равно нужно идти в items. Первый элемент — обычный успех, status: 201, result: created, ровно та же форма, что и у одиночного POST /_doc из начала статьи. Второй — status: 400, и внутри error два уровня: type: mapper_parsing_exception — верхнеуровневая причина (поле не удалось разобрать под объявленный тип), и caused_by.type: number_format_exception — конкретика (строку "oops" не удалось привести к числу). Поле _id присутствует у обоих элементов, включая отказавший, — значит, по нему можно однозначно сопоставить ошибку с исходным документом в пакете и, например, положить его в отдельную очередь на ручной разбор или дедлеттер.

Проверять errors вместо status: 200 ответа — не стилистическая придирка, а единственный способ не терять данные молча. Все четыре клиента статьи дают эту информацию в разборе ответа, просто в разной форме: Go — булев br.Errors и br.Items[i]["index"].Status/.Error, Java — br.errors() и br.items().get(i).error(), Rust — сырые out["errors"] и out["items"][i]["index"]["status"] из serde_json::Value, Python — не поле ответа, а второй элемент кортежа из helpers.bulk(..., raise_on_error=False). Форма разная, суть одна: код, который читает только код ответа HTTP и не разворачивает items, для partial failure одинаково слеп во всех четырёх экосистемах.

Retry с exponential backoff. Не любая ошибка в items заслуживает повторной попытки. mapper_parsing_exception из примера выше — не тот случай: документ со строкой в числовом поле не станет валиднее, если отправить его ещё раз без изменений, — сначала нужно починить сами данные (или тип поля) на стороне приложения, а потом уже повторять запись. Retry имеет смысл для ошибок другой природы — временных: обрыв соединения, таймаут, кластер занят и не успел ответить. Для них стандартный паттерн — повтор с растущей задержкой между попытками (1s → 2s → 4s → 8s, обычно с ограничением сверху и джиттером, чтобы много клиентов не забились в кластер синхронно после одного и того же сбоя), а не мгновенный повтор в цикле, который в момент перегрузки только усугубляет проблему.

Backpressure. Когда кластер не успевает принимать записи с той скоростью, с которой их присылают, OpenSearch не ставит лишние запросы в бесконечную очередь — он явно отклоняет их кодом 429 Too Many Requests, что соответствует rejected execution во внутренних очередях индексации. Это не ошибка в данных и не повод для немедленного повтора того же объёма с той же скоростью — это сигнал кластера «сбавь темп»: уменьшить размер батча, снизить число параллельных bulk-запросов или просто выдержать паузу перед следующей попыткой. Клиент, который в ответ на 429 начинает слать ещё агрессивнее (в расчёте «пробить» перегрузку), только продлевает её; корректная реакция — тот же exponential backoff, что и для прочих временных ошибок, но с уменьшением конкурентности до следующего успешного запроса.

Типичные ошибки

  • Отправка по одному документу вместо bulk. Не архитектурная тонкость, а прямые потери в пропускной способности — раздел «Почему bulk» показал разницу в 45-50 раз на одном и том же стенде и объёме данных. Если поток документов больше единиц в секунду, _bulk — не опциональная оптимизация, а необходимость с самого начала.
  • Игнорирование partial failures. HTTP 200 OK от /_bulk принимается за «всё записалось», хотя это только «сервер обработал запрос» — реальный результат в errors/items, разобранных в разделе «Обработка ошибок» выше. Без проверки каждого item часть данных теряется молча, и обнаруживается это обычно не в момент записи, а позже — когда в поиске не находится документ, который, как казалось, был отправлен.
  • Незакрытые соединения и клиенты. Клиент во всех четырёх экосистемах статьи держит пул TCP-соединений (или, для Go, http.Transport) — если создавать его на каждый запрос вместо одного разделяемого экземпляра на процесс, соединения и файловые дескрипторы копятся быстрее, чем освобождаются. Клиент создаётся один раз при старте приложения и переиспользуется на всё время его жизни; при штатном завершении — закрывается явно (transport/client Close(), где API это предоставляет), а не оставляется сборщику мусора.
  • Mapping-конфликты на dynamic-полях. Строка вместо числа в поле, которое проиндексировано первым документом как integer (или наоборот), — та же ситуация, что и в примере partial failure выше, только источник обычно один и тот же: непредсказуемая структура входных данных при dynamic mapping. Если схема известна заранее, явный маппинг с проверенными типами полей и dynamic: strict/strict_allow_templates отсекает такие документы предсказуемо на входе, а не роняет случайные items пакетами — подробно об этом в разделе «Контроль над схемой» предыдущей статьи серии.

Короткий чек-лист перед продакшеном — собран из того, что разобрано выше и в соседних статьях серии:

  • Маппинг и index template — до первой записи, а не после: тип поля у существующего индекса не переопределить без reindex (см. «Индексы, маппинги и шаблоны»).
  • Всегда _bulk при заметном потоке; размер батча — по метрикам, а не на глаз.
  • Разбирать каждый item bulk-ответа, а не только HTTP-код; неуспешные документы — в метрики и в dead-letter queue, чтобы не терять их молча.
  • Retry только для transient-ошибок (сеть, 429) с exponential backoff и jitter; неисправимые (mapping-конфликт) — не повторять, а чинить документ или схему.
  • Отдельное право на запись для сервисного пользователя: cluster-level indices:data/write/bulk в его роли (см. security plugin в статье про установку кластера).
  • Один клиент на процесс и явное закрытие при остановке — иначе утечка соединений и файловых дескрипторов.

Что дальше

Данные пишутся, ошибки разбираются поэлементно, скорость — на уровне тысяч документов в секунду с одного клиента. Дальше в серии — то, что происходит с этими данными после записи: ISM и retention целиком (rollover, состояния warm/cold, снапшоты и удаление по расписанию — мост к теме был показан в предыдущей статье), а также OpenSearch Dashboards — как смотреть на те же индексы глазами, а не только через API и curl.

Отдельная тема на будущее — агентный сбор вместо клиентского. Всё в этой статье писало данные из кода приложения напрямую, но на реальном стенде логи и метрики чаще собирают агентом в стороне от процесса: например, контур Vector → NATS → OpenSearch, где Vector читает логи с диска или сокета, NATS буферизует поток между источником и хранилищем, а в OpenSearch данные попадают уже пакетами через тот же _bulk, что и в этой статье, — просто с другой стороны сборки.

Официальные источники

  • Index document — Document API: POST/PUT _doc, _create
  • Bulk — Bulk API: формат NDJSON, разбор ответа, errors/items
  • opensearch-go — официальный Go-клиент
  • opensearch-java — официальный Java-клиент
  • opensearch-rs — официальный Rust-клиент
  • opensearch-py — официальный Python-клиент

Обсуждение в Telegram

Присоединиться →

Комментарии