Why disaster happens at the edges: An introduction to queue theory
thenewstack.io
thenewstack.io
Queue theory encourages people to think in terms of flow rates and has the technically correct de-rates to nominal capacity to account for variance. But normally that is overkill except in highly engineered systems.
The battle in an organisation is convincing people to look at flows at all. In my experience people love metrics that track stock (we have X widgets or can handle Y orders) and not flows (we built A widgets and sold B widgets). Anyone who has training in queue theory knows when it is appropriate to look at flows rather than stocks and that is where the value is. The formulas and variances tend to just scare people away from monitoring flows because they don't understand what is happening.
Of course there are systems where not being able to process everything right when it comes in means that you're "down". It's not a given though and queuing systems are actually perfect for use cases where not being able to process at the speed of incoming requests is completely fine. Eventual consistency is a thing.
I have the same experience though that it seems to be hard for folks to accept that yes, 100k messages in _that_ queue is totally fine. Some process dumped 100k messages to then be processed over the next hour and that's totally fine. That way we don't have to actually spend a small fortune to process them at the speed at which they can be generated.
Of course it depends on the problem at hand. If that particular flow in the system is one that is highly user interactive, then yes, this means 'system down' for your users. Never mind that the system will eventually process it all. I always smile when I see another system than ours and I can spot distributed processing w/ queues in between in how the system behaves in certain circumstances.
Granted, depending on how well behaved your clients are, the number of events each second may as the failed requests are retried.
Send the feedback as far up the source as possible.
But you can also redirect cargo to different ports, or increase the price you charge for servicing cargo. Even if you dump containers you don’t have to do so randomly. Customers could even mark containers as rejectable, or non critical and therefore first to go to a backlog for a discount.
I found the article really interesting. For most IT use cases jobs are generic but thanks for raising the port example.
See https://www.cips.org/knowledge/procurement-topics-and-skills... for more.
Look at the port. Is ingress consistently larger than egress? Then limit ingress.
Now look at the ocean. Is ingress consistently larger than egress? Then limit ingress.
Now look at the exporting ports. Is ingress consistently larger than egress? Limit ingress.
Now look at trucks going to the port. Is ingress consistently larger than egress? Limit ingress.
Now look at the factories that load trucks. Is ingress consistently larger than egress? Limit ingress.
----
Back pressure builds on looking at one queue and asking if arrivals come faster than departures go -- exactly what the top-level comment suggested.
I am talking about a bazillion other use cases in which there is variance in the number of events produced at any given time and where the immediate processing of those events is not of utmost importance. "Slap a queue in between" and suddenly you are no longer bound by what your systems can handle, you are only bound by what the queuing system can handle and that is usually much easier to scale (or can already handle this kind of load).
E.g. think about a system needing to synchronize with another external system. Back in the old days (and well I'm sure those systems are still out there today) synchronizing two systems may well have been an overnight job. A nightly sync job. Nowadays many such systems probably try near realtime synchronization but you don't want to make it part of the regular flow either and block the actual operation.
But sometimes there's too many synchronizations that need to happen at once, so a queue builds up over some period and goes down again a bit later. It's still better user experience than a nightly job. E.g. your regular load might be at 10 e/s. During some particularly busy time or specific bulk operations being executed in parallel by multiple people, you might have an influx of 200 e/s for a minute or so.
If you have a queue or something in the middle then the failure is potentially recoverable or never happens in the first place.
However they do allow you to smooth or average spikes that could otherwise be problematic.
In your articles analogy I am able to drain a pot of pasta down the sink even though I might not be able to pour the same pot down the drain.
For a numerical example of queuing theory 101 and why variance is very important rather than something that “scares people away” [0]. Flow tracking is not enough.
Also this discussion is relevant [0]:https://www.johndcook.com/blog/2008/10/21/what-happens-when-...
Which, in practice, is a model of queuing systems that often gets most of the value of queue theory and is easy to explain without jargon. Then apply an 80% fudge factor. Unless the situation needs care.
I didn’t understand GPs comment about flows - perhaps that’s where the dispute is.
Genuine question, looking forward to learning more :)
In the link I shared there is a simple worked example of 10.3 capacity/hr for 10 events/hr that piles up 28 events. It also shows the problem can be mitigated by analysing it in queuing/probabilistic terms. In real life it is even more complicated obviously!
I guess that's the point... an organization needs enough people to know about queue theory "to look at flow" and understand that queues are ubiquitous and anticipate that these queues experience a latency phase change at some point depending on utilization (and job duration variance). If an organization doesn't have anyone who understands queue theory, then it's probably more likely to have queue-related failures.
In order to look at flow - "should learn about queue theory"!
"queue theory" itself would be useful to design queuing related systems.
Queuing theory deals with the stochastic behavior of a discreet event in at a single queue, and possibly multiple ones connected in some structured manner.
It's importance is that it can determine before hand what are the stochastic boundary of the system.
Your whole idea of 99 events and 101 events are so strawman that I don't know where to start poking at the holes...
For one thing, the rate of something has its definition, there is jitter or busyness. To say something has a rate is plainly a impractical concept. As there is nothing in the network that had steady rate...
Edit: I should stand corrected that the parent does not really miss the meaning of queuing theory. It was just that I only read the first paragraph.
I guess that's the OP's point that you can't escape from talking about queues and which you sort of validated by bringing up Little's law.
Queues are literally everywhere around us in the physical world. It needs to be understood at a basic level -- Little's law is good enough and only requires primary school math. If execs aren't able to understand such a ubiquitous and simple concept then I don't know what to say.
Average flow rate is indeed useful but please explain that along with a queue and show whoever needs to know how queues can grow unboundedly if output rate is consistently less than input rate. And if needed explain further the upstream consequences of unbounded queue growth.
Not really. You need actual theory to inform you whether 99 events per second is okay. See https://www.johndcook.com/blog/2008/10/21/what-happens-when-...
Most of the time feeling it out in production is a lot cheaper than keeping an in house specialist. If it isn't then sure, maintain a specialist to monitor the situation. It is a relatively rare role though.
And keeping your safety margins nice and wide isn't just inefficiency. It's also building in resilience for when something unexpected does happen. Yes, if you need your six-nines performance on a 1% margin then maybe you should try to optimize this to the nth degree, but realistically almost none of us have to sail that close to the wind.
This reminds me of an entrepreneurship class I took in college where the professor emphasized the importance of cash flow, something I hadn't considered much and found somewhat counterintuitive.
I'd love to learn more about this perspective on when to focus on the flows over stocks, are there any resources you'd recommend to help me better understand it?
Basically all you are doing when you "study queue theory" is training your brain to go very quickly from looking at some random metric to thinking "oh, this is a queue with X distribution, I need to work out the arrival rate and service rate then I know everything".
My personal reference is "Fundamentals of Queuing Theory" by Gross, Shortle, Thomson & Harris. But to learn, just pick some interesting queues (ie, queues with variance in the rates) and start calculating (if cars arrive at a light with X mean rate, what variance means that there will sometimes be a queue extending around the street? Would adding another server at this shop have a big or small difference on how long it takes to prepare my order? Based on the number of checkouts at this supermarket, what rate are their planners expecting - is this a busy period?).
No individual idea in queue theory is that powerful, the advantage is really that when you see a queue you don't waste hours trying to link variables in O(n^2) ways and instead work with arrival/service rates. And you realise that the system works by phase changes rather than anything else, which is in a sense obvious but not actually very intuitive if you haven't identified the queue.
Do you want to provision your system for the worst case scenario volume, and what is this worst case scenario you are willing to pay to insure your system against.
It’s an easy question, but I don’t see how queueing theory can provide answers to. It is at its root a subjective business decision.
Cold start is an important scenario that lots of people overlook - “we won’t restart it during busy times”, etc. however if you have a bug or outage then that might be precisely when you restart. During a cold start you will often find your system behaves as it does at the tails of the performance curve, either because you are getting 100% cache misses, or because all your queues are full due to a client pile-on or queued upstream work, or some combination.
The article described back pressure; I can’t emphasize that enough as part of any resilient queuing system. Chains of queues are only as performant as their weakest link.
https://github.com/joelparkerhenderson/queueing-theory
Queueing theory is the mathematical study of waiting lines, or queues. We use queueing theory in our software development, for purposes such as project management kanban boards, inter-process communication message queues, and devops continuous deployment pipelines.
Dμ = Delivery service rate. Devops teams may say "deployment frequency" or "we ship X times per day".
Rτ = Restore lead time. Site reliability engineers may say "time to restore service" or "mean time to restore (MTTR)".>Every distribution can be described as a curve. The width of this curve is known as its variance.
Not every distribution has a variance. Some notable examples include the Cauchy distribution (or Lorentz lineshape in physics),
https://en.m.wikipedia.org/wiki/Cauchy_distribution
the Power law distribution for powers lesser or equal to 2,
https://en.m.wikipedia.org/wiki/Pareto_distribution
and the Levy distribution,
https://en.m.wikipedia.org/wiki/Lévy_distribution
These are not mathematical curiousities but actually describe real physical systems.
Additionally, the variance works as a "width" only for unimodal distributions. Any variance based metrics or analysis should only be used after checking for multimodality and, for multimodal data, with extreme caution thereafter.
From the article: "It’s tempting to focus on the peak of the curve. That’s where most of the results are. But the edges are where the action is. Events out on the tails may happen less frequently, but they still happen. In digital systems, where billions of events take place in a matter of seconds, one-in-a-million occurrences happen all the time. And they have an outsize impact on user experience."
EDIT: IMO, the title is still a little annoying in this respect. I think everyone would agree if a request to your site fails 5% of the time, that is unacceptable, even though it "usually works." The discussion of the distribution curve simply to make the point that spikes in usage cause backed up queues which impact performance isn't necessarily helpful as far as I can tell, and it seems done largely in service to the title. In my mind while reading this, I'm thinking, "Okay, cool, but how does the fact that this interesting issue exists at the edge of the curve help me identify it?" Answer: It doesn't. If you see errors occurring, you will investigate them once they are noticed. Being at the edge of the curve may mean it takes longer to notice, but like, what kind of alerting system are you using that discriminates against rare issues in favor of common ones?
Discussing queues, over provisioning, back pressure, etc. are all super interesting and helpful.
> As a rule of thumb, target utilization below 75%
This is one good reason to rely on serverless. Simply outsource the problem to a system that knows how to handle this better.
> Steer slower workloads to paths with lower utilization
This is fraught with all sorts of perils [0], so take caution going down this valid but tricky route.
> Limit variance as much as possible when utilization is high
This is key, and one of the most elegant solutions to this problem I know of comes from a Facebook talk on CoDel + Adaptive LIFO. [1]
> Implement backpressure in systems where it is not built-in
Making downstream dependencies behave is never an option. Selectively isolating noisy neighbours, if possible, from other well behaved clients, tends to work well for multi-tenant systems [2], in addition to monitoring long work (by dropping work that takes forever) [3], or better yet, doing constant amount of work [4] (aka eliminating modes) [5].
> Use throttling and load shedding to reduce pressure on downstream queues.
One way is to impl admission control (ala Token Bucket) to minimise the impact of thundering herds. Though, clients do not take kindly to being throttled. [6]
A great article; many have been written at this point. Very many still have been burned by queues. Though, in my experience, without an exception, some component somewhere was always building up that backlog! [7][8][9]
[0] Interns with Toasters: How I taught people about Load Balancers, https://news.ycombinator.com/item?id=16894946
[1] Fail at scale: Controlling queue delays, https://blog.acolyer.org/2015/11/19/fail-at-scale-controllin...
[2] Worload isolation, https://aws.amazon.com/builders-library/workload-isolation-u...
[3] Avoiding insurmountable queue backlogs, https://aws.amazon.com/builders-library/avoiding-insurmounta...
[4] Constant work, https://aws.amazon.com/builders-library/reliability-and-cons...
[5] Cache, modes, and unstable systems, https://news.ycombinator.com/item?id=28344561
[6] Fairness in multi-tenant systems, https://aws.amazon.com/builders-library/fairness-in-multi-te...
[7] Using load shedding to avoid overload, https://aws.amazon.com/builders-library/using-load-shedding-...
[8] Treadmill: Precise load testing, https://research.fb.com/publications/treadmill-attributing-t...
[9] Cascading failures, https://sre.google/sre-book/addressing-cascading-failures/
Surely they mean 'HTTP protocol'? TCP doesn't do back pressure but IP does at the router level, IIRC.
Doesn't TCP typically use sliding window flow control? I'm not sure what "code" is referring to in this context.
But yes, TCP uses windows for back-pressure, but that isn't really useful for application level backpressure as the OS controls the queues sizes, so pretty much most systems have their own backpressure on top.