Skip to content

The finished design, decision by decision

How to design a Product Analytics System

Not the one correct diagram, but a design you can defend under these constraints: the finished architecture, then every stage's question with the reasoning that answers it, the tradeoffs it accepts, and where another engineer could land differently.

Product analytics over billions of events

The short answer

8 parts, each with one job. The map below shows how requests and data move between them; the stages after it explain why each part is there.

Customers' apps
Send events in small batches through the analytics snippet.
Capture API
Checks the API key, stamps the team and a unique event ID, and appends to the stream; nothing else.
Event stream
Durable log of accepted events, partitioned by team and person ID, kept for several days.
Ingestion workers
Resolve each event's person, copy person properties onto it, and insert into the column store in large batches.
Events (column store)
One wide table, sorted by team, event name and time; hot properties stored as their own columns.
Postgres
Customers, dashboards, and the mapping from anonymous and known IDs to people, with their current properties.
Query service
Turns a chart into SQL, runs it, and caches the result for a short time.
Customers' dashboards
Build charts and read them.
12345678CLIENTCustomers' appsEDGECapture APILOG / STREAMEvent streamWORKERIngestionworkersDATABASEEvents(column store)DATABASEPostgresSERVICEQuery serviceCLIENTCustomers'dashboards

Select a component to see what it is responsible for and which state it owns.

  1. 1Customers' apps → Capture API: Batches of events
  2. 2Customers' dashboards → Query service: Chart request
  3. 3Query service → Postgres: Teams and saved charts
  4. 4Query service → Events (column store): Aggregate by column
  5. 5Capture API → Event stream: Append, then acknowledge
  6. 6Ingestion workers → Event stream: Read a partition
  7. 7Ingestion workers → Postgres: Who is this ID?
  8. 8Ingestion workers → Events (column store): Batch insert

Why does this design work?

Customers' apps send events to a capture API that only checks the key, stamps the team and a unique event ID, and appends to a durable event stream partitioned by team and person. Ingestion workers read each partition in order, resolve the person (linking anonymous IDs at signup), copy the person's current properties onto each event, and insert large batches into a column store, committing their position only after each insert, with a deduplication token so retries never double-count.

Events sit in one wide table sorted by team, event name and time, sharded by a hash of team and person so no customer becomes a hotspot. Frequently filtered properties are real columns; the rest stay in JSON. Charts read only the columns they need and skip everything outside the team and date range, so a 90-day chart over hundreds of millions of events takes seconds. Person merges are recorded in a small override table applied at query time and folded in periodically.

Invariants, and where they are enforced

  • An event the capture API accepted is eventually stored, even after hours of database downtime.

    Capture acknowledges only after the append; workers move their position in the stream only after a batch is safely inserted; retention outlasts the longest outage. Enforced by Capture API, Event stream, Ingestion workers.

  • A retried batch never counts an event twice.

    Every event gets a unique ID at capture; a retried batch reuses the same deduplication token, and the table collapses rows with the same event ID. Enforced by Capture API, Ingestion workers, Events (column store).

  • One person's events are processed in order, by one worker at a time.

    The stream is partitioned by team and person ID, and each partition has a single consumer. Enforced by Event stream, Ingestion workers.

What does it rely on?

  • Stream retention outlasts the longest database outage plus catch-up time.
  • Queries almost always filter by team and time, matching the sort key.
  • A small set of properties accounts for most filters.
  • Customers accept that copied person properties describe the moment of the event.

What tradeoffs does it make?

ChoiceGainsCosts
Column storeFast scans and compression.Batch-only inserts; expensive updates.
Durable stream in frontNo loss through outages; batching.Another system; charts lag during catch-up.
Person data copied onto eventsNo joins on reads.Point-in-time semantics; merges need corrections.

What are the reasonable alternatives?

Postgres with indexes and rollups
Better when volumes are modest, or the charts are a fixed set.
Managed warehouse (BigQuery, Snowflake)
Better when freshness of minutes is fine and the team does not want to run a database.
Pre-aggregated time-series store
Better when questions are metrics over fixed dimensions rather than per-person funnels.

When does it stop working?

  • Queries routinely span all teams or ignore time, so the sort key cannot skip data.
  • Customers need person properties as they are now for every chart.
  • Merges become so frequent that the override table grows faster than it can be folded in.

