Chapter 11 · Ingest Pipelines, Processors, Enrichment, Failure Handling, and Data Quality

Build and Benchmark an Ingest Pipeline Against Upstream ETL Alternatives and Define Ownership Boundaries

Benchmark the complete AtlasMart ingest path against upstream ETL, measure latency and throughput rather than guessing, and define explicit ownership boundaries for parsing, schema enforcement, enrichment, redaction, routing, and replay.

Intermediate105–125 minutesIngest pipeline & data-quality labElasticsearch 9.5.3 · OpenSearch 3.8.0Last reviewed: September 2026

Learning outcomes

AtlasMart can parse and enrich logs in several places: application code, an agent/shipper, Data Prepper/Logstash/stream processing, or the search cluster's ingest pipeline. The easiest place to add one processor is not automatically the safest architecture. This lesson makes the ownership decision measurable.

01

Build a repeatable direct-vs-ingest-vs-upstream benchmark with correctness gates.

02

Separate transformation CPU from indexing, refresh, network and storage effects.

03

Define ownership boundaries for parsing, enrichment, redaction, schema validation and routing.

04

Use replayability and failure isolation as architecture criteria, not afterthoughts.

05

Publish a release gate with measured p95/p99 latency, throughput, failure and resource evidence.

Chapter baseline reviewed 11 September 2026

The reproducible examples target self-managed Elasticsearch 9.5.3 and OpenSearch 3.8.0 using the established AtlasMart lab conventions: Elasticsearch on https://localhost:9200 with the copied CA certificate, OpenSearch on https://localhost:9201 with the disposable demo certificate explicitly treated as local-only, pinned server versions, and no moving latest tags. OpenSearch documentation explicitly recommends Data Prepper for larger/complex preprocessing and describes in-cluster ingest pipelines as suitable for simpler preprocessing/smaller workloads. Elastic likewise documents that ingest pipelines run per document and cannot serve as a general multi-document transformation engine; the enrich processor is a deliberate exception for lookup data.

Execution and measurement note

The generation environment does not run the two search servers. Commands are reviewed deterministic lab specifications and expected invariants, not fabricated captured output. Execute them against disposable AtlasMart resources and record your own processor timings, CPU, throughput, error counts, response bodies, and security behavior before making production decisions.

Execution and safety note

Run mutating, destructive, security, lifecycle, snapshot, failure-injection, and load-test commands only in the disposable AtlasMart lab or an equivalently isolated environment. Verify the target cluster, index, tenant, credentials, and rollback path before execution; treat shown output as an expected invariant unless the lesson explicitly labels it as captured evidence.

1. Define three comparable paths

Path Search-cluster work Best questions to ask
A · direct normalized write Index already-normalized JSON What is the indexing baseline without transformation?
B · in-cluster ingest Parse/normalize/enrich before indexing Does convenience justify CPU/tail-latency/backpressure cost?
C · upstream ETL External worker produces normalized JSON; search cluster indexes it Does moving CPU out improve isolation/replay enough to justify another component?

All three paths must produce the same canonical document fields. A faster path that changes semantics has failed the benchmark.

2. Correctness fixture first

Canonical AtlasMart expected document
{
  "event_id":"evt-bench-001",
  "@timestamp":"2026-09-11T08:30:00Z",
  "service":"catalog-api",
  "level":"INFO",
  "tenant_id":"tenant-a",
  "client_ip":"192.0.2.10",
  "http_method":"GET",
  "url_path":"/products/P-1001",
  "status_code":200,
  "duration_ms":12.4,
  "ingest":{
    "pipeline_version":"v1"
  }
}

For each path, index a deterministic corpus containing valid events, malformed timestamps, missing optional fields, lookup misses, high Unicode, long messages and duplicate event IDs. Compare canonical fields, quarantine classification and stable IDs before recording any throughput number.

3. Benchmark protocol

Measurement protocol
# Before each trial:
GET _nodes/stats/ingest,os,process,jvm,indices,thread_pool
GET atlasmart-telemetry-v11/_stats

# Run fixed corpus at fixed concurrency and fixed refresh policy.
# Collect client-side:
#   total docs, elapsed time, docs/s,
#   per-request p50/p95/p99 latency,
#   HTTP/Bulk per-item failures and retries.
# Collect cluster-side deltas:
#   ingest count/time/failed by pipeline and processor,
#   process CPU, JVM GC, indexing time, rejected thread-pool work,
#   store growth, refresh/merge activity, indexing lag.

# Repeat warm trials; do not compare one cold run with one warm run.

Use the same document IDs and target topology. If Path C includes a network hop to an external ETL service, include its CPU, queue time and failure/retry metrics; otherwise you are merely moving cost off the dashboard.

4. Ownership matrix

