Transformation & Processing → Stream Processing

Event Deduplication

Removes duplicate events caused by retries, multi-path delivery, or at-least-once semantics — so one customer action is never counted twice.

High-Level Design

Deduplication runs first in the Stream Worker, before any other processing.

Data Source
Queue & Retry
Hands off at-least-once delivered events
→
Ingestion
Azure Event Hubs
Source stream the Stream Worker consumes
→
Processing
Event Deduplication
.NET Core Stream Worker — idempotency check via event_id
→
Foundation
Sessionization
Runs next on de-duplicated events
→
Intelligence
Azure Cache for Redis
Dedup window storage
→
Activation
Every Downstream Consumer
Guaranteed effectively-once processing

💼 Business Context

  • Duplicate events silently inflate metrics (double-counted revenue, inflated engagement) — this is the control that keeps numbers trustworthy
  • Lets every upstream SDK/connector retry aggressively for reliability without worrying about double-processing downstream
  • Owned by Platform Engineering

🔌 Technical Overview

Every event carries a unique event_id assigned at the edge. The Stream Worker checks each incoming event_id against a short-lived deduplication window maintained in Azure Cache for Redis (typically 24-48 hours, covering realistic retry windows) before processing it further. A duplicate is logged and dropped, not reprocessed — this is what makes the Queue & Retry layer's at-least-once delivery guarantee safe to build on.

Dedup Window

24-48 hour Redis TTL event_id as dedup key At-least-once → effectively-once

💾 Dedup Check

var seen = await _redis.StringGetAsync($"dedup:{evt.EventId}");
if (seen.HasValue) { _metrics.Increment("dedup.dropped"); return; }
await _redis.StringSetAsync($"dedup:{evt.EventId}", 1, TimeSpan.FromHours(48));

🔗 Integration Points

  • Azure Cache for Redis — dedup window storage
  • .NET Core Stream Worker — where the check runs, first step after consuming from Event Hubs
  • Queue & Retry layer — the reason this control exists (at-least-once delivery upstream)
  • Azure Monitor — dedup-rate metric, a useful early signal of upstream retry storms

🧰 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: Redis lookups are sub-millisecond even at high event throughput
  • Latency: adds negligible per-event latency
  • Reliability: dedup window is sized generously beyond the longest realistic retry delay
  • Security/Privacy: dedup keys are event IDs only, no payload content is cached

🎯 Enterprise Example

A mobile client loses connectivity mid-send and retries the same add_to_cart event three times once back online. Deduplication ensures the customer's cart analytics show one addition, not three — keeping conversion-funnel math accurate.

← Back to Stream Processing