Every stage, decided and explained

Spoilers, for the whole investigation: each stage's question and its answer, the reasoning behind it, and the tradeoffs it accepts. If you have not worked through the stages yet, you may want to do that first.

Work through the stages

Stage 1 of 9 · Model

What one chart has to read

Use 86,400 seconds a day and 1 KB an event. A large team sends 300 million events a month, and its daily-signups chart filters on two properties over 90 days.

What you need to know first

2 billion events a day. About how many events a second on average?

About 23,000 per second.

2,000,000,000 ÷ 86,400 ≈ 23,000 a second; peaks at 4× are about 93,000. A busy write rate, but an ordinary one if writes are batched.

At about 1 KB per event, how many terabytes of raw events arrive per day?

About 2 TB.

2,000,000,000 × 1 KB = 2 TB a day, about 180 TB over 90 days. Far too much to answer charts from memory.

A row store like Postgres keeps each row's fields together on disk. To count rows, it reads whole rows, every field, even if the query uses three of them.

An index helps find rows. It doesn't make reading them cheaper once they're found.

A chart needs 900 million events, using 3 small fields from each 1 KB event. Reading whole rows, about how many gigabytes are read?

About 900 GB.

900,000,000 × 1 KB = 900 GB read, for perhaps 10 to 20 GB of fields actually used. The query isn't slow at finding rows; it's slow because it reads everything else in them.

What the stage asks

Which statements follow?

  1. Holds

    2 billion events a day is about 23,000 a second on average, so peaks near 90,000 a second.

    2,000,000,000 ÷ 86,400 ≈ 23,000; four times that is about 93,000. A busy but ordinary write rate if writes are batched.

  2. Holds

    That is about 2 TB of raw events a day, or roughly 180 TB for 90 days.

    2 billion × 1 KB = 2 TB a day. The volume, more than the rate, rules out answering from memory.

  3. Holds

    The big team's 90-day chart has to consider about 900 million events.

    300 million a month × 3 months. Even if each row takes a microsecond, a single chart is minutes of work on one core unless most of each row is never read.

  4. Fails

    An index on (team, timestamp) in Postgres would make this chart fast.

    The index finds the 900 million rows quickly, but then each full row, properties and all, must be read to count them. The problem is not finding the rows; it is reading 900 GB to use three fields from each.

The reasoning

  1. Analytics queries are wide in rows and narrow in columns.
  2. Row stores read whole rows, so a three-field question pays for every field.
  3. Indexes find rows; they don't make reading them cheaper.

Analytical questions are wide in rows and narrow in columns: they touch a large share of a team's events but only a few fields of each. A store that keeps whole rows together must read everything to answer them. That is the bottleneck to design around, not the write rate.

Stage 2 of 9 · Decide

Where the events live

Charts are built by customers, so the questions are not known in advance. They nearly always filter by team and time range, then by event name and some properties.

What you need to know first

A column store keeps each column's values together on disk: all timestamps in one place, all event names in another. A query reads only the columns it uses. See Columnar storage.

Values in one column are similar (the same few event names, increasing timestamps), so they compress very well, often 5 to 10×.

Column stores also keep data sorted by a chosen key, and record the min and max of each block. A query whose filter matches the sort key can skip whole blocks without reading them.

With the sort key (team, event, time), "team 42's signups last quarter" reads a thin slice and skips every block from other teams, events and dates.

Why not precompute daily counts for every event and property instead?

Customers ask questions nobody planned (combined filters, funnels), and rollups only answer the ones you prepared.

Rollups are great for fixed dashboards. A funnel needs individual events in order, which a count can't give you.

What does a column store do badly?

Updates and deletes (they rewrite compressed blocks), fetching one whole row (every column is a separate read), and many tiny inserts (each creates a small part to merge later). It wants large, append-only batches.

What the stage asks

Where should events be stored for querying?

  1. Defensible

    Keep Postgres: partition the table by month, add read replicas and more indexes

    Fine for smaller volumes, and it keeps one database. At this scale every chart still reads whole rows, so big teams stay slow however many replicas share the pain.

  2. Defensible

    Pre-compute daily counts for each event and property as events arrive

    Known charts become instant. But customers ask new questions, combine filters, and build funnels that need individual events in order. Rollups are a good addition for fixed dashboards, not a replacement for the events.

  3. Sound

    A column store with one wide events table, sorted by team, event name and time

    A chart reads only the columns it uses, compressed several times over, and the sort order lets it skip every block belonging to other teams, other events and other dates. This is the layout most product analytics systems settle on.

  4. Defensible

    A search engine, indexing every property of every event

    Excellent at finding individual events by any property. Aggregating hundreds of millions of them, and holding an index on thousands of keys, is expensive compared with a column store.

