Skip to content

Design a Product Analytics System, 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.

System so far· 8 parts
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

What you need to know

0 of 1 checks done
  1. 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.

  2. Check

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