Interview framing
You are in a mid-level data engineering system design interview. Budget about 45 to 55 minutes. The interviewer will not grade you on naming a specific vendor. They want to hear how you think under load: clarify, estimate, draw a path that can fail safely, and defend trade-offs with numbers.
The prompt you will hear is short on purpose:
Design a system that ingests 2 million events per second from mobile apps, processes them in real time, and serves dashboards with a 10-second refresh SLA. Events are about 1 KB each.
That is enough to fail if you jump straight to boxes. Good candidates slow down, do byte math out loud, and only then sketch architecture.
Background from first principles
Start with what an event is. A mobile app fires a small JSON payload when something happens: screen view, button tap, purchase, crash. Each payload is roughly 1 KB. At two million such payloads every second, you are not talking about a web form. You are talking about a firehose.
Why does that matter? Write path and read path are different animals. Dashboards refresh every few seconds and serve a modest number of analysts. The write path must absorb millions of inserts per second without melting. If you treat both paths as one database, you lose.
Postgres and most OLTP databases shine at transactions, joins, and strong consistency for application state. They are not built to ingest 2 GB of new rows every second forever. Indexes, WAL, vacuum, and disk IOPS will collapse long before you hit the dashboard SLA. Saying that out loud is interview gold: we cannot write raw events into Postgres at this rate.
You need three ideas before the whiteboard:
1. A durable buffer that absorbs spikes (a log such as Kafka or a managed equivalent). 2. Stream processing that turns raw events into small aggregates on the fly. 3. A serving store optimized for fast dashboard reads, plus a cheap archive for raw history.
Also understand event time versus processing time. Event time is when the phone says the action happened. Processing time is when your job sees the record. Mobile networks make those diverge. Any serious real-time design must name a late-event policy.
Expanded problem statement and constraints
Build a first production version of a real-time analytics pipeline:
- Peak ingest: 2,000,000 events/sec from mobile clients.
- Payload size: about 1 KB per event.
- Dashboard freshness: charts should not show data older than about 10 seconds at the serving layer. State your interpretation of end-to-end latency versus refresh interval.
- Metrics: rolling counts, rates, approximate uniques, simple funnels suitable for ops and product.
- Replay: keep a recent window of raw events so you can fix bugs and backfill.
- Poison messages: bad payloads must not stall the hot path; route them to a dead-letter path.
- Region: single-region v1 is acceptable if you name the DR gap.
- Cost: you must defend storing and moving roughly 2 GB/sec of raw data.
What good looks like in the room
You restate the goal in one sentence. You write the byte math on the board before drawing Kafka. You ask which metrics truly need 10-second freshness versus nightly batch. You separate hot aggregates from cold raw archive. You name one late-event strategy and one lag alert. You leave time for failure modes. You do not promise exactly-once everywhere without an idempotent sink story.
The interviewer should leave believing you felt the scale, protected the product under spikes, and knew where truth lives when the hot path is approximate.
Clarifying questions to ask
- Which metrics must be real-time vs can wait for a nightly batch?
- Is uniqueness exact or approximate (HLL)? Does cross-device identity matter?
- Retention for raw events vs aggregates?
- Is the 10-second SLA end-to-end (phone to UI) or chart refreshes every 10s with whatever is already stored?
- Multi-tenant customer dashboards or one internal BI surface?
- Peak multiplier over average (2x to 3x is common)?
- PII in the event payload (user ids, device tokens)?
- Which dimensions are allowed on hot tiles (event_type, country) versus only in the lake (user_id)?
Scale and estimation prompts
Do this math out loud:
- 2e6 events/sec x 1 KB = 2 GB/sec raw ingest.
- Per day: 2 GB/sec x 86,400 sec is about 173 TB/day uncompressed.
- With 3x to 5x columnar compression on cold storage, plan roughly 35 to 60 TB/day retained.
- Kafka with 3x replication moves about 6 GB/sec of network traffic to brokers at peak.
- Dashboard QPS is low (tens to hundreds). Optimize the write path into the serving store, not a giant scan of raw events every refresh.
Ask the interviewer whether peak is 2x or 3x average, and size the buffer and stream jobs for peak.
A useful verbal check: if your design still works if event size were 200 bytes or 5 KB, say which hop changes. That shows you tied architecture to numbers.
Out of scope for v1
Say these out loud so the interviewer knows you can bound the problem:
- Multi-region active-active serving.
- Exact distinct users over 30 days at full fidelity for every tile.
- Real-time personalization or online feature serving for ML (different problem).
- Guaranteed end-to-end exactly-once with zero cost (prefer at-least-once plus idempotent aggregates).
- Storing every raw byte forever without tiering.
- Full-text search over every raw event on the hot path.