Start by asking two things: how fresh must the dashboards be (seconds or a minute), and what do they show (counts, funnels, per-user detail). Then draw the path from the browser to the dashboard, and say where each risk sits. Here is a design that fits "millions of users, near real-time".
Web/app SDK -> collector API -> Kafka (key = session_id)
|-> stream processor -> OLAP store -> dashboards
|-> raw sink -> object storage (lake)Collection
An SDK batches events on the device and sends them to a thin, stateless collector behind a load balancer. The collector validates the shape, stamps a server receive time, and writes to Kafka (or Pub/Sub). It does no heavy work, so it can scale out and never blocks on downstream systems. Each event carries a client-generated event_id, an event_ts and a schema version, and the schemas live in a registry.
Streaming layer
Key the topic by session_id or user_id, so the events of one user stay in order within a partition. A stream processor (Flink, Spark Structured Streaming or Dataflow) does three jobs. It drops duplicates using event_id inside a window. It filters bots, for example by user agent or by rate per client. And it computes aggregates in event-time windows with a watermark, so a phone that was offline for 10 minutes still lands in the right minute. Choose how late is too late (say 15 minutes), and send anything later to a correction path.
Serving
Dashboards need fast, filterable aggregates, so write them to an OLAP store built for that (Druid, Pinot, ClickHouse, or BigQuery with streaming ingestion). A warehouse can serve it, but concurrency and latency get expensive. Write the raw events in parallel to object storage as Parquet, partitioned by date. That is the history for replays, backfills and ad-hoc analysis.
The hard parts to mention
- A 10x spike (a sale, a push notification): partitions set above current need, autoscaling consumers, and a buffer in Kafka so the processor can catch up.
- Lag: monitor consumer lag and the gap between event time and processing time, and alert on both.
- Corrections: because raw data lands in the lake, you can recompute a day's aggregates and overwrite the dashboard table when the logic changes.
- Privacy: hash or drop personal fields early.
Trade-off to say out loud
True exactly-once across all of this is hard. Aim for at-least-once delivery with dedup by event_id, and say that small count differences between the real-time and the batch numbers are expected, with batch being the source of truth.