What a strong answer covers

  • Reads only the columns the query uses, and those compress well.
  • Sorting by team, event and time lets queries skip most of the data.
  • Works for questions nobody planned, unlike rollups.
  • Names the costs: updates, deletes and single-row lookups are expensive; inserts must be batched.Supporting

The reasoning

  1. Store by column: queries read only the columns they use, and columns compress well.
  2. Sort by team, event and time so queries skip most blocks.
  3. Column stores dislike updates, single-row reads and tiny inserts; batch your writes.

The choice follows from the first stage: if queries use few columns of many rows, store data by column. The sort key is the other half: with (team, event, time) first, a query for one team's signups in one quarter reads a thin slice of a few columns. See Columnar storage.

The costs are real and shape the rest of the design: a column store wants large batches of inserts, dislikes updates, and is slow at fetching one whole row.

Tradeoffs

ChoiceGainsCosts
Row store with indexesOne database, easy updates.Reads whole rows; slow for big teams.
RollupsInstant known charts.Cannot answer new questions.
Column storeFast scans over few columns.Batch inserts only; updates are expensive.

Stage 3 of 9 · Break it

The database is down for twenty minutes

Today the capture API inserts each request's events straight into the database. Twenty minutes at peak is about 100 million events.

What you need to know first

Events arrive at a rate you don't control; the database absorbs them at a rate that varies, sometimes zero during an upgrade. A durable log (a stream like Kafka) between the two lets each run at its own pace. See Append-only logs.

The capture API's job shrinks to: check, stamp, append to the stream, acknowledge.

The database is down for 20 minutes at a peak of about 90,000 events a second. Roughly how many million events must be held?

About 108 million events.

90,000 × 1,200 s ≈ 108 million events, about 100 GB at 1 KB each. A stream holding days of data absorbs that easily; capture servers' memory doesn't.

Capture returns an error during the outage and relies on the snippet to retry later. What happens to many of those events?

They're lost: tabs close and phones go offline before retrying.

Client retries are best effort. Accepting events into a durable buffer is the only way not to depend on them.

After an outage, the workers have the backlog plus live traffic. Size them, and the database, for the catch-up, not only for the average. See Backpressure and capacity.

What the stage asks

How should events get from the capture API into the column store?

  1. Flawed

    Keep inserting directly, and return an error so the snippet retries later

    Browsers close tabs and phones go offline, so many retries never happen and events are lost. And on a normal day, thousands of tiny inserts a second create thousands of small parts for the column store to merge.

  2. Flawed

    Buffer events in the capture servers' memory and insert them every few seconds

    Batching fixes the small inserts, but a capture server that restarts loses its buffer, and twenty minutes of peak traffic will not fit in memory.

  3. Sound

    Capture appends each batch to a durable, partitioned stream and acknowledges; workers read the stream and insert large batches

    Capture's only dependency is the stream, so an outage of the database delays charts but loses nothing: the stream holds days of events, and workers catch up when the database returns. Workers can also make batches as large as the column store likes.

  4. Defensible

    Write events as files to object storage, and load them into the database every hour

    Durable and cheap, and common for warehouses. But charts would lag by an hour, and customers expect events within a minute or so.

What a strong answer covers

  • Capture depends only on the stream, so database outages do not lose or reject events.
  • Retention must outlast the longest outage, plus time to catch up.
  • Workers insert large batches, which the column store needs.
  • Workers record their position only after a batch is inserted.Supporting

The reasoning

  1. A durable stream between capture and the database lets each run at its own pace.
  2. Capture depends only on the stream, so database outages delay charts but lose nothing.
  3. Workers insert large batches and must be sized for catching up, not just the average.

Ingestion has two speeds: the speed events arrive, which you do not control, and the speed the database can take them, which varies. A durable log between them lets each run at its own pace. The capture API becomes simple and very hard to break: check, stamp, append, acknowledge.

