Instagram launched in 2010 with three engineers, one Postgres database on a rented EC2 machine, and photos in S3. Eighteen months later Facebook paid a billion dollars for a product with 27 million users sitting on a 2 terabyte Postgres instance. No Kafka, no microservices, no exotic storage engine. Just one relational database doing its job while everyone around them was being told that relational databases do not scale.

The part worth studying is not that Postgres is great. It is the order in which things broke, and the fact that each wall got a fix instead of a migration. Reddit, Notion, Discord, Strava and Heroku all made a similar call at some point. They knew the alternatives existed, they evaluated them, and they stayed.

Why did Instagram stay on Postgres instead of switching to NoSQL?

By 2012 the single database was cracking. Two terabytes was close to the memory limit of the biggest instance Amazon would rent them, and disk IO was saturated. Vertical scaling had run out of road. This is exactly the moment where the standard advice kicks in: time to move to Cassandra, DynamoDB or Mongo, because relational databases cannot scale horizontally.

They took that advice seriously and evaluated the options. The conclusion was that the problems they were hitting were not Postgres problems, they were scale problems. Any database running that workload would hit the same walls. Moving to Cassandra does not make horizontal partitioning disappear, it hides it behind a different abstraction and adds a migration you cannot walk back from. Once your data lives in a new engine, you have burnt the boats.

So they chose the harder and more boring option: make Postgres itself horizontal. Everything below is what that actually required.

What breaks first in Postgres at scale: connections, not data

The first wall had nothing to do with data volume. Postgres uses a process per connection, and each one costs roughly a megabyte of memory before it does anything useful. Application servers make this worse, because every process keeps its own pool.

The arithmetic is unpleasant once you write it down:

# Rough memory cost of idle Postgres connections held by an app tier.
APP_SERVERS = 50
PROCESSES_PER_SERVER = 6
POOL_PER_PROCESS = 5
MB_PER_CONNECTION = 1.3

connections = APP_SERVERS * PROCESSES_PER_SERVER * POOL_PER_PROCESS
puts "connections: #{connections}"
puts "memory just to hold them: #{(connections * MB_PER_CONNECTION / 1024).round(1)} GB"
# connections: 1500
# memory just to hold them: 1.9 GB

Fifteen hundred connections, about two gigabytes of RAM, and the database has not planned a single query yet. That memory should be caching pages and sorting results instead of holding sockets open.

The fix is a connection pooler, and the classic one is PgBouncer. It is a thin proxy that sits between the application and Postgres. Your app thinks it is talking to a database, PgBouncer multiplexes all those client connections onto a small set of real backend connections. In transaction pooling mode a few dozen real connections can serve thousands of clients, because a backend is only held for the duration of a transaction rather than the lifetime of a request.

Transaction pooling is the mode that gives you the big win, and it comes with a rule: no session state. Session-level SET, LISTEN/NOTIFY, session-scoped advisory locks and server-side prepared statements can all break, because the next transaction may land on a different backend. In Rails you either set prepared_statements: false in database.yml or run PgBouncer 1.21 or newer with max_prepared_statements set to a non-zero value, otherwise you will chase confusing "prepared statement already exists" errors.

If you run Postgres at any meaningful traffic and you have nothing in front of it, this is the highest leverage afternoon of work available to you this week. It is also the bottleneck teams discover last, usually during an incident.

When is it actually time to shard a Postgres database?

Later than you think. Instagram did not shard until they were at tens of millions of users and had genuinely exhausted one machine. Before sharding there is a long list of cheaper moves: a pooler, read replicas for read-heavy traffic, moving cold or write-heavy side data (sessions, analytics events, logs) out of the primary, fixing the queries that dominate pg_stat_statements, and partitioning large tables by time inside a single database.

Sharding is different in kind from all of those, because it changes the shape of your application code permanently. Cross-shard joins stop existing. Transactions stop being free. Every query needs to know where it is going. You reach for it when one machine cannot hold the working set, not when a dashboard looks busy.

