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.

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

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.

01

Explain where ingest processing runs and why processor order is part of the data contract.

02

Distinguish processor-level ignore_failure/on_failure from pipeline-level failure handling.

03

Use conditions and tags to make optional branches and processor cost observable.

04

Simulate good and bad events before indexing and interpret _ingest metadata correctly.

05

Use node ingest statistics to separate transformation cost from downstream indexing/search cost.

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. 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.

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.

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

Dev Tools · disposable indices
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

Dev Tools · 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.

Failure handling changes control flow

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

Dev Tools · good and poison samples
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

Conditional normalization example
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

Dev Tools · write and observe
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

  1. Why is processor order part of the schema contract?
  2. What is the safest first execution environment for a changed pipeline?
  3. Why tag processors?
  4. Why is ignore_failure dangerous on required parsing?
  5. 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

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.