Chapter 04 · Indexing and CRUD: Bulk APIs, Refresh, Concurrency Control, and Idempotent Ingestion
Bulk Request Format, Batch Sizing, Backpressure, Error Inspection, and Retry Classification
Treat a bulk HTTP request as a container of independent document operations: preserve NDJSON framing, inspect every item, classify failures, and apply bounded backpressure-aware retries.
Learning outcomes
AtlasMart now needs to ingest tens of thousands of product changes efficiently. Sending one HTTP request per mutation wastes connection and parsing overhead, but simply making one enormous request creates a different failure mode: long queues, large memory pressure, timeouts, and expensive retries. The Bulk API is therefore a transport envelope for many independent actions—not an all-or-nothing transaction.
Construct valid newline-delimited JSON (NDJSON) for index, create, update, and delete actions and preserve the final newline.
Inspect every item in the bulk response instead of using
only the HTTP status or top-level errors flag.
Classify mapping/validation failures, conflicts, throttling, unavailable shards, authentication failures, and ambiguous network timeouts by retry policy.
Size batches by bytes, action mix, latency, concurrency, and server pressure rather than a universal document count.
Apply bounded concurrency and exponential backoff without blindly replaying successful items.
Examples are written against
Elasticsearch 9.5.3 and
OpenSearch 3.8.0. The portable core uses document
and index APIs verified on both products; product-specific
behavior is labeled rather than normalized away. Elasticsearch
examples assume the default self-managed distribution with its
bundled JVM. OpenSearch examples assume the upstream 3.8.0
distribution with the Security plugin present. Keep the
earlier course endpoints: Elasticsearch at
https://localhost:9200 with
ELASTIC_PASSWORD, and OpenSearch at
https://localhost:9201 with
OPENSEARCH_INITIAL_ADMIN_PASSWORD.
The generation environment does not run these two search
servers. Requests were checked against current official
documentation but were not executed here, so example responses
describe expected fields, status classes, and invariants
rather than fabricated captured output or benchmarks. Use only
disposable atlasmart-* indices, preserve the
CA/certificate paths established in Chapter 01, and never
point cleanup, force-merge, or failure-injection commands at
unrelated or production data.
1. NDJSON is a framing contract, not ordinary JSON
Bulk requests alternate action metadata lines with optional
source/update lines. delete consumes only its
metadata line; index, create, and
update require an additional line. Each JSON object
must remain on one physical line, and the payload ends with a
newline. When using curl with a file,
--data-binary preserves the newlines; plain
form-style data options can corrupt the framing.
{"index":{"_index":"atlasmart-products-v4-write-lab","_id":"P-710"}}{"product_id":"P-710","name":"USB Hub","price":29.90,"stock":15,"updated_at":"2026-09-11T05:20:00Z","event_id":"evt-710"}{"create":{"_index":"atlasmart-products-v4-write-lab","_id":"P-711"}}{"product_id":"P-711","name":"Desk Lamp","price":49.00,"stock":8,"updated_at":"2026-09-11T05:20:02Z","event_id":"evt-711"}{"update":{"_index":"atlasmart-products-v4-write-lab","_id":"P-701"}}{"doc":{"price":35.90,"updated_at":"2026-09-11T05:20:03Z"}}{"delete":{"_index":"atlasmart-products-v4-write-lab","_id":"P-does-not-exist"}}
At an index-specific bulk endpoint the _index value
can be omitted from each action; explicit per-item targets are
useful in teaching because the routing of each operation remains
visible. Do not use the same pattern to mix unrelated tenants
without an authorization model.
2. HTTP 200 can contain failed documents
The server can parse and execute a Bulk request successfully at
the HTTP layer while individual actions fail. The response
includes an items array in request order. Each item
has an action key, status, result metadata, and—on failure—an
error object. The top-level
errors flag is only a fast signal that at least one
item failed; it is not a substitute for per-item inspection.
POST _bulk?refresh=false{"create":{"_index":"atlasmart-products-v4-write-lab","_id":"P-711"}}{"product_id":"P-711","name":"Desk Lamp","price":49.00,"stock":8,"updated_at":"2026-09-11T05:21:00Z","event_id":"evt-711-duplicate"}{"index":{"_index":"atlasmart-products-v4-write-lab","_id":"P-712"}}{"product_id":"P-712","name":"Cable","price":"NOT-A-NUMBER","stock":20,"updated_at":"2026-09-11T05:21:01Z","event_id":"evt-712-bad"}{"index":{"_index":"atlasmart-products-v4-write-lab","_id":"P-713"}}{"product_id":"P-713","name":"Mouse Pad","price":12.50,"stock":30,"updated_at":"2026-09-11T05:21:02Z","event_id":"evt-713"}
The expected shape is one conflict for the duplicate create, one
permanent mapping/parsing failure for the invalid numeric field,
and one success for P-713. Retrying the whole
request repeats P-713 unnecessarily and may be
unsafe if some successful actions are non-idempotent updates.
3. Classify failures before you retry
| Signal | Typical class | Default application response |
|---|---|---|
| 2xx item | Success | Record success; never resend merely because a sibling item failed. |
| 409 conflict | Semantic/concurrency conflict | Re-read/reconcile or treat duplicate create as expected dedupe; do not blind-backoff forever. |
| 400 mapping/parse/illegal argument | Permanent payload/schema defect | Dead-letter/quarantine with reason; fix data or schema. |
| 401/403 | Authentication/authorization defect | Stop or route to operator/config remediation; repeated document retry is pointless. |
| 429 rejection/throttling | Transient capacity/backpressure signal | Retry only failed items with jittered exponential backoff and reduced concurrency. |
| 502/503/unavailable shard | Potentially transient service/topology condition | Retry idempotent failed items after bounded backoff; preserve attempt history. |
| Network timeout / connection loss | Ambiguous outcome | Assume the server may have committed; use stable IDs/OCC/reconciliation before replay. |
| 413/request too large | Batch construction defect | Split the batch; retry smaller envelopes, not the same oversized request. |
Error codes are context, not a universal theorem. A
409 from create-only deduplication can be business
success, while a 409 from stale optimistic
concurrency can require a new read and merge. A
404 delete can mean “desired absence achieved” or
“wrong ID,” depending on the operation contract.
4. Batch sizing is an experiment, not a magic count
Start with a modest byte/action envelope and measure. Increase until indexing throughput stops improving, p95/p99 latency rises sharply, memory/queue pressure appears, or retries become expensive. Document size, ingest pipelines, scripts, mappings, shard fan-out, replicas, refresh policy, storage, JVM resources, network, and concurrent workers all change the optimum. The correct unit is often bytes plus observed service time, not “5,000 documents.”
| Measure | Why it matters |
|---|---|
| Request bytes and action count | Controls parsing/network cost and retry blast radius. |
| Bulk p50/p95/p99 latency | Shows tail growth hidden by average throughput. |
| Per-item 429/5xx rate | Direct backpressure/failure signal. |
| Indexing pressure / thread-pool rejection evidence | Shows the server is protecting itself from more work. |
| Client in-flight requests and queue depth | Prevents the producer from hiding overload in its own memory. |
| Refresh/merge/segment behavior | Separates ingestion pressure from background search-structure work. |
Rejections are safety signals. A much larger queue can turn fast backpressure into long tail latency, heap pressure, and a larger ambiguous-retry window. Reduce producer concurrency, batch size, or ingest cost; add measured capacity; or reshape the workload before changing protective limits.
5. Lab: partition a mixed response into three queues
Run the mixed bulk payload on both local engines. Save the entire response. For each request position, join the original action with the corresponding response item and emit exactly one classification: success, retryable, or permanent/reconcile. Keep the original stable ID and event ID with the classification so a later retry can prove it is operating on the same logical mutation.
for each (request_item, response_item) in order: status = response_item.status if 200 <= status < 300: success.append(request_item) elif status in {429, 502, 503, 504}: retryable.append(request_item) elif status == 409: reconcile.append(request_item) # duplicate-create or OCC conflict else: permanent.append({request_item, response_item.error})retry only retryable items after bounded jittered backoffnever replay success items because a sibling failed
Force at least one 409 and one mapping failure as above. Capacity-related 429s are environment-sensitive; do not overload a workstation merely to manufacture one. If no safe local rejection occurs, use the recorded deterministic 429 response shape as a simulation and test the classifier against it.
Check your understanding
- Why is a successful Bulk HTTP status insufficient?
- Why is blind whole-batch retry unsafe?
- What makes 429 different from a mapping error?
- Why should batch size be measured in more than document count?
- What is the correct default after a network timeout?
Review the answers
1. Because individual actions are executed independently and failures are reported per item inside the response.
2. Successful non-idempotent actions may execute again, and even idempotent successes create unnecessary load.
3. 429 normally signals temporary backpressure/capacity, while a mapping/parsing 400 is a payload or schema defect that backoff alone cannot fix.
4. Documents and actions vary in bytes and server work; bytes, action mix, scripts/pipelines, shard routing, latency and concurrency drive pressure.
5. Treat the outcome as ambiguous and reconcile or safely replay using idempotent identity/OCC rather than assuming failure.
Production judgment
Bulk throughput is useful only when the application can explain every item. Preserve correlation IDs, original stable IDs, attempt counts, error types, and latency so an incident can reconstruct which mutations committed. Tune producer concurrency against server backpressure and tail latency, not maximum queue depth. Keep permanent data defects out of retry loops, and ensure dead-letter handling is observable and recoverable rather than a silent data sink.
Summary and next step
The Bulk API reduces transport overhead but increases the importance of per-item accounting. The next lesson separates refresh, flush, and force merge so the loader does not buy visibility or disk cleanup by accidentally destroying indexing efficiency.
Authoritative references
- Elastic Bulk API — Current bulk action, per-item result, refresh, versioning, routing, and OCC reference.
- Elastic refresh parameter — Current visibility semantics for index/update/delete/bulk requests.
- Elastic optimistic concurrency control — Sequence-number and primary-term concurrency semantics.
- OpenSearch Bulk API — Current NDJSON, per-item failure, OCC, versioning, and refresh behavior.
- OpenSearch Document APIs — Current document operation and sequence-number/primary-term overview.