JB logo
CoffeeyOUTUBE
Blog
PreviousNext

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:

  1. Symptom: what you'll see when you need this stage
  2. Diagnosis: how to confirm it's really this problem
  3. Implementation: the concrete steps and code
  4. Cost: what you're paying (money, complexity, correctness)
  5. 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

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

MetricWhy it mattersWhere to get it
Request latency (p50, p95, p99)Averages hide pain; p99 shows your worst usersAPM tool, Nginx logs, app middleware
Requests per secondTells you how close you are to capacityLoad balancer / platform metrics
Server CPU & memoryIs the app server the bottleneck?Host metrics
DB CPU, memory, disk I/OIs the database the bottleneck?Managed DB dashboard
DB active connections vs max_connectionsConnection exhaustion (Stage 5)pg_stat_activity
Slowest queriesMissing indexes (Stage 6)pg_stat_statements
Error rate by typeTimeouts vs connection errors vs bugsLogs / 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 scattered new 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.

  1. In your hosting dashboard, move to the next instance size up (more vCPUs, more RAM).
  2. 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: 3 in 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 app to 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

StrategyHow it worksUse when
Round robin (default)Each server in turnRequests are roughly equal in cost
least_connServer with fewest active connectionsSome requests are slow, others fast
ip_hashSame client IP → same serverAvoid; 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 stateWhere it lives nowWhere it should live
Login sessionsIn-memory session storeRedis (or stateless tokens)
Uploaded filesLocal diskObject storage (S3/R2/GCS)
In-memory caches (const cache = new Map())Process memoryRedis (Stage 7), or accept per-instance caching for non-critical data
Rate limit countersProcess memoryRedis
WebSocket connectionsTied to one serverRedis pub/sub adapter to broadcast across servers
Cron jobs / setInterval tasksRun on every instanceA 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 sessionsSigned tokens (JWT)
Extra lookup per requestYes (fast)No
Instant logout / revokeEasy: delete the keyHard: needs a denylist (which is... a shared store)
InfrastructureRedis requiredNone
Good forWeb apps with cookiesAPIs, 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 connect

Servers 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: SET without LOCAL, advisory locks held across transactions, LISTEN/NOTIFY, and (on older PgBouncer versions) named prepared statements. Prisma needs ?pgbouncer=true in 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: 1 or 2).
  • 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 WHERE clauses
  • Foreign keys used in JOINs (Postgres does not auto-index these)
  • Columns used in ORDER BY with LIMIT (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 primaryReplica is fine
Account balances, payments, inventory counts during checkoutPublic feeds, search results
Reads inside a transactionProfile pages of other users
Reads immediately after the same user's writeAnalytics, 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_count column updated on follow/unfollow. Then the "expensive count" disappears entirely, and the cache just protects an already-cheap read.

Invalidation strategies

StrategyHowTrade-off
TTL onlyLet entries expire after N secondsSimple; data is stale up to N seconds
Delete on writeinvalidate() after every write that affects the keyFresher; you must remember every key a write affects
Versioned keysPut a version in the key (v2:user:...) and bump it on deploy/schema changeClean way to discard a whole class of entries at once
Both TTL + deleteRecommended defaultDelete 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 carefullyNever cache
Follower / like countsEvent listings with pricesAccount balances
Public profile dataSeat maps (display only)Seat availability at purchase time
Homepage / trending feedsUser's own notification countPayment / order status during checkout
Static config, feature flagsSearch resultsAuth / 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

  1. 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.
  2. Pass IDs, not whole objects. Enqueue { userId } and load fresh data in the worker; payloads can go stale while waiting.
  3. Set timeouts on outside calls so one hung API call doesn't block a worker slot forever.
  4. 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.
  5. 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 outbox table 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)

OptionWhat it does
Bigger primary / faster storageBuys years for many apps
Archive old dataMove 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 outAnalytics to a warehouse, search to a search engine, logs/events elsewhere
Distributed databasesCitus (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 typeGood shard keyWhy
Social appuser_idMost queries are "this user's posts/profile/feed"
Multi-tenant SaaStenant_id / org_idA company's data stays together; almost no cross-tenant queries
Ticketingevent_idSeats, 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_id lookup table, then route by user_id.

Step 5: Plan the migration (the genuinely hard part)

Moving a live database onto shards without downtime typically looks like:

  1. Dual write: new writes go to both the old database and the correct shard.
  2. Backfill: copy historical data to the shards in batches.
  3. Verify: compare old vs new (row counts, checksums, sampled reads).
  4. Shadow read: read from shards in the background, compare with the old DB, log mismatches.
  5. Cut over reads, then stop writing to the old DB.
  6. 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 observeLikely causeGo to
App server CPU pegged, DB fineApp server capacityStage 2, then 3
Server crash = full outageSingle point of failureStage 3
Random logouts / lost carts after scalingHidden per-server stateStage 4
Cron jobs running multiple timesScheduled work on every instanceStage 4 + 8
too many clients / connection timeout errors during spikesConnection exhaustionStage 5
One query is slow, shows Seq ScanMissing indexStage 6, step 1
DB CPU maxed, queries indexed, mostly readsRead loadStage 6, step 2
User's own post missing right after postingReplication lagStage 6 (read-your-writes)
Same expensive query run thousands of timesRepeated computationStage 7
Users see outdated numbersCache invalidation / TTL too longStage 7
Endpoints slow due to emails/APIs/processingSynchronous slow workStage 8
Feature fails when a 3rd-party service is downTight couplingStage 8
Data too big for any single machineDataset sizeStage 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:

UseDesired eviction policyWhy
Cacheallkeys-lruWhen memory fills up, drop the least-used entries. That's the point of a cache.
Sessionsnoeviction (or volatile-* with TTLs)Evicting a session silently logs someone out
Queues (BullMQ)noevictionEvicting 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 separate REDIS_URL and QUEUE_REDIS_URL variables.

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
  • /health endpoint

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 proxy set; graceful shutdown on SIGTERM
  • No per-server state: sessions in Redis or stateless tokens
  • Shared SESSION_SECRET across 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.