TT Lab
Get started
Learn Learning paths Courses

Load Testing

Rejection Is Protection, Not Failure

Continue in TT Lab

In one line

The design question is not 'should we have a queue' but which queue reaches its limit first and what you do then.

Why this matters

On the road a request travels before it is processed, queues line up. Socket buffer → thread pool queue → connection pool wait → DB lock wait → processing. Once any one of them saturates first, everything after it turns into waiting. So bottleneck investigation is the job of sweeping this ladder in order.

The upper bound of a queue length can be calculated. It is the allowed wait time divided by the average processing time. With a target of 200ms, processing of 20ms, and 10 workers, the bound is roughly 100. If you go beyond this bound and keep an unbounded queue, you only turn an outage into latency. Rejection is not failure but protection. It is better to return a 503 quickly.

How it works

The most famous war story is the collision of two defaults. HikariCP maximum-pool-size defaults to 10, and Tomcat threads.max defaults to 200. When traffic grows, most threads get stuck waiting for a connection and Connection is not available, request timed out after 30000ms pours out. The fix is to change three things together. Raise the pool to 30, reduce connection-timeout from 30 seconds to 5 seconds to make it fail fast, and actually shrink Tomcat threads.max to 50. The alert to prevent recurrence fires when hikaricp.connections.pending > 0 persists for 30 seconds or more.

The procedure for finding the saturation point has five steps.

  1. Set an initial value with a formula. Pool size = (number of cores × 2) + number of disk spindles. With 8 cores it is around 20.
  2. Record the pool wait p99 and throughput at the target load.
  3. Try halving it. If throughput holds, the original was too large.
  4. Increase it until throughput stops growing, and slightly below the value where it starts to flatten is the right size.
  5. If the wait is still long even at that value, you must fix the queries, not the pool.

A long connection wait means a long hold time, and that is usually a slow query or a long transaction. The pool size only hides that symptom for a while. So the principle is this: the pool size should be the number of requests the database can handle well at the same time, not the number the application wants to send.

What it looks like in the field

If you increase concurrency beyond the limit, throughput X does not grow and only response time W grows proportionally. Soon context switching, lock contention approaching the square of the concurrency, a drop in the cache hit rate, and retry load pile on and throughput actually decreases. The message this equation gives is not the number itself but the order of magnitude. The right answer is tens, not hundreds.

The measuring tool hides the bottleneck

The most common reason you fail to find a bottleneck is on the measuring side, not the server. If the load tool sends the next request only after it has received a response, then while the server is stalled, no request is fired at all. That stall is not recorded as any request's response time, so even if the server is frozen for a whole second, it does not show up in the metrics. This is called coordinated omission. Real users do not hold back their requests depending on the server's condition, so a p99 measured this way is far more optimistic than reality.

There are two ways to avoid it. First, use a tool that supports fixing the target arrival rate. It fires requests at set intervals without waiting for responses, and records the delay added to the response time. Second, if the tool cannot do that, instrument the queue wait time separately on the server side. The difference between the moment a request arrived at the socket and the moment a worker picked it up is that value, and the point when this value grows is the real saturation point.

The average deceives people in a similar way. The response time distribution is not symmetric but stretches long to the right, so the average becomes a value that most users do not experience. That is why you look at percentiles. However, if you average the p99s of several instances, that value means nothing. That is because percentiles are not values that can be added or divided. You must collect the histogram buckets, rebuild the whole distribution, and then compute the percentile.

One more thing to point out. If a single screen consists of ten internal calls and each call's p99 is 100ms, the p99 of the whole screen is not 100ms. The probability of landing on the slow side at least once in ten is much larger, so the latency users actually experience rises well above that. A decision to add one more internal call should be seen not as adding that call's average but as adding one more tail. So the most effective way to protect performance is not to make each call a little faster but to reduce the number of calls itself or make them overlap in parallel.

What you will do in the next lab

You deliberately build and start a Python HTTP server that handles only one request at a time. You measure at concurrency 1 and 10 to confirm the serialization signal in which throughput stays the same and only latency grows, set up three bottleneck hypotheses by layer, and then switch to a server that can handle requests concurrently and compare how many times throughput jumps. At the end, you calculate the number of instances needed to handle the target traffic.