In short
- Spikes keep finding the same five weak points: queries per page, the connection pool, an expiring cache key, slow work inside the request and the number of network hops.
- Each one has a simple calculation behind it. Do the calculation before launch and not after.
- Splitting a system into services adds hops, and every hop adds latency and a way to fail.
- Start with one application, one database, a queue and a cache. Add boundaries when they buy something.
The failures that make the news are usually boring ones
When Ticketmaster opened the Eras Tour presale in November 2022, it said its systems took 3.5 billion requests, about four times its previous peak, with bots adding to the load. When Segment looked at its own system it had grown to more than 140 services, with three engineers spending most of their time keeping it alive, and the team merged everything into one service.
Those are opposite problems. One is too much traffic and the other is too much architecture. Both come down to two questions: what happens when more people show up than you planned for, and what does each piece of your design cost when they do?
I see the same handful of causes in startups at a far smaller scale. A product that works for a hundred users breaks at ten thousand, and the framework is almost never the culprit. It is one of five places in the request path.
One page, a hundred and one queries
An ORM makes it easy to load a list of orders and then, inside the loop, load each order's customer. On a laptop with ten rows nobody notices. In production the page runs one query for the list and one more for every row, and each of them is a network round trip to the database.
Swipe sideways to see the whole figure
Assume a round trip costs four milliseconds. A page of 100 rows then spends 404 milliseconds waiting on the database, one query after another, before it renders anything. Fetching the related rows with one query, using a join or an IN list, takes two round trips for the same page: about eight milliseconds. The data is identical and only the number of trips changed.
Catch it early. Count queries per request in development, and fail a test when a page runs more than a budget you chose.
The connection pool is a hidden queue
Postgres starts one process for each connection and, by default, allows 100 of them. Application pools are much smaller, often 10 or 20 per instance, and every request that touches the database waits for one. The ceiling on throughput is pool size divided by how long each request holds its connection.
Startups lose most of that ceiling by holding a connection while doing something slow. A checkout handler opens a transaction, calls the payment provider, waits 300 milliseconds and then commits. With 10 connections the ceiling is 10 divided by 0.3 seconds, about 33 requests a second, whatever the CPU is doing. Call the provider before or after, hold the connection only for the 15 milliseconds of real queries, and the same pool passes about 667 a second.
Swipe sideways to see the whole figure
A pooler such as PgBouncer shares a small number of real connections between many clients, which helps once you run many app instances. It does not rescue a handler that holds a connection while it waits on someone else.
When a popular cache key expires
A cache hides an expensive query until the entry expires. At that moment every request that arrives sees a miss, and each one starts the same expensive recomputation. If the query takes 1.5 seconds and 20 requests a second are arriving, 30 copies of the query run at once against a database that was only ever sized for one. The database slows down, the recomputation takes longer, more requests pile in and the cache never refills.
Swipe sideways to see the whole figure
Facebook's paper on scaling memcache describes this as a thundering herd and introduces leases, so that only one client at a time may refill a key. On a set of keys that were prone to herds, the authors measured a peak database load of 17,000 queries a second without leases and 1,300 with them. You do not need Facebook's infrastructure to do the same. Coalesce concurrent misses so one request refills and the others wait or serve slightly stale data, and add random jitter to expiry times so that keys written together do not expire together.
Work the user does not need to wait for
Sending a confirmation email, calling a partner's webhook and rendering a PDF all feel like part of placing an order, so they end up inside the request. Each adds its own latency and its own way to fail. If the email provider is down, the customer sees an error for an order you already saved.
Swipe sideways to see the whole figure
Save the order, put a job on a queue and answer. A worker does the rest, retries on failure and can be scaled separately. The customer sees their confirmation in tens of milliseconds, and the provider's bad afternoon becomes a retry queue's problem. Give every job an idempotency key so that a retry cannot send the email twice.
Every network hop has a price
Splitting a system into services adds network calls, and every call is a place to be slow or to fail. Availability multiplies along a chain. Twenty services in a row at 99.9% each give about 98%, which is roughly 14 hours of trouble in a 30-day month, against 43 minutes for one service on its own. Fan-out hurts latency in the same way. Dean and Barroso point out in The Tail at Scale that if a request waits on 100 servers and each one is slow one time in a hundred, 63% of requests will be slow.
Swipe sideways to see the whole figure
Segment's move shows the cost from the other side. Its destinations had grown into more than 140 separate services, and the post describes shared libraries drifting into different versions across them, autoscaling that had become more art than science, and operational overhead that grew linearly with every destination added. After merging them into one service, the team reports that on-call pages for load spikes on low-traffic destinations went away.
Services earn their keep when the boundary buys something larger than the hops it adds: a team that can deploy alone, a component that scales alone, a failure you can fence off. For most early products a well-organised single application gives you the option to split later without paying for it now.
What I build first
My default for an early product is one deployable application with clear internal modules, PostgreSQL, a queue for slow work and a cache for hot reads. Connection pooling goes in from the first day. Every request logs how many queries it ran and how long it held a connection, so the first two problems in this article show up as numbers on a dashboard before they show up as complaints. Before a launch I run a load test against staging with traffic shaped like the launch, which means a sudden jump and not an average day.
A checklist before a launch
- How many queries does your busiest page run, and what is the budget?
- How long does a request hold a database connection, and does any outside call happen while it does?
- What happens to the database when your most popular cache key expires?
- Which work in the request could happen after the response, and what retries it?
- How many services does one user action touch, and what is their combined availability?
- When did you last load test with a spike and not a steady rate?
Sources
- Dark Reading, Ticketmaster blames bots in the Eras Tour presale (Ticketmaster's figure of 3.5 billion requests)
- Twilio Segment, Goodbye Microservices
- Scaling Memcache at Facebook (NSDI 2013)
- Dean and Barroso, The Tail at Scale (Communications of the ACM, 2013)
- PostgreSQL documentation, connection settings
The animated figures are simplified illustrations that use the numbers stated in the text. Where a figure is modelled on a real incident, the sources above are the full accounts.
