LearnThatStack Ace your next interview

System Design Concepts · API Design

Scale and availability: what statelessness makes possible, and what still fails

Holding no session is what lets the servers scale out. For availability, the other half, assume parts fail and stop one failure from spreading.

statelessSession in memory breaks the balancerany request, any boxboxes24can answer24holding you0browserGET /cartBearer eyJhbGciOiany box can check itwhat the user gotGET /cart200GET /cart200GET /cart200load balancerpicks a boxround robinbox Aholds nothinga.internalin memory: nothingbox Bholds nothingb.internalin memory: nothingbox Cholds nothingc.internalin memory: nothingbox Djust startedd.internalin memory: nothingwhere the session livesa shared store, or the tokenthe two homes that are not a box's memorypick either one and every box is interchangeablea box's memorya shared storethe tokenadding a box is a number in a config file, not a project
statelessSession in memory breaks the balancerboxes24can answer24holding you0browserBearer eyJhbGciOi200load balancerround robinbox Aholds nothingbox Bholds nothingbox Cholds nothingbox Djust startedwhere the session livesa shared store, or the tokena shared store or the token, never the boxadding a box is a number, not a project

Session in memory breaks the balancer

  1. Box A is the only server, and it keeps the session sid 7f3a for user 124 in its memory, along with 811 other sessions. Every request from user 124's browser gets a 200 reply, and nothing is shared or wrong yet.
  2. Traffic grows, so a second server, box B, is added, and a load balancer picks one box for each request. Box B is running but has never seen sid 7f3a, so only box A of the two running boxes can answer this user.
  3. The load balancer sends the next request to box B, which has no session for sid 7f3a in its memory. So box B answers 401 and signs the browser out, because of the load balancer's choice and nothing the user did.
  4. Sticky sessions make the load balancer send every request with this cookie to box A, so the error goes away. Then box A is deployed and loses the 812 sessions in its memory, so every user pinned to box A is signed out at once.
  5. The session moves out of box A's memory into Redis, a shared session store that keeps sid 7f3a with a 30 minute expiry. Both boxes read sid 7f3a from Redis, so both answer 200 without holding anything that belongs to a user.
  6. Another option is to give the state to the client, in a signed token that holds user 124 and an expiry. Each server checks the signature itself, so there is no session store to run, but nothing can revoke the token before it expires.
  7. None of the four server boxes holds a session, so any request can go to any box. The boxes are stateless, so adding a fifth box means changing a number in a config file, not starting a project.

© LearnThatStack - diagrams may not be republished without permission.

Most systems start with one server, box A, that keeps each session in its own memory. Box A holds sid 7f3a, the session ID for user 124, plus 811 other sessions. Nothing is wrong yet, and every request from that browser comes back with status 200. But watch two counts: how many servers are running, and how many servers can answer you.

The whole chapter is about one gap: the servers that are running, and the servers that can answer you. Traffic grows, so you add a second server and put a load balancer in front of both. Now two servers are running, but only the first server, box A, can answer you. The gap opened as soon as there were two of something.

The request was correct, and it still came back 401. The load balancer sent the request to box B, a server that had never seen session sid 7f3a. The user did nothing, but the load balancer signed them out. The balancer will sign them out again on roughly half of their requests.

Sticky sessions are popular because the cookie pins each user to box A, the server holding the session, so the error goes away. Then box A is deployed and loses everything in its memory, so 812 people are signed out at once. Sticky routing still sends those 812 people to box A, the server that lost their sessions. So the sticky-session patch fails at exactly the moment you needed it.

In practice, stateless does not mean there is no state. It means no single server owns the state. The session moves out of the server process into Redis, which keeps sid 7f3a with a thirty minute expiry. Both servers read the session from Redis, so both return 200.

The other place a session can live is the client itself. The client holds a signed token that carries user 124 and an expiry time. Each server checks the signature on its own, so there is no session store to run. The trade-off is revocation, which means cancelling a session early. A shared store lets you delete a session right away. But a token works until it expires, so you keep the expiry short.

This is what you mean when you say a service scales horizontally. None of the four servers holds a session, so any request can go to any server. Four servers are running, and at last all four can answer requests. A fifth server needs only a new number in a config file, because the fourth server holds nothing worth copying.