How to pick a shard key you will not regret

The shard key is the one decision you get to make once. Whatever column you split on shapes every query the application will ever run.

For Instagram the answer was the user. A user's photos, likes and profile data live together on one shard, so "show me my own profile" hashes the user ID, hits exactly one shard, and comes back fast.

The problem is that a social network's defining query is not "show me my photos". It is "show me photos from the 200 people I follow", and those 200 people are scattered across 50 shards. That query becomes 50 partial reads merged and sorted in the application, and it is as slow as the slowest shard in the set. This is the trade every sharded system makes in some form: single-entity reads get fast, cross-entity reads get expensive. There is no clever schema that removes it. Instagram accepted the trade and paid for it in the application layer with caching and precomputed feeds.

Two things are worth writing on a wall before you commit. Resharding a few hundred terabytes means rewriting every row, so treat the key as permanent. And pick the key that matches your most frequent query, not the one that looks most balanced on paper.

How to shard Postgres with logical shards instead of physical servers

Most sharding implementations in 2010 mapped data straight onto machines. Sixteen servers, one sixteenth of the users each, and growth meant reshuffling data. Instagram put a layer between the two. They created several thousand logical shards, each one a Postgres schema holding the same set of tables, and they mapped those schemas onto a much smaller number of physical machines.

-- Every logical shard is a schema with an identical table layout.
CREATE SCHEMA IF NOT EXISTS shard_0042;

CREATE TABLE shard_0042.photos (
  id          bigint PRIMARY KEY,
  user_id     bigint      NOT NULL,
  caption     text,
  created_at  timestamptz NOT NULL DEFAULT now()
);

-- The mapping lives in one small table, not in application constants.
CREATE TABLE shard_map (
  logical_shard int PRIMARY KEY,
  physical_node text NOT NULL,
  moved_at      timestamptz NOT NULL DEFAULT now()
);

The application never knows about machines. It hashes the user ID into a logical shard, looks up which node currently owns that shard, and queries that schema:

require 'digest'

class ShardRouter
  LOGICAL_SHARDS = 4096

  def initialize(shard_map)
    @shard_map = shard_map # { logical_shard => "db-07.internal" }
  end

  def logical_shard_for(user_id)
    # Stable hash: the same user must always land on the same logical shard,
    # so never use Ruby's #hash here, it is randomized per process.
    Digest::MD5.hexdigest(user_id.to_s)[0, 8].to_i(16) % LOGICAL_SHARDS
  end

  def route(user_id)
    shard = logical_shard_for(user_id)
    { node: @shard_map.fetch(shard), schema: format("shard_%04d", shard) }
  end
end

Note what the router does not contain: the number of machines. That number is free to change.

In the beginning all of those thousands of schemas sat on a single physical machine, and the application could not tell. When a node fills up, say it is at 95 percent, you bring up a new machine, clone the whole node with built-in streaming replication, wait for it to catch up, promote it, update the mapping table, and drop on each side the schemas it no longer owns. Shards 1920 to 2047 now answer on node 16. No data was rewritten, no schema changed, no application code was deployed. Only a row in a lookup table changed.

If you hardcode "we have 16 shards" into your application, every growth step is a migration project. If you have 4096 logical shards mapped onto however many nodes you happen to own today, growth is a config change and a replication job.

How to generate unique IDs across shards without a coordinator

Sharding breaks auto increment immediately. If shard_0000 hands out 1, 2, 3 and shard_0001 hands out 1, 2, 3, you now have two different photos with the same ID and an application that cannot tell them apart.

The obvious replacement is a UUID, and Instagram ruled it out. Version 4 UUIDs are 128 bits, twice the width in every index, and they carry no ordering, so "the newest 20 photos" cannot be answered by the primary key and random inserts scatter writes across the B-tree. A dedicated ticket server, the approach Flickr used, gives you ordered IDs at the cost of a single point of failure that stops every insert in the company when it dies. Twitter's Snowflake was the closest fit, but it meant running and monitoring another distributed service next to the database.

