Skip to content
Reliable Data Engineering
Practice problem medium streamingkafkalakehousesessionizationolap
Practise with timer, notes and rubric

Design a Real-Time Clickstream Analytics Platform

Try it first. Set a 45-minute timer, draw on paper, talk out loud. Then compare with the reference answer below.

Problem

An e-commerce company wants to track every user interaction on web and mobile (page views, searches, product clicks, add-to-cart, checkout) to:

  1. Power live dashboards: active users per minute, trending products (last 15 min), conversion funnel.
  2. Give analysts SQL access to full history.
  3. Provide sessionized data for ML (recommendations) and product analytics.

Clarifying questions (and the answers we’ll assume)

QuestionAssumed answer
Event volume?500 M events/day, 3× peak during sales
Dashboard freshness?< 1 minute
Exactness?Dashboards may be approximate (±1%); analyst tables exact & deduplicated
Late events?Mobile can buffer offline: up to 48 h late
Retention?Raw 2 years queryable, aggregated forever
PII?user_id, IP, device ids, all subject to GDPR deletion
Consumers?Internal dashboards (~200 users), analysts (SQL), ML pipelines

1. Requirements

Functional: collect events from web/iOS/Android; real-time metrics; sessionized history; funnel analytics; self-serve SQL. Non-functional: < 1 min dashboard freshness; no data loss; dedup for analytics; GDPR deletion within 30 days; cost-conscious.

2. Estimates

500M/day ÷ 10⁵ ≈ 5k/s avg → 15k/s peak
1 KB JSON → 15 MB/s peak (≈ 4 MB/s Avro)
Kafka: 24–48 partitions (parallelism + growth), RF=3, 7-day retention ≈ 9 TB disk
Lake: 500 GB/day raw → ~70 GB/day Parquet → 2 yrs ≈ 50 TB (~$1.2k/month storage)

Conclusion: medium scale. Kafka + Spark Structured Streaming + Delta is enough. Don’t over-engineer.

3. High-level architecture

flowchart LR
    subgraph Clients
        WEB[Web SDK]
        MOB[Mobile SDK<br/>offline buffer]
    end
    WEB --> COL[Collector service<br/>validate, enrich, batch]
    MOB --> COL
    COL --> K[[Kafka: events.raw<br/>key = user_id]]
    COL -.->|invalid| DLQ[[events.dlq]]
    K --> RT[Streaming job:<br/>real-time aggregates]
    K --> ING[Streaming job:<br/>bronze ingest]
    RT --> OLAP[(Pinot / Druid<br/>or Redis)]
    OLAP --> DASH[Live dashboards]
    ING --> BR[(Bronze Delta<br/>raw events)]
    BR --> SIL[Silver: dedup,<br/>sessionize, conform]
    SIL --> SV[(silver.events<br/>silver.sessions)]
    SV --> GO[(Gold: funnels,<br/>daily metrics, features)]
    GO --> BI[Analysts / BI]
    GO --> ML[ML training]

Walk one event: user taps “add to cart” → SDK assigns event_id (UUID) + client_ts → collector adds server_ts, geo from IP (then drops raw IP), validates schema → Kafka partition by user_id → (a) real-time job updates per-minute counters in the OLAP store, (b) ingest job appends to bronze → silver job dedups by event_id, assigns session_id → gold funnel tables.

4. Data model

Event (Avro, schema registry):

{ "event_id": "uuid", "event_type": "add_to_cart", "schema_version": 3,
  "user_id": "u_123", "anonymous_id": "a_987", "session_hint": "s_55",
  "client_ts": "2026-10-01T10:15:02.123Z", "server_ts": "2026-10-01T10:15:03.010Z",
  "platform": "ios", "app_version": "8.2.1", "page": "/p/shoes-42",
  "product_id": "p_42", "properties": {"price": "59.90", "currency": "EUR"},
  "geo": {"country": "DE", "city": "Berlin"} }
TableGrainKey columnsLayout
bronze.events1 row per received event (may have dupes)+ _kafka_partition, _kafka_offset, _ingested_atpartition by ingest_date
silver.events1 row per unique event_idevent_date, user_id, session_idcluster by (event_date, user_id)
silver.sessions1 row per sessionsession_id, user_id, start, end, n_events, convertedcluster by (session_date)
gold.funnel_dailyday × platform × country × stepusers_at_stepsmall
rt.metrics_minute (OLAP)minute × dimensionactive_users (HLL), eventsTTL 7 days

5. Deep dives

5.1 Real-time aggregates

events = (spark.readStream.format("kafka").option("subscribe", "events.raw")
          .option("maxOffsetsPerTrigger", 2_000_000).load()
          .select(from_avro("value", schema).alias("e")).select("e.*")
          .withColumn("event_time", col("client_ts").cast("timestamp")))

active = (events.withWatermark("event_time", "5 minutes")
          .groupBy(window("event_time", "1 minute"), "platform", "country")
          .agg(approx_count_distinct("user_id").alias("active_users"),
               count("*").alias("events")))