Catching up after an outage is its own load: the workers have twenty minutes of backlog plus live traffic. Size the workers and the database for that, not only for the average. See Backpressure and capacity.

Stage 4 of 9 · Break it

Slow again, for a different reason

Properties are stored as one JSON string column, because every customer sends different keys. Select the lines that point at the cause.

What you need to know first

A column store only helps when the data is actually in columns. A single JSON column holding every property is, for those fields, a row store again: to read one property, the query reads and parses the whole blob for every row.

A query profile shows 98% of blocks skipped by the sort key, but 37 of 38 GB read are the properties column. What's the problem?

The sort key works; the waste is inside each row, in the JSON blob.

Skipping found the right rows efficiently. Then, to use two properties, each row's whole JSON was read and parsed.

The usual fix is a hybrid: keep the JSON for the long tail of rare keys, and promote the most-queried keys to real, typed columns, extracted from the JSON as rows are inserted. Which keys? Look at what queries actually filter on. Old data needs a backfill for the new columns to be useful over long ranges.

What the stage asks

Select the lines that explain the slowness.

LogQuery profile: daily signups where plan = 'pro' and country = 'DE', team 5150, last 90 days
  1. 1blocks skipped by sort key (team, event, timestamp): 98.6%
  2. 2rows read: 41,200,000
  3. 3columns read: team_id, event, timestamp, properties

    To use two properties, the query reads the whole properties column, which is most of each event's bytes.

  4. 4bytes read: 37.9 GB (properties: 37.1 GB)

    Almost all the bytes are the JSON blob. The column layout is wasted on a column that holds everything.

  5. 5CPU profile: 81% JSONExtractString

    Every row's JSON is parsed on every query to pull out two values. The same parsing is repeated for every chart that uses those keys.

  6. 6parallel threads: 32; peak memory: 1.4 GB

What the fix has to do

  • One JSON column holding all properties defeats reading only what you need.
  • Parsing JSON per row per query dominates CPU.
  • Store the most used properties as their own typed columns, filled at insert time (and backfilled for recent data).
  • Notes that the sort key is working; the waste is inside each row.Supporting

The reasoning

  1. One JSON column for all properties turns a column store back into a row store for those fields.
  2. Promote frequently queried keys to typed columns computed at insert time.
  3. Choose which keys to promote from real query patterns, and backfill recent data.

A column store only helps if the data is actually in columns. Free-form properties tempt you into one big JSON column, which turns the column store back into a row store for exactly the fields people filter on.

The usual answer is a hybrid: keep the JSON for the long tail of rare keys, and promote the keys that queries use most into real columns, computed from the JSON as rows are inserted. Find them by looking at which keys queries actually filter on. Old data needs a backfill for the new column; recent data is what most charts read, so start there.

Stage 5 of 9 · Decide

Filtering by who did it

Customers filter charts by person properties too: "signups from people on the Pro plan", where plan belongs to the person and changes when they upgrade. People live in Postgres: about 2 billion of them across all teams, keyed by ID, with their current properties.

What you need to know first

Denormalisation copies data to where it's read, so queries don't have to join. It costs storage and write-time work, and the copy is a snapshot: it records what was true when it was copied.

In analytics, reads are the expensive part and writes are append-only, so the trade usually pays.

Why not fetch Pro users' IDs from Postgres and filter events with WHERE person_id IN (…)?

A big team can have millions of Pro users; shipping millions of IDs into a query is slow and fragile.

It works for small teams and breaks as customers grow, which is exactly when it matters.

Person properties are copied onto each event at ingestion. A user upgraded from Free to Pro last week. A chart asks 'signups by plan'. Which plan does their signup event show?

Free: the plan they were on when they signed up. That's often the right answer ("what plan were people on when they did X?"), but it differs from "people who are Pro now". Write down which questions use snapshots, and offer current properties through a join for the ones that need them.

What the stage asks

