System designCase study 7 of 15
System design · Case study 7 of 15
Design a Real-Time Analytics Pipeline
Design a pipeline that turns application events into business metrics (orders per minute, revenue, conversion) visible on a dashboard within one minute of the events happening.
On this page
Functional requirements
- Ingest order and page-view events from many application servers
- Compute per-minute metrics by region and product category
- Serve the last 24 hours of metrics to a dashboard
- Make raw events available for later batch analysis
Non-functional requirements
- End-to-end latency under 60 seconds for 99% of minutes
- Correct counts despite duplicates and events arriving up to 10 minutes late
- No data loss on component restarts
- Dashboard queries return in under a second
Scale assumptions
- Average 20,000 events per second, peaks of 100,000
- About 1 KB per event
- Dashboard used by around 100 people concurrently
Technologies
Kafka, Spark Structured Streaming or Flink, Delta Lake (raw events), A low-latency OLAP store or key-value store for serving
Producers publish events to Kafka; a stream processor aggregates by event time with watermarks and writes results to a low-latency store for the dashboard, while raw events also land in a lakehouse table.
Approach
Clarify what “real time” means (here, one minute), then design around event time, late data and duplicates, which are where real-time systems usually go wrong.
Architecture
- Producers publish events with an event id and event timestamp to Kafka, keyed by user or order id.
- Stream processor parses, deduplicates by event id within the watermark, and aggregates into one-minute event-time windows.
- Serving store receives upserted window results keyed by window and dimensions.
- Dashboard reads the last 24 hours from the serving store.
- Raw sink: a second query appends raw events to a lakehouse table for batch use.
Event time and late data
Aggregate by the timestamp in the event, not arrival time. A 10-minute watermark tells the engine how long to keep each window open for late events; results for a window are upserted again as late events arrive, then the window’s state is dropped.
Duplicates
Producers can retry and consumers can reprocess after failure. Deduplicate on event id within the watermark, and make sinks idempotent (upsert by window key).
Storage
Kafka retains events for several days for replay. Raw events land in a partitioned lakehouse table. Aggregates live in a serving store sized for 24 hours of minute-level rows.
Reliability
Checkpoint stream state and offsets so restarts resume exactly where they stopped. Monitor consumer lag; if lag grows, the dashboard silently falls behind.
Reprocessing
To fix a logic bug, rebuild historical aggregates in batch from the raw lakehouse table and overwrite the affected windows.
Observability
Track input rate, processing rate, consumer lag, watermark delay, end-to-end latency and the number of late events dropped.
Security
Events may contain personal data: restrict access to raw topics and tables, and aggregate before exposing data to the dashboard.
Cost
Streaming compute runs continuously. Right-size the cluster for typical load with autoscaling headroom for peaks, and compact the raw table’s small files.