Transformation & Processing → Stream Processing

Derived Event Generation

Synthesizes new, higher-level events from patterns in the raw stream — turning 'what happened' into 'what it means'.

High-Level Design

Derived Event Generation computes business-meaningful moments once, for everyone.

Data Source
Real-time Enrichment
Hands off an enriched event
→
Ingestion
Azure Cache for Redis
Holds in-flight windowed state per customer
→
Processing
Derived Event Generation
.NET Core Stream Worker — pattern-based synthetic events
→
Foundation
Data Lakehouse
Derived events are stored like any other event
→
Intelligence
Azure Event Hubs
Derived events are republished here
→
Activation
Activation API
Cart-abandonment recovery and similar triggers

💼 Business Context

  • Business logic like 'cart abandoned' or 'browsing spree' isn't a single raw event — it's a pattern the platform computes once, so every team doesn't reinvent it
  • Standardizes business-meaningful moments across the company instead of each team defining its own version of 'abandonment'
  • Owned by Platform Engineering, with derivation rules defined in partnership with the business teams that consume them

🔌 Technical Overview

The Stream Worker maintains short-lived windowed state per customer to detect patterns and emit derived events: e.g., cart_abandoned fires when an add_to_cart event has no matching order_paid within a configurable window (default 60 minutes); browsing_spree fires after N product views within M minutes. Derived events are published back onto Azure Event Hubs with the same schema as any other event, so they flow through the rest of the pipeline identically.

Example Derived Events

cart_abandoned browsing_spree repeat_visitor price_drop_relevant

💾 Derived Event

{
  "event": "cart_abandoned",
  "event_id": "derived-3c8e...",
  "user_id": "cust_004821",
  "derived_from": ["evt_add_to_cart_991"],
  "properties": { "cart_value": 6998, "minutes_since_add": 62 }
}

🔗 Integration Points

  • .NET Core Stream Worker — windowed pattern detection and derivation logic
  • Azure Cache for Redis — holds in-flight windowed state per customer
  • Azure Event Hubs — derived events are republished here, same as primary events
  • Activation API — the most common consumer of derived events (e.g., cart-abandonment recovery)

🧰 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: windowed state is sharded by customer ID and scales with the Stream Worker's partition count
  • Latency: derived events fire as soon as their defining window/condition is met, typically within minutes
  • Reliability: derivation rules are versioned; a rule change doesn't retroactively alter already-emitted derived events
  • Security/Privacy: derived events carry the same classification and consent handling as the raw events they're computed from

🎯 Enterprise Example

A customer adds an item to cart and closes the tab. 62 minutes later, having seen no order_paid event, the Stream Worker emits cart_abandoned — which the Activation API picks up to trigger a recovery email, entirely without the marketing team having to build their own abandonment-detection logic.

← Back to Stream Processing