How do charts filter events by person properties?

  1. Flawed

    Fetch the matching person IDs from Postgres, then filter events with an IN list

    A big team can have millions of Pro users. Shipping millions of IDs from one database into a query on another is slow and fragile, and it gets worse as customers grow.

  2. Defensible

    Copy people into a table in the column store, and join events to people at query time

    Correct, always using people's current properties. But joining hundreds of millions of events to millions of people on every chart is heavy, and the persons table changes constantly, which a column store handles poorly.

  3. Sound

    When an event is ingested, copy the person's ID and current properties onto it

    Charts filter on columns that are already on each event: no join. The meaning changes, though: an event records the person's properties as they were when it happened. For many questions ('what plan were people on when they did this?') that is the better answer; for others it is a surprise to explain.

  4. Defensible

    Pre-compute, for each person, counts of each event they have done

    Fast for 'people who did X at least 3 times', but it cannot answer time-based charts or funnels, which need the events themselves.

What a strong answer covers

  • Removes the join from every chart; filters read columns on the event.
  • Properties become as of the event, not current; says when that is right and when it surprises.
  • Costs storage and ingestion work, softened by column compression.
  • Mentions offering current properties through a join for the queries that need them.Supporting

The reasoning

  1. Copy person properties onto events at ingestion to remove the join from every chart.
  2. Copied properties are snapshots as of the event, not current values.
  3. Offer current properties through a join for the questions that need them.

This is denormalisation on purpose: pay once at write time to avoid paying on every read. In an analytics system reads are the expensive part and writes are append-only, so the trade usually wins.

The cost is semantic, not just storage. Copied data is a snapshot: it tells you what was true when the event happened. Write down which questions that answers differently, because customers will ask.

Tradeoffs

ChoiceGainsCosts
Join at query timeAlways current; no duplication.Heavy joins on every chart.
Copy onto eventsNo joins; history is preserved.Point-in-time meaning; more bytes; later changes do not apply to old events.

Stage 6 of 9 · Break it

Two visitors who were one person

Anonymous events carry a random device ID. At signup, the snippet sends an identify call linking that device ID to the new account. Person data is copied onto events at ingestion, as decided in the last stage.

What you need to know first

Before signup, a visitor's events carry a random device ID. At signup, an identify call links that device ID to the new account. But earlier events were already stored with the device ID, and in a column store nothing rewrites stored events cheaply.

So after the link, the same human appears as two people in older data.

What's a cheap way to make funnels count them as one person without rewriting events?

Keep a small table mapping old IDs to merged ones and apply it at query time, folding it in periodically

Merges are rare compared with events, so the table stays small and the join is cheap. A background job can rewrite events in bulk occasionally.

The identify call and the user's next event are processed by different workers, and the event is processed first. What goes wrong?

The event is tagged with the old anonymous person, because the link didn't exist yet when it was enriched. Partitioning the stream by team and person ID puts both on one worker, in order. See Ordering.

What the stage asks

Which statements hold?

  1. Holds

    Linking the two IDs at signup does not change the person on events already stored.

    The person was copied onto each anonymous event when it was ingested, and nothing rewrites stored events. The link affects only events ingested afterwards.

  2. Depends

    The fix is to update the person ID on every past event whenever two people are merged.

    It gives the right answer, but updates rewrite whole blocks of a column store, and merges happen constantly. Done one merge at a time it can overwhelm the database. Batching them, or keeping a small table of overrides that queries apply on top, are the usual compromises.

  3. Fails

    Ingestion can process events for the same person on any worker, in any order.

    If the identify call and the next event are processed out of order, the event is tagged with the wrong person. Partition the stream by team and person ID so one worker sees each person's events in order. See Ordering.

  4. Holds

    A small table mapping old person IDs to merged ones, applied at query time, keeps funnels right without rewriting events.

    Merges are rare compared with events, so the table is small and the join is cheap. A background job can fold it into the events table periodically and empty it.

The reasoning

  1. Identity merges change the meaning of old events, which a column store can't cheaply rewrite.
  2. Keep a small table of ID overrides applied at query time, folded in periodically.
  3. Partition the stream by person so identify calls and events are processed in order.

Identity is the place where "events are never edited" meets reality: the meaning of old events changes when you learn who they belonged to. You can rewrite history (expensive in a column store), or keep a small, separate record of corrections and apply it when querying, folding it in from time to time.

Ordering matters too. Processing one person's events in order requires that they land in one partition, which is a Partitioning decision made at the capture API.

Where another engineer could land differently

Some products choose never to merge retroactively and document it. That is simpler, and fine if customers understand that funnels start counting at identification.

Stage 7 of 9 · Decide

