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.
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.
Build a repeatable direct-vs-ingest-vs-upstream benchmark with correctness gates.
Separate transformation CPU from indexing, refresh, network and storage effects.
Define ownership boundaries for parsing, enrichment, redaction, schema validation and routing.
Use replayability and failure isolation as architecture criteria, not afterthoughts.
Publish a release gate with measured p95/p99 latency, throughput, failure and resource evidence.
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.
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.
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
{
"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
# 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.
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
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
- Why must correctness precede throughput in a pipeline benchmark?
- What cluster metric helps attribute in-cluster transformation cost?
- When does upstream ETL become preferable?
- Why keep prior pipeline versions during rollout?
- 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
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
- Elastic ingest pipelines — Pipeline execution, conditionals, failure handling, node statistics, default/final pipeline concepts and operational limitations.
- Elastic simulate pipeline API — Test an existing or inline pipeline; verbose mode exposes per-processor intermediate results.
- Elastic ingest error handling — Processor-level and pipeline-level on_failure semantics and ingest failure metadata.
- Elastic enrich processor setup — Source data, enrich policy execution, generated enrich index and ingest processor workflow.
- Elastic enrich processor reference — policy_name, match field, target field, max_matches and lookup behavior.
- Elastic script processor — Painless ingest context and script processor behavior.
- OpenSearch ingest pipelines — Pipeline model, ingest node prerequisite and guidance on using Data Prepper for larger or more complex preprocessing.
- OpenSearch ingest processors — Current core processor inventory and Nodes Info inspection endpoint.
- OpenSearch pipeline failures — Failure handling and ingest-pipeline node statistics.
- OpenSearch access data in a pipeline — Source, metadata such as _index/_routing, and _ingest.timestamp access.
- OpenSearch grok processor — Pattern-based parsing and debug controls.
- OpenSearch dissect processor — Delimiter-based parsing for stable log formats.
- OpenSearch user-agent processor — User-agent parsing and target field configuration.
- OpenSearch Nodes Stats API — Per-node, pipeline and processor ingest counts, time, current work and failures.