Transformation & Processing → Stream Processing

Real-time Enrichment

Adds computed, business-meaningful context to an event as it streams — beyond the raw geo/device enrichment done at the edge.

High-Level Design

Real-time Enrichment joins events against fast lookup sources in-stream.

Data Source
Identity Resolution
Hands off an identity-resolved event
→
Ingestion
Product Catalog / Profile Caches
Fast lookup sources
→
Processing
Real-time Enrichment
.NET Core Stream Worker — profile and catalog lookups
→
Foundation
Derived Event Generation
Runs next, using enriched context
→
Intelligence
Azure Cache for Redis
Low-latency lookup cache
→
Activation
In-session Personalization
The use case this enables

💼 Business Context

  • Lets downstream consumers get rich context (customer tier, product category, current cart value) without each one re-querying multiple systems
  • Powers real-time use cases (in-session personalization) that can't wait for a batch join
  • Owned by Platform Engineering, with enrichment sources owned by the relevant domain teams

🔌 Technical Overview

Real-time Enrichment joins each event, in-stream, against fast lookup sources: the Identity & Profile Service (customer tier, LTV band), a product catalog cache (category, price tier), and current session state (running cart value). These lookups are backed by Azure Cache for Redis to keep per-event latency low, with the underlying Data Lakehouse tables as the periodically refreshed source of truth for the cache.

Enrichment Sources

Customer profile (tier, LTV band) Product catalog (category, price) Session state (cart value)

💾 Enriched Event Excerpt

"enrichment": {
  "customer_tier": "gold",
  "product_category": "audio",
  "session_cart_value": 6998
}

🔗 Integration Points

  • Azure Cache for Redis — low-latency lookup cache for profile/catalog/session data
  • Identity & Profile Service — source of customer-tier/LTV enrichment
  • Product catalog service — source of category/price enrichment
  • .NET Core Stream Worker — where the joins are performed inline

🧰 Services Consumed

  • Owning microservice — Cxos.Processing.Api (see the Full Application Service Map)
  • Database — ADLS Gen2 (Iceberg) + Cosmos DB Table API (stream checkpoints)

⚠️ Non-Functional Considerations

  • Scale: cache-backed lookups keep enrichment cost roughly constant regardless of event volume
  • Latency: adds low-single-digit milliseconds per lookup; enrichment failures degrade gracefully (event proceeds with partial enrichment)
  • Reliability: cache staleness is bounded by a refresh interval appropriate to each source's rate of change
  • Security/Privacy: only non-sensitive, business-purpose fields are added — enrichment is scoped, not an unrestricted profile dump

🎯 Enterprise Example

A gold-tier customer adds a fourth item to their cart. Real-time enrichment attaches their tier and running cart value to the event inline, letting the Activation API trigger a tier-appropriate free-shipping nudge within the same session — a use case that would be impossible if enrichment only happened in a nightly batch job.

← Back to Stream Processing