So they built Snowflake-style IDs into Postgres itself. A 64 bit integer, split into a 41 bit millisecond timestamp counted from a custom epoch, a 13 bit shard ID, and a 10 bit per-millisecond sequence. That gives about 35 years of timestamps before a signed bigint overflows, up to 8192 shards, and 1024 IDs per millisecond per shard.

Here is the same idea for an orders table, written as a plain function each shard owns:

CREATE SEQUENCE shard_0042.order_id_seq;

CREATE OR REPLACE FUNCTION shard_0042.next_order_id(OUT new_id bigint) AS $$
DECLARE
  epoch_ms  bigint := 1735689600000; -- custom epoch: 2025-01-01, keeps the timestamp delta small
  now_ms    bigint;
  shard_id  int    := 42;            -- hardcoded per schema, never passed in by the app
  seq_value bigint;
BEGIN
  SELECT floor(extract(epoch FROM clock_timestamp()) * 1000)::bigint INTO now_ms;
  SELECT nextval('shard_0042.order_id_seq') % 1024 INTO seq_value;

  -- timestamp in the high bits, then shard, then sequence in the low bits
  new_id := ((now_ms - epoch_ms) << 23) | (shard_id << 10) | seq_value;
END;
$$ LANGUAGE plpgsql;

ALTER TABLE shard_0042.orders
  ALTER COLUMN id SET DEFAULT shard_0042.next_order_id();

No external service, no coordination between shards, no single point of failure, as long as a single shard stays under 1024 inserts per millisecond. Two properties fall out of it for free. The ID tells you which shard produced it, so a bare ID is enough to route a lookup. And because the timestamp sits in the high bits, IDs sort by creation time, which means "latest first" is ORDER BY id DESC LIMIT 20 with no separate timestamp index and no extra sort.

Discord, Slack and a long list of others now use the same layout. If you are building something that might shard one day, adopting this on day one costs you an afternoon. Retrofitting it onto a billion existing rows costs you a quarter.

Which Postgres features do most teams never use?

Three of them did a lot of work here, and they are already installed on your server.

Partial indexes index only the rows matching a condition. If you keep a billion photos but only ever query the recent ones, there is no reason to carry a billion index entries.

-- Only recent, still-visible rows get an index entry.
-- The date is frozen at CREATE INDEX time, so rebuild it as it ages.
CREATE INDEX index_photos_on_created_at_recent
  ON photos (created_at DESC)
  WHERE deleted_at IS NULL
    AND created_at > '2026-01-01';

The index stays a fraction of the size, so more of it lives in memory and scans touch fewer pages. The catch is that the planner only uses it when it can prove your WHERE clause implies the index predicate, so the condition has to be one your queries actually write.

Functional indexes index the result of an expression instead of the raw column. Instagram had 64 character random tokens and no reason to duplicate all 64 characters inside an index when the first eight are already unique enough.

CREATE INDEX index_api_tokens_on_prefix
  ON api_tokens (substr(token, 1, 8));

-- The query has to use the same expression for the index to be used.
SELECT id, user_id
FROM api_tokens
WHERE substr(token, 1, 8) = 'a91f7c02'
  AND token = 'a91f7c02c4e1...';

The prefix narrows the search to a handful of rows, the full comparison confirms the match, and the index is roughly three times smaller. Irrelevant at ten thousand rows, decisive at ten billion.

Logical replication is the most underrated of the three. It streams every insert, update and delete on selected tables to downstream consumers as they happen, as long as every table has a replica identity, which a primary key already gives you.

-- On the primary
CREATE PUBLICATION photo_stream FOR TABLE photos, photo_tags;

-- On the consumer
CREATE SUBSCRIPTION search_index_sync
  CONNECTION 'host=10.0.0.11 dbname=gallery user=replicator'
  PUBLICATION photo_stream;

