Menu
System designCase study 22 of 40

System design · Case study 22 of 40

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.

  • Advanced
  • 2 min read
  • Updated Oct 2026
On this page
  1. Approach
  2. Architecture
  3. Event time and late data
  4. Duplicates
  5. Storage
  6. Reliability
  7. Reprocessing
  8. Observability
  9. Security
  10. Cost

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

  1. Producers publish events with an event id and event timestamp to Kafka, keyed by user or order id.
  2. Stream processor parses, deduplicates by event id within the watermark, and aggregates into one-minute event-time windows.
  3. Serving store receives upserted window results keyed by window and dimensions.
  4. Dashboard reads the last 24 hours from the serving store.
  5. Raw sink: a second query appends raw events to a lakehouse table for batch use.
Two outputs from one stream: fresh aggregates for the dashboard and raw events for everything else.

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.

By Data Career Hub Editorial · Last reviewed Oct 2026

Search
Filter by type