Chapter 04 · Indexing and CRUD: Bulk APIs, Refresh, Concurrency Control, and Idempotent Ingestion
Build a High-Throughput Bulk Loader that Handles Partial Failure Without Silent Data Loss or Duplication
Assemble the chapter into a dual-platform AtlasMart ingestion worker that classifies every bulk result, retries only safe work, proves one forced conflict, and reconciles timeout ambiguity without silent loss or duplication.
Learning outcomes
The chapter capstone is not “send Bulk requests fast.” AtlasMart needs an ingestion worker that can prove every input event ended in exactly one documented state: committed, deduplicated, retry scheduled, quarantined, or awaiting reconciliation. It must preserve stable IDs, inspect per-item responses, respect backpressure, force one deterministic conflict, and measure search visibility separately from write acknowledgement.
Define an ingestion record with stable event ID, target document ID, operation intent, attempt metadata, and reconciliation state.
Build NDJSON batches, correlate responses by position, and partition successes, retryable failures, permanent failures, and semantic conflicts.
Implement bounded exponential backoff with jitter and a concurrency limit that responds to 429/5xx pressure.
Force and resolve one sequence-number conflict and one duplicate create without duplicating a non-idempotent effect.
Verify source count/IDs, event identity, acknowledgement-versus-visibility timing, and final dead-letter/reconciliation queues before declaring success.
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. Define the loader’s state machine before writing code
Every input event gets a stable event_id from the
upstream domain, not a random retry ID. The loader derives the
target product _id and an operation contract:
authoritative snapshot (index), create-once event, conditional
update, or delete. The local ledger records attempt number, last
status/error type, next eligible retry time, and final
disposition. This makes a retry observable rather than recursive
“try again” code.
| Disposition | Meaning | May be resent automatically? |
|---|---|---|
| Committed | Item succeeded and its response metadata is stored | No. |
| Deduplicated | Create conflict matches the same previously accepted event | No. |
| Retry scheduled | Transient/capacity failure on an idempotent item | Yes, after bounded backoff. |
| Quarantined | Permanent mapping/validation/auth/config defect | No; operator/data repair required. |
| Reconcile | Outcome or business conflict is ambiguous | No blind retry; read current state/event history first. |
2. A small dual-platform loader core
The following Python uses raw HTTPS deliberately so the classification core is identical for the two local endpoints; production applications should normally use the current official Elasticsearch or OpenSearch client for their chosen platform and carry the same state machine into the client helper. The example requires a CA path and credentials from Chapter 01 and never disables certificate verification.
import json, os, ssl, base64, urllib.request, urllib.error, time, randomdef auth_header(user, password): token = base64.b64encode(f"{user}:{password}".encode()).decode() return "Basic " + tokendef send_bulk(base_url, ca_file, user, password, ndjson): ctx = ssl.create_default_context(cafile=ca_file) req = urllib.request.Request(base_url + "/_bulk?refresh=false", data=ndjson.encode(), method="POST", headers={"Content-Type":"application/x-ndjson", "Authorization":auth_header(user, password)}) with urllib.request.urlopen(req, context=ctx, timeout=20) as r: return json.load(r)def classify(action, item): op, detail = next(iter(item.items())) status = detail.get("status", 0) if 200 <= status < 300: return "committed" if status in (429, 502, 503, 504): return "retry" if status == 409: return "reconcile" return "quarantine"def backoff(attempt, cap=30.0): return random.uniform(0, min(cap, 0.5 * (2 ** attempt)))
The transport itself can fail before an item array is available. A timeout or connection loss after bytes were sent is an ambiguous batch outcome. Do not mark every action “failed.” Reconcile by stable IDs/event IDs or replay only operations whose semantics are explicitly idempotent.
3. Deterministic test batch with one failure of each important kind
Use the strict Chapter 03 mapping. Create one event twice to generate a duplicate conflict, send one invalid numeric price to generate a permanent mapping failure, include one valid index operation, and construct one stale OCC action using metadata captured before an intervening update.
{"create":{"_index":"atlasmart-products-v4-write-lab","_id":"event-9001"}}{"product_id":"event-9001","name":"Dedupe Marker","price":0.0,"stock":0,"updated_at":"2026-09-11T06:30:00Z","event_id":"event-9001"}{"create":{"_index":"atlasmart-products-v4-write-lab","_id":"event-9001"}}{"product_id":"event-9001","name":"Dedupe Marker","price":0.0,"stock":0,"updated_at":"2026-09-11T06:30:00Z","event_id":"event-9001"}{"index":{"_index":"atlasmart-products-v4-write-lab","_id":"P-730"}}{"product_id":"P-730","name":"Keyboard","price":"bad-price","stock":6,"updated_at":"2026-09-11T06:30:01Z","event_id":"event-9002"}{"index":{"_index":"atlasmart-products-v4-write-lab","_id":"P-731"}}{"product_id":"P-731","name":"Webcam","price":59.9,"stock":4,"updated_at":"2026-09-11T06:30:02Z","event_id":"event-9003"}{"update":{"_index":"atlasmart-products-v4-write-lab","_id":"P-701","if_seq_no":S,"if_primary_term":T}}{"doc":{"stock":8,"updated_at":"2026-09-11T06:30:03Z"}}
Before sending the final update item, deliberately mutate
P-701 so the saved S/T pair is stale.
The classifier should produce: one committed create, one 409
reconcile/deduplicate candidate, one permanent/quarantine
mapping error, one committed product index, and one 409 OCC
reconcile item. Exact item status text can differ, but the
categories must be evidence-driven.
4. Retry only what is safe, and reduce pressure when the server pushes back
For retryable 429/5xx items, place only those failed items on a retry queue. Use a maximum attempt count, exponential backoff with full jitter, and a producer concurrency limit. If rejection rate or tail latency rises, reduce in-flight requests rather than growing the client queue indefinitely. Record per-attempt bytes, item count, response time, statuses, and server-side rejection/indexing evidence where available.
MAX_ATTEMPTS = 6for attempt in range(MAX_ATTEMPTS): response = send_bulk(...) retry_items = [] for original, result in zip(batch, response["items"]): state = classify(original, result) ledger.record(original.event_id, attempt, state, result) if state == "retry": retry_items.append(original) if not retry_items: break time.sleep(backoff(attempt)) batch = retry_itemselse: move_remaining_to_reconcile_or_dead_letter()
A loader that achieves higher throughput by ignoring failed items or retrying without a ledger is not reliable. Measure accepted actions per second together with p95/p99 latency, retry rate, permanent failure rate, reconciliation backlog, indexing lag, and search-visibility lag.
5. Verify no silent loss, no silent duplication
After the loader reaches a terminal state, run a reconciliation
query. Compare the set of input event IDs with
committed/deduplicated/quarantined/reconcile ledger states; no
input may be missing. Search the disposable index only after an
explicit refresh or a wait_for visibility boundary,
and separately preserve the original write acknowledgement
timestamps. Count stable document IDs and verify the successful
product sources.
POST atlasmart-products-v4-write-lab/_refreshGET atlasmart-products-v4-write-lab/_countGET atlasmart-products-v4-write-lab/_search{ "size":100, "_source":["product_id","event_id","price","stock","updated_at"], "sort":[{"product_id":"asc"}], "query":{"exists":{"field":"event_id"}}}GET atlasmart-products-v4-write-lab/_stats?filter_path=indices.*.total.indexing,indices.*.total.refresh,indices.*.total.flush,indices.*.total.translog
The acceptance test is set-based: every input event has a known disposition, every retried event retains the same logical identity, committed final-state documents match expected source values, duplicate-create events are not duplicated, the forced OCC conflict is explicitly reconciled, and no permanent failure remains hidden inside a successful Bulk HTTP response.
Check your understanding
- What is the minimum reliable unit of Bulk accounting?
- Why must a timeout batch go to reconciliation rather than be marked failed?
- What should happen to a mapping error after backoff?
- How does the loader react to sustained 429s?
- How do you prove “no loss/duplication”?
Review the answers
1. The individual action/item correlated with its original stable event/document identity.
2. Because some or all actions may have committed before the response was lost.
3. Nothing useful; it is a permanent data/schema defect and should be quarantined for repair.
4. Retry only rejected idempotent items with jittered backoff and reduce concurrency/batch pressure instead of increasing queues blindly.
5. Reconcile the complete input event-ID set against terminal ledger dispositions and verify the final document/event state after an explicit visibility boundary.
6. Cleanup and reset
Only after preserving the lab evidence, delete the dedicated write-lab/event indices. Do not delete shared Chapter 01–03 fixtures if you still need them for later chapters.
DELETE atlasmart-products-v4-write-labDELETE atlasmart-product-events
Production judgment
A production loader is a controlled queueing system, not merely a fast client. Capacity must include retry headroom; security must restrict which indices and scripts the ingestion identity can use; observability must expose per-item outcomes and age of retry/reconcile queues; and deployment must preserve the ledger across process restarts. Stable IDs, OCC, and external versions solve different problems. Document the exact combination used by each event class, then test server restarts, partial failures, credential expiration, throttling, and ambiguous timeouts before trusting the pipeline.
Summary and next step
Chapter 04 now separates operation intent, Bulk envelope success, search visibility, recovery/segment maintenance, concurrency preconditions, and idempotent retry behavior. Chapter 05 builds on these reliable writes to show how analyzers transform text into token streams and therefore define what full-text matching can mean.
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.
- OpenSearch Refresh API — Visibility controls used by final reconciliation.
- Elastic Python client bulk helpers — Official Elasticsearch client helper patterns for production applications.
- OpenSearch Python client — Official OpenSearch Python client guidance.