The interview setup
You are designing streaming analytics for 1 million IoT devices. Each device sends about 200 bytes every 5 seconds. You must compute per-device rolling averages, detect anomalies, and store time-series data for historical charts.
This interview is a trap for candidates who only flex Kafka knowledge. Bandwidth is modest. The real bosses are series cardinality, keyed state size, downsample policy, and reconnect storms after outages.
What the interviewer wants:
- You notice cardinality before you brag about throughput
- You design ingest that survives thundering herds
- You pick storage that fits time-series access patterns
- You prevent alert fatigue
- You can explain how rolling windows live in stream state
Background concepts from first principles
What is an IoT telemetry pipeline?
Devices emit readings (temperature, vibration, GPS, voltage). A cloud pipeline ingests them, computes short-horizon analytics, alerts on bad behavior, and stores history for dashboards and forensics.
Unlike a clickstream ad-tech firehose, many IoT fleets are "chatty keys, medium bytes." One million devices means one million (or more) active series.
Why cardinality hurts
Every distinct series key (device_id, or device_id × metric_name) needs indexing and often memory for recent state. A warehouse table that is fine for 10k products can become painful with 1M hot devices if every query filters one device across billions of rows without the right layout.
Time-series databases and careful lake clustering exist because of this access pattern: append by time, fetch by series and range.
Rolling averages need state
A 15-minute rolling average is not a stateless map. The job must remember recent points or a compact summary per device. In Flink or similar engines, that is keyed state. RocksDB state backends are common. Unbounded state for dead devices is a slow-moving outage.
Downsampling is a product decision
Ops debugging yesterday wants fine resolution. A monthly capacity chart wants coarse aggregates. Keeping 5-second raw forever at fleet scale becomes a storage and compaction tax. Healthy designs tier: raw for days, minutes for months, hours for years.
Burst reconnect
After a network partition, devices come back and flush buffers. Your gateway and Kafka must absorb a herd. Without admission control, you create a second outage during recovery.
The expanded problem
Build a system that:
- Ingests about 200K events/sec (1M devices × one event every 5 seconds)
- Computes per-device rolling averages over windows such as 5 to 15 minutes
- Detects anomalies or threshold breaches and alerts
- Stores time-series for charts
- Handles offline devices and reconnect bursts
Devices --> MQTT/HTTP gateway --> Kafka --> stream jobs
| |
v v
alerts TSDB / lake historyConstraints and scale prompts
Do the math out loud:
- 200K events/sec × 200 bytes ≈ 40 MB/s raw ingest
- Per day raw ≈ 40 MB/s × 86,400 ≈ 3.5 TB/day before compression and downsample
- Cardinality ≈ 1M series baseline; multiply by metrics per device
- Alert fanout must be bounded; noisy sensors can page humans to death
What good looks like
- Lead with cardinality and reconnects, not only MB/s
- Partition or key by device_id carefully, with hotspot controls
- State TTL for inactive devices
- Explicit retention tiers
- Alert hysteresis and cooldown
- Optional edge aggregation if devices can do it
Clarifying questions
- How many metrics per device?
- Alert latency SLA: seconds or minutes?
- Retention at raw resolution versus downsampled?
- Can devices aggregate on the edge?
- How long may devices buffer offline?
- Single region or global fleet with regional gateways?
Out of scope
- Designing sensor firmware in hardware detail
- Publishing a research-grade anomaly model paper
- Building the entire device twin UI
- Multi-cloud networking deep dive unless invited
A concrete story to keep in your head
A cold-chain pharmacy fridge fleet loses connectivity for twenty minutes in one city. When the radios return, fifty thousand devices dump buffered temperature samples. Without gateway admission control, Kafka request latency spikes, stream jobs fall behind, and alert rules think the entire city is on fire because catch-up data looks like sudden change.
Your design must make recovery boring: absorb, meter, compute, and alert with cooldown. That story is what this problem is really about.