A search index, a cache invalidator and an analytics warehouse can each read that stream and stay current without the application publishing events by hand. One thing to check before you plan around it: the primary needs wal_level = logical, which is not the default and needs a restart, so it is worth setting early rather than during an incident. Plenty of teams build exactly this with dual writes and a message broker, carry the dual write consistency bugs that come with it, and never notice the database was already offering the same stream. Change data capture tools like Debezium sit on top of the same mechanism.

Does Instagram still run on Postgres in 2026?

Partly, and it is worth being honest about that. Everything above is the story of roughly 2010 to 2015. After the acquisition, Meta folded Instagram into its own infrastructure, and the social graph now lives in TAO, a distributed graph store built for Meta's scale, while Postgres kept the user data it always held.

That is still the strongest version of the argument. Postgres carried the product from three engineers to a nine figure user count, and it gave up the graph only when Instagram had a genuinely different problem and a company with thousands of infrastructure engineers to solve it with. Almost nobody reading this is in that position. There is a complementary case worth reading here too: OpenAI serves over 800 million weekly ChatGPT users on a single primary Postgres, with no sharding at all. Same database, opposite architecture, both working.

A checklist before you reach for a new database

  1. Put a connection pooler in front of Postgres and switch it to transaction mode.
  2. Read pg_stat_statements and fix the five queries that dominate total time.
  3. Move read-heavy traffic to replicas and non-critical writes out of the primary.
  4. Add partial and functional indexes where full indexes are wasting memory.
  5. Partition your biggest tables by time inside the one database you already have.
  6. Only then design a shard key, and design it around your most frequent query.
  7. Create thousands of logical shards, map them to nodes in a table, and keep machine counts out of the code.
  8. Adopt Snowflake-style IDs before you have a billion rows, not after.

FAQ

How many connections can Postgres handle?

Postgres itself will let you set max_connections to several thousand, but the practical ceiling is much lower because each connection is an OS process with its own memory. Most deployments start suffering well before a thousand. The number that matters is how many concurrent active transactions your cores can serve, usually somewhere between two and four times the CPU count, and a pooler is how you serve thousands of clients with that many backends.

Is PgBouncer still the right choice, or should I use pgcat or Supavisor?

PgBouncer is still the safe default, it is small, boring and everywhere. The newer poolers such as pgcat and Supavisor add query routing, read/write splitting and multi-threading, which matter if you want the pooler to also handle load balancing across replicas or shards. Start with PgBouncer, move only when you have a specific feature you need.

Can Postgres shard itself without application changes?

Partially. Declarative partitioning splits a table across partitions within one database, which helps with data volume but not with running out of machine. Extensions like Citus do distribute tables across nodes and handle routing for you, which is closer to what Instagram built by hand. Both still need you to choose a distribution key, and that decision carries the same weight.

Do I need Snowflake IDs if I already use UUIDv7?

Not really. UUIDv7 is time-ordered, which removes the write scatter and the sorting problem that made UUIDv4 painful. What it does not give you is a shard ID encoded in the value or a 64 bit width. If you shard, the embedded shard number is genuinely useful for routing. If you do not, UUIDv7 is the simpler choice.

How many logical shards should I create?

More than you think you need, and a power of two, because the cost of a spare logical shard is one empty schema while the cost of running out is a full reshard. A few thousand is a common answer, and 13 bits in an ID layout caps you at 8192, which is a reasonable place to stop.

The takeaway from all of this is not "use Postgres because Instagram used Postgres". It is that every team has a limited budget for complexity, and every new piece of infrastructure spends some of it on quirks to learn, tooling to build and people to hire. Spend that budget when you have a genuinely new problem, like vector search at scale or graph traversal across a planet. Having a lot of users is not a new problem. It is the most documented problem in our industry, and the database you already run has been solving it for twenty years.

Happy sharding!