2 billion users, 1 boring database
Instagram gets brought up constantly as proof that an ordinary relational database will carry you further than you'd think. I wanted to see what that looked like underneath, so I spent a while reading about how they actually got there. These are my notes.
What stayed with me is this. They hit three separate ceilings on the way up, and Postgres was not the cause of any of them.
the setup was boring on purpose
At the start there were three engineers.
Django ran on EC2. Images went into S3. Everything else, who posted what and the comments and the likes, lived in one Postgres database.
That was it. No queue, no separate services, no cache sitting in front.
People move past this part too quickly. Running something this plain wasn't a shortcut they were lucky to survive. It's the reason three people could build anything at all.
the first wall wasn't the database
Response times started climbing, and nothing about the queries had changed.
Each Django process was holding a small set of open connections of its own. On its own that sounds harmless. The arithmetic is what stings.
Picture 50 application servers. Give each of them roughly 30 connections. You are now sitting at 1,500 connections held open against Postgres.
Postgres hands every connection its own process, and one costs somewhere near 1.3 MB simply to exist. Multiply it out and you land around 2 GB.
All of that is committed while the database sits idle. None of it holds a cached page or works on a query. It buys nothing except the right to be connected.
That's what I found odd. Load wasn't the problem. Introductions were.
pgbouncer
What you want is pooling, and pgbouncer is the usual answer.
It's a thin proxy you drop in front of Postgres. Your application connects to pgbouncer and cannot tell the difference. Behind it, pgbouncer folds all those application connections down onto a small number of genuine connections to the database.
Fifteen hundred going in. Around thirty coming out. The amount of actual work never changes.
The memory that comes back goes where it belonged in the first place: caching, sorting, planning.
Almost every application that grows runs into this at some point. It's easy to miss, because nothing in your code looks wrong.
2 TB, and the detour they didn't take
A few years along, the database is roughly 2 TB. It's on the largest instance AWS sold at the time, with memory and disk both pinned. pgbouncer bought them room, but a single box only goes so far.
Plenty of teams would go database shopping at this point, and Instagram did look around. Mongo, Cassandra, the usual list.
They stayed. What they concluded was that Postgres wasn't the thing failing. The real limit was keeping all of it on one machine, and a different database doesn't change that. Splitting the data does.
sharding: the question you have to answer first
Sharding gets tossed around like it's a switch you flip. Reading through this, it's really a single decision made early, and every query you write afterwards is shaped by it.
You have to choose the column you split on.
They chose user id.
the trade-off you don't get out of
That choice makes some queries lovely and others miserable.
Pulling up one account? Run the user id through a hash, that gives you a shard number, and you go straight there. One shard, one lookup. No slower than if the data had never been divided.
Building a feed out of the 300 accounts someone follows? Those accounts are scattered across all of them. The query now has to visit every shard and stitch the pieces back together. Add more shards and it only gets worse.
No column avoids this. Something always ends up on the slow side of the split.
Instagram's move was to take the question away from the database. Feeds get assembled ahead of time and served from a cache.
the trick: logical shards, not physical ones
I hadn't seen anyone lay this out properly before.
My picture of sharding was a row of servers with the data divided between them. Shard 5 means box 5.
Instagram slid a layer in between. They defined thousands of shards up front, with no connection to how many machines they owned. Each one is a Postgres schema, a named container holding a set of tables. The table definitions are identical everywhere. What differs is which rows sit inside.
So a machine isn't "a shard." It's home to a group of schemas.
When a box starts filling up, you bring in another one and shift some of those schemas across.
Shard 5 was on N1 yesterday and is on N3 today. No table was altered. No application code was touched. The only thing rewritten was the lookup that says which machine holds which shard.
Keeping how the data is divided apart from the hardware it sits on is what makes this survivable. You grow by moving things instead of rebuilding them.
then the ids broke
New problem. Auto-increment.
Every shard runs its own counter, alone. Shard 3 issues 1001. Shard 7 issues 1001 as well. Two unrelated photos now answer to the same id, and nothing downstream can separate them.
They worked through a few options.
UUIDs fixed uniqueness on the spot, with nothing to coordinate. The costs were size and lost ordering. Asking for recent photos stops being cheap, because sorting on the id tells you nothing and you fall back to a timestamp instead.
Twitter's Snowflake was built for exactly this job. But it's one more service to deploy and look after, and they wanted fewer pieces to keep alive, not more.
So they wrote a smaller version that lives inside Postgres.
64 bits, three parts
An Instagram id is one 64-bit integer, carved up like this:
41 bits, milliseconds counted from a chosen starting point
13 bits, naming the logical shard that produced it
10 bits, a counter that resets every millisecond, per shard
Every section pulls its weight.
Time sits in the high bits, so later photos always get larger numbers. Sorting by id sorts by time. You skip the timestamp column and the index that would come with it.
The shard bits mean the id already tells you which shard to go to. Nothing to look up beforehand.
The counter allows 1,024 ids per shard inside any one millisecond, produced right there, without asking anyone.
All of it is a Postgres function running inside whichever shard is doing the insert. Nothing extra to deploy.
The design has aged well. Discord, Slack and plenty of others arrived at nearly the same layout.
where they actually are now
Postgres isn't carrying Instagram on its own today. It runs alongside Cassandra and a large amount of Meta's internal infrastructure. Nobody is serving two billion people from a stock install.
The takeaway I'd actually keep is narrower and more useful than the headline. For years, the database wasn't the thing standing in the way. Connection handling was. Fitting onto one machine was. Producing unique ids was. Every one of those had a fix that left Postgres exactly where it was.