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.

Intermediate → Advanced135–160 minutesBulk-loader integration labElasticsearch 9.5.3 · OpenSearch 3.8.0Last reviewed: September 2026

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.

01

Define an ingestion record with stable event ID, target document ID, operation intent, attempt metadata, and reconciliation state.

02

Build NDJSON batches, correlate responses by position, and partition successes, retryable failures, permanent failures, and semantic conflicts.

03

Implement bounded exponential backoff with jitter and a concurrency limit that responds to 429/5xx pressure.

04

Force and resolve one sequence-number conflict and one duplicate create without duplicating a non-idempotent effect.

05

Verify source count/IDs, event identity, acknowledgement-versus-visibility timing, and final dead-letter/reconciliation queues before declaring success.

Chapter baseline reviewed 11 September 2026

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.

Execution and safety note

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.

python · platform-neutral bulk classifier core
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.

bulk test · illustrative NDJSON with explicit correlation
{"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.

python · bounded retry loop sketch
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()
Do not benchmark by deleting the evidence.

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.

portable REST · final verification
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

  1. What is the minimum reliable unit of Bulk accounting?
  2. Why must a timeout batch go to reconciliation rather than be marked failed?
  3. What should happen to a mapping error after backoff?
  4. How does the loader react to sustained 429s?
  5. 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.

portable REST · cleanup only named Chapter 04 resources
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

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.