Write the batch inserter

Write the loop each ingestion worker runs for one stream partition. Inserts can fail, and workers can crash at any point. The database deduplicates an insert whose dedupToken it has seen recently.

What you need to know first

A stream consumer has two steps per batch: insert the batch, then commit its position ("I've processed up to offset 5,000"). The order decides what a crash does:

  • Commit then insert: a crash in between skips the batch. Lost.
  • Insert then commit: a crash in between re-reads the batch. Duplicated, unless the database can recognise it.

To let the database recognise a retried batch, what should its deduplication token be?

Derived from the batch's partition and offset range, so the same work always gets the same token

A retry of the same batch covers the same offsets and produces the same token, which the database recognises as a repeat.

"Exactly once" in practice is at least once plus deduplication. See Delivery guarantees. And batches must be large (by size or time), because a column store turns every insert into a part it later merges.

What the stage asks

Implement runPartition.

Reference implementation

const MAX_ROWS = 50_000;
const MAX_WAIT_MS = 2_000;
const sleep = (ms: number) => new Promise((r) => setTimeout(r, ms));

export async function runPartition(partition: number): Promise<never> {
  let next = await stream.committedOffset(partition);
  for (;;) {
    const batch: Event[] = [];
    const started = Date.now();
    while (batch.length < MAX_ROWS && Date.now() - started < MAX_WAIT_MS) {
      const events = await stream.read(partition, next + batch.length, MAX_ROWS - batch.length);
      if (events.length === 0) await sleep(100);
      for (const e of events) batch.push(await enrich(e)); // in order: identify calls first
    }
    if (batch.length === 0) continue;

    const first = batch[0]!.offset;
    const last = batch[batch.length - 1]!.offset;
    const token = `${partition}:${first}-${last}`;

    for (let attempt = 0; ; attempt++) {
      try {
        await db.insert(batch, token);
        break;
      } catch {
        await sleep(Math.min(30_000, 500 * 2 ** attempt) * (0.5 + Math.random() / 2));
      }
    }
    next = last + 1;
    await stream.commit(partition, next);
  }
}
  • The position is committed after the insert. A crash in between means the batch is read and inserted again, so the token must be the same next time.
  • The token comes from the offsets, not from a random value, so a restarted worker produces the same one for the same batch, provided batches are cut the same way. Each event's unique ID gives a second line of defence if they are not.
  • Batches close on size or time, so quiet partitions still reach charts within seconds.
  • A failed insert is retried with backoff rather than skipped: skipping would commit past events that were never stored.

What a strong answer covers

  • Commits the stream position only after the insert succeeds.
  • Accumulates a large batch (by size or time) instead of inserting each event.
  • Uses a deduplication token derived from the batch's offsets, so a retried batch is recognised.
  • Retries a failed insert with backoff instead of skipping the batch.
  • Enriches events in order within the partition.Supporting

The reasoning

  1. Insert, then commit the stream position: crashes then cause repeats, never loss.
  2. Derive the dedup token from the batch's offsets so retries are recognised.
  3. Accumulate large batches, and retry failed inserts with backoff instead of skipping.

"Exactly once" here is really at least once, plus deduplication: the stream may hand you a batch twice, and the database recognises the repeat. The two rules that make it work are the order of operations (insert, then commit) and a deduplication key that is the same every time the same work is retried. See Delivery guarantees.

Stage 8 of 9 · Change it

One customer is a third of the traffic

The column store runs on several servers. Each table is split into shards across them, and a query runs on every shard that holds relevant data, in parallel.

What you need to know first

How you split data across servers (sharding) decides where load lands. The natural choice is to shard by what queries filter on. But if one value of that key gets most of the traffic, one server gets most of the work. See Partitioning.

Events are sharded by team. One team becomes a third of all traffic. What happens?

Its server takes a third of all writes and runs its heaviest charts alone, while others idle.

Sharding by team keeps each team on one machine, so the biggest team can never use more than one.

Hashing (team, person) spreads each team over every server, so its charts run in parallel everywhere. Keeping one person's events together keeps per-person questions (funnels, retention) on one shard.

Spreading data isn't isolation, though: per-team ingestion quotas and query concurrency limits stop one team's success from becoming everyone else's outage.

What the stage asks

