Interview Setup
Interview Prompt
Design a user analytics pipeline like Google Analytics. Track pageviews, clicks, and custom events from websites/apps, process them for real time dashboards and batch reporting.
Clarifying Questions (ask before designing)
| Question | Why it matters |
|---|---|
| Event volume: what are the expected daily events and peak ingestion rates? | 100B events daily (~1.2M/sec) drives Kafka partition planning, edge collector horizontal scaling, and peak capacity headroom. |
| Real time vs batch latency requirements? | Active user counts in 1-minute windows require Flink streaming, whereas daily executive reports tolerate hours of batch latency, leading to a Lambda architecture. |
| Query patterns: do clients require real time counters or ad hoc SQL exploratory queries? | Use ClickHouse for OLAP aggregations, Redis for low latency live counters, and an S3 data lake for long term raw retention. |
| Privacy requirements: are PII fields present, and what are the retention and GDPR erasure rules? | Directly influences schema normalization, hashing rules at ingestion, and partition key strategy for customer erasure requests. |
Scope
In scope
- Lambda vs Kappa architecture
- Event schema evolution
- Data warehouse star schema
- ETL vs ELT
- Data lake partitioning
- Capacity estimation with shown math
Out of scope (state explicitly)
- Detailed frontend/UI pixel implementation
- Org structure, staffing, and hiring plan
Functional Requirements
Clarify the event schema, data freshness targets, and reporting capabilities with your interviewer. The design requires client side event capture, high throughput ingestion, sessionization, and analytical dashboards for funnels and retention, confirming whether reporting requires real time streaming, batch processing, or both.
In the room: ask whether they need per user funnels or only aggregate page views, because that distinction fundamentally changes storage engine selection and partition key strategy.
- Event tracking: Capture client side page views, link clicks, scroll events, and custom events across web and mobile platforms.
- Real time dashboard: Display active current users, events per second, and top pages for the latest completed window with sub minute freshness.
- Historical reports: Run analytical aggregations, such as pageviews over 30 days grouped by country or device category.
- Funnel conversion: Trace multi step user conversion journeys, such as progressing from Landing Page to Registration and completing Check Out.
- Cohort retention: Group users by registration cohorts to measure user retention curves over days, weeks, and months.
- Custom dimensions: Allow developers to attach arbitrary key-value metadata to custom events.
- Session management: Stitch independent event streams into discrete user sessions using a 30-minute inactivity threshold.
Non-Functional Requirements
The ingestion tier must withstand sharp client traffic spikes without dropping accepted events while balancing near real time dashboard updates with petabyte-scale durability and cost efficiency.
- High Scale Ingestion: Sustain approximately 1.2 million incoming events per second as the average design rate, with additional peak headroom for bursty client traffic across millions of active client domains.
- Low Ingestion Latency: Enriched events must reach the durable event log and data lake within 5 seconds of receipt. Real time dashboard freshness is governed separately by the streaming SLO, while hourly batch reporting follows its own latency target.
- Low-Latency Queries: Set a product objective below 2 seconds at p99 for supported dashboard query patterns, with an operational SLO below 3 seconds at p99.
- Petabyte Scale: Scale to process and index over 100 billion events daily, totaling approximately 4.5 petabytes of raw data over 90 days.
- Tiered Retention: Retain high-fidelity raw event data for 90 days in hot storage and pre-aggregated rollups for 2 years in cold archives.
- At-Least-Once Delivery: Durably append accepted events before acknowledging HTTP clients. Duplicates are allowed by the delivery model and are removed from business aggregates using tenant-scoped (site_id, event_id) deduplication.
Capacity Estimations
With 100 billion events generated daily, sizing Kafka ingestion pipelines, network bandwidth, and columnar storage tiers is essential before designing the query layer:
| Metric | Calculation | Value |
|---|---|---|
| Events / day | Given (typical workload assumption) | 100 Billion |
| Events / sec | 100B / 86400 = ~1.157M/sec average. 1.2M/sec is the rounded design rate | ~1.2 Million |
| Avg event size | Given (typical workload assumption) | 500 bytes |
| Ingestion throughput | Derived | 600 MB/s |
| Storage / day (raw) | 100B x 500 bytes = 50 TB rounded. 1.2M/sec x 500 bytes = 51.84 TB/day | 50 TB planning figure |
| Storage / 90 days | 50 TB/day x 90 = 4.5 PB rounded. 51.84 TB/day x 90 = ~4.66 PB | 4.5 PB planning figure |
| Unique users tracked | Derived | 1 Billion |
I/O and Bandwidth Derivations: - Global Ingestion rate: 100,000,000,000 events / 86,400 sec = ~1,157,400 events/sec average. - Network Bandwidth: 1.2M events/sec x 500 bytes = 600 MB/s network ingestion rate. - Storage Calculations (Raw Events): - 600 MB/s x 86,400 sec = 51.84 TB raw logs per day. - 90-Day Raw Storage Footprint: 51.84 TB x 90 = ~4.66 Petabytes.
Architecture Diagram
Client SDKs batch telemetry events locally and transmit them to stateless edge collectors. Kafka buffers the stream, an S3 lake writer persists the durable raw and enriched event archive, Flink handles low latency windowed metrics in Redis, Spark manages batch sessionization from S3, and ClickHouse powers interactive OLAP queries.
Component Deep Dives
1. Edge Collector and Enrichment Architecture
Tracing data movement from client side beacon capture to columnar warehouse queries reveals critical optimizations in session stitching, memory bounds, and storage partitioning.
The client SDK batches events locally before forwarding them over HTTP to minimize mobile radio wakeups and network overhead. Where consent or another legal basis is required, edge collectors apply that decision before enrichment or durable storage. Stateless collectors then parse, validate, and enrich the incoming JSON payload before publishing to Kafka.
{
"site_id": "site-123",
"client_id": "anon-uuid",
"session_id": "sess-uuid",
"events": [
{
"event_id": "event-192a-001",
"type": "pageview",
"url": "/pricing",
"timestamp_ms": 1710320000000,
"referrer": "google.com"
},
{
"event_id": "event-192a-002",
"type": "click",
"element": "#signup-btn",
"timestamp_ms": 1710320005000
},
{
"event_id": "event-192a-003",
"type": "custom",
"name": "video_play",
"props": {
"video_id": "v123"
},
"timestamp_ms": 1710320010000
}
]
}- GeoIP Lookup: Resolve client IP addresses to geographic country and city using an in memory MaxMind database without blocking network I/O, then discard or minimize raw IP retention according to the privacy policy.
- User-Agent Resolution: Parse client browser families, operating system versions, and device categories into structured LowCardinality fields.
- Bot and Abuse Filtering: Match client fingerprints against known crawler signatures, enforce per-site quotas and per-IP rate limits, and reject suspicious traffic before event publication.
2. Real Time Processing (Apache Flink)
Apache Flink consumes the enriched event stream from Kafka, deduplicates the tenant-scoped (site_id, event_id) pair within a bounded retention horizon, and executes tumbling and sliding windows to maintain live counters. Processing-time windows are used for live freshness metrics, while event time windows with watermarks and allowed lateness are reserved for metrics where historical timestamp accuracy matters:
// Apache Flink Stream Processing Aggregations
// Consumes enriched-events topic and writes low latency metrics to Redis
// 1. Active Users (1-minute tumbling processing time window for live freshness)
// Use an HLL-style sketch to keep state bounded for high cardinality tenants.
events
.keyBy(event -> event.site_id)
.window(TumblingProcessingTimeWindows.of(Time.minutes(1)))
.aggregate(new HyperLogLogCounter())
.addSink(new RedisCommandSink("SETEX active_users:{site_id} 120 {count}"));
// 2. Event Rate (10-second sliding processing time window with 1-second slide)
events
.keyBy(event -> event.site_id)
.window(SlidingProcessingTimeWindows.of(Time.seconds(10), Time.seconds(1)))
.aggregate(new EventRateCounter())
.addSink(new RedisTimeSeriesSink("TS.ADD events_per_sec:{site_id} {ts_ms} {count}"));
// 3. Top Pages in the Latest Completed Minute (event time)
events
.assignTimestampsAndWatermarks(
WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofMinutes(5)))
.filter(event -> "pageview".equals(event.type))
.keyBy(event -> event.site_id + ":" + event.url)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.allowedLateness(Time.minutes(5))
.aggregate(new PageviewCounter())
.addSink(new LatestWindowTopPagesSink("top_pages:{site_id}", "latest"));
// The sink stages the completed window, trims it to the top 100 members, and
// atomically replaces top_pages:{site_id} so the API exposes the latest window only.3. Batch Processing and Session Stitching (Apache Spark)
Apache Spark pulls raw event partitions from S3 on scheduled hourly cadences, executing complex multi pass calculations that require broader temporal context:
1. Read raw event Parquet partitions already landed in S3, partitioned by time and site bucket. 2. Reconstruct user sessions by grouping events by site_id and client_id, ordering events by timestamp, splitting sessions when inactivity exceeds 30 minutes, and calculating total session duration and page depth. 3. Compute funnel analysis aggregations by calculating retention steps across configured conversion journeys. 4. Materialize daily pre-aggregations by generating daily pageview rollups grouped by country and device type. 5. Execute archival compaction jobs that merge small hourly files into large 128 to 256 MB Parquet blocks in cold S3 storage.
4. ClickHouse OLAP Columnar Store
Row-oriented relational databases struggle when executing analytical aggregations across billions of records. ClickHouse stores data in columnar format with extreme data compression using the MergeTree engine:
CREATE TABLE events (
event_date Date,
site_id String,
session_id String,
client_id String,
event_id String,
event_type LowCardinality(String),
url String,
referrer String,
country LowCardinality(String),
device_type LowCardinality(String),
browser LowCardinality(String),
custom_props Map(String, String), -- primitive custom values normalized to strings
timestamp DateTime64(3)
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(event_date)
ORDER BY (site_id, event_date, session_id, timestamp, event_id);Event Bus Design (Kafka)
A distributed event log decouples stateless ingestion collectors from downstream stream processors and storage sinks. It buffers traffic surges while preserving order for records that share the configured Kafka key. Selective hot tenant salting trades global tenant ordering for additional parallelism when necessary. A second-stage aggregation restores the logical site_id metric.
# Central Kafka Cluster Topic Topology and Consumer Groups
topics:
raw_events:
purpose: "Collector output following schema validation"
partition_key: "site_id, or site_id + stable hash(client_id) for hot tenants"
hot_tenant_strategy: "apply salting only to hot tenants, followed by a second-stage aggregation to restore site_id"
partitions: 128
retention_hours: 48
replication_factor: 3
min_insync_replicas: 2
enriched_events:
purpose: "Payloads enriched with GeoIP, user-agent parsing, bot filtering, and PII scrubbing"
partition_key: "site_id, or site_id + stable hash(client_id) for hot tenants"
hot_tenant_strategy: "apply salting only to hot tenants, followed by a second-stage aggregation to restore site_id"
ordering_note: "Salting a hot tenant trades global per-site Kafka ordering for parallelism. Downstream aggregation restores the logical site_id result"
partitions: 128
retention_days: 7
replication_factor: 3
min_insync_replicas: 2
session_events:
purpose: "Optional Flink output for low latency sessionized events when streaming session views are required"
partition_key: "client_id"
identity_note: "Use resolved user_id only when an authenticated identity is available, and otherwise partition by client_id. Preserve site_id in the record for tenant isolation."
partitions: 128
retention_days: 7
replication_factor: 3
min_insync_replicas: 2
deletion_requests:
purpose: "Erasure requests keyed by tenant and client identifier"
partition_key: "site_id + client_id"
partitions: 32
retention_days: 30
replication_factor: 3
min_insync_replicas: 2
producers:
edge_collector:
semantics: "Idempotent Kafka producer retries plus business-level deduplication by tenant-scoped event_id"
payload_fields:
- "event_id"
- "site_id"
- "user_id"
- "event_type"
- "properties"
- "timestamp_ms"
- "geo"
- "device"
acknowledgment: "all (acks=-1) with enable.idempotence=true and retries enabled"
consumer_groups:
s3_lake_writer:
pipeline: "Persist validated and enriched Kafka events to the S3 Parquet data lake"
sink: "S3 Parquet landing"
flink_realtime:
pipeline: "One-minute active users, events per second, and the latest completed top-page window written to Redis"
sink: "Redis in memory cluster"
alert_pipeline:
pipeline: "Automated anomaly evaluation and threshold breach detection"
sink: "Incident response systems including PagerDuty and Slack"
erasure_pipeline:
pipeline: "Apply deletion tombstones to ClickHouse and schedule S3 Parquet rewrites for affected identifiers"
sink: "ClickHouse tombstones and S3 rewrite workflow"
batch_jobs:
spark_reporting:
input: "S3 Parquet landing"
pipeline: "Hourly session reconstruction, conversion funnels, and cohort retention"
sink: "ClickHouse OLAP columnar storage"
pipeline_paths:
synchronous_ingest: "Client SDK batch to collector validation to durable Kafka append, targeting HTTP 204 within 50 ms under normal load"
asynchronous_processing: "Flink streaming and Spark batch pipelines materializing views without blocking ingest"
dead_letter_queue: "raw-events-dlq receives unparseable messages. Monitor DLQ rate and oldest-message age independently of consumer lag"API Design
Service Contract and Data Types
The analytics API separates high throughput ingestion from low latency querying, ensuring write traffic does not compete with interactive report generation.
Strong domain types define the batch collection request, real time metrics payload, and report query filter:
// Domain models for analytics ingestion and query APIs
// All event timestamps use Unix epoch milliseconds and represent client event time.
export interface AnalyticsEvent {
event_id: string;
user_id?: string;
type: "pageview" | "click" | "custom" | (string & {});
url?: string;
element?: string;
name?: string;
props?: Record<string, string | number | boolean>;
referrer?: string;
timestamp_ms: number; // Unix epoch milliseconds representing client event time
}
export interface CollectBatchRequest {
site_id: string;
client_id: string;
session_id: string;
events: AnalyticsEvent[];
idempotency_key?: string; // Optional request-level deduplication key for batch retries, scoped by site_id
}
export interface RealtimeMetrics {
site_id: string;
active_users: number; // HLL estimate, not an exact distinct count
events_per_sec: number;
top_pages: Array<{
url: string;
views: number;
}>;
}
export interface ReportQueryFilter {
site_id: string;
metric: "pageviews" | "sessions" | "bounce_rate" | "conversions";
from_date: string;
to_date: string;
group_by?: "country" | "device" | "browser" | "url";
}
export interface AnalyticsPipelineService {
collectBatch(batch: CollectBatchRequest): Promise<void>;
getRealtimeMetrics(siteId: string): Promise<RealtimeMetrics>;
queryReport(filter: ReportQueryFilter): Promise<{ data: Record<string, string | number>[] }>;
}1. Collect Events (Ingestion API)
The edge ingestion endpoint accepts batched client events and returns HTTP 204 No Content after durable Kafka append for seamless browser beacon integration.
POST /collect HTTP/1.1
Host: analytics.example.com
Content-Type: application/json
Idempotency-Key: req-20260314-0001
{
"site_id": "site-882",
"client_id": "anon-uuid-2918",
"session_id": "sess-uuid-9921",
"events": [
{
"event_id": "event-192a-001",
"type": "pageview",
"url": "/pricing",
"timestamp_ms": 1710320000000,
"referrer": "google.com"
}
]
}
HTTP/1.1 204 No Content2. Query Real-Time and Historical Reports
Query interfaces power interactive user analytics dashboards across active visitor metrics and historical dimensional rollups.
GET /api/v1/analytics/site-882/realtime HTTP/1.1
Host: analytics.example.com
HTTP/1.1 200 OK
Content-Type: application/json
{
"site_id": "site-882",
"active_users": 1523,
"events_per_sec": 342,
"top_pages": [
{ "url": "/pricing", "views": 842 },
{ "url": "/docs", "views": 321 }
]
}
GET /api/v1/analytics/site-882/report?metric=pageviews&from=2026-03-01&to=2026-03-14&group_by=country HTTP/1.1
Host: analytics.example.com
HTTP/1.1 200 OK
Content-Type: application/json
{
"data": [
{ "country": "US", "pageviews": 1523000 },
{ "country": "DE", "pageviews": 381000 }
]
}Common Error Responses
Standardized error responses cover payload validation failures and queue throttling.
400 Bad Request: invalid input, missing required fields, or malformed JSON payload
401 Unauthorized: missing or invalid authentication token or API key
403 Forbidden: authenticated caller lacks required permissions for this resource
404 Not Found: requested resource ID does not exist
409 Conflict: duplicate write or version conflict, retry with a unique idempotency key
422 Unprocessable Entity: syntactically valid request failed semantic business validation
429 Too Many Requests: rate limit quota exceeded, client should honor Retry-After header
500 Internal Error: unexpected server failure, retry safely with an idempotency key
503 Service Unavailable: downstream dependency is unavailable or overloaded, retry with exponential backoff
202 Accepted: asynchronous job queued successfully, poll GET /jobs/{id} for completion status
408 Request Timeout: background job is still executing, continue polling status endpointData Model
Raw Enriched Event Schema (Kafka and S3 Parquet)
Analytics data structures are split across an immutable event log in S3 Parquet, in memory cache structures in Redis, and a columnar ClickHouse schema partitioned primarily by time.
Enriched event payload containing core tracking identifiers, including the tenant-scoped event identity (site_id, event_id), normalized device and geographic dimensions, and arbitrary custom properties. The edge timestamp_ms value is converted to ClickHouse DateTime64(3).
{
"event_id": "event-192a-3312",
"site_id": "site-123",
"session_id": "sess-abc",
"client_id": "anon-xyz",
"user_id": null,
"event_type": "pageview",
"url": "/pricing",
"referrer": "google.com",
"country": "US",
"device": "mobile",
"browser": "Chrome",
"timestamp_ms": 1710320000000,
"custom": {
"plan": "pro",
"source": "ad_campaign_123"
}
}Redis Real-Time Cache Topology
In-memory data structures supporting sub second dashboard telemetry, separating atomic counters, time-series velocity, and sorted URL leaderboards:
# Active users key: string counter with 120-second expiration
active_users:{site_id}:
data_structure: "string"
value: "HLL estimate for the latest completed one-minute window"
ttl_seconds: 120
# Real time event velocity: rolling time-series window
events_per_sec:{site_id}:
data_structure: "time-series"
retention: "60 data points (1 per second)"
# Top URLs for the latest completed minute
top_pages:{site_id}:
data_structure: "sorted-set"
score: "pageview count in latest completed window"
member: "page URL path"
max_members: 100
ttl_seconds: 120
update: "stage completed window, trim to top 100, then atomically replace top_pages:{site_id}"Fault Tolerance
Failure Matrix
High volume analytics pipelines must preserve event continuity across network drops, stream worker crashes, and skewed customer workloads. The failure matrix covers durability, stream recovery, client reliability, abuse control, and partition skew.
| Ingestion Component | System Failure Solution |
|---|---|
| Accepted event durability | Kafka replication factor of 3 with idempotent producers using acks=all and HTTP 204 only after Kafka append acknowledgment. Downstream consumers deduplicate by the tenant-scoped (site_id, event_id) pair |
| Flink failure | Asynchronous state checkpointing to S3 every 60 seconds with state recovery resuming from the latest checkpoint |
| ClickHouse failure | ReplicatedMergeTree tables deployed across multiple availability zones with automated replica failover |
| Client offline | Client SDK buffers unsent events in browser storage or app-local durable storage and flushes payloads after connectivity returns |
| Bot traffic | Filter known bot signatures at edge collectors, enforce token-bucket IP rate limits, and issue CAPTCHAs on burst patterns |
| Data skew | Partition Kafka topics by site_id for the baseline workload and use adaptive salting with a stable client hash for hot tenants. Partition ClickHouse primarily by time and use site_id in the sort key |
1. Real Time Session Stitching and Watermarks
Out-of-order events from mobile devices and spotty networks make it difficult to determine when a user session concludes. The pipeline coordinates Flink session windows with event time, watermarks, and an explicit allowed-lateness policy.
Session Definition: A sequence of user interactions within a site containing no inactivity gap greater than 30 minutes. Flink Session Windowing: 1. Group incoming events by site_id and client_id using Flink session windows with a 30-minute inactivity gap. 2. Advance event time watermarks with a 5-minute bounded-out-of-orderness tolerance, then keep window state available for the separately configured allowed lateness. Late Event Edge Case: - Condition: Mobile events arriving after the watermark has advanced beyond the session window. - Streaming impact: Events arriving within the configured 5-minute allowed-lateness period can update or merge an existing session. Events arriving after window state has been cleaned up may form a separate session or be routed to a late-event stream. In this scenario, late events beyond the allowed lateness are expected to cause approximately 1% real time session count inflation. - Batch mitigation: Hourly Spark jobs re-stitch raw S3 events using true event timestamps to reconcile accurate historical session counts.
2. Browser Exit Loss and Beacon Delivery
When a user abruptly closes a browser tab, pending asynchronous HTTP requests are often cancelled. The client SDK deploys background delivery methods and persistent local buffers to reduce telemetry loss during page lifecycle transitions:
Problem: Client browsers often terminate or navigate away before asynchronous requests complete, causing lost exit events. Delivery Mechanisms and Fallbacks: 1. Primary transport (navigator.sendBeacon): Queues an HTTP POST asynchronously without blocking page navigation. A successful return means the browser accepted the payload for transfer, not that the analytics server acknowledged it. Browser implementations also limit the queued payload size to roughly 64 KiB. 2. Page lifecycle guidance: Flush on visibilitychange when the document becomes hidden, with pagehide as a fallback, rather than relying on unload or beforeunload. 3. Fallback transport 1 (Pixel Tracking): Sends an HTTP GET request to /pixel.gif with encoded event parameters. It can work in some browser contexts but may be blocked by content blockers and privacy tooling. 4. Fallback transport 2 (Durable Client Buffer): The web SDK saves unconfirmed events in browser storage and the mobile SDK uses app-local durable storage. Buffered batches retry when the application becomes active and connectivity returns. 5. Tolerance threshold: Treat 2% to 5% event loss as an explicit product-level tolerance for this scenario rather than an industry-wide guarantee. Stronger requirements should use durable client queues, event IDs, and retry semantics.
3. High Scale Unique Counting with HyperLogLog
Maintaining exact distinct user sets in memory across tens of thousands of client websites creates severe memory bottlenecks. The system deploys Redis HyperLogLog for probabilistic cardinalities:
Problem: Counting distinct active users across 50,000 corporate websites in real time creates massive memory bottlenecks if tracking raw sets.
Flink + HyperLogLog (HLL) Solution:
1. Write: The Flink keyed window maintains an HLL-style sketch per site and one-minute bucket. Each client_id updates the sketch instead of materializing the full distinct-user set.
2. Read: The window aggregate emits an estimated distinct-user count to Redis as active_users:{site_id}, with a 120-second TTL so stale live metrics expire naturally.
3. Accuracy: HyperLogLog provides an approximate distinct count with a standard relative error of approximately 0.81% for the canonical configuration.
4. Memory efficiency: Each HyperLogLog sketch consumes at most about 12 KB regardless of cardinality.
5. Scale footprint: One active minute bucket for 50,000 customer websites multiplied by 12 KB yields roughly 600 MB before allocator and replication overhead. Retaining 60 minute buckets would require roughly 36 GB before overhead, so retention policy directly affects state sizing.
Alternative Redis-native implementation:
PFADD active_users:{site_id}:{minute} {client_id}
EXPIRE active_users:{site_id}:{minute} 120
PFCOUNT active_users:{site_id}:{minute}4. Event Identity and Deduplication
The ingestion contract treats the tenant-scoped pair (site_id, event_id) as the stable idempotency key because the pipeline uses at least once delivery. Retries can therefore produce the same logical event more than once. Flink streaming aggregates the (site_id, event_id) pair within a bounded retention period before updating live counters, while batch jobs deduplicate immutable S3 events before producing historical rollups. The retention horizon must cover the expected retry and replay window. ClickHouse retains site_id and event_id in the fact schema so reconciliation jobs can identify duplicate records and audit lineage.
5. Late Arriving Event Windowing Accuracy
Balancing dashboard update immediacy against historical analytical precision requires clear delineation between processing time and event time:
Dashboard Query: SELECT COUNT(*) FROM events WHERE timestamp > NOW() - INTERVAL 1 HOUR Case 1 (Processing Time): Count events at the moment they reach the ingestion collector. - Advantage: Stable output because historical numbers never shift retroactively. - Disadvantage: Inaccurate because delayed mobile events are attributed to the wrong time window. Case 2 (Event Time): Count events using the client side recorded creation timestamp. - Advantage: Provides more faithful historical attribution based on when the user action occurred, subject to timestamp quality, clock skew, duplicates, and the configured lateness policy. - Disadvantage: Dynamic shifting occurs because metrics for the latest hour update as delayed events arrive. - Operational convention: Display a 5-minute settling window warning in user interfaces indicating that current hour counts remain provisional until the configured 5-minute settling period has elapsed.
Additional Considerations
1. GDPR Right to Erasure and De-identification
Enterprise analytics architectures must address privacy compliance, long term retention economics, and practical interview execution.
Privacy regulations require systems to erase individual user histories upon request, presenting significant challenges for append only columnar files stored in S3:
GDPR, CCPA, and Erasure Compliance: 1. Consent and Opt-Out: Suppress client SDK telemetry collection where consent is the required legal basis until the applicable consent signal is granted. 2. Right to Erasure: Publish user deletion requests keyed by site_id and client_id to the dedicated deletion_requests Kafka topic so downstream storage and analytical consumers can apply the same erasure decision. Raw Kafka topics use bounded retention, so directly identifying fields should be minimized before publication or retained only for the documented short retention period. 3. Managing Immutable S3 Parquet Storage: - Challenge: Deleting a single client record directly in Parquet requires rewriting affected partition files. - Solution: Mark deleted identifiers in ClickHouse immediately so analytical queries exclude them, then run background S3 rewrite jobs on a cadence that satisfies the applicable erasure SLA. Weekly is an illustrative cadence, not a universal compliance requirement. Tracking Alternatives: 1. First-Party Cookies: Set tracking identifiers strictly on the host domain, which remains durable within the browser but clears if the user switches browsers. 2. Server-Side Collection: Route analytics payloads through the customer origin server instead of direct browser endpoints, reducing client side ad-blocker suppression. 3. Privacy-First Aggregate Analytics: Record aggregate pageview counts without storing unique client fingerprints or IP addresses. This can reduce privacy exposure and may change consent obligations depending on jurisdiction and implementation, but it does not universally eliminate legal requirements.
Related Problems and Concepts
Real time stream processing, sliding window aggregations, and top trending algorithms share architectures with Trending Topics, Top-K Rankings, and Like Count (High-Profile). Foundational stream mechanics connect directly to Stream Processing Basics, Kafka Architecture and Guarantees, Back-of-the-Envelope Estimation, and System Design Interview Patterns.
Interview Walkthrough
Structuring the discussion during a 45 minute architectural interview:
- 25 minute cut
Skip deep multi region and staff-level compaction mechanics unless requested.
- Separate ingest from query paths clearly on the whiteboard (5 min)
- Cover client reliability including beacon delivery, local buffering, and batching to reduce event loss during page unloads (6 min)
- Detail the Lambda architecture with a Flink speed tier for real time counters and a Spark batch tier for reconciliation (5 min)
- Explain session stitching using Flink session windows and identity resolution across anonymous client IDs within each site (5 min)
- Justify choosing event time over processing time, emphasizing how watermarks and allowed lateness bound event time completeness for funnel accuracy (4 min)
- Separate the ingestion path from client SDK to Kafka from the analytical query path in ClickHouse, as they have opposing latency and optimization goals.
- Use
navigator.sendBeacon()and client side local storage buffering to reduce page unload event loss when the browser tab closes. - Apply a Lambda architecture with Redis as the speed layer for live counters and ClickHouse as the analytical engine for deep cohort and funnel queries.
- Stitch user sessions using Flink session windows per site_id and client_id with a 30-minute inactivity threshold, a bounded-out-of-orderness tolerance of 5 minutes, and a separately configured allowed-lateness policy.
- Choose event time over processing time for accuracy, displaying a 5-minute settling window banner so dashboard numbers avoid confusing dynamic shifts.
- Use HyperLogLog for unique visitor estimation, capping memory at 12 KB per customer site regardless of visitor cardinality while maintaining an error rate of roughly 0.81%.
- Plan GDPR erasure by scrubbing personal identifiers at ingest and issuing tombstone markers rather than executing destructive row deletes on immutable S3 Parquet files.
- Avoid the common pitfall of counting raw pageview events as unique users without deduplication, which can materially inflate daily active user metrics, with the impact depending on traffic and deduplication behavior.
Engineering Trade-offs
1. Columnar OLAP Database Analysis
Key architectural trade-offs evaluate columnar storage efficiency, real time vs batch reconciliation, and multi step conversion funnel calculations.
Comparing analytical storage engines across query latency, ingestion throughput, compression efficiency, and operational overhead:
| Feature | ClickHouse ⭐ | Druid | BigQuery | TimescaleDB |
|---|---|---|---|---|
| Query Speed | Can be sub second for suitable filtered and aggregated queries | Can be sub second with pre-aggregation | Often seconds for analytical scans | Often seconds on large analytical scans |
| Ingestion Rate | Can exceed 1M rows/sec depending on schema and hardware | Can exceed 1M rows/sec with suitable ingestion design | Managed streaming or batch ingestion with service and cost limits | Can reach high write rates with tuned schema and hardware |
| Storage Cost | Low (extreme compression) | Moderate | High (pay-per-query scan) | Low |
| Real time Lag | Near real time, workload dependent | Near real time, workload dependent | Latency depends on ingestion and query mode | Near real time for suitable workloads |
| SQL Support | Rich SQL dialect | SQL with engine-specific limitations | Standard SQL | PostgreSQL compatible |
| Operational Complexity | Moderate, depending on topology | Higher operational complexity | Low infrastructure operations because it is managed | Lower operational complexity when PostgreSQL is familiar |
2. Conversion Funnel Query Optimization
Evaluating multi step conversion funnels across billions of events requires vectorized analytical engines rather than iterative self-joins. We leverage ClickHouse windowFunnel functions:
-- Conversion Funnel Analysis across User Journeys
-- ClickHouse windowFunnel() evaluates multi step conversion chains in a vectorized scan
SELECT
countIf(step >= 1) AS landing,
countIf(step >= 2) AS signup,
countIf(step >= 3) AS payment,
countIf(step >= 4) AS purchase
FROM (
SELECT
site_id,
client_id,
windowFunnel(86400)(
timestamp,
url = '/landing',
event_type = 'signup_complete',
event_type = 'payment_added',
event_type = 'purchase_complete'
) AS step
FROM events
WHERE site_id = {site_id:String}
AND timestamp > now() - INTERVAL 7 DAY
GROUP BY site_id, client_id
);
-- Query latency depends on filtering, grouping cardinality, data volume, and cluster resourcesReview
How helpful was this walkthrough?
Click a star to rate. We actively use this feedback to refine and update our system design content.
Discussion
Share your thoughts, ask questions, or help others.