advanced 2 min answer Multiple choice

Figma's Postgres stack grew roughly 100x between 2020 and 2024. The team had already partitioned tables vertically and now needed horizontal scale for the largest tables. They evaluated distributed SQL engines including CockroachDB and TiDB and rejected all of them, then spent about nine months building their own sharding layer with a query proxy written in Go. Which fact best explains that choice?

figmapostgresshardingtechnology-selectionoperability
Pick one
Show the full answer Hide the answer

The situation they were in

Vertical partitioning had already moved table groups onto separate hosts, which buys time and stops at the first table too large for one machine. The 2024 write-up describes the next step as horizontal sharding of those tables, under schedule pressure, by a small databases team.

What they chose, and why it fit their constraints

They stayed on Postgres and built the routing themselves. Three pieces carry it: colos, colocation groups of tables that share a sharding key so joins and transactions still work within one key; logical shards implemented as Postgres views, so the application could be routed as if sharded before any data moved; and DBProxy, a Go service that parses queries and routes them to the right physical shard.

The published reasoning for rejecting the distributed SQL options is that the team would have had to rebuild its domain expertise from scratch, under an aggressive timeline. The scarcest resource was operational knowledge of one database, not a feature in another one, and the decision follows from that rather than from any benchmark.

The views step is the part worth stealing. Logical sharding was rolled out by percentage against a single database, which made the risky part reversible: routing bugs surfaced under real traffic while the data was still in one place, and only then did the physical failover happen. Roughly nine months to shard the first table end to end is the honest price.

Why the other options fail

"Distributed SQL cannot support transactions within a shard key." It can; that is the point of those engines. The constraint Figma engineered around is the same one everyone has: cross-shard joins and transactions are the expensive case, which is why colos exist.

"Their access patterns needed schemaless documents." This inverts the story. A broad relational query surface with joins and transactions is precisely why the destination was a sharded relational database rather than a document or wide-column store, and it is the fact most teams get wrong when they cite "scale" as the reason to leave SQL.

"A managed service would have cost more per gigabyte." Storage price is rarely the deciding term at this size, and if it had been, a nine-month engineering programme with a permanent proxy to operate is not the cheaper answer. Comparing infrastructure prices while ignoring engineer-years is the standard error in this decision.

When this is the wrong answer, and where copying it would be a mistake

The same reasoning points the other way for almost everyone. A team without a dedicated databases group, without deep Postgres operational depth, and without the traffic to justify a permanent query proxy should read this as an argument for buying the managed distributed engine: the deciding question is which capability you already run well, and the answer differs by team. Build the proxy when you have both the expertise and a workload that a managed engine's tail latency or transaction model genuinely cannot serve, and remember that the proxy becomes a piece of critical infrastructure you now own forever.