Programmable RevenueBuilder Edition Subscribe
Issue 01 Builder Spotlight Reported

The Postgres scale-out problem, and how Figma crossed it in ten seconds

Every product that wins eventually outgrows one Postgres box, and every path out is expensive: NewSQL migration, middleware, extensions, or sharding logic smeared through the codebase. The standout answer is Sammy Steele's work at Figma: shard logically first, route in one proxy, rehearse until the irreversible step is seconds long.

Thomas Cornelius August 22, 2026 7 min read
Technical schematic: one postgres with users, files, and orgs table groups, an app sending sql through dbproxy with hash(file_id), routing to one of eight shards, with the steps views first then move data.
The whole move: one postgres, a routing proxy, eight shards, views before data.
Engraved line-art portrait of Sammy Steele Photo: graph8
Builder Spotlight Sammy Steele Senior Staff Engineer, databases tech lead , Figma

Relational databases scale up beautifully and stop. A growing product rides one Postgres instance through years of bigger boxes, and then the physics arrive: tables of several terabytes and billions of rows, vacuums that cannot keep pace, IOPS at the provider’s ceiling. Scaling out, splitting one database into many, is the only door left, and it is the single most feared migration in infrastructure, because it has to happen underneath a product that cannot stop.

Builder Spotlight takes one problem like this per profile and points at the builder whose public work is the standout answer. For Postgres scale-out, that is Sammy Steele, senior staff engineer and databases tech lead at Figma, previously a builder of petabyte-scale metadata systems at Dropbox. Her team’s horizontal sharding of Figma’s Postgres, after load grew almost 100x from 2020, is the model of how to cross this wall without the product noticing.

The problem, across the board

There are five known doors out of a full Postgres box, and each has a price. Managed NewSQL (Spanner, CockroachDB) means migrating semantics and leaving your provider in one jump, the riskiest possible sequencing for a live product. Vitess is proven sharding middleware, but it grew up in the MySQL world. Citus extends Postgres itself, at the price of coupling your upgrade path to an extension. Application-level sharding, the route Notion documented in its own sharding write-up, works, but spreads routing logic through the codebase where every feature pays for it forever. And the fifth door, the one most teams fear, is building your own routing layer.

Figma took the fifth door and published the map.

What Figma built

Three decisions carry the design.

They sharded around the data model instead of forcing one universal key. Related tables that share a sharding key move together as a group, which the team calls colos, and keys like user, file, and organization were chosen per table group. Hashing the keys prevents hotspots. The classic sharding trap, contorting every table onto one key, never happens.

They separated logical sharding from physical sharding. Views made tables behave as sharded before any data physically moved, so production traffic could run through the sharded layout while rollback stayed cheap. The risky step was rehearsed until it was routine.

They built DBProxy, a Go service between the application and Postgres that parses each query’s AST, extracts the shard key, and routes it. Hundreds of application call-sites did not change, because the sharding knowledge lives in one owned service at the SQL boundary.

SELECT ... FROM files WHERE file_id = 8123 dbproxy (go) parse AST -> shard key: file_id hash(8123) -> shard 3 shard 1 · colo: files+ shard 3 · colo: files+ shard 5 · colo: users+ shard 7 · colo: orgs+ colos: tables sharing a shard key move together. views made tables logically sharded before any data physically moved, so rollback stayed cheap.
DBProxy: the application sends ordinary SQL; the routing layer knows where the row lives.

The project took nine months, with the first sharded table live in September 2023. That first failover cost about ten seconds of partial availability on the primaries, no impact on replicas, and no latency or availability regressions after.

Why this was the right technical direction

The deep idea is that reversibility was a design requirement, not a hope. Every step before the physical failover could be undone cheaply, because the logical layout ran in production while the data stayed put. The project’s risk was concentrated into a window measured in seconds, and that window was rehearsed. Most sharding projects fail precisely in the step Figma made small.

The second idea is about where knowledge lives. Sharding is permanent; once the data splits, something must route every query forever. Putting that knowledge in one proxy at the SQL boundary, rather than in application code or a vendor’s middleware, means Figma owns its hardest dependency and its applications stay ordinary Postgres clients. The honest cost, which the team acknowledges, is maintaining a query-parsing proxy forever.

What applies today

The five doors have not changed, and neither has the trade at the center: who owns the routing layer. What Figma’s project settled is the sequencing question. Logical-before-physical, rehearse-then-cut is now the reference pattern for changing a database under live traffic, and it generalizes beyond sharding: the same shape de-risks engine swaps, provider moves, and schema splits. Any table with a natural entity key, contacts, accounts, workspaces, events, can run this play; the hard cases are the tables that join across keys, which is exactly what colos exist to contain, and cross-shard queries, where every team still pays the scatter-gather tax her write-up compresses past.

Run the same play

Nine months of work, ten seconds of impact. The ratio came from a sequence any team can copy.

  1. Choose shard keys per table group from real query patterns, not one universal key. Tables that join constantly share a key and move together.
  2. Make the sharding logical first: views (or a routing flag) that treat tables as sharded while the data has not moved. Run production traffic through the logical layout until it is boring.
  3. Concentrate routing in one service at the SQL boundary, so application call-sites stay untouched and the sharding knowledge has one home.
  4. Rehearse the physical failover until the irreversible window is measured in seconds, and keep the rollback alive until the moment you cross it.

A contact graph or an event store, the workloads under the graph8 platform, hits the same wall Figma hit, and this sequence is the calmest public account of the way through it.

About this section: Builder Spotlight profiles an engineer outside graph8, from their own published work only.

Source: Sammy Steele, “How Figma’s Databases Team Lived to Tell the Scale”, Figma engineering blog, March 2024. All Figma numbers are from that post. More from Steele: her Postgres.fm episode on the road to 100TB and the Software Engineering Daily interview. Background: Notion on application-level Postgres sharding and the Vitess documentation for the middleware path.

Build with graph8