How should events be spread across the database servers?

  1. Flawed

    Shard by team, so each team's data sits on one server

    The viral team's server takes a third of all writes and runs its heaviest queries alone, while the others sit idle. It also cannot grow beyond one machine.

  2. Sound

    Shard by a hash of team and person ID, so every team's events spread across all servers; queries run on all shards in parallel

    Writes and reads for the big team are spread evenly, and its charts run in parallel on every server. One person's events stay together, which keeps per-person questions (funnels, retention) local to a shard.

  3. Defensible

    Move the largest customers to their own cluster

    Strong isolation, and worth it for a handful of very large accounts. It is more to operate, and it does not tell you how to spread data within any one cluster.

  4. Defensible

    Sample the viral team's events at 10% until the spike passes

    A legitimate tool if the customer agrees and charts scale numbers back up, but it changes their data. It is a product decision, not a sharding scheme.

What a strong answer covers

  • Hashing within a team spreads one team's writes and reads across all servers.
  • Keeping a person's events together keeps funnel and retention work local.
  • Adds per-team limits (ingestion quotas or query concurrency) so one team cannot take all capacity.
  • Notes that the stream absorbs the burst; ingestion lag rises for a while but nothing is lost.Supporting

The reasoning

  1. Shard by what spreads load, not by what you filter on.
  2. Hashing team and person spreads a big team across all servers while keeping each person local.
  3. Add per-team quotas and query limits for isolation; spreading data isn't the same thing.

Shard by what spreads the load, not by what you filter on. Every query filters by team, but sharding by team puts the biggest team's load on one machine. Hashing within the team spreads it, and the sort key still lets each shard skip other teams' data.

Spreading data is not the same as isolating customers. Per-team quotas at capture and limits on concurrent queries keep one team's success from becoming everyone else's outage. See Partitioning.

Stage 9 of 9 · Defend it

Defend the pipeline

Your interviewer: "That is a stream, a fleet of workers and a separate database. Why not keep Postgres with good rollups, or send everything to a managed warehouse?"

What you need to know first

A good defence ties every component to a requirement:

ComponentRequirement it serves
Column storead hoc questions over huge volumes
Streamno lost events through outages and spikes
Copied person propertiesfast filters without joins
Hashed sharding and quotasfairness between customers

A component you can't tie to a requirement shouldn't be in the design.

When would Postgres with rollups be the better choice?

At smaller volumes, or when the set of charts is fixed and known in advance

Then rollups cover the questions, and one database is far simpler to run.

What the stage asks

Defend the design, concede what the alternatives do better, and say when you would pick them.

Reference answer

The requirement that decides it is that customers build their own charts. Rollups answer the questions you planned; they cannot answer a funnel someone thought of this morning. So the raw events must be queryable, fast, at the scale of the largest customer.

Why a column store. Those questions read a few fields from a very large number of events. Storing by column, compressed and sorted by team, event and time, turns a 900 GB scan into a few gigabytes. Promoting hot properties to columns and copying person properties onto events keeps it that way.

Why the stream. It separates accepting events from storing them. The capture API stays simple and up; outages and spikes become lag rather than loss; and the database gets the large batches it needs.

What I concede. Under tens of millions of events a month, Postgres with sensible indexes and a few rollups is simpler and good enough. A managed warehouse removes the operational work, and is the right call when charts can be minutes old or query costs per scan are acceptable. This design earns its complexity when interactive, ad hoc questions over billions of events are the product.

What a strong answer covers

  • Customers ask questions nobody planned, so rollups alone cannot serve them.
  • Explains why a column store fits: few columns, many rows, compression, sort-key skipping.
  • Explains what the stream buys: outages and spikes do not lose events, and inserts are batched.
  • Concedes when Postgres or a managed warehouse is better: smaller volumes, or freshness and cost measured differently.

The reasoning

  1. Tie each component to a requirement it serves.
  2. Rollups alone can't answer questions nobody planned.
  3. Concede when Postgres or a managed warehouse is the better trade.

A good defence ties each part to a requirement: ad hoc questions (column store), no lost events (stream), fast person filters (copied properties), fairness between customers (hashing and quotas). If a part cannot be tied to a requirement, it should not be in the design.

How you did

Now try it as an interview question

  • “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?”

The interview mode mixes stages from this and other investigations with concept recall and questions about your own projects.

Back to the last stage