trending = (events.filter("event_type in ('product_click','add_to_cart')")
            .withWatermark("event_time", "5 minutes")
            .groupBy(window("event_time", "15 minutes", "1 minute"), "product_id")
            .agg(count("*").alias("interactions")))

5.2 Deduplication and sessionization (silver)

Duplicates come from SDK retries and at-least-once delivery. Dedup on event_id:

(spark.readStream.table("bronze.events")
   .withWatermark("event_time", "48 hours")              # max lateness we accept in streaming
   .dropDuplicatesWithinWatermark(["event_id"])
   .writeStream.foreachBatch(merge_events)               # MERGE ... WHEN NOT MATCHED INSERT (catches older dupes)
   .option("checkpointLocation", "/chk/silver_events").start())

Sessionization (30-min inactivity gap) is a classic gaps-and-islands problem. In batch SQL:

WITH ordered AS (
  SELECT *, LAG(event_time) OVER (PARTITION BY user_id ORDER BY event_time) AS prev_time
  FROM silver.events WHERE event_date BETWEEN :d - 1 AND :d
), flagged AS (
  SELECT *, CASE WHEN prev_time IS NULL
                  OR event_time > prev_time + INTERVAL 30 MINUTES THEN 1 ELSE 0 END AS new_session
  FROM ordered
)
SELECT *, user_id || '-' || SUM(new_session) OVER (PARTITION BY user_id ORDER BY event_time) AS session_id
FROM flagged;

In streaming: Spark session_window(event_time, "30 minutes") or transformWithState for custom logic. Decision: batch (hourly) sessionization is enough for analysts/ML; real-time session metrics aren’t a requirement. Say this explicitly.

5.3 Late data and restatement

flowchart LR
    LATE[Event 30h late] --> BR[(bronze: ingest_date = today)]
    BR --> S[silver MERGE by event_id<br/>event_date = 2 days ago]
    S --> R[Daily job restates gold<br/>for the last 3 event_dates]
    R --> G[(gold funnel tables)]

Partition bronze by ingestion date (writes never touch old partitions), model silver/gold by event date, and re-compute the trailing 3 days of gold daily. Dashboards mark the last 48 h as preliminary.

5.4 Collector & schema governance

6. Trade-offs

DecisionChosenAlternativeWhy
BufferKafkaKinesis / Pub/SubReplay, fan-out, ecosystem; managed alternatives fine if cloud-native
Partition keyuser_idround-robinPer-user ordering for sessionization; risk: bots/hot users → monitor, filter bots
Stream engineSpark SSFlinkSeconds latency suffices, unified with batch/Delta; Flink if sub-second
Real-time storePinot/DruidRedis, query Delta directlySub-second slice-and-dice + top-N; Delta SQL would be seconds and costly at high concurrency
Distinct usersHLLexact setsState size; exact numbers come from batch
Sessionizationhourly batchstreaming session windowsNo real-time requirement; simpler, cheaper

7. Failure modes

FailureDetectionHandling
Collector overload on flash salep99 latency, 5xxAutoscale, SDK retries with backoff, client-side buffering
Streaming job crashLag alertRestart from checkpoint; idempotent sinks
Bad SDK release sends malformed eventsDLQ rate spike, volume per app_versionReject at collector, alert app team, fix + replay DLQ
Bot traffic inflates metricsAnomaly on events/userBot filtering rules in silver; flag in real-time path
DuplicatesUniqueness check on silverDedup by event_id (stream + MERGE)

8. Scaling 10× (5 B events/day)

Kafka partitions 128+, tiered storage; pre-aggregate at the collector or in a first streaming stage (per-minute partials) before the OLAP store; Pinot real-time tables with star-tree index; compaction/liquid clustering on silver; consider Flink for the real-time tier if state grows; cost: move bronze older than 90 days to cold storage.

9. What separates a senior answer

10. Follow-up questions

How would you compute "users who viewed a product then purchased within 1 hour" in real time?

Keyed stateful processing by user_id: store recent product views in state with a 1h TTL/timer; on purchase, check for a matching view and emit. In Spark: transformWithState/flatMapGroupsWithState; in Flink: keyed process function with timers, or CEP pattern. In batch: self-join on user with time bounds.

Kafka partition 7 has 10× the lag of others. Why?

Hot key: a bot or a load-test account keyed to that partition, or a misbehaving client sending a flood. Check top keys on that partition. Mitigate: filter/rate-limit bots at the collector, salt the key for the aggregation path, or re-key by event_id for paths that don’t need per-user order.

How do you support GDPR deletion here?

Deletion request table → daily batch DELETE from silver/gold by user_id (and anonymous_id links), deletion vectors + purge + VACUUM within the retention SLA; bronze either crypto-shredded (per-user key for PII fields) or limited retention; OLAP store has short TTL (7 days) so it ages out; ML feature tables rebuilt.

Why not have clients write straight to Kafka?

Security (no broker credentials on devices), validation, enrichment, rate limiting, batching/compression, and protocol (HTTP from browsers/mobile). The collector is a thin, horizontally scalable, stateless layer.


Self-assessment rubric