advanced 3 min answer

An interviewer says - we shard the way Pinterest documented it: the shard id is packed into every object id, routing needs no lookup table, and an object never leaves the shard it was created on. Tell me what that decision has already decided for us, then tell me what you do about the first query that has to span shards.

pinterestshardingidentifiersscatter-gatherdenormalisation
Show the full answer Hide the answer

What the interviewer is testing

Whether you can read a design's consequences off its data model. Pinterest described this scheme publicly in a 2015 engineering post and an earlier conference talk — a 64-bit object id packing a shard id, a type id and a local id, so any id tells the client which database to talk to; the id's shard field allows 65,536 shards and the first few thousand were opened across a handful of hosts, each running hundreds of databases; queries are primary-key or index lookups with no joins; data never leaves the shard it was created on, although a whole shard can be moved to another host. The weak candidate says "clever, no lookup table". The strong one says what has become impossible.

The clarifying questions that change the answer

Is the new query user-facing or analytical · is its result set bounded or unbounded · how stale may it be · what is the fan-out distribution, because a feature where the median user follows 50 accounts and the top 0.1% follow 500,000 is two different designs.

What the design has already decided

  1. The shard count is fixed, by the bit width and by every id already issued. Growth happens by moving logical shards onto more hosts, never by splitting a shard. That is a feature — it is why there is no resharding project — and it means capacity planning is "how many shards will we ever want", answered once.
  2. The shard key is identity. Re-keying means re-issuing ids, and ids are in URLs, caches, other services' foreign keys and users' bookmarks. This is the least reversible decision in the system.
  3. There are no cross-shard transactions and no joins. Anything spanning shards is application-level scatter-gather, written by you.
  4. Co-location is the whole bet. Everything a user owns has to be reachable from the user's shard, or the scheme buys nothing.

A strong answer's arc for the cross-shard query

Three options, each with the condition that selects it:

  • Denormalise at write time into the reader's shard. Right when the query is user-scoped and the fan-out is bounded. The cost is write amplification equal to the average audience size, and a tail you must special-case: above roughly tens of thousands of recipients, push becomes pull.
  • Build a purpose-made derived store — a search index or a time-ordered secondary store fed from the shards. Right when the query is global, ad hoc, or ordered by something that is not the shard key. The cost is a second consistency boundary with its own staleness window and rebuild time.
  • Bounded scatter-gather. Acceptable over tens of shards, never thousands. The arithmetic that kills it is availability, before latency: 64 shards at 99.9% each give the query 0.999^64, about 94% — one request in 16 meets a shard that is unavailable, and you must then decide between a slow answer and an incomplete one.

Common weak answers

  • "Add a global secondary index." In a shared-nothing fleet of independent databases there is no such thing; you would be building option two and calling it an index.
  • "Reshard on a better key." The ids are the key. That is a re-identification programme touching every consumer of an id, not a migration.
  • "Scatter-gather with a timeout." A timeout converts a slow query into a wrong answer with missing shards, which for a user-facing list is worse than slow.

What a strong answer adds

Where copying Pinterest would be a mistake. Their domain partitions cleanly by one entity, and the scheme was chosen by a small team that needed routing to cost nothing. If your access patterns are many-to-many in both directions — an ads system queried by campaign and by audience — baking the shard into identity locks you out of the cheap path for half your queries, and a routing table you can change, at the cost of one lookup, is worth its price.