PostgreSQLdatabase scalingread replicasPgBouncerconnection poolingcache lockingAzure Cosmos DBOpenAIsystem designbackend engineering
TL;DR OpenAI scaled ChatGPT to 800 million users without abandoning PostgreSQL. Their secret? Read replicas for global read distribution, cache locking to prevent thundering herd, Azure Cosmos DB for polyglot persistence, and PgBouncer for connection pooling. This deep-dive reveals every strategy with production-grade code you can use today.
Every database engineer's nightmare: your system goes from zero to 800 million users in two years and you're still running PostgreSQL. OpenAI didn't abandon their primary database for a distributed system — they mastered the disciplined art of scaling what they had. This is exactly how they did it, and what you can copy.
50 Read Replicas Cache Locking Azure Cosmos DB PgBouncer PostgreSQL Primary Post ExcerptOpenAI scaled ChatGPT to 800 million users without abandoning PostgreSQL. Their secret? Read replicas for global read distribution, cache locking to prevent thundering herd, Azure Cosmos DB for polyglot persistence, and PgBouncer for connection pooling. This deep-dive reveals every strategy with production-grade code you can use today.
The ProblemPicture this: November 2022. ChatGPT launches. In five days, it hits one million users. In two months, one hundred million — the fastest-growing consumer application in history. The engineering team is staring at a single PostgreSQL primary database and making a choice that most architects would call reckless: they're going to scale it, not replace it. No rush to Cassandra. No panic migration to DynamoDB. No week-one rewrite into microservices with separate data stores. Just disciplined, methodical PostgreSQL scaling under conditions no database team had faced before.
Today, ChatGPT serves over 800 million users. Somewhere between "that's impressive" and "that should be impossible" lies the answer to one of the most instructive database engineering stories of the decade. OpenAI's approach wasn't magic — it was a deliberate sequence of four tactical decisions, each addressing a specific bottleneck without introducing unnecessary complexity. They resisted the urge to over-engineer. They made the right move at the right scale, not the "enterprise" move that looks good in architecture diagrams.
The Myth This Story Busts"You need a distributed database to serve hundreds of millions of users." This is the most expensive misconception in database engineering. Distributed systems introduce consistency complexity, operational overhead, cross-shard query nightmares, and dramatically higher engineering costs. PostgreSQL, properly configured and surrounded by the right infrastructure layer, can scale far beyond what most engineers believe — if you have the discipline to address bottlenecks one at a time rather than preemptively rewriting everything.
Read ReplicasThe first and most impactful decision OpenAI made was straightforward in principle and powerful in execution: route all read traffic away from the primary database. In a chat application like ChatGPT, the overwhelming majority of database operations are reads — users loading their conversation history, browsing past exchanges, reviewing previous outputs. Writes are comparatively rare: new messages, user settings updates, authentication tokens. This asymmetry creates the perfect case for read replicas.
PostgreSQL's streaming replication allows you to maintain synchronized copies of the primary database — replicas that apply the same write-ahead log (WAL) stream and stay within milliseconds of the primary's state. OpenAI deployed 50 of these replicas across multiple global regions, using them as the target for all user-facing read queries. The primary database became dedicated almost exclusively to writes: new messages, conversation creation, settings changes. This single architectural decision gave them roughly 50× the read capacity compared to a single database setup, while keeping the primary database load-free for the writes that absolutely had to land there.
The regional distribution matters as much as the count. A user in Tokyo hitting a read replica hosted in US-East introduces 150–200ms of network latency on every chat history load — perceptible, frustrating, cumulative. By distributing replicas across Azure regions (reflecting their cloud partnership), OpenAI ensured that read traffic was served from geographically proximate infrastructure. A Tokyo user reads from an Asia-Pacific replica; a Berlin user reads from a European one. The primary in US-East processes writes, replicates to all 50 replicas, and stays completely shielded from direct user read traffic.
Real-World AnalogyThink of the primary database as a master author's original manuscript. Only the author (writes) can edit it. But you can print 50 copies of it across the world — one in each major city's library. Readers can visit their local library to access the content without bothering the author or disrupting the editing process. The content might be a few minutes behind the author's latest revisions, but for reading last week's chapter, that's perfectly fine. That's read replicas.
# Django / SQLAlchemy read replica routing — production pattern
# Route reads to replicas, writes to primary
import random
from sqlalchemy import create_engine
from sqlalchemy.orm import sessionmaker
REPLICAS = [
"postgresql://user:pass@replica-us-east.db/chatgpt",
"postgresql://user:pass@replica-eu-west.db/chatgpt",
"postgresql://user:pass@replica-ap-tokyo.db/chatgpt",
# ... 47 more regional replicas
]
PRIMARY = "postgresql://user:pass@primary.db/chatgpt"
class DatabaseRouter:
def __init__(self):
self.primary = create_engine(PRIMARY, pool_size=20, max_overflow=40)
self.replicas = [create_engine(url, pool_size=20) for url in REPLICAS]
def get_read_session(self, region: str = None):
# Route to nearest regional replica if specified
replica = random.choice(self.replicas) # or geo-route
return sessionmaker(bind=replica)()
def get_write_session(self):
# All writes go to primary only
return sessionmaker(bind=self.primary)()
# Usage
router = DatabaseRouter()
# Reading chat history (replica)
with router.get_read_session(region="ap-tokyo") as session:
history = session.query(Conversation).filter_by(user_id=user_id).all()
# Writing new message (primary)
with router.get_write_session() as session:
session.add(Message(content=new_message))
session.commit()
⚡ Pro Tips — Read Replicas
pg_stat_replication view to monitor lag in production.Read replicas solve the steady-state read problem beautifully. But they introduce a failure mode that catches most engineers off-guard: cache invalidation stampedes, better known as the thundering herd problem. Here's how it plays out. You cache a user's conversation list in Redis with a 10-minute TTL. Everything works great. Then the TTL expires for 10,000 users simultaneously (say, all users who signed up during a viral growth spike at the same hour). Every one of those cache misses causes the application to reach past the read replicas and directly query the primary database to reload the data. In milliseconds, you've fired 10,000 simultaneous queries at the one database you've been so carefully protecting. This is how companies take down their primary database when they think they're being careful.
OpenAI solved this with cache locking: a distributed lock mechanism that ensures only one thread or process reloads a given cache key from the database when it expires. All other concurrent requests that hit the same expired cache key see the lock, recognize that a reload is in progress, and either wait for the result or return the stale cached value temporarily. The lock holder queries the database, populates the cache, releases the lock, and all waiting requests get the freshly cached result. Net database queries for a mass expiration event: one per unique cache key, not one per user.
Here's the thing most tutorials miss: the locking mechanism itself must be implemented carefully. A naive implementation using Redis SETNX can deadlock if the lock holder crashes before releasing the lock. Production implementations use lock TTLs (the lock expires after a maximum time regardless of whether it was explicitly released) and, for resilience, a fallback strategy that serves the stale cached value when the lock can't be acquired within a timeout threshold. Stale data for 50ms is infinitely better than no data for 30 seconds while your database recovers from a thundering herd.
Imagine a restaurant that posts a "Today's Menu" board. When the menu runs out (cache expiry), 500 waiting customers all simultaneously rush the kitchen asking "What's on today?" The kitchen collapses. Cache locking is like appointing one customer as the spokesperson — they go to the kitchen, everyone else waits in line, the spokesperson returns with the new menu and announces it to the whole queue. One kitchen trip. Five hundred customers served.
import redis
import json
import time
r = redis.Redis(host="redis-primary", decode_responses=True)
def get_with_lock(cache_key: str, db_fetch_fn, ttl: int = 600, lock_timeout: int = 5):
"""
Get from cache with distributed lock on cache miss.
Prevents thundering herd: only ONE process reloads from DB.
"""
# Happy path: cache hit
cached = r.get(cache_key)
if cached:
return json.loads(cached)
lock_key = f"lock:{cache_key}"
lock_acquired = r.set(lock_key, "1", nx=True, ex=lock_timeout)
if lock_acquired:
try:
# We hold the lock — reload from database
data = db_fetch_fn()
r.setex(cache_key, ttl, json.dumps(data))
return data
finally:
r.delete(lock_key) # always release
else:
# Another process is reloading — wait and retry
deadline = time.time() + lock_timeout
while time.time() < deadline:
time.sleep(0.05) # 50ms polling
cached = r.get(cache_key)
if cached:
return json.loads(cached)
# Fallback: fetch from DB directly (lock holder may have crashed)
return db_fetch_fn()
# Usage
def get_user_conversations(user_id: str):
cache_key = f"conversations:{user_id}"
return get_with_lock(
cache_key=cache_key,
db_fetch_fn=lambda: db.query_conversations(user_id),
ttl=300 # 5 minutes
)
⚡ Pro Tips — Cache Locking
conversations:{user_id}, conversation:{conv_id}, user_settings:{user_id} separately.ttl = 600 + random.randint(-60, 60). This spreads cache misses over time and prevents synchronized expiration waves.Here's a subtle but critical point that separates good database architecture from great database architecture: the goal isn't to force every feature into the same database. It's to use the right database for each data access pattern. OpenAI's team was disciplined enough to recognize when a new feature's data model genuinely required different capabilities than PostgreSQL's relational model could provide efficiently — and to use a different database for that specific use case rather than retrofitting their existing schema.
The clearest examples are ChatGPT's newer features: group chats and image storage. Group chat data has a fundamentally different access pattern than two-party conversations. It involves multiple simultaneous participants, real-time message fan-out, presence tracking, and participant roster management — operations that relational databases handle awkwardly and document databases handle naturally. Image storage involves binary data at potentially large sizes (uploaded screenshots, generated images) with metadata — a pattern better suited for blob storage with document-style metadata, not relational rows with byte arrays. For both, Azure Cosmos DB provided a natural fit: a multi-model NoSQL database with native support for document storage, flexible sharding (partitioning), and horizontal scalability built in from the start.
The architectural principle OpenAI was applying is called polyglot persistence: using multiple, specialized databases within one application, each handling the data it's best suited for. The critical discipline is keeping the boundary clean. PostgreSQL owns conversational metadata, user accounts, billing, and authentication. Cosmos DB owns group chat messages and media metadata. The application layer knows which store to query for which resource type. This isn't microservices for the sake of complexity — it's choosing the right tool for each job, implemented cleanly at the service layer.
Real-World AnalogyA law firm doesn't keep client contracts, physical evidence samples, and audio deposition recordings all in the same filing cabinet. Contracts go in fireproof legal files, physical evidence in locked storage, recordings in a digital archive. Each storage type is optimized for what it holds. Your application data deserves the same thoughtfulness — not everything belongs in one relational database just because that's where you started.
from azure.cosmos import CosmosClient, PartitionKey
import psycopg2
# Polyglot persistence: two databases, clean routing
class DataLayer:
def __init__(self):
# PostgreSQL: user accounts, billing, conversation metadata
self.pg = psycopg2.connect("host=primary dbname=chatgpt")
# Cosmos DB: group chats, images, flexible schema data
cosmos = CosmosClient("https://chatgpt.documents.azure.com",
credential="your_key")
db = cosmos.get_database_client("chatgpt")
self.group_chats = db.get_container_client("group_chats")
self.media = db.get_container_client("media_metadata")
def get_user(self, user_id: str):
# User accounts → PostgreSQL (ACID, relational)
cursor = self.pg.cursor()
cursor.execute("SELECT * FROM users WHERE id = %s", (user_id,))
return cursor.fetchone()
def get_group_chat(self, chat_id: str):
# Group chats → Cosmos DB (flexible, sharded by chat_id)
return self.group_chats.read_item(item=chat_id, partition_key=chat_id)
def store_image_metadata(self, image_data: dict):
# Image metadata → Cosmos DB (schema-flexible, large blobs)
return self.media.create_item(body=image_data)
def create_conversation(self, user_id: str, title: str):
# Conversation metadata → PostgreSQL (ACID guarantees)
cursor = self.pg.cursor()
cursor.execute(
"INSERT INTO conversations (user_id, title) VALUES (%s, %s) RETURNING id",
(user_id, title)
)
self.pg.commit()
return cursor.fetchone()[0]
⚡ Pro Tips — Polyglot Persistence
chat_id is natural — all messages for a chat are co-located on the same partition, making conversation reads a single-partition operation. A bad partition key (like a timestamp) creates hot partitions under high write load.ALTER TABLE migrations with tools like Alembic or Flyway. Cosmos DB schema changes are schema-less by nature but need application-layer backward compatibility. Document your data models carefully when mixing databases — the implicit schema lives in your application code, not in the database.Here's a problem that surprises most engineers the first time they hit it: PostgreSQL has a hard connection limit, and it's surprisingly low relative to the number of users you might want to serve. Each PostgreSQL connection is a heavyweight OS process — it allocates memory (5–10MB per connection in a typical workload), spawns a backend process, and maintains a dedicated socket. Configuring PostgreSQL to accept 10,000 simultaneous connections would require ~100GB of RAM just for connection overhead, leaving little for actual query execution. In practice, production PostgreSQL instances run with connection limits in the hundreds — often 200–500 for a single node.
But OpenAI's application tier might have thousands of web server processes, each wanting a database connection. And millions of simultaneous active users. The naive solution — give every application server process its own persistent PostgreSQL connection — immediately hits the connection ceiling and crashes the database with "FATAL: sorry, too many clients already." The elegant solution is PgBouncer: a lightweight PostgreSQL connection pooler that sits between the application servers and the database. Application servers connect to PgBouncer (which can handle thousands of simultaneous incoming connections), and PgBouncer maintains a small, fixed pool of actual PostgreSQL connections that it multiplexes across all incoming requests.
PgBouncer operates in three modes with very different semantics. Session mode: an application server gets a PostgreSQL connection for its entire session duration. Transaction mode: a connection is assigned only for the duration of a transaction, then returned to the pool — the most efficient mode for typical web applications. Statement mode: a connection is assigned only for the duration of a single SQL statement. OpenAI almost certainly runs in transaction mode, which allows a pool of 100 PostgreSQL connections to service thousands of simultaneous application requests, since each only holds the connection during an active transaction (typically milliseconds).
Real-World AnalogyThink of PostgreSQL connections like checkout lanes in a supermarket — each one requires a dedicated cashier. You have 50 cashiers (connections). Ten thousand customers want to check out simultaneously. Without pooling: 9,950 customers can't get a lane and crash the store. PgBouncer is the store manager who says: "Most of you are done in 30 seconds. I'll give you a lane the moment one becomes free and free it the moment you're done." Same 50 cashiers, dramatically higher throughput.
# PgBouncer configuration — pgbouncer.ini
[databases]
chatgpt = host=primary-postgres.internal port=5432 dbname=chatgpt
[pgbouncer]
listen_port = 6432
listen_addr = 0.0.0.0
# Pool mode: transaction (most efficient for web apps)
pool_mode = transaction
# Connection pool to PostgreSQL (the actual DB connections)
default_pool_size = 100 # 100 real PostgreSQL connections
max_client_conn = 10000 # up to 10k app connections to PgBouncer
reserve_pool_size = 10 # emergency reserve connections
# Timeouts
server_connect_timeout = 3 # seconds to get a DB connection
client_idle_timeout = 60 # disconnect idle clients after 60s
server_idle_timeout = 300 # reclaim idle DB connections after 5min
# Auth
auth_type = md5
auth_file = /etc/pgbouncer/userlist.txt
# Monitoring
stats_period = 60 # emit stats every 60 seconds
# Run PgBouncer with Docker
docker run -d \
--name pgbouncer \
-p 6432:6432 \
-v $(pwd)/pgbouncer.ini:/etc/pgbouncer/pgbouncer.ini \
-v $(pwd)/userlist.txt:/etc/pgbouncer/userlist.txt \
edoburu/pgbouncer
# Monitor pool stats (connect to PgBouncer admin console)
psql -h 127.0.0.1 -p 6432 -U pgbouncer pgbouncer
# SHOW POOLS; — see connection pool utilization
# SHOW STATS; — see request rates and latency
# SHOW CLIENTS; — see active client connections
⚡ Pro Tips — PgBouncer
SET LOCAL, advisory locks, LISTEN/NOTIFY, and prepared statements. If your application uses any of these, you need to either use session mode (lower efficiency) or handle these cases outside the pool (e.g., use a separate direct connection for advisory locks).default_pool_size based on your PostgreSQL max_connections setting, leaving headroom. If PostgreSQL allows 200 connections, set PgBouncer's pool to 180 — reserve 20 for monitoring, admin queries, and emergency access. Running the pool at 100% capacity means any spike overflows to the wait queue.Let's trace the complete journey of a ChatGPT request through OpenAI's database architecture to see every component working in concert. A user in Singapore opens ChatGPT and asks to see their conversation history. The request hits an application server. Before touching any database, the application checks Redis for the cached conversation list using the cache-lock pattern — if the cache is warm, the response returns in under 10ms without any database involvement.
If the cache is cold (miss), the cache lock is acquired and a read query is dispatched — not to the primary PostgreSQL, but to the nearest Asia-Pacific read replica. The replica, synchronized within milliseconds of the primary, returns the conversation list. The result is cached with a TTL plus jitter. If 500 other Singapore users hit the same miss simultaneously, 499 of them wait for the lock holder, then receive the cached result from that single database query. The primary database saw zero traffic from this entire scenario.
Now the user sends a new message. The write goes through PgBouncer, which pulls a connection from its transaction pool, executes the INSERT against the primary PostgreSQL, commits in milliseconds, and returns the connection to the pool. The WAL stream asynchronously replicates this write to all 50 replicas within milliseconds. The next time any user reads this conversation from any regional replica, the new message is already there. The user then uploads an image — this routes entirely to Cosmos DB's media container, completely bypassing PostgreSQL. PgBouncer, the replicas, and the cache layer are never involved.
Four independent strategies, each solving one specific problem, composing into a system that serves 800 million users without a distributed primary database. The lesson isn't "copy OpenAI's exact stack." The lesson is: understand your bottlenecks before you build solutions for them. Reads outpace writes → replicas. Cache stampedes kill the primary → locking. New data models don't fit relational → polyglot. Too many connections → pool them. One bottleneck at a time. One tool at a time.
Getting Started# Step 1: Set up PostgreSQL with streaming replication
# On primary - enable WAL streaming
# Edit postgresql.conf:
wal_level = replica
max_wal_senders = 10
wal_keep_size = 512MB
# Create replication user
psql -c "CREATE USER replicator REPLICATION LOGIN PASSWORD 'secret'"
# Step 2: Create a read replica
# On replica server:
pg_basebackup -h primary-host -U replicator -D /var/lib/postgresql/data \
-P -Xs -R
# -R flag creates standby.signal and recovery config automatically
pg_ctl start -D /var/lib/postgresql/data
# Step 3: Install and configure PgBouncer
apt install pgbouncer
# Edit /etc/pgbouncer/pgbouncer.ini with your configuration (see above)
systemctl start pgbouncer
systemctl enable pgbouncer
# Step 4: Deploy Redis for caching
docker run -d --name redis -p 6379:6379 redis:7-alpine \
redis-server --maxmemory 2gb --maxmemory-policy allkeys-lru
# Step 5: Verify everything
# Check replication status on primary:
psql -c "SELECT client_addr, state, sent_lsn, write_lag FROM pg_stat_replication;"
# Check PgBouncer pool:
psql -p 6432 -U pgbouncer pgbouncer -c "SHOW POOLS;"
# Check Redis:
redis-cli ping # should return PONG
# Step 6: Application setup — complete database layer
import redis
from sqlalchemy import create_engine
# Connect through PgBouncer (port 6432, not 5432)
PRIMARY = create_engine("postgresql://app:pass@pgbouncer:6432/mydb",
pool_pre_ping=True) # health-check connections
REPLICA = create_engine("postgresql://app:pass@replica:5432/mydb",
pool_pre_ping=True)
cache = redis.Redis(host="redis", decode_responses=True)
# Monitor replication lag in your health check endpoint
def check_replication_lag():
with REPLICA.connect() as conn:
result = conn.execute("SELECT EXTRACT(EPOCH FROM (now() - pg_last_xact_replay_timestamp()))")
lag_seconds = result.scalar()
if lag_seconds > 5:
# Alert: replica is falling behind
alert(f"Replica lag: {lag_seconds}s")
return lag_seconds
FAQ
Simulate read replicas under load, watch cache locking prevent database overload, configure PgBouncer pools, and test your knowledge — all in your browser.
Simulate traffic across the primary and 6 regional replicas. Watch the primary stay near-zero while replicas handle all reads.
Traffic Distribution — 50k req/s 0 req/s routed Routing LogSimulate 20 concurrent requests hitting an expired cache key. Compare what happens with and without locking.
Cache Expiry Event — 20 Concurrent Requests ❌ Without Lock — Thundering Herd ✅ With Cache Lock — Controlled DB QUERIES (no lock) 0 DB QUERIES (with lock) 0 PRIMARY LOAD (no lock) 0% PRIMARY LOAD (lock) 0%Click each data type to see which database handles it and why. Green = PostgreSQL, Purple = Cosmos DB.
Data Routing Map Click any data type above to learn why it lives in that database. Query RoutingAdjust the slider to change incoming request volume. Watch how PgBouncer multiplexes thousands of app connections through a small pool of real PostgreSQL connections.
Connection Pool Monitor Incoming App Requests 105000 500 req/s Pool Size (PG connections) 10200 100 connections App Connections to PgBouncer (capped display: 100) ↓ PgBouncer multiplexes to PostgreSQL pool ↓ Real PostgreSQL Connections UTILIZATION - QUEUED 0 STATUS OKToggle each strategy on/off and see how system health metrics respond under a 100k req/s load.
System Simulator Strategies: Read Replicas Cache Locking PgBouncer Cosmos DB PRIMARY LOAD 0% write only target P99 LATENCY - ms DB ERRORS/s 0 connection refused SYSTEM STATUS IDLE overall health Primary Database Load Cache Hit Rate Connection Pool Utilization System Log8 questions covering all four OpenAI scaling strategies.