The servers are now replaceable, so you could add one, but try these options first. Their order matters more than the list itself. At three thousand requests a second, checkout answers in 950 ms at the 95th percentile. Each row names one real part of that wait, and the rows add up to 950 ms. A bigger server is already used up, because doubling the machine is a real first move that stops at one machine.

Caching the hot reads comes first, because an afternoon of work removes 560 ms, more than half the wait. The same product page is requested 82 times a second, and a cache answers that page once instead of 82 times. The cost is stale data: a read can return a page up to a minute old. You also have to handle invalidation now, which means removing a cached page when the product data changes.

The primary database is answering reads that it never had to answer. Reads outnumber writes about ten to one here, so sending the reads to replicas (copies of the database) saves 80 ms. The cost has a name that interviewers listen for: replication lag, the delay before a write reaches the replicas. A read straight after a write can miss that write, so any path that must see its own write stays on the primary.

Connection pooling needs no new machines and no new service, only one setting. Opening a new database connection for every request costs more than most of the queries run on that connection. Reusing open connections from a pool takes 27 ms off the wait. But pool size becomes a limit you have to tune, and one slow query now holds a connection that another request wanted.

Moving work off the request saves time by changing the contract with the client. The confirmation email goes out inside the request, so the request waits for the email. Put the email on a queue instead, and the request stops waiting, which cuts 60 ms from the response time. An asynchronous API is this trade: the client gets a 202 now and finds out the result later.

Sharding (splitting the data across databases) comes last, and the numbers show why. Sharding cuts 25 ms from this p95 latency, the smallest saving of any option here. And sharding adds costs: every query now needs a shard key, you give up cross-shard joins, and you need a plan to rebalance the data. So shard when one database runs out of capacity, not to make a request faster.

Most of performance work is knowing which part of the wait nothing can remove. Field selection and gzip shrink the JSON from 180 kB to 24 kB, so the wait drops by 90 ms here. Sending less saves more over mobile, because bandwidth and battery both cost the user something there. The 108 ms left after this last optimization is the endpoint doing its work, and no optimization removes that time.

All of that assumed the parts work, but next, one of them gets slow. Checkout calls pricing, pricing calls inventory, and the page returns in 60 ms. Each service has a pool of forty threads, with no more than three busy. Remember that all three health checks return 200 up.

The sentence to say out loud here is that slow is worse than down. Inventory is not down. But its database is saturated, so each call takes 9 seconds instead of 25 milliseconds. Every pricing thread that calls inventory waits, so pricing's pool of forty threads fills up. A call to a service that is down fails at once, but a slow service holds your threads.

The failure now spreads up the chain of calls because each caller waits. Checkout's threads wait on pricing, and pricing waits on inventory, so all forty of checkout's threads are busy too. The user gets a 503 after thirty seconds of waiting. Three services are unusable, but not one process crashed, so the CPU graphs look idle through the whole incident.

The one thing to say about availability is that every call over a network has a deadline. A missing timeout is the usual root cause. With a deadline on the call, the pricing service stops waiting for the inventory service after 800 milliseconds instead of 9 seconds. Each thread returns to the pool eleven times sooner, so the pricing service's thread pool drops from forty busy threads to twelve.

Twenty failures in ten seconds move the breaker on pricing's calls to inventory from closed to open. A call to inventory now fails in 1 millisecond without leaving pricing's process, so pricing holds four threads instead of twelve. In half-open, the breaker's third and last state, one probe call every thirty seconds checks whether inventory has recovered. Inventory falls to six busy threads, because you stopped sending calls to a service that was already struggling.

The pricing service answers from its last good price, which is three minutes old, so the page still renders. Accepting a three minute old price is a product question, not an engineering decision. Ask that question now rather than during the outage. Decide in advance what a degraded answer looks like on each screen.

The Inventory service answered 200 through this whole outage. A health check that only proves the process is alive gives a false answer, and the platform trusts that answer. So the platform keeps sending traffic to an instance that cannot do its job. Check the dependency and let the readiness check fail, so the platform stops sending traffic to that instance.

Answer in two halves. Holding no session is what makes scaling out possible. Session state lives in a shared store or a signed token, so any server can handle any request. Then scale in this order: cache hot reads, add read replicas, and pool connections. Next, move slow work out of the request. Shard last, because sharding saves the least and costs the most. Availability is the second half. Every call gets a deadline. A circuit breaker stops calls to a service that is already failing. Slow is worse than down.

In interviews · 4 questions

Related Questions

Also helps with

← All concepts