Chapter 11 · Ingest Pipelines, Processors, Enrichment, Failure Handling, and Data Quality
Dead-Letter/Failure Indices, Observability, Retry, and Preventing Poison Documents from Blocking Pipelines
Keep malformed or poison events observable without silently dropping them or blocking healthy ingestion: classify failures, quarantine records, preserve evidence, retry only when safe, and monitor processor-level failure cost.
Learning outcomes
A poison record is an input that repeatedly violates the ingestion contract: malformed timestamp, invalid IP, incompatible type, oversized field, or an enrichment key that violates policy. Retrying it blindly can create an infinite failure loop and starve healthy traffic. AtlasMart therefore needs a failure lane with evidence, ownership and replay rules.
Distinguish transient transport/indexing failures from deterministic transformation failures.
Quarantine poison records without losing the original payload or failure context.
Use pipeline and processor metrics to detect failure-rate and latency changes.
Design safe replay with a corrected pipeline version and stable document identity.
Prevent failure indices from becoming an ungoverned sensitive-data dump.
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. Both products support on_failure and ingest statistics, but failure metadata details and processor catalogs are version specific. The portable requirement is stronger: every quarantined event must preserve original evidence, a stable failure class, pipeline version, destination intent and replay identity.
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. Classify before retrying
| Failure class | Example | Retry policy |
|---|---|---|
| Deterministic data-quality failure | timestamp=not-a-date, impossible conversion | Do not blind-retry; quarantine, correct producer/parser, then replay. |
| Transient cluster/backpressure failure | 429/rejection, temporary unavailable shard | Retry with bounded exponential backoff and idempotent identity. |
| Authentication/authorization failure | expired credential, denied pipeline/index permission | Stop and repair security/configuration; repeated retries add noise. |
| Deployment mismatch | pipeline references unavailable processor/plugin | Rollback or restore capability; do not mutate data to hide deployment drift. |
| Schema contract failure | strict mapping rejects unexpected field/type | Decide whether producer is wrong or schema must evolve; never coerce blindly. |
2. Quarantine through pipeline-level failure handling
PUT _ingest/pipeline/atlasmart-log-ingest-v2
{
"description":"AtlasMart parser with explicit quarantine",
"processors":[
{"set":{"field":"ingest.pipeline_version","value":"v2"}},
{"set":{"field":"ingest.received_at","value":"{{{_ingest.timestamp}}}"}},
{"dissect":{"tag":"parse-line","field":"message","pattern":"%{event_time}|%{service}|%{level}|%{tenant_id}|%{status_code}|%{duration_ms}"}},
{"date":{"tag":"parse-date","field":"event_time","target_field":"@timestamp","formats":["ISO8601"]}},
{"convert":{"tag":"status-int","field":"status_code","type":"integer"}},
{"convert":{"tag":"duration-double","field":"duration_ms","type":"double"}}
],
"on_failure":[
{"set":{"field":"_index","value":"atlasmart-ingest-failures-v1"}},
{"set":{"field":"ingest.failed","value":true}},
{"set":{"field":"ingest.failure_class","value":"TRANSFORM"}},
{"set":{"field":"ingest.failed_pipeline","value":"atlasmart-log-ingest-v2"}}
]
}
Elastic also exposes failure metadata such as the failing
processor type/tag/message inside on_failure.
OpenSearch documentation emphasizes failure handling and metrics
but does not require the course's portable contract to depend on
identical failure-metadata names. If you use product-specific
metadata, normalize it into stable application fields before
indexing the quarantine record.
3. Prove healthy and poison events separate cleanly
POST _bulk?pipeline=atlasmart-log-ingest-v2&refresh=true
{"index":{"_index":"atlasmart-telemetry-v11","_id":"evt-good-4"}}
{"event_id":"evt-good-4","message":"2026-09-11T08:00:00Z|checkout-api|INFO|tenant-a|200|18.2"}
{"index":{"_index":"atlasmart-telemetry-v11","_id":"evt-bad-4"}}
{"event_id":"evt-bad-4","message":"2026-99-99|broken"}
GET atlasmart-telemetry-v11/_doc/evt-good-4
GET atlasmart-ingest-failures-v1/_doc/evt-bad-4
GET _nodes/stats/ingest?filter_path=nodes.*.ingest
Inspect every Bulk API item even when the HTTP request itself
succeeds. A pipeline handler that changes
_index may turn a transformation exception into a
successful quarantine index operation; that is expected only if
your application records and monitors quarantine counts as
failures at the business level.
4. Preserve replay identity and immutable evidence
{
"event_id":"evt-bad-4",
"original_target":"atlasmart-telemetry-v11",
"original_message":"2026-99-99|broken",
"ingest":{
"failed":true,
"failure_class":"TRANSFORM",
"failed_pipeline":"atlasmart-log-ingest-v2",
"received_at":"...",
"replay_count":0
}
}
Do not overwrite the raw event during triage. Replay should create or update the canonical destination using the same stable event ID so an ambiguous network retry does not duplicate the logical event. Add an audit trail for replay attempts and the pipeline version that finally succeeded.
5. Retry algorithm: bounded and classification-driven
for item in failed_items:
if item.failure_class == "TRANSIENT":
retry_with_exponential_backoff_and_jitter(item, max_attempts=bounded)
elif item.failure_class == "TRANSFORM":
quarantine(item) # human/config fix first
elif item.failure_class == "AUTH":
stop_pipeline_and_page_owner(item)
else:
quarantine_and_investigate(item)
# Replay only after a new parser/schema version passes fixtures.
# Use the original stable event_id as the canonical destination _id.
The retry worker must preserve ordering requirements where they exist, enforce tenant boundaries, and cap concurrent retries so a downstream incident does not become a retry storm.
6. Wrong approach: “on_failure means the problem is solved”
An on_failure branch can prevent loss, but it can
also hide an outage if nobody monitors the failure lane. A
pipeline with 100% quarantined events can report successful HTTP
indexing to the failure index. Repair by treating quarantine
rate, oldest-unreplayed age, replay backlog and processor
failure deltas as first-class SLO signals.
7. Failure-index governance
Quarantine records are often more sensitive than successful normalized documents because they contain raw headers, malformed payloads and unexpected fields. Give the failure index a restrictive mapping where possible, separate write/read roles, encrypt transport, audit access, cap retention, and redact secrets before persistence when feasible. Never expose the failure index directly to tenant search users.
GET _nodes/stats/ingest?filter_path=nodes.*.ingest
GET atlasmart-ingest-failures-v1/_count
GET atlasmart-ingest-failures-v1/_search
{
"size":0,
"aggs":{
"by_class":{"terms":{"field":"ingest.failure_class.keyword","size":20}},
"oldest":{"min":{"field":"ingest.received_at"}}
}
}
If you keep ingest.failure_class as a keyword in
production, query that keyword directly; this lab failure index
is intentionally permissive, so inspect its dynamic mapping
first.
Check your understanding
- Why should a deterministic parse failure not be retried immediately forever?
- Can a Bulk API HTTP 200 still contain failed work?
- Why preserve the original event ID?
- What must accompany a quarantine mechanism?
- Why is the failure index security-sensitive?
Review the answers
1. The same input/pipeline combination will keep failing and can starve healthy traffic.
2. Yes. Inspect each item result; pipeline quarantine can also be a successful write to the failure destination.
3. It enables idempotent replay into the canonical destination and prevents duplicate logical events.
4. Metrics, ownership, bounded retention, replay procedure, security controls and alert thresholds.
5. It retains raw/malformed data that can contain fields or secrets removed from normal documents.
Production judgment
Budget failure handling as part of normal capacity planning. Monitor failures per processor/source/tenant, p95/p99 ingest latency, retry rate, quarantine backlog and replay age. Test failure injection with small isolated batches. A failure path that has never been replayed is unproven. Backups protect stored quarantine and successful indices; source-system replay or durable queues are needed when the business requires reconstruction after pipeline defects.
Summary and next step
You can now keep poison records from blocking healthy ingestion without hiding them. The final lesson measures where transformation belongs—inside the cluster or upstream—and turns that choice into an explicit ownership decision.
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.