Scaling Your App, One Break at a Time
A practical, stage-by-stage guide to scaling a web app without over-engineering — measure first, then add exactly one thing when something actually breaks: vertical scaling, a load balancer, stateless servers, connection pooling, indexes and read replicas, caching, queues and background jobs, and finally sharding. Node.js/TypeScript, PostgreSQL, Redis and Nginx examples, each with its symptom, diagnosis, cost and 'done when'.
Scaling Your App, One Break at a Time
A practical implementation guide to load balancing, stateless servers, connection pooling, read replicas, caching, queues, and sharding
How to Use This Guide
This guide follows one rule: don't add a piece of infrastructure until something is actually breaking, and then add exactly one thing. Each stage below is structured the same way:
- Symptom: what you'll see when you need this stage
- Diagnosis: how to confirm it's really this problem
- Implementation: the concrete steps and code
- Cost: what you're paying (money, complexity, correctness)
- Done when: how you know the stage is complete
The code examples use Node.js / TypeScript, PostgreSQL, Redis, and Nginx because they're common and the concepts transfer directly. If you use a different stack (Python/Django, Go, Rails, MySQL), the patterns are identical; only the library names change.
The golden rule: Every box you add makes your system faster or more resilient and more complex. Most apps should stop at Stage 1 or 2 for a long time. Skipping ahead doesn't make you a better engineer; knowing when to move forward does.
Table of Contents
- Stage 0: Measure Before You Scale
- Stage 1: The Starting System (One Server, One Database)
- Stage 2: Vertical Scaling
- Stage 3: Horizontal Scaling + Load Balancer
- Stage 4: Stateless Servers (Shared Session Store)
- Stage 5: Connection Pooling
- Stage 6: Indexes First, Then Read Replicas
- Stage 7: Caching
- Stage 8: Queues & Background Jobs
- Stage 9: Sharding
- Appendix A: Symptom → Stage Diagnosis Table
- Appendix B: The Redis Question
- Appendix C: Final Architecture & Checklist
Stage 0: Measure Before You Scale
You cannot tell which stage you need without data. Before touching architecture, put these in place. They're cheap and they'll tell you where the bottleneck is instead of letting you guess.
What to measure
| Metric | Why it matters | Where to get it |
|---|---|---|
| Request latency (p50, p95, p99) | Averages hide pain; p99 shows your worst users | APM tool, Nginx logs, app middleware |
| Requests per second | Tells you how close you are to capacity | Load balancer / platform metrics |
| Server CPU & memory | Is the app server the bottleneck? | Host metrics |
| DB CPU, memory, disk I/O | Is the database the bottleneck? | Managed DB dashboard |
DB active connections vs max_connections | Connection exhaustion (Stage 5) | pg_stat_activity |
| Slowest queries | Missing indexes (Stage 6) | pg_stat_statements |
| Error rate by type | Timeouts vs connection errors vs bugs | Logs / error tracker |
A minimal timing middleware (Express)
import express from "express";
const app = express();
app.use((req, res, next) => {
const start = process.hrtime.bigint();
res.on("finish", () => {
const ms = Number(process.hrtime.bigint() - start) / 1e6;
console.log(
JSON.stringify({
method: req.method,
path: req.route?.path ?? req.path,
status: res.statusCode,
ms: Math.round(ms),
})
);
});
next();
});Enable query statistics in Postgres
-- Requires shared_preload_libraries = 'pg_stat_statements' (most managed DBs have it on)
CREATE EXTENSION IF NOT EXISTS pg_stat_statements;
-- Top 10 queries by total time spent
SELECT query, calls, round(total_exec_time) AS total_ms, round(mean_exec_time, 2) AS mean_ms
FROM pg_stat_statements
ORDER BY total_exec_time DESC
LIMIT 10;Load test before your users do
Don't wait for the concert tickets to go on sale. Simulate it with a tool like k6:
// load-test.js — run with: k6 run load-test.js
import http from "k6/http";
import { sleep } from "k6";
export const options = {
stages: [
{ duration: "1m", target: 100 }, // ramp up
{ duration: "3m", target: 1000 }, // spike
{ duration: "1m", target: 0 }, // ramp down
],
};
export default function () {
http.get("https://staging.yourapp.com/api/events");
sleep(1);
}Done when: you can answer "what is my slowest endpoint, my slowest query, and my peak requests per second?" from a dashboard instead of a guess.
Stage 1: The Starting System
[ Users ] ──► [ App Server ] ──► [ Database ]This handles thousands of users. Most apps live here forever, and that's a success, not a failure.
Make Stage 1 as good as it can be
Before scaling anything, make sure the basics are solid, because every later stage is easier if these are true:
- Configuration comes from environment variables, not hard-coded values. You'll be pointing at different databases, Redis instances, and replicas later.
- Database access goes through one module (e.g.
db.ts), not scatterednew Client()calls. This single decision makes Stages 5, 6, and 9 dramatically easier. - No local file storage for user uploads. Use object storage (S3, R2, GCS) from day one. Local files are state that will bite you in Stage 4.
- A health check endpoint exists. You'll need it in Stage 3.
// db.ts — the single place the app talks to the database
import { Pool } from "pg";
export const pool = new Pool({ connectionString: process.env.DATABASE_URL });
export async function query<T = any>(text: string, params?: unknown[]) {
const result = await pool.query(text, params);
return result.rows as T[];
}// health.ts
app.get("/health", async (_req, res) => {
try {
await pool.query("SELECT 1");
res.status(200).json({ status: "ok" });
} catch {
res.status(503).json({ status: "db_unreachable" });
}
});Done when: your app runs, config is externalized, DB access is centralized, and uploads don't touch the local disk.
Stage 2: Vertical Scaling
Symptom
Server CPU or memory is pegged during peak traffic. Requests queue up and time out. Your code has no obvious bugs or bad queries.
Diagnosis
Server CPU sits at 80–100% during slow periods while the database is comfortable. Latency climbs in step with traffic.
Implementation
Buy a bigger machine. That's it.
- In your hosting dashboard, move to the next instance size up (more vCPUs, more RAM).
- If you're on Node.js, remember that one Node process uses one CPU core for your JavaScript. A 4-core machine running a single process wastes 3 cores. Use all of them:
// cluster.ts — run one worker process per CPU core
import cluster from "node:cluster";
import os from "node:os";
if (cluster.isPrimary) {
const cores = os.availableParallelism();
for (let i = 0; i < cores; i++) cluster.fork();
cluster.on("exit", (worker) => {
console.log(`Worker ${worker.process.pid} died, restarting`);
cluster.fork();
});
} else {
await import("./server.js");
}Or with PM2: pm2 start dist/server.js -i max.
Notice something: the moment you run multiple processes, even on one machine, you already have multiple copies of your app. Everything in Stage 4 about state applies here too.
Cost
Money, and a hard ceiling. There's a biggest machine you can rent, and pricing gets steep long before you get there. Also, one machine is still one point of failure.
Done when: peak CPU sits comfortably below ~70% with headroom for spikes. If you've climbed several sizes and you're still pegged, or the next size is painfully expensive, move to Stage 3.
Stage 3: Horizontal Scaling + Load Balancer
┌──► [ App Server 1 ] ──┐
[ Users ] ──► [ LB ] ├──► [ App Server 2 ] ──┼──► [ Database ]
└──► [ App Server 3 ] ──┘Symptom
You've outgrown vertical scaling, or you need to survive a single server dying without an outage.
Key idea
These are identical copies of the same app, not different parts of it. You deploy once and run N instances. Any instance can handle any request.
Implementation
Option A: Managed platforms (easiest)
On Vercel, Railway, Render, Fly.io, Cloud Run, ECS, or Kubernetes, you set an instance count or autoscaling rule and the platform provides the load balancer. For example:
- Railway / Render: set replicas / instance count in service settings
- Fly.io:
fly scale count 3 - Kubernetes:
replicas: 3in your Deployment, plus a Service/Ingress - Vercel: serverless functions scale automatically; you already have this
If you're on one of these, skip to the "Before you move on" checklist below.
Option B: Self-hosted with Nginx + Docker Compose
# docker-compose.yml
services:
app:
build: .
environment:
- DATABASE_URL=${DATABASE_URL}
- REDIS_URL=${REDIS_URL}
expose:
- "3000"
# run with: docker compose up --scale app=3
nginx:
image: nginx:stable
ports:
- "80:80"
volumes:
- ./nginx.conf:/etc/nginx/conf.d/default.conf:ro
depends_on:
- app# nginx.conf
upstream app_servers {
least_conn; # send to the server with fewest active connections
server app:3000 max_fails=3 fail_timeout=30s;
# With explicit hosts instead of Docker DNS:
# server 10.0.0.11:3000 max_fails=3 fail_timeout=30s;
# server 10.0.0.12:3000 max_fails=3 fail_timeout=30s;
# server 10.0.0.13:3000 max_fails=3 fail_timeout=30s;
keepalive 32;
}
server {
listen 80;
location / {
proxy_pass http://app_servers;
proxy_http_version 1.1;
proxy_set_header Connection "";
proxy_set_header Host $host;
proxy_set_header X-Real-IP $remote_addr;
proxy_set_header X-Forwarded-For $proxy_add_x_forwarded_for;
proxy_set_header X-Forwarded-Proto $scheme;
# If one server fails, retry the request on another
proxy_next_upstream error timeout http_502 http_503;
proxy_connect_timeout 5s;
proxy_read_timeout 30s;
}
}Note: with Docker's internal DNS, Nginx resolves
appto all container IPs when it starts. If you scale up later, reload Nginx (docker compose exec nginx nginx -s reload) so it picks up the new containers.
Load balancing strategies
| Strategy | How it works | Use when |
|---|---|---|
| Round robin (default) | Each server in turn | Requests are roughly equal in cost |
least_conn | Server with fewest active connections | Some requests are slow, others fast |
ip_hash | Same client IP → same server | Avoid; this is "sticky sessions," a crutch that hides Stage 4 problems |
Tell your app it's behind a proxy
// Without this, req.ip is the load balancer's IP and secure cookies may break
app.set("trust proxy", 1);Graceful shutdown
When the platform removes an instance (deploys, scale-down), it should finish in-flight requests instead of dropping them:
const server = app.listen(3000);
process.on("SIGTERM", () => {
console.log("SIGTERM received, draining connections");
server.close(async () => {
await pool.end();
process.exit(0);
});
setTimeout(() => process.exit(1), 25_000).unref(); // hard stop
});Cost
- The load balancer is now a single point of failure. Managed load balancers (AWS ALB, GCP LB, Cloudflare, platform-provided) handle redundancy for you. Self-hosted, you'd need a second Nginx with a floating IP (e.g. keepalived), which is real operational work and a good reason to use a managed one.
- Your servers must be interchangeable. Right now they probably aren't. That's Stage 4, and you should do it at the same time as this stage.
Done when: you can kill any single app instance during a load test and users see no errors.
Stage 4: Stateless Servers
Symptom
After adding servers, users get randomly logged out, carts empty themselves, or flash messages disappear. With 3 servers, roughly 2 out of 3 requests land on a server that doesn't "know" the user.
Diagnosis
The problem vanishes when you scale back to 1 instance. Search your code for anything stored in process memory that must survive between requests.
Implementation
Step 1: Audit your code for hidden state
Sessions are the obvious one, but check for all of these:
| Hidden state | Where it lives now | Where it should live |
|---|---|---|
| Login sessions | In-memory session store | Redis (or stateless tokens) |
| Uploaded files | Local disk | Object storage (S3/R2/GCS) |
In-memory caches (const cache = new Map()) | Process memory | Redis (Stage 7), or accept per-instance caching for non-critical data |
| Rate limit counters | Process memory | Redis |
| WebSocket connections | Tied to one server | Redis pub/sub adapter to broadcast across servers |
Cron jobs / setInterval tasks | Run on every instance | A single scheduler or a queue with repeatable jobs (Stage 8) |
That last row is sneaky: with 3 instances, your "send daily digest" cron sends three emails.
Step 2: Move sessions to Redis
pnpm add express-session connect-redis redis
// session.ts
import session from "express-session";
import { RedisStore } from "connect-redis";
import { createClient } from "redis";
export const redis = createClient({ url: process.env.REDIS_URL });
redis.on("error", (err) => console.error("Redis error", err));
await redis.connect();
export const sessionMiddleware = session({
store: new RedisStore({ client: redis, prefix: "sess:" }),
secret: process.env.SESSION_SECRET!, // same secret on every instance!
resave: false,
saveUninitialized: false,
name: "sid",
cookie: {
httpOnly: true,
secure: process.env.NODE_ENV === "production",
sameSite: "lax",
maxAge: 1000 * 60 * 60 * 24 * 7, // 7 days
},
});// server.ts
app.set("trust proxy", 1);
app.use(sessionMiddleware);
app.post("/login", async (req, res) => {
const user = await verifyCredentials(req.body.email, req.body.password);
if (!user) return res.status(401).json({ error: "Invalid credentials" });
req.session.regenerate((err) => {
// prevent session fixation
if (err) return res.status(500).end();
req.session.userId = user.id; // written to Redis, visible to all servers
res.json({ ok: true });
});
});⚠️ Every instance must share the same
SESSION_SECRET. If each instance generates its own, sessions signed by server 1 will be rejected by server 2, and you'll get the exact same random-logout bug.
Alternative: Stateless tokens (JWT)
Instead of storing sessions anywhere, the server signs a token containing the user ID and any server can verify it with a shared key.
| Redis sessions | Signed tokens (JWT) | |
|---|---|---|
| Extra lookup per request | Yes (fast) | No |
| Instant logout / revoke | Easy: delete the key | Hard: needs a denylist (which is... a shared store) |
| Infrastructure | Redis required | None |
| Good for | Web apps with cookies | APIs, mobile clients, service-to-service |
Many teams use short-lived access tokens plus a refresh token stored server-side, getting most of both benefits. If you're using an auth provider (Clerk, Auth0, Supabase Auth, NextAuth with a DB adapter), check its docs: it may already be stateless-compatible.
Cost
- One extra network hop to Redis on every authenticated request (usually well under a millisecond on the same network).
- Redis becomes critical infrastructure: if it's down, nobody is logged in. Use a managed Redis with replication/failover for production.
Done when: you can run 3+ instances without sticky sessions, log in, click around rapidly, and never get logged out; and cron-style jobs run exactly once.
Stage 5: Connection Pooling
Symptom
Not slowness, errors, and they appear exactly when traffic spikes:
FATAL: sorry, too many clients already
FATAL: remaining connection slots are reserved for non-replication superuser connections
Error: timeout exceeded when trying to connectServers look healthy. The database looks healthy. Things fail anyway.
Diagnosis
-- How many connections are open right now vs the limit?
SELECT count(*) AS open_connections,
current_setting('max_connections') AS max_connections
FROM pg_stat_activity;
-- Who's holding them?
SELECT application_name, client_addr, state, count(*)
FROM pg_stat_activity
GROUP BY 1, 2, 3
ORDER BY 4 DESC;If open_connections is at or near max_connections during spikes, this is your stage.
The math you must do
total connections = (number of app instances) × (pool size per instance)This must stay below your database's max_connections, with headroom for migrations, admin tools, and replicas. Example: DB allows 100, reserve 20 → 80 available. With 4 instances, each pool gets max: 20. Scale to 8 instances and you must drop to max: 10, or add a pooler.
Implementation
Step 1: Configure your app-level pool correctly
// db.ts
import { Pool } from "pg";
export const pool = new Pool({
connectionString: process.env.DATABASE_URL,
max: Number(process.env.DB_POOL_MAX ?? 10), // per instance!
idleTimeoutMillis: 30_000, // close idle connections
connectionTimeoutMillis: 5_000, // fail fast instead of hanging
});
// Always release connections. pool.query() does this for you.
// If you check out a client manually, use try/finally:
export async function withClient<T>(
fn: (c: import("pg").PoolClient) => Promise<T>
) {
const client = await pool.connect();
try {
return await fn(client);
} finally {
client.release(); // forgetting this leaks connections until the pool is empty
}
}For ORMs: Prisma uses connection_limit in the connection string; Drizzle/Knex/TypeORM/Sequelize accept pool options similar to the above.
Step 2: Add an external pooler (PgBouncer) when instances multiply
An app-level pool limits connections per instance. A pooler limits connections globally: hundreds of app connections share a small fixed set of real database connections.
[ App 1 ] ─┐
[ App 2 ] ─┼─► [ PgBouncer: 20 real connections ] ──► [ Postgres ]
[ App N ] ─┘Check first: most managed databases (Supabase, Neon, AWS RDS Proxy, Azure, DigitalOcean, Crunchy) include a pooler you can switch on. You just use a different connection string. Prefer that over running PgBouncer yourself.
If self-hosting:
; pgbouncer.ini
[databases]
appdb = host=postgres-primary port=5432 dbname=appdb
[pgbouncer]
listen_addr = 0.0.0.0
listen_port = 6432
auth_type = scram-sha-256
auth_file = /etc/pgbouncer/userlist.txt
pool_mode = transaction ; connection returned after each transaction
max_client_conn = 1000 ; app-side connections allowed
default_pool_size = 20 ; real DB connections per db/user pair⚠️ Transaction mode gotchas: because consecutive transactions may run on different real connections, session-level features break:
SETwithoutLOCAL, advisory locks held across transactions,LISTEN/NOTIFY, and (on older PgBouncer versions) named prepared statements. Prisma needs?pgbouncer=truein the URL. Run migrations against a direct connection, not the pooler.
Step 3: Serverless-specific rules
On Vercel, Netlify, or AWS Lambda, each cold-started function instance opens its own connections, so a spike creates hundreds of connections in seconds.
- Always connect through a pooler URL (the provider's pooled connection string).
- Create the client outside the handler so warm invocations reuse it.
- Keep per-instance pool size tiny (
max: 1or2). - Consider HTTP-based database drivers (e.g. Neon's serverless driver) designed for this.
// Reused across warm invocations of the same function instance
const pool =
globalThis.__pool ??
new Pool({
connectionString: process.env.DATABASE_POOLED_URL,
max: 1,
});
globalThis.__pool = pool;Cost
Low. One more component (unless managed), plus the transaction-mode caveats above. This is one of the cheapest, highest-value stages.
Done when: a load test spike produces no connection errors, and pg_stat_activity stays flat regardless of instance count.
Stage 6: Indexes First, Then Read Replicas
Symptom
Database CPU is high and queries are slow even though connections are under control.
Step 1: The free fix (always do this first)
A missing index looks exactly like a capacity problem. Check before spending money.
-- See how a query actually runs
EXPLAIN ANALYZE
SELECT * FROM users WHERE email = 'bob@example.com';If you see Seq Scan on users on a large table, the database is reading every row. Add an index:
-- CONCURRENTLY avoids locking the table on a live app (can't run inside a transaction)
CREATE INDEX CONCURRENTLY idx_users_email ON users (email);
-- Composite index for common filter + sort patterns
CREATE INDEX CONCURRENTLY idx_orders_user_created
ON orders (user_id, created_at DESC);Index checklist:
- Columns in frequent
WHEREclauses - Foreign keys used in
JOINs (Postgres does not auto-index these) - Columns used in
ORDER BYwithLIMIT(feeds, lists) - Don't index everything: each index slows down writes and uses disk
Also look for N+1 queries (a loop that runs one query per item). Fetching 50 posts then running 50 separate author queries is 51 round trips that should be 1 or 2.
Re-run the pg_stat_statements query from Stage 0. If the top queries are now fast and DB CPU is still maxed, continue.
Step 2: Read replicas
writes ──► [ Primary DB ]
[ App Servers ] ──┤ │ replication
reads ──► [ Replica 1 ] [ Replica 2 ]Most apps read far more than they write, so offloading reads removes most of the load.
Create the replicas
On managed databases this is a button or one CLI command (RDS, Cloud SQL, Supabase, Neon, DigitalOcean, Azure). You get a separate read-only connection string (or one per replica).
Route queries in your data layer
This is where centralizing DB access in Stage 1 pays off.
// db.ts
import { Pool } from "pg";
const primary = new Pool({
connectionString: process.env.DATABASE_URL,
max: 10,
});
const replicas = (process.env.DATABASE_REPLICA_URLS ?? "")
.split(",")
.filter(Boolean)
.map((url) => new Pool({ connectionString: url, max: 10 }));
let next = 0;
function pickReplica(): Pool {
if (replicas.length === 0) return primary; // graceful fallback
const r = replicas[next % replicas.length];
next++;
return r;
}
export const db = {
/** INSERT/UPDATE/DELETE, transactions, and anything that must be fresh */
write: primary,
/** Reads that can tolerate being a fraction of a second stale */
read: () => pickReplica(),
};// Usage
await db.write.query("INSERT INTO posts (user_id, body) VALUES ($1, $2)", [
userId,
body,
]);
const feed = await db
.read()
.query("SELECT * FROM posts ORDER BY created_at DESC LIMIT 50");Many ORMs support this natively (Prisma's read replicas extension, Rails connects_to, Django database routers, Sequelize replication option). Prefer the built-in feature when it exists.
Handle replication lag: read-your-own-writes
The classic bug: a user posts, the feed loads from a replica that hasn't received the post yet, and their own post is missing. Fix it by routing that user's reads to the primary for a short window after they write:
const STICKY_MS = 5_000;
export function markWrite(req: Express.Request) {
req.session.lastWriteAt = Date.now();
}
export function readerFor(req: Express.Request) {
const recent = req.session.lastWriteAt &&
Date.now() - req.session.lastWriteAt < STICKY_MS;
return recent ? db.write : db.read();
}
// In a route:
app.post("/posts", async (req, res) => {
await db.write.query("INSERT INTO posts ...", [...]);
markWrite(req);
res.status(201).end();
});
app.get("/feed", async (req, res) => {
const { rows } = await readerFor(req).query("SELECT ...");
res.json(rows);
});Decide what must always read from the primary
This is a product decision, not a technical one. Write it down for your team:
| Always primary | Replica is fine |
|---|---|
| Account balances, payments, inventory counts during checkout | Public feeds, search results |
| Reads inside a transaction | Profile pages of other users |
| Reads immediately after the same user's write | Analytics, dashboards, reports |
| Anything used to make a decision that writes (e.g. "is this seat still free?") | Product listings, comments |
For the ticket app: "is this seat available?" must hit the primary (ideally with a row lock like SELECT ... FOR UPDATE inside the purchase transaction). Showing a slightly stale seat map is fine; selling based on a stale one is not.
Monitor lag
-- Run on a replica: how far behind is it?
SELECT now() - pg_last_xact_replay_timestamp() AS replication_lag;Alert if lag exceeds a threshold (e.g. 5 seconds) and consider temporarily routing reads to the primary when it does.
Cost
- More database instances (money).
- Your app is now slightly wrong on purpose. Replicas can return stale data, and you've chosen which features may tolerate that.
Done when: primary CPU drops substantially, reads are spread across replicas, and you have a documented list of what always reads from the primary.
Stage 7: Caching
Symptom
The same expensive query runs thousands of times and returns nearly the same answer (follower counts, homepage feeds, event listings, leaderboards, config).
Diagnosis
In pg_stat_statements, look for queries with a huge calls count and meaningful mean_exec_time. Those are your caching candidates.
Implementation: the cache-aside pattern
Request ──► Check Redis ──hit──► return cached value
│
miss
▼
Query database ──► store in Redis with TTL ──► return value// cache.ts
import { redis } from "./session.js"; // reuse your existing Redis client (or a dedicated one — see Appendix B)
export async function cached<T>(
key: string,
ttlSeconds: number,
loader: () => Promise<T>
): Promise<T> {
try {
const hit = await redis.get(key);
if (hit !== null) return JSON.parse(hit) as T;
} catch (err) {
console.warn("Cache read failed, falling back to DB", err); // cache down ≠ app down
}
const value = await loader();
try {
// Add jitter so many keys don't all expire at the same second
const jitter = Math.floor(Math.random() * ttlSeconds * 0.1);
await redis.set(key, JSON.stringify(value), { EX: ttlSeconds + jitter });
} catch (err) {
console.warn("Cache write failed", err);
}
return value;
}
export async function invalidate(...keys: string[]) {
if (keys.length) await redis.del(keys);
}// Usage: follower count
const followerKey = (userId: string) => `v1:user:${userId}:followers`;
app.get("/users/:id", async (req, res) => {
const count = await cached(followerKey(req.params.id), 30, async () => {
const { rows } = await db
.read()
.query("SELECT count(*)::int AS n FROM follows WHERE followee_id = $1", [
req.params.id,
]);
return rows[0].n;
});
res.json({ followers: count });
});
app.post("/users/:id/follow", async (req, res) => {
await db.write.query(
"INSERT INTO follows (follower_id, followee_id) VALUES ($1, $2) ON CONFLICT DO NOTHING",
[req.session.userId, req.params.id]
);
await invalidate(followerKey(req.params.id)); // next read recomputes
res.status(204).end();
});Tip: for extremely hot counters, consider also storing a denormalized
follower_countcolumn updated on follow/unfollow. Then the "expensive count" disappears entirely, and the cache just protects an already-cheap read.
Invalidation strategies
| Strategy | How | Trade-off |
|---|---|---|
| TTL only | Let entries expire after N seconds | Simple; data is stale up to N seconds |
| Delete on write | invalidate() after every write that affects the key | Fresher; you must remember every key a write affects |
| Versioned keys | Put a version in the key (v2:user:...) and bump it on deploy/schema change | Clean way to discard a whole class of entries at once |
| Both TTL + delete | Recommended default | Delete keeps it fresh; TTL is a safety net for missed invalidations |
Cache stampede protection
When a hot key expires, thousands of concurrent requests can all miss and hammer the database at once. For your hottest keys, let only one request rebuild the value:
export async function cachedWithLock<T>(
key: string,
ttl: number,
loader: () => Promise<T>
): Promise<T> {
const hit = await redis.get(key);
if (hit !== null) return JSON.parse(hit);
const gotLock = await redis.set(`lock:${key}`, "1", { NX: true, EX: 10 });
if (gotLock) {
try {
const value = await loader();
await redis.set(key, JSON.stringify(value), { EX: ttl });
return value;
} finally {
await redis.del(`lock:${key}`);
}
}
// Someone else is rebuilding; wait briefly and retry
await new Promise((r) => setTimeout(r, 100));
return cachedWithLock(key, ttl, loader);
}The decision that matters: what may be stale?
Make this list explicitly with your team or product owner. An AI can write the caching code; it can't decide this for your users.
| Cache hard (seconds to minutes stale is fine) | Cache briefly or carefully | Never cache |
|---|---|---|
| Follower / like counts | Event listings with prices | Account balances |
| Public profile data | Seat maps (display only) | Seat availability at purchase time |
| Homepage / trending feeds | User's own notification count | Payment / order status during checkout |
| Static config, feature flags | Search results | Auth / permission checks after role changes |
Don't forget the cheaper layers too: HTTP caching (Cache-Control headers) and a CDN for static assets and public pages can remove load before it ever reaches your servers.
Cost
The sharpest correctness cost in this guide: cached data will be wrong for a while. Plus Redis memory and invalidation bugs, which are some of the hardest bugs to track down.
Done when: your top repeated queries hit the cache (track hit rate, aim for 80%+ on cached keys), DB load drops, and you have a written list of what is and isn't cacheable.
Stage 8: Queues & Background Jobs
Symptom
Endpoints are slow because they do work the user doesn't need to wait for: sending emails, calling third-party APIs, processing images/video, analytics, webhooks. Worse, an outside service being down makes your core feature fail (signup fails because the email provider is down).
Diagnosis
Look at slow endpoints and list every step. Ask of each: "Does the user need this done before we respond?" If not, it belongs in a queue.
Implementation with BullMQ (Redis-based)
pnpm add bullmq ioredis
// queue.ts — shared by web servers (producers) and workers (consumers)
import { Queue } from "bullmq";
import IORedis from "ioredis";
export const queueConnection = new IORedis(process.env.QUEUE_REDIS_URL!, {
maxRetriesPerRequest: null, // required by BullMQ workers
});
export const emailQueue = new Queue("email", {
connection: queueConnection,
defaultJobOptions: {
attempts: 5, // retry up to 5 times
backoff: { type: "exponential", delay: 30_000 }, // 30s, 60s, 120s...
removeOnComplete: 1000, // keep last 1000 for debugging
removeOnFail: false, // keep failures for humans to inspect
},
});Producer: the web server answers fast and enqueues the slow work
// routes/signup.ts
app.post("/signup", async (req, res) => {
const user = await createUser(req.body); // the ONLY thing the user must wait for
await emailQueue.add(
"verify-email",
{ userId: user.id, email: user.email },
{ jobId: `verify-email:${user.id}` } // dedupe: same job won't be queued twice
);
res.status(201).json({ id: user.id }); // respond immediately
});Consumer: a separate worker process
// worker.ts — deployed as its own service/process, NOT inside the web server
import { Worker } from "bullmq";
import { queueConnection } from "./queue.js";
const worker = new Worker(
"email",
async (job) => {
switch (job.name) {
case "verify-email":
await sendVerificationEmail(job.data.email, job.data.userId);
break;
case "password-reset":
await sendPasswordReset(job.data.email, job.data.token);
break;
default:
throw new Error(`Unknown job: ${job.name}`);
}
},
{ connection: queueConnection, concurrency: 10 }
);
worker.on("failed", (job, err) => {
console.error(
`Job ${job?.id} failed (attempt ${job?.attemptsMade})`,
err.message
);
// After all attempts are used, the job stays in the "failed" set.
// Alert on it so a human can look (this is your dead-letter queue).
});
process.on("SIGTERM", async () => {
await worker.close(); // finish current jobs before exiting
process.exit(0);
});Run it: node dist/worker.js as a separate service (its own container, Railway service, Render background worker, etc.). Scale workers independently of web servers.
Replacing cron with repeatable jobs
Remember the "cron runs on every instance" problem from Stage 4? Queues fix it:
await reportsQueue.upsertJobScheduler(
"daily-digest",
{ pattern: "0 8 * * *" }, // every day at 08:00
{ name: "send-daily-digest" }
);
// Only one worker picks up each scheduled run, no matter how many instances exist.Rules for reliable jobs
- Make jobs idempotent. Retries mean a job may run more than once. Sending the same verification email twice is fine; charging a card twice is not. Store a "done" marker or use idempotency keys with payment providers.
- Pass IDs, not whole objects. Enqueue
{ userId }and load fresh data in the worker; payloads can go stale while waiting. - Set timeouts on outside calls so one hung API call doesn't block a worker slot forever.
- Watch the failed set and queue depth. A growing backlog means workers can't keep up (add workers); a growing failed set means something is broken.
- Know the gap: if the database write succeeds but enqueueing fails (Redis blip), the job is lost. For critical jobs, use the transactional outbox pattern: write the job to an
outboxtable in the same DB transaction as the user record, and have a small process move outbox rows into the queue.
Also good queue candidates: image/video processing, PDF generation, webhook delivery, search indexing, analytics events, notification fan-out.
Alternatives to BullMQ: Celery/RQ (Python), Sidekiq (Ruby), Oban (Elixir), or managed queues like AWS SQS, Google Cloud Tasks, Inngest, Trigger.dev, and QStash (especially handy on serverless, where long-running workers aren't available).
Cost
- "Done" now means "promised," not "finished." Design your UI for that ("Check your inbox in a moment").
- Two more things to run and monitor: the queue and the workers.
- Failure handling is your job: retries, backoff, dead letters, alerts.
Done when: slow endpoints respond fast, signup succeeds even with the email provider down, failed jobs are retried and visible, and scheduled tasks run exactly once.
Stage 9: Sharding
Read this first: This is the most powerful and most expensive tool here. Most apps never need it. Teams rightly postpone it for years. Only shard when your data genuinely cannot fit on the largest reasonable single machine, or write volume exceeds what one primary can handle.
Symptom
The dataset is so large (or writes so heavy) that a single primary can't hold or handle it, even after all previous stages. Read replicas don't help because each replica is a full copy.
Try these first (seriously)
| Option | What it does |
|---|---|
| Bigger primary / faster storage | Buys years for many apps |
| Archive old data | Move cold rows (e.g. tickets for past events) to cheaper storage |
| Table partitioning (Postgres native) | Splits one huge table into pieces on the same server; keeps queries simple |
| Move specific workloads out | Analytics to a warehouse, search to a search engine, logs/events elsewhere |
| Distributed databases | Citus (Postgres), Vitess (MySQL), CockroachDB, YugabyteDB, PlanetScale handle sharding for you |
If you still need to shard yourself, here's how.
Step 1: Choose the shard key (the most important decision)
Pick the column that most of your queries already filter by, so each query touches exactly one shard.
| App type | Good shard key | Why |
|---|---|---|
| Social app | user_id | Most queries are "this user's posts/profile/feed" |
| Multi-tenant SaaS | tenant_id / org_id | A company's data stays together; almost no cross-tenant queries |
| Ticketing | event_id | Seats, orders, and holds for one event live together |
Bad shard keys: timestamps (all new writes hit one shard), low-variety values like country (uneven shards), or anything your common queries don't include.
Keep related data together: if you shard by user_id, a user's posts, likes, and settings must also carry user_id and live on the same shard.
Step 2: Use globally unique IDs
Auto-increment IDs collide across shards (every shard has a user #1). Use UUIDs (v7 is time-ordered and index-friendly), ULIDs, or Snowflake-style IDs.
Step 3: Route with logical shards (not raw id % N)
The simple rule hash(id) % numberOfDatabases works, but when you go from 3 to 4 databases, almost every user's shard changes, forcing a massive data move.
Instead, hash into a fixed, large number of logical shards and map those to physical databases. To add capacity later, you move whole logical shards between databases, and the hash never changes.
// shard-router.ts
import { createHash } from "node:crypto";
import { Pool } from "pg";
const LOGICAL_SHARDS = 256; // fixed forever. choose generously
// Physical databases. Start with a few; add more later.
const physical: Record<string, Pool> = {
db1: new Pool({ connectionString: process.env.SHARD_DB1_URL, max: 10 }),
db2: new Pool({ connectionString: process.env.SHARD_DB2_URL, max: 10 }),
db3: new Pool({ connectionString: process.env.SHARD_DB3_URL, max: 10 }),
};
// Map each logical shard to a physical database.
// In production, store this map in config or a small metadata DB.
function physicalFor(logical: number): string {
if (logical < 86) return "db1";
if (logical < 171) return "db2";
return "db3";
}
// Same key → same number, every time, on every server
export function logicalShard(key: string): number {
const digest = createHash("sha1").update(key).digest();
return digest.readUInt32BE(0) % LOGICAL_SHARDS;
}
export function shardFor(key: string): Pool {
return physical[physicalFor(logicalShard(key))];
}
export function allShards(): Pool[] {
return Object.values(physical);
}// Single-shard query: goes straight to one database, no searching
async function getUser(userId: string) {
const { rows } = await shardFor(userId).query(
"SELECT * FROM users WHERE id = $1",
[userId]
);
return rows[0];
}
async function createPost(userId: string, body: string) {
await shardFor(userId).query(
"INSERT INTO posts (id, user_id, body) VALUES (gen_random_uuid(), $1, $2)",
[userId, body]
);
}Step 4: Handle cross-shard queries (scatter-gather)
Questions that span all users now have to ask every shard and combine the answers:
async function totalUserCount(): Promise<number> {
const results = await Promise.all(
allShards().map((pool) =>
pool.query("SELECT count(*)::int AS n FROM users")
)
);
return results.reduce((sum, r) => sum + r.rows[0].n, 0);
}
async function latestPostsGlobally(limit = 20) {
const results = await Promise.all(
allShards().map((pool) =>
pool.query("SELECT * FROM posts ORDER BY created_at DESC LIMIT $1", [
limit,
])
)
);
return results
.flatMap((r) => r.rows)
.sort((a, b) => b.created_at - a.created_at)
.slice(0, limit);
}These get slower with every shard you add. Strategies to avoid them:
- Keep global counts elsewhere: maintain a counter in Redis or a metadata table.
- Precompute: background jobs (Stage 8) that aggregate into summary tables.
- Send analytics to a data warehouse instead of querying shards.
- Look up by non-shard keys via an index table: e.g. login by email needs a small global
email → user_idlookup table, then route byuser_id.
Step 5: Plan the migration (the genuinely hard part)
Moving a live database onto shards without downtime typically looks like:
- Dual write: new writes go to both the old database and the correct shard.
- Backfill: copy historical data to the shards in batches.
- Verify: compare old vs new (row counts, checksums, sampled reads).
- Shadow read: read from shards in the background, compare with the old DB, log mismatches.
- Cut over reads, then stop writing to the old DB.
- Keep a rollback plan until you're confident.
Budget weeks to months for this, not days.
Cost
- Cross-shard queries become expensive and complex.
- Transactions across shards are effectively gone (design so they're never needed).
- Operations multiply: backups, migrations, and monitoring per shard.
- Hard to undo. Once code is written around shards, going back is its own migration.
Done when: each common query touches exactly one shard, cross-shard work runs only in background jobs or dedicated systems, and you can add physical capacity by moving logical shards.
Appendix A: Symptom → Stage Diagnosis Table
| What you observe | Likely cause | Go to |
|---|---|---|
| App server CPU pegged, DB fine | App server capacity | Stage 2, then 3 |
| Server crash = full outage | Single point of failure | Stage 3 |
| Random logouts / lost carts after scaling | Hidden per-server state | Stage 4 |
| Cron jobs running multiple times | Scheduled work on every instance | Stage 4 + 8 |
too many clients / connection timeout errors during spikes | Connection exhaustion | Stage 5 |
One query is slow, shows Seq Scan | Missing index | Stage 6, step 1 |
| DB CPU maxed, queries indexed, mostly reads | Read load | Stage 6, step 2 |
| User's own post missing right after posting | Replication lag | Stage 6 (read-your-writes) |
| Same expensive query run thousands of times | Repeated computation | Stage 7 |
| Users see outdated numbers | Cache invalidation / TTL too long | Stage 7 |
| Endpoints slow due to emails/APIs/processing | Synchronous slow work | Stage 8 |
| Feature fails when a 3rd-party service is down | Tight coupling | Stage 8 |
| Data too big for any single machine | Dataset size | Stage 9 (after alternatives) |
Appendix B: The Redis Question
Redis shows up three times in this guide: sessions (Stage 4), cache (Stage 7), and queues (Stage 8). Using one piece of technology for all three is great for simplicity, but be aware of one real conflict:
| Use | Desired eviction policy | Why |
|---|---|---|
| Cache | allkeys-lru | When memory fills up, drop the least-used entries. That's the point of a cache. |
| Sessions | noeviction (or volatile-* with TTLs) | Evicting a session silently logs someone out |
| Queues (BullMQ) | noeviction | Evicting job data silently loses jobs |
Recommendation:
- Small apps: one Redis with
maxmemory-policy noeviction, TTLs on all cache keys, and memory alerts. Simple and safe. - Growing apps: split into two instances: one for cache (
allkeys-lru) and one for sessions + queues (noeviction, persistence enabled). That's why the examples above use separateREDIS_URLandQUEUE_REDIS_URLvariables.
Use managed Redis (or a compatible alternative like Valkey) with replication and automatic failover in production, since sessions and queues make it critical infrastructure.
Appendix C: Final Architecture & Checklist
┌─────────────┐
│ CDN / edge │ (static assets, public pages)
└──────┬──────┘
│
┌──────▼──────┐
│Load Balancer│ (managed, redundant)
└──────┬──────┘
┌───────────────────┼───────────────────┐
┌──────▼─────┐ ┌──────▼─────┐ ┌──────▼─────┐
│ App (state-│ │ App (state-│ │ App (state-│
│ less) │ │ less) │ │ less) │
└──┬───┬───┬─┘ └──┬───┬───┬─┘ └──┬───┬───┬─┘
│ │ │ │ │ │ │ │ │
┌─────────┘ │ └──────┐ │ │ │ │ │ │
▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼ ▼
┌───────┐ ┌──────────┐ ┌──────────────────┐ ┌──────────────┐
│ Redis │ │ Redis │ │ Connection Pooler│ │ Queue (Redis)│
│ cache │ │ sessions │ └────────┬─────────┘ └──────┬───────┘
└───────┘ └──────────┘ │ │
┌────────────┼──────────┐ ┌─────▼─────┐
▼ ▼ ▼ │ Workers │
┌─────────┐ ┌──────────┐ ┌──────────┐ └───────────┘
│ Primary │─► Replica 1│ │ Replica 2│
│ (writes)│ │ (reads) │ │ (reads) │
└─────────┘ └──────────┘ └──────────┘
(sharded into multiple primaries only if data outgrows one)Implementation checklist
Foundation
- Metrics: latency percentiles, RPS, CPU/memory, DB connections, slow queries
- Config via environment variables
- All DB access through one module
- Uploads in object storage, not local disk
-
/healthendpoint
Scaling the front
- Vertical scaling exhausted or too expensive before going horizontal
- All CPU cores used (cluster / PM2 / multiple workers)
- Multiple instances behind a (managed) load balancer
-
trust proxyset; graceful shutdown on SIGTERM - No per-server state: sessions in Redis or stateless tokens
- Shared
SESSION_SECRETacross instances - Cron/scheduled tasks run exactly once
Scaling the database
- Pool sizes calculated: instances × pool max <
max_connections - Pooler enabled (managed or PgBouncer); migrations use direct connection
- Slow queries reviewed with
EXPLAIN ANALYZE; indexes added - N+1 queries eliminated
- Read replicas with read/write routing
- Read-your-own-writes handled
- Documented list of "always read from primary"
- Replication lag monitored
Doing less work
- Cache-aside for expensive, repeated reads
- TTL + delete-on-write invalidation
- Documented list of what may and may not be cached
- Cache failure falls back to DB instead of breaking the app
- Slow/external work moved to a queue
- Jobs idempotent, retried with backoff, failures visible and alerted
- Workers deployed and scaled separately
Only if truly needed
- Alternatives to sharding evaluated (archiving, partitioning, distributed DB)
- Shard key chosen based on real query patterns
- Global IDs (UUIDv7/ULID/Snowflake)
- Logical → physical shard mapping
- Cross-shard queries moved to background jobs or other systems
- Zero-downtime migration plan with rollback
The Takeaway
Every stage trades something away: money, simplicity, or correctness. Replicas and caches make your app faster and slightly wrong on purpose, and queues make "done" mean "promised." The engineering skill isn't knowing the names of these boxes; AI can write the code for every one of them. The skill is knowing which ones you actually need, in what order, and what each one costs, and deciding which parts of your product are allowed to be a little wrong for a little while.
Start simple. Measure. Break it. Fix exactly one thing. Repeat.


