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
Select a component to see what it is responsible for and which state it owns.
- 1Customers' apps → Capture API: Batches of events
- 2Customers' dashboards → Query service: Chart request
- 3Query service → Postgres: Teams and saved charts
- 4Query service → Events (column store): Aggregate by column
- 5Capture API → Event stream: Append, then acknowledge
- 6Ingestion workers → Event stream: Read a partition
- 7Ingestion workers → Postgres: Who is this ID?
- 8Ingestion workers → Events (column store): Batch insert
What you need to know
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.
Check
Events are sharded by team. One team becomes a third of all traffic. What happens?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.