Product analytics over billions of events
Design a Product Analytics System
Take in billions of product events a day and answer questions nobody planned for in seconds: choose where events live, keep ingestion alive through spikes and outages, make filters on people fast, and decide what happens when two anonymous visitors turn out to be one person.
Intermediate, about 45 minutes, 9 stages
The situation
An analytics product for software teams. Customers add a snippet to their app, and it sends events: a page was viewed, a button clicked, a plan upgraded. Each event has a name, a timestamp, the ID of whoever did it, and a bag of properties (browser, country, plan, and whatever else the customer's developers decide to send).
Customers build charts from these events without writing SQL: daily signups from Germany, the share of visitors who sign up and then pay within a week, how many people come back a month later. Each chart is a new question; the product cannot know in advance which ones will be asked.
There are about 6,000 customer teams, sending about 2 billion events a day in total, and the biggest few teams send a large share of it. Everything currently lives in one Postgres table with the properties in a JSON column. Charts for small teams are fine. Charts for big teams time out.
What it has to do
Functional
- Accept events from customers' apps and servers.
- Answer trends, funnels and retention questions, filtered by any event or person property.
- Link a visitor's anonymous events to their account once they sign up.
- Show events in charts within a minute or so of being sent.
Non-functional
- A chart over 90 days of a large team's data returns in a few seconds.
- An accepted event is never lost, even if the analytics database is down.
- A retried batch never counts an event twice.
- One customer's traffic spike does not delay everyone else's charts.
Constraints and assumptions
- About 2 billion events a day, with peaks around four times the average.
- The largest teams each send hundreds of millions of events a month.
- Properties are free-form: thousands of different keys across all customers.
- An event is about 1 KB as sent, most of it properties.
- Queries almost always filter by one team and a time range.
- Events are never edited after they arrive; people's properties change over time.
Interview questions it prepares you for
- “Design Google Analytics, Mixpanel or Amplitude.”
- “Design a system that counts events and shows real-time dashboards.”
- “Design an ad-click aggregation system.”
- “Why would you use a column store for analytics, and what is it bad at?”
Read and practise next
How PostHog built it · Analytics on billions of events, in their engineers' own words
Concepts to know first: Message queues, Partitioning.
Similar systems: Design a URL Shortener, Design a Distributed Job Queue, Design a Distributed Cache (Memcache).