Chapter 11 · Ingest Pipelines, Processors, Enrichment, Failure Handling, and Data Quality
Ingest Nodes/Pipelines, Processor Execution, on_failure, Conditional Processing, and Simulation
Follow one AtlasMart log event through an ingest node, ordered processors, conditions, simulation, failure handling, metadata, and pipeline statistics before allowing the transformed document to become searchable.
Learning outcomes
AtlasMart receives access logs from multiple services. The raw line is useful evidence, but search and alerting need typed fields such as timestamp, service, tenant, status code and latency. An ingest pipeline is a named, ordered sequence of processors executed before indexing. This lesson makes that execution path visible so a failed parser is not mistaken for a storage or search failure.
Explain where ingest processing runs and why processor order is part of the data contract.
Distinguish processor-level ignore_failure/on_failure from pipeline-level failure handling.
Use conditions and tags to make optional branches and processor cost observable.
Simulate good and bad events before indexing and interpret _ingest metadata correctly.
Use node ingest statistics to separate transformation cost from downstream indexing/search cost.
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 ordered ingest processors, conditions, simulation, on_failure and ingest statistics. Elasticsearch and OpenSearch also evolve their processor catalogs independently, so production code should discover/verify available processors instead of assuming identical catalogs.
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.
1. The write path has a transformation stage
A client sends an index or bulk operation, optionally naming a pipeline. The coordinating path sends the document through an ingest-capable node. Each processor sees the document resulting from the previous processor. Only after the pipeline finishes does the transformed document proceed to normal indexing, routing, primary-shard execution and replication. Therefore an acknowledgement is about the transformed document, not the original wire payload.
| Term | Mechanism | Operational consequence |
|---|---|---|
| Ingest node | Node allowed to execute ingest pipelines | CPU-heavy parsing can compete with write coordination if capacity is not budgeted. |
| Processor | One ordered transformation/validation action | Changing order can change output or failure behavior. |
| Condition | Expression controlling whether a processor runs | Optional fields should be modeled explicitly rather than hidden by broad ignore_failure. |
| on_failure | Processors invoked after a failure | Can annotate/quarantine a poison record or make a failure observable. |
| Simulation | Pipeline execution without committing the document | Use it as a deterministic contract test before deployment. |
| Pipeline stats | Cumulative count/time/current/failed metrics | Useful for deltas and bottleneck localization, not a benchmark substitute by itself. |
2. Build the disposable AtlasMart targets
PUT atlasmart-telemetry-v11
{
"settings": {"number_of_shards":1,"number_of_replicas":0},
"mappings": {
"dynamic":"strict",
"properties": {
"@timestamp":{"type":"date"},
"event_id":{"type":"keyword"},
"service":{"type":"keyword"},
"level":{"type":"keyword"},
"tenant_id":{"type":"keyword"},
"client_ip":{"type":"ip"},
"http_method":{"type":"keyword"},
"url_path":{"type":"keyword"},
"status_code":{"type":"integer"},
"duration_ms":{"type":"double"},
"message":{"type":"text"},
"ingest":{"properties":{
"pipeline_version":{"type":"keyword"},
"received_at":{"type":"date"},
"failed":{"type":"boolean"},
"error":{"type":"text","index":false}
}}
}
}
}
PUT atlasmart-ingest-failures-v1
{
"settings":{"number_of_shards":1,"number_of_replicas":0},
"mappings":{"dynamic":true}
}
The success index is strict so the pipeline cannot accidentally invent fields forever. The failure index is intentionally permissive in this small lab because its job is preservation of diagnostic evidence; in production, quarantine schemas should also be governed and access controlled because failed records often contain malformed or sensitive input.
3. Define an ordered portable pipeline
PUT _ingest/pipeline/atlasmart-log-ingest-v1
{
"description":"AtlasMart deterministic access-log parser v1",
"processors":[
{"set":{"tag":"stamp-pipeline-version","field":"ingest.pipeline_version","value":"v1"}},
{"set":{"tag":"stamp-received-at","field":"ingest.received_at","value":"{{{_ingest.timestamp}}}"}},
{"dissect":{
"tag":"parse-access-line",
"field":"message",
"pattern":"%{event_time}|%{service}|%{level}|%{tenant_id}|%{client_ip}|%{http_method}|%{url_path}|%{status_code}|%{duration_ms}"
}},
{"date":{"tag":"parse-event-time","field":"event_time","target_field":"@timestamp","formats":["ISO8601"]}},
{"convert":{"tag":"status-to-int","field":"status_code","type":"integer"}},
{"convert":{"tag":"duration-to-double","field":"duration_ms","type":"double"}},
{"remove":{"tag":"remove-temp-time","field":"event_time"}}
],
"on_failure":[
{"set":{"field":"_index","value":"atlasmart-ingest-failures-v1"}},
{"set":{"field":"ingest.failed","value":true}},
{"set":{"field":"ingest.error","value":"pipeline atlasmart-log-ingest-v1 failed"}}
]
}
Notice the tags. A processor tag is not decoration: it gives
operators a stable stage name in diagnostics and per-processor
statistics. The pipeline-level on_failure changes
_index so a poison record is quarantined instead of
being silently discarded. The lab stores a stable pipeline
version and ingest receipt time so later replays can prove which
transformation contract ran.
A processor-specific on_failure handles that processor and permits later processors to continue. A pipeline-level on_failure is the fallback for an unhandled processor exception and ends the normal pipeline path. Do not use ignore_failure simply to make dashboards green: ignored parse errors become data-quality defects.
4. Simulate before writing
POST _ingest/pipeline/atlasmart-log-ingest-v1/_simulate?verbose=true
{
"docs":[
{"_index":"atlasmart-telemetry-v11","_id":"evt-good-1","_source":{
"event_id":"evt-good-1",
"message":"2026-09-11T07:30:00Z|catalog-api|INFO|tenant-a|192.0.2.10|GET|/products/P-1001|200|12.4"
}},
{"_index":"atlasmart-telemetry-v11","_id":"evt-bad-1","_source":{
"event_id":"evt-bad-1",
"message":"malformed|record"
}}
]
}
Expected invariant for the good sample: the final source
contains typed status_code=200,
duration_ms=12.4, an ISO timestamp, stable
service/tenant fields and
ingest.pipeline_version=v1. Expected invariant for
the bad sample: the normal parse path fails and the
pipeline-level handler changes its target to
atlasmart-ingest-failures-v1 while retaining the
original message. Exact verbose response shape is version
dependent; assert fields and destination, not incidental
ordering.
5. Conditions should express intentional optionality
PUT _ingest/pipeline/atlasmart-log-ingest-v1-cond
{
"processors":[
{"uppercase":{
"tag":"normalize-level",
"field":"level",
"if":"ctx.level != null"
}},
{"set":{
"tag":"mark-server-error",
"field":"is_server_error",
"value":true,
"if":"ctx.status_code != null && ctx.status_code >= 500"
}}
]
}
A condition says “this field is optional and this branch is
defined.” ignore_failure says “this processor
failed and we decided to continue.” Those are different
operational facts and should not be conflated.
6. Index through the pipeline, then inspect evidence
POST atlasmart-telemetry-v11/_doc/evt-good-1?pipeline=atlasmart-log-ingest-v1&refresh=true
{
"event_id":"evt-good-1",
"message":"2026-09-11T07:30:00Z|catalog-api|INFO|tenant-a|192.0.2.10|GET|/products/P-1001|200|12.4"
}
POST atlasmart-telemetry-v11/_doc/evt-bad-1?pipeline=atlasmart-log-ingest-v1&refresh=true
{
"event_id":"evt-bad-1",
"message":"malformed|record"
}
GET atlasmart-telemetry-v11/_doc/evt-good-1
GET atlasmart-ingest-failures-v1/_doc/evt-bad-1
GET _nodes/stats/ingest?filter_path=nodes.*.ingest
Do not infer throughput from one request. The node statistics are cumulative counters: capture a before snapshot, run a controlled batch, capture an after snapshot, then compute deltas. Pair that with client-side elapsed time and system CPU so you can distinguish processor time from networking, indexing, refresh and storage effects.
7. Deliberately wrong approach: hide every error
Setting ignore_failure:true on the parser can make
writes appear successful while required fields remain absent.
AtlasMart then gets a worse failure: documents are searchable
but semantically unusable for dashboards, alerting or tenant
filters. Repair by making required parsing fail closed,
quarantine the raw event with a stable error classification, and
reserve optional conditions for genuinely optional fields.
Check your understanding
- Why is processor order part of the schema contract?
- What is the safest first execution environment for a changed pipeline?
- Why tag processors?
- Why is ignore_failure dangerous on required parsing?
- What does _nodes/stats/ingest prove?
Review the answers
1. Each processor receives the previous processor output, so reordering can change both values and failures.
2. The simulate API with deterministic good and bad fixtures.
3. Tags make failures and per-processor performance attributable to stable named stages.
4. It can turn visible ingestion failures into silently malformed indexed documents.
5. Cumulative pipeline/processor execution counts, time, current work and failures; use deltas rather than treating it as an end-to-end benchmark.
Production judgment
Keep pipeline logic bounded, versioned and testable. Measure p95/p99 end-to-end indexing latency and sustained throughput with and without the pipeline, plus processor time, CPU, failure rate, quarantine growth and indexing lag. Protect pipeline-management privileges, because a malicious or accidental pipeline change can rewrite tenant fields, routing, destination indices or sensitive data. Snapshot/recovery protects indexed state; it does not replace retaining replayable raw events when your business requires deterministic reprocessing.
Summary and next step
You can now trace a document through an ingest pipeline and distinguish transformation failure from indexing failure. Next we choose specific processors from the structure of the source data and compare their correctness and CPU tradeoffs.
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.