Concern Preferred owner Reason
Producer-semantic validation Producer or upstream contract layer Reject bad domain data as close to source as practical.
Simple field rename/type normalization Either upstream or ingest pipeline Choose from operational ownership and measured cost.
Secrets/token redaction Before durable indexing; preferably before leaving trusted boundary Search cluster should not become the first place a secret is noticed.
Complex parsing/multi-source joins Upstream ETL/stream processor Better isolation, state, replay and scaling model.
Elastic stable reference lookup Elastic enrich can be appropriate Use only with explicit freshness/version and benchmarked cost.
OpenSearch generic business lookup Upstream/Data Prepper/application by default Core 3.8 processor catalog is not the same as Elastic enrich policies.
Final index-specific normalization/routing Ingest pipeline can be appropriate Closest to index contract; keep bounded and observable.
Historical repair/backfill Replay/reindex job Ingest pipelines only run when documents are indexed; existing data is not retroactively changed.

5. Backpressure and failure domains

If transformation runs inside the search cluster, a CPU-heavy parser can slow ingestion and consume headroom needed for search, merges and recovery. If it runs upstream, queues can absorb bursts and isolate cluster CPU, but you now own another availability domain. The correct choice depends on measured saturation behavior, replay requirements and team boundaries.

Controlled load acceptance criteria
correctness fixtures: PASS
quarantine fixtures: PASS
no silent field loss: PASS
stable event IDs / replay idempotency: PASS
p95 indexing latency: MEASURED <= SLO
p99 indexing latency: MEASURED <= SLO
sustained docs/s: MEASURED >= target
pipeline processor time/doc: MEASURED
search-cluster CPU headroom: MEASURED >= policy
indexing lag: MEASURED <= SLO
failure/quarantine rate: MEASURED and explained
replay drill: PASS
security/tenant authorization tests: PASS

6. Deliberately wrong benchmark: compare only average docs/s

Average throughput hides the failure mode that users feel first: tail latency and backlog growth. It also rewards pipelines that silently drop or misparse records. Repair by gating on correctness, then reporting p50/p95/p99 latency, sustained throughput, CPU/GC, retries, indexing lag, quarantine rate and output equivalence. Preserve the test corpus and environment metadata so future versions can reproduce the decision.

7. Deployment and rollback contract

Pipeline release sequence
1. Version pipeline/ETL configuration in source control.
2. Run simulation/unit fixtures for good + bad inputs.
3. Deploy under a new ID/version, do not mutate blindly.
4. Shadow/canary a bounded traffic slice.
5. Compare transformed output + performance against baseline.
6. Promote producers/index settings to the new version.
7. Monitor failure rate, p95/p99, CPU, lag, quarantine age.
8. Keep previous version and replay compatibility during rollback window.
9. Retire only after rollback/replay evidence is complete.

Pipeline updates are schema changes in disguise. Coordinate them with mappings/templates from Chapters 03 and 09 and with data-stream/backing-index creation from Chapter 10. A new processor output field can fail a strict mapping even when simulation itself succeeds unless the target mapping is part of the test.

8. Security, tenancy and auditability

Separate privileges for pipeline management, lookup-source writes, normal ingest, quarantine replay and failure-index reads. Never let a tenant-controlled field select arbitrary destination indices or routing without validation. Log pipeline-version changes and keep raw/replay evidence according to retention policy. Managed services may expose different processor/plugin sets and permission models; deploy only what the target service documents and test the real account.

Check your understanding

  1. Why must correctness precede throughput in a pipeline benchmark?
  2. What cluster metric helps attribute in-cluster transformation cost?
  3. When does upstream ETL become preferable?
  4. Why keep prior pipeline versions during rollout?
  5. Why is a pipeline change tied to mapping/template tests?
Review the answers

1. A fast pipeline that changes semantics or drops failures is not a valid candidate.

2. Ingest node/pipeline/processor count-time-failure deltas, paired with CPU/JVM and indexing metrics.

3. When transformation is complex/stateful, needs joins/replay/isolation, or consumes unacceptable search-cluster headroom.

4. For rollback and deterministic replay of quarantined/raw events.

5. Processor output must still satisfy the target schema; transformation and mapping are one end-to-end contract.

Production judgment

Use in-cluster ingest for bounded, index-adjacent transformations whose cost and failure behavior are well understood. Move complex ETL upstream when it needs durable queues, multi-record state, joins, network lookups or independent scaling. Re-evaluate after version upgrades because processor implementations and managed-service limits change. Do not publish universal batch sizes or CPU thresholds; derive them from representative load and explicit SLOs.

Chapter 11 release gate

Evidence manifest
server baseline: Elasticsearch 9.5.3 / OpenSearch 3.8.0
pipeline definitions versioned: PASS
processor availability verified on target nodes: PASS
good/bad simulation fixtures: PASS
strict target mapping compatibility: PASS
quarantine route + access controls: PASS
lookup enrichment freshness/version test: PASS
bulk per-item error inspection: PASS
replay idempotency drill: PASS
p95/p99 + sustained throughput: MEASURED
processor CPU/time deltas: MEASURED
indexing lag/backpressure: MEASURED
rollback to previous pipeline: PASS

Replace MEASURED with actual lab/production-like results. The chapter intentionally does not fabricate throughput or latency values.

Summary and next step

Chapter 11 turns ingest into an observable, owned transformation system: deterministic fixtures, processor choice, explicit lookup freshness, quarantine/replay, resource measurements and deployment versions. Chapter 12 uses the same discipline for zero-downtime index changes through aliases, rollover, reindexing and task control.

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.