Chapter 11 · Ingest Pipelines, Processors, Enrichment, Failure Handling, and Data Quality
Enrich Policies and Lookup Data: Refresh, Consistency, and Deployment Workflows
Design lookup enrichment as a versioned data product: implement Elastic enrich policies correctly, contrast OpenSearch 3.8 alternatives, prove freshness boundaries, and deploy reference-data changes without assuming lookups are live.
Learning outcomes
AtlasMart wants every log to carry a stable support-team name
derived from service. The lookup changes
infrequently, but operators must know exactly when a change
becomes effective. This is an enrichment-consistency problem,
not merely a “join” problem.
Explain Elastic enrich source indices, policies, executed enrich indices and ingest lookup behavior.
Prove that editing source lookup data does not automatically refresh an already executed Elastic enrich index.
Distinguish OpenSearch 3.8 core ingest enrichment capabilities from Elastic generic enrich policies.
Design a portable upstream lookup path with an explicit enrichment version and deterministic fixture.
Deploy lookup changes with tests, freshness SLOs, rollback and performance 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. Elastic 9.5.3 provides native enrich policies/processors. OpenSearch 3.8 core ingest processors include specific enrichment processors such as geoip/ip2geo/user_agent but its documented core processor catalog does not include an Elastic-style generic enrich-policy processor. For generic business lookups, this lesson uses upstream/Data Prepper/application enrichment on the OpenSearch path unless your deployed distribution/plugin explicitly provides another capability.
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. Enrichment is a consistency contract
| Question | Why it matters |
|---|---|
| What is authoritative? | The lookup source—not the denormalized copy embedded in each indexed event. |
| When does a lookup change take effect? | Determines how long newly indexed documents can carry stale reference data. |
| Is historical data rewritten? | Usually not automatically; decide whether old enriched values remain historical truth or require reprocessing. |
| How is the enrichment version recorded? | Allows support and replay jobs to explain which lookup snapshot produced a field. |
| What happens on lookup miss? | Must be explicit: null/unknown, quarantine, fail closed or continue without enrichment. |
2. Elastic native enrich: build, execute, use
PUT atlasmart-service-directory-v1
{
"mappings":{"properties":{
"service":{"type":"keyword"},
"team":{"type":"keyword"},
"tier":{"type":"keyword"},
"lookup_version":{"type":"keyword"}
}}
}
PUT atlasmart-service-directory-v1/_doc/catalog-api?refresh=true
{"service":"catalog-api","team":"search-platform","tier":"1","lookup_version":"2026-09-11.1"}
PUT /_enrich/policy/atlasmart-service-owner-v1
{
"match":{
"indices":"atlasmart-service-directory-v1",
"match_field":"service",
"enrich_fields":["team","tier","lookup_version"]
}
}
POST /_enrich/policy/atlasmart-service-owner-v1/_execute
Executing the policy creates a read-only, optimized internal enrich index. The lookup used at ingest is therefore derived materialized state. Treat policy execution as a deployment action and record its completion before sending traffic that depends on the new reference data.
PUT _ingest/pipeline/atlasmart-enrich-v1
{
"processors":[
{"enrich":{
"policy_name":"atlasmart-service-owner-v1",
"field":"service",
"target_field":"service_owner",
"max_matches":1
}},
{"set":{"field":"ingest.enrichment_contract","value":"service-owner-v1"}}
]
}
POST _ingest/pipeline/atlasmart-enrich-v1/_simulate
{"docs":[{"_source":{"service":"catalog-api"}}]}
3. Staleness test: source edit is not a live join
# Change authoritative source lookup data.
PUT atlasmart-service-directory-v1/_doc/catalog-api?refresh=true
{"service":"catalog-api","team":"commerce-platform","tier":"1","lookup_version":"2026-09-11.2"}
# Simulate BEFORE re-executing the enrich policy.
POST _ingest/pipeline/atlasmart-enrich-v1/_simulate
{"docs":[{"_source":{"service":"catalog-api"}}]}
# Rebuild the enrich materialization.
POST /_enrich/policy/atlasmart-service-owner-v1/_execute
# Simulate AFTER execution and compare lookup_version/team.
POST _ingest/pipeline/atlasmart-enrich-v1/_simulate
{"docs":[{"_source":{"service":"catalog-api"}}]}
The acceptance criterion is not a hard-coded response body; it is the state transition: before policy re-execution, the pipeline can still resolve from the prior enrich materialization; after successful execution, new ingests should carry the new lookup version. Existing indexed documents are not retroactively rewritten.
An enrich policy is not an always-live distributed join. If the source index changes, operators must execute the policy again to build refreshed enrich data. Treat freshness as a deployment/SLO question.
4. OpenSearch path: do not invent an Elastic-compatible API
OpenSearch 3.8 supports several enrichment-like ingest processors, including GeoIP/IP2Geo and user-agent parsing, but its current documented core processor inventory does not expose a generic Elastic-style enrich-policy processor. For AtlasMart's service-owner lookup, keep the cross-platform contract at the application/data-pipeline level: enrich before indexing using a versioned lookup cache/table, or use a supported Data Prepper/application transformation stage. Then send the already enriched document through the search ingest pipeline for validation and final normalization.
{
"event_id":"evt-3001",
"service":"catalog-api",
"service_owner":{
"team":"search-platform",
"tier":"1",
"lookup_version":"2026-09-11.1"
},
"ingest":{
"enrichment_contract":"service-owner-v1"
}
}
This makes business enrichment portable even though the mechanism is not. It also avoids hiding a product-specific API behind a fake common abstraction.
5. Deployment workflow
1. Update lookup source under change control.
2. Validate uniqueness/cardinality of the match key.
3. Build/refresh derived lookup materialization.
- Elasticsearch: execute enrich policy.
- OpenSearch path: refresh/version upstream lookup cache or Data Prepper source.
4. Simulate canonical hit, miss and malformed-key fixtures.
5. Verify expected lookup_version in transformed output.
6. Run throughput/latency comparison against previous version.
7. Promote pipeline/configuration.
8. Monitor miss rate, ingest failures, CPU and freshness age.
9. Keep prior lookup materialization/configuration until rollback window closes.
6. Wrong approach: assume reference edits are instantly visible
If AtlasMart changes ownership from Team A to Team B and assumes every new log immediately carries Team B, incident routing can be wrong for hours or days depending on how enrichment is materialized. Repair by defining a lookup-freshness SLO, stamping the version into every enriched document, monitoring version distribution, and making refresh/execution a release step.
7. Validate performance and miss semantics
GET _nodes/stats/ingest?filter_path=nodes.*.ingest
# For a fixed corpus record:
# - lookup hit count / miss count
# - output lookup_version distribution
# - pipeline failure count
# - ingest time delta per document
# - end-to-end docs/s and p95/p99 indexing latency
# - CPU and indexing lag
# - behavior when source lookup is unavailable during deployment
Enrichment can dominate ingest cost if every document performs expensive lookup work. Elastic recommends benchmarking enrich processors and notes they fit reference data that changes infrequently rather than real-time rapidly changing state. The same architecture principle applies to upstream lookup mechanisms.
Check your understanding
- Does updating an Elastic enrich source index automatically update the existing enrich index?
- Are existing indexed events rewritten after enrichment data changes?
- What should every business enrichment carry when reproducibility matters?
- Why not expose one fake enrich API for both products?
- What must happen on lookup miss?
Review the answers
1. No. Re-execute the policy to build refreshed enrich data.
2. No, not automatically; reprocessing/reindexing is a separate decision.
3. A stable lookup/enrichment version or equivalent provenance field.
4. The native capabilities differ; portability belongs in the data contract, not invented API equivalence.
5. A documented policy such as unknown value, quarantine or failure, backed by monitoring and tests.
Production judgment
Use ingest enrichment for small, stable reference data with explicit freshness. If the lookup changes continuously, requires multi-row joins, remote network calls, transactional consistency or fan-out, move it upstream. Protect lookup sources and enrichment-management privileges from tenant users. Monitor lookup miss rate and freshness age alongside throughput; stale enrichment is a correctness incident even when indexing remains green.
Summary and next step
You can now deploy enrichment without confusing reference snapshots with live joins. Next we turn failure handling into a first-class subsystem: quarantine, retry classification, replay and poison-record observability.
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.