Queues should be empty
joshvoigts.com
joshvoigts.com
The way I see it is that there are two opposing things you can optimise for that depend on queue depth: utilisation or latency.
If you care about processing each message as quickly as possible then queues should be empty. This often requires a higher cost and lower utilisation as you inevitably have idle workers waiting for new messages.
If you care about utilisation, then you never want your queue to be empty while something is polling it. This might be some background task that runs on a GPU - every second that sits idle is wasted cash, and those tasks usually benefit heavily from batching inputs together.
In this case you want it to always have something to read from the queue and shut it down the moment this isn’t the case.
Also please don’t accuse someone of not “responding to the content of the article” when the content is a bog standard list. If you want my response to that: “yes, that is indeed a list of 4 things queues help with”. Not sure it adds much to the discussion.
I guess I should really ask, why do people care about utilization instead of throughput, latency, and cost? We could increase utilization by rewriting everything to be more wasteful, for example, but it doesn't seem like we should do it.
People care about utilization when, eg, processing petabytes of data — something that takes hours to days, where you’re optimizing batch throughout and cost. By making sure you efficiently utilize your hardware.
For a given rate of incoming work, increasing throughput decreases utilization. So I guess it's not choosing throughput too much.
> People care about utilization when, eg, processing petabytes of data — something that takes hours to days, where you’re optimizing batch throughout and cost. By making sure you efficiently utilize your hardware.
If you have fixed resources and lots of work, it would be good to have nonempty queues to prevent them from idling. That is maybe the big unwritten part of the comment upthread. It doesn't seem worth saying "no, the article is wrong in all cases because of this particular case!" though, so I'm not sure.
If you have fixed resources and fixed work, wanting your queues to be empty is a wish, not a strategy.
You might argue that a queue with many items in it which can still accept more items is not really "full", but it is clearly not empty either.
(1) Obviously all queues are bounded by the capacity of the hardware, but since most queueing software can persist to disk and/or lives in the cloud with elastic storage this is not usually a problem.
> always have something to read from the queue
In the first case, you've gone a bit far. To maximise utilisation, you want the queue to never be empty when polled.
Obviously you don’t want it to be 100% full, but you do want it to be the kind of full you get after a nice meal but you still have room for desert and a coffee.
You really don't want your queue more than 70-80% utilized. Once it gets to 75% you should consider that "red zone" and you need to fix it somehow. Once the queue is 80% you're really hitting into some issues. Above 80% and you will have significant problems.
If you’re dealing with something like SQS and are receiving variable sized batches of messages, the consumers can be classed as “100% utilised” (always processing a batch of messages with no empty receives) whilst the queue would not be considered 100% utilised (the average delivered batch size is always less than the max).
Obviously there is a line where you do not want a bunch of expensive consumers hanging around - but this isn't the place IMO to micro-optimize your costs. You shouldn't have like 20x the consumers you need but you shouldn't worry about having EXACTLY the correct amount.
In your example, are you implying that if more jobs come in more workers get spawned? Is that specific to the Amazon service?
A process could receive N tasks in a batch and process them at the same time.
In this scenario it is possible for all workers to be utilised whilst the queue is not “100% utilised” - which is a failure mode the link you shared explained:
> And by induction, we can show that a system that’s forever at 100% utilization will exceed any finite queue size you care to pick
The inflows to the queue are less than the outflows, despite every worker spending 100% of its time processing tasks.
It is specific on the utilisation type though, I’m mainly thinking about ML inference that has a fairly fixed initial cost per inference but scales quite nicely with >1 input.
> The inflows to the queue are less than the outflows, despite every worker spending 100% of its time processing tasks.
Right, but at some point they will have to have SOME idle time. If they are NEVER idle the queue by definition is growing, not shrinking, as you said.
Unless you have some system that monitors queue length and doesn't add jobs if it is like >n, which is what the article suggests (maximum queue size). But then you have to figure out what you do with those jobs that can't be added, which probably means just creating some kind of queue higher up in the pipe.
If they are processing more than is being added, eventually things will run out of the queue.
Any queue that fills up, say halfway, is an indication the system is in a precarious state.
Its simple. Queues have limited space if the velocity of data going in is faster then data going out the queue will fill up until it's full.
A half full queue indicates the velocity is in such a state and such a state is unsustainable. You need to build your system such that it stays out of unsustainable states as much as possible.
If you're finding that your queue fills up half way a lot that means you're building your system in such a way that it operates at the border of sustainability. In other words the system often ingests data at rates that cannot be sustained indefinitely.
If a system with a queue half full operates in that same state twice as long, you now have a full queue and you're losing data.
In essence your queues should mostly be empty most of the time. At best you can account for occasional spikes as the article stated but a queue even half full should never be the operational norm.
Not all queues are “realistically bounded”. Let’s take a hypothetical silly “order queue” where each entry represents a €10 item that’s been purchased. The queue feeds into a boxing machine that puts the item in a box and posts it to the user.
The point at which a queue hits an upper bound would represent billions of euros of orders. That’s not really a realistic failure mode to hit without noticing long beforehand.
Meanwhile, if you always want your queue empty then your boxing machine is going to be shipping thousands of small boxes because the moment an order arrives, it’s processed into a box.
It’s much more efficient to have a small, managed and time-bounded backlog of orders that the boxing machine uses to batch orders together in a smaller number of larger boxes.
There is no system on the face of the earth that can operate at an average positive velocity. That is categorical failure.
This is mathematically and logically true no matter how big your queue is.
Even for the situation your describing... A really big ass queue, the average velocity of that system must still be negative. Unless your queue is so big that it can withstand an average negative velocity over the lifetime of the system which is kind of absurd.
Most of the time your queues should be empty. If your queues are observationally often filled with data but not full it means you're still operating at a negative velocity but you've tuned your system to operate at the border of failure.
It means your system hits unsustainable velocities often. You may want this for cost reasons but but such a system is not robust enough for my comfort. The queue should be empty most of the time with occasional spikes at most.
I agree. But that’s not the argument the post is making, nor does it support the position that “queues should be empty”. No, they shouldn’t. It’s fine to let a queue build up (and thus not be empty), then process all the items in the queue in a short time frame. There are efficiency savings if you are able to do this.
This is true even when you have a big ass queue. The system must on average be more empty then it is filled.
This is not to say that your system has to never have spikes in data. Your system can often have spikes, but this cannot ever be the norm as in the amount of spikes in data production rate cannot push the average production rate to exceed the average consumption rate.
If I have 1 message per minute being added to the queue, I can safely run 2 of these batch processes once a day to clear the queue. Or more generally can launch (queue_size / 1000) tasks once a day.
Because surely it doesn’t matter about the queue being empty, what matters is the rate at which messages can be processed.
With 2 processes running once a day the rate would be 2000/24/60 = 1.3 messages per minute in aggregate.
The queue can be filled 23 hours per day without issue and without being on the brink of a failure state.
Edit: you edited your comment whilst I was replying to specifically mention consumption rate. The aggregate consumption rate is all that matters, but it doesn’t follow that the queue “has to be more empty than not empty”. Those are different issues.
You're talking about a single batch job that runs once a day and empties the queue at the same rate no matter how much data is in the queue. Thus the more data you grab in batch the more efficient the processing so you grab data once per day instead of continuosly. This is a clear exception from the assumed norm that needs to be mentioned.
That is not to say that your case is invalid. Many analytic databases possess your ingestion parameters. But the case is exceptional enough that if not explicitly mentioned it can be assumed that such a batch job is not a primitive in the model being discussed.
Let's say you flew from one city to another. People assume you took a plane. If you took a helicopter, though helicopters are common, it should be mentioned otherwise it won't be assumed.
For analytics databases the batch ingestion requirement are usually so big that you have to save it to the file system hence the need for an external queue.
Heap allocated queues are different, at that stage I would say that it's in "processing" already and popped out of the external queue.
It's more similar to streaming... what you're doing here. In that case my velocity measurements are more applicable. You want your queues to be empty in general. A batch job is run every hour or something like that which is not what you're doing here.
If you ran your load test for 5 minutes and you see your queues are 50 percent full. Well that means 10 minutes in you'll hit OOM. Assuming your load tests are at a constant rate.
If your queues are mostly empty then it can handle the load you gave it and have room for spikes. It's just math.
Very few technical problems don’t have any elasticity in outflows, but a lot of problems do benefit from a form of batching. Heck, even something as simple as inserting records into a non-analytic database benefits greatly from batching (inserting N rows rather than 1) and thus reducing TPS, or calling an external service, or writing files to some storage, or running inferences, or writing to a ledger, or anything.
Batching, in all forms, exists absolutely everywhere from super-scalar CPU uops to storage to networking to the physical world. The exception is to do something without batching, even if the whole stack you’re using from the ground up is built on it.
Given the common use case of batching, which inherently requires items to be queued for a period of time, how can you say that queues should always be empty? They would statistically rarely be empty. Which is a good thing. Because it enables batching.
If you want to queue and insert 1000 records into a database with 1000 transactions across 1000 workers, then say your workload is highly optimised because there are never any records in your queue, cool. Not very true though.
They have elasticity but you don't care for it. You run it at the maximum. As you mentioned the only case where you wouldn't is if you have some transaction that doesn't scale linearly with with frames/events.
>Given the common use case of batching, which inherently requires items to be queued for a period of time, how can you say that queues should always be empty? They would statistically rarely be empty. Which is a good thing. Because it enables batching.
In my experience, while batching is not uncommon, streaming is the much more common use case. So we can agree to disagree here.
Now, that is an extreme example, but similar behavior exists. It’s fine to have queues partially filled most of the time, as long as your average rate of processing exceeds the rate of incoming items by some margin.
Spiking the processing is not what people think about with queues and is most likely NOT what the author of the article is talking about.
Typically there's no reason why you would spike your processing unless such processing doesn't scale linearly with the amount of items processed (as the author in a sibling post to yours mentioned). Such a case must be deliberately introduced into the discussion as an exception to the model because it's not what people think about when discussing this. Or maybe it is given that 3 people brought the same exact exception up. Welp if that's the case then I can only say the author of this article and I certainly aren't referring to this use case.
edit: I'm rate limited. So I'm referencing another post in my response that you won't see until the rate limit dies and I post it.
Your comfort is secondary to the business needs that the system is designed to fulfill. In practice some risks are better to accept than mitigate
Either way you have admitted that you tuned the system to operate with an acceptance of risk. Likely for cost reasons. This is fine, if you deliberately chose to do this.
Yet you bravely asserted that
> At best you can account for occasional spikes as the article stated but a queue even half full should never be the operational norm.
Suddenly this thing that you said "should never be" changed to "is fine"
You want to sky dive while deliberately choosing not to use a parachute? That's fine. But don't expect me to adjust my communication to account for your exception. You chose to do this.
But let's not pretend that you don't know about this fact in human communication. It's an obvious thing you know. Yet you chose to characterize your post in a way that negatively makes my statement look unreasonable and your actions look normal. Please don't do this. If you choose to sky dive without a parachute, that's fine, but admit you are operating with risk for the sake of cost. Do not try mischaracterize my intentions, that is just rude.
You assert that without justification. Typical counter-example would be system where you have predictable daily load-pattern, it can be perfectly reasonable to build the system that accumulates items during peak hours and clears queues during quiet hours.
You're maybe referring to what the other person mentioned... processing that does not scale with the amount of data processed so you want the processing to operate on the biggest batch of data possible. This is, while a common use case, not common enough that it should be assumed without be mentioned.
However, in your example it seems to be happening on the same computer (you mentioned mic), (and thus not web, which is what's the assumed context here) therefore data production rate doesn't seem to spike. If your processing is faster then data production you don't even need a queue, just stream each frame straight to processing in a single thread. But this also assumes the nature of your processing. The details for the nature of the processing can change things to require a queue for something like the averages of frames at discreet intervals... but that's different from the use case of queues like kafka or zmq for which I think is the main context here.
Edit: I reference other responses here, but you won't see them immediately. I'm currently rate limited and I have those replies in a queue so those will be posted later today.
Since we don't have infinite space, we can expect eventually to lose some messages in this scenario.
If there is a 1 hour queue for the checkout at a grocery store, you probably won't join the queue - you'll go to a different store...
Sure, that's technically correct but applies to basically everything. It is very likely your users/orders/whatever table in your database is also "trending to infinity" over time, except this doesn't mean databases should always be empty or that we should expect to eventually lose some users/orders/whatever.
Or, more succinctly, if your queue capacity is 10 million messages and your queue messages represent "houses purchased through our website", then in reality your queue capacity is infinite because nobody will ever purchase 10 million homes through your application per small unit of time.
But this is also where the disagreement comes in with flexible serverless workers. They generally do not cost more, so you'll end up spending $24 in an hour rather than $1 an hour all day. If the serverless costs more then you are likely not working through your backlog.
I agree it should be a business decision, so you catch up on weekends instead for example, but non-technical folks don't always grasp those numbers and that they are getting a lag with their data.
It depends what you mean by "approaching empty". If you mean "empty most of the time", then no, that's wasted resources. If you mean "its size is decreasing most of the time", then no, not possible. Even if you mean "it is empty some of the time", then maybe, but that is not a strong requirement either.
Of course more average input than average output is bad, but stating it as "should be empty" or "should be approaching empty" seem to suggest the wrong thing entirely. It should be smaller than the target latency of the process.
70-75% is just a rough number, the real number depends highly on several factors related to your specific use of the queue.
It may not make sense but it is the reality.
So I disagree entirely with "it's either shrinking or growing". It is both growing, shrinking, and stable, at the same time, over different horizons, and I have no idea what law you are trying to state.
The link you provided above is a really good primer based on serious math, I don't understand how you think it supports your vague claims.
Queues are either growing (adding items) or getting smaller (consuming items). There is not a state where they stay the same length, unless you're doing something really weird like monitoring to make sure there are always 10 items in the queue, or something, and adding if that number drops, but then you'd need a second queue-type structure to support that, and the conversation kind of goes off the rails.
I literally cannot understand how you think time works in your area.
You're also being really rude for no good reason. My polite suggestion is that you reflect on yourself here about why you are being rude on the internet to strangers. It could lead to growth and hopefully help you with whatever hurt you're experiencing.
[1] https://en.m.wikipedia.org/wiki/Little%27s_law
[2] broadly construed. For example a service that gets a mass of requests at the top of every hour displays seasonality.
This is literally the slippery slope fallacy. You aren't accounting for minor fluctuations or finite timescales. If you were driving a car, you might say "if you aren't steering into oncoming traffic then you're steering towards the ditch" and conclude that no cars will ever safely reach their destination.
If the only valid state for a queue is empty, then why waste time implementing a queue?
If you feel like you have enough spare capacity and any given item isn’t taking too long to process then it doesn’t matter if you are ever empty.
Corner conditions are where I start to care:
* What is the worst case latency? How long until the work can be done? How long until the work can finish from the point at which it enters the queue?
* What is the worst case number of items in the queue? How large does the queue need to be?
* What is the maximum utilization over some time period with some granularity? How spiky is the workload? A queue is a rate-matching tool. If your input rate exactly equals your output rate at all times, you don't need a queue.
* What is the minimum utilization over some time period? Another view on how spiky the workload is.
Minimums and maximums I find are much more illustrative than averages or even mean-squared-errors. Minimums and maximums bound performance; an average doesn't tell you where your boundary conditions are.
In general you don't want your queue to fully fill up, but like the other poster said, it's some tradeoff between utilization and latency, and the two are diametrically opposed.
The unifying aspects of using queues for us is that it allows us to load-balance jobs across workers, and allows us to monitor throughput in a centralised fashion. We can set alerts on the time-sensitive queues and then use different thresholds for the background queues, but we're using the same metric data and same alerting system.
In such a case, you'd expect at least N messages in the queue at all times — and often a small multiple of N. Not because the consumer can't consume those messages in a timely manner, but because, given a reliable queue, messages can't be consumed until they've been 2PCed into the queue, and that means that there are always going to be some messages sitting around in the "acknowledged to the producer, but not yet available to the consumer" state. Which still takes up resources of the queueing system, and so should still be modelled as the messages being in the queue.
I think the reason you get the responses you do is people are used to systems with high variation in flow rates that can only process a limited number of items at the same time, and for these, queues can soak up a lot of variability by growing long and then being worked off quickly again.
Love the downvotes! People should read this: https://blog.danslimmon.com/2016/08/26/the-most-important-th...
I understand this is a little counterintuitive but it is true. You should not have a queue that always has a job ready to be processed.
You either have to be consuming faster than the queue is growing, or you are adding faster than you consume. There is no equilibrium.
If you are consuming faster than it is growing, at some point you'll consume all jobs, meaning your queue is sometimes empty. If you are adding faster than you consume, your wait time approaches infinity (tho in reality you'll fill up your queue storage first, of course).
You may need that or want to add this buffer - but you're not solving any computation resource problem directly. If your jobs take XYZ CPU cycles they take XYZ CPU cycles - it doesn't matter if you make them wait a little bit.
You could even be spending more money, since it does take some resources and cost to scale things.
It doesn't mean it's not valid to do things this way, I'm just saying it's still the exact same problem.
Because, no, nothing you said changes the fact that queues only do anything when kept mostly empty. And the only thing they do is some short term trade of latency for reliability. They actually can't change utilization.
This airport seems busy but the planes seem full, how can we improve utilisation?
Which refutes the original point about queues not doing much unless they are empty and being unable to effect utilisation. If you have a somewhat inelastic output that benefits from batching, but has variable inputs, they are absolutely crucial to utilisation.
This is no longer the case using the latest in serverless IaaS architecture, where the pool of queue workers dynamically expand to meet demand.
However, yes, you don’t pay for this or care about it much. It’s good.
A good comparison is a restaurant's order queue.
If the order queue is always empty, that means no one is ordering food at the restaurant. That is an undesirable state. We want people to be ordering food, an empty queue shows that we have a problem in getting people to eat at our restaurant.
We don't want our queue to be always growing. If it's always growing, that means we're understaffed, and people waiting for their meals will get frustrated. Some will wait and leave a bad review, others might leave before they ever get their meal--and won't pay for it, either. Lost time and effort on our part.
But an order queue that almost always has two meals in the making is golden. It means we're getting consistent clientele to dine at our restaurant, that meals are being cooked and served in good time, diners are happy with the experience, and our system operations are healthy.
I think the author of TFA understands this, but stopped a little short of the optimal expression of their ideas. It's not about what you have, it's about what you should always be working towards.
You need to intentionally design your bottlenecks not pretend that you can avoid them completely.
And really, we need more context to say whether a queue should be empty, trend towards being empty, or towards always having tasks. I don't think it's possible to make a general prescription for all queues.
Then you can employ statistics to plan what resources you will allocate to processing the queue.
Though having capacity for Poisson event spikes isn’t nearly the only way to ensure queues trend to zero elements. For example, you can also scale the workload for each queue element depending on the length of the queue.
It’s very popular in video game graphics to allocate a given amount of frame time for ray tracing or anti aliasing, and temporally accumulate the results. So we may process a frame in a real time rendering pipeline for more or less time, depending on the state of the pipeline. Say, if the GPU is still drawing the previous frame, we can accumulate some ray traced lighting info on the CPU for the next frame for longer. If the CPU can’t spit out frames fast enough, we might accumulate ray traced lighting kernels/points for less time each frame and scale their size/impact.
Statistical distributions of frame timings don’t figure here much as workloads are unpredictable and may spike at any time - you have to adjust/scale rendering quality on the fly to maintain stable frame rates. This example assumes that we’re preparing one frame on the CPU while drawing the previous one on the GPU. But there may also be other threads that handle frames in parallel.
Other workload spikes could be predicted by things like Poisson distributions. But once again, the key takeaway is that you can’t generalise this stuff in any practical application. It all depends on system requirements and context.
The right statement is not "they should be empty" but "they should be under the target latency", which might be 0 (then "the queues should be empty"), or might be some other target (then empty means overprovisioned resources).
Everyone in this thread seem to be talking over each other because they imagine vastly different services (e.g. shipping company vs hospital emergency rooms).
(Imagine sending a TCP packet for every individual byte that was enqueued on a socket!)
Same applies if you have any sort of load shedding on the consumer side. If the queue is empty, it means you've shed load you didn't need to; when the next brief drop in load comes, you can't take advantage of it.
(Fun analogy: this is why Honda hybrids typically don't charge their battery all the way while driving, to allow the chance to do so for free while braking.)
Say, you have 100 requests coming in per second. Each consumer of your queue requires 0.3 seconds per request. Now, you can optimise for the number of consumers. Would you choose 40 consumers, then the probability of serving the request with an empty queue in 0.3 seconds would be 17% (and > 0.3 seconds would be 84%). However, the total waiting time per request would be 0.1 second.
With 80 consumers, the probability of an empty queue would be higher at 58%. However, most of those consumers are doing nothing, while your waiting time is only reduced by 0.1 second! Only a 30% improvement.
(edit: probability of empty queue in second example is higher, not lower)
Like, suppose you had a specific service level agreement that requires X% of requests to be served within Y seconds over some time interval, and breaking the SLA costs Z dollars. Each provisioning level would generate some distribution of latencies, which can get turned into likelihood of meeting the SLA, and from there you can put a dollar value on it. And crucially, this allows the "ideal" amount of provisioning to vary based off the relative cost of over and under-provisioning; if workers are cheap and breaking the SLA is costly, you would have more workers than if the SLA is relatively unimportant and workers are expensive.
In the end, the distribution is not really poisson, of course. So, you might be interested in low pass filtering to elastically scale your provisioned workers. There is quite some theory about this, including sophisticated machine learned models to predict future load. But I digress.
The thing that queues are good for is enforcing order: when I get to the front of the queue and it’s time to do my thing, I have a guarantee that anyone who entered the queue before me has already done their thing, and anyone who entered the queue after me has not yet had a chance to do their thing. This is a useful guarantee for tasks that have an order dependence, for example database operations looking to guarantee consistency.
You could conceivably have a non-zero queue size that asymptotically reaches some stable size over time. In such a case you are keeping up with messages but with some approximately static amount of added latency.
I don’t think I’ve seen this happen in practice but I think some systems can be engineered to be approximately stable (grows and shrinks but within predictable bounds)
That said, consider a physical machine. The "hopper" on many is essentially a queue to feed in items to work on. If you let that go empty, you messed up. Hard.
If a queue is just being used to decouple a producer/consumer pair, then "it should be [almost] empty" is a reasonable notion, assuming consumer/producer are compute bound and there's no fundamental reason why one should be faster than the other...
Which brings us to the other type of use where things are being enqueued specifically as a means of waiting for resources that are in limited supply: maybe just compute if these are computationally heavy jobs, or maybe I/O bandwidth, etc. In that case we expect the queue to often be non-empty, and there's nothing wrong with that, even if there's potential to throw money/hardware at it to get the queue size down to zero!
The system should be designed to handle as many messages as possible ensuring that your resources don't get over burdened. Empty queues where your messages keep getting rejected because your resources are being hammered is not necessarily an efficient system.
Can you expand on this?
That is - initially you can decrease latency by sacrificing utilization and just adding more workers that usually sit idle. This increases throughput during heavy load, and decreases latency. Until you can't, because the overhead from all the extra workers causes a slowdown up the stack (either the controller that's handling messages, or the db itself that's got the queue state, or whatever is up there). Then when you throw more overhead, you increase latency and decrease throughput, because of the extra time wasted on overhead.
It depends on a lot of different variables.
If you have work items that vary significantly in cost, it can add even more problems (i.e. if your queue items can vary 3+ orders of magnitude in processing cost).
1. Limits concurrency based on latency gradient or queue depths etc.
2. Weighted fair queueing scheduler - prioritized load shedding based on workload priorities and latency estimation.
Anyone interested in Little’s Law, token buckets, network schedulers, control theory (PID controllers) should definitely check it out! [1]
Instead of letting the system fail under heavy load.
You put the tasks in a holding place.
They will be processed later.
Computers have spoiled us with instant and now.... But even a 1 hour processing time is amazing if that task would take 3 days manually.
The runtime acts as an event loop and if a specific program is hogging too much time on the CPU, it will get paused and thrown to the back of the execution stack. This is all transparent (relatively) for the developer.
Here’s a video if you want to see an example: https://youtu.be/JvBT4XBdoUE
Basically everything that processes data has a queue attached to it, from your router to your scheduler. Not knowing anything about queue theory is such an insult to the profession of SWE. It's like being a Mech engineer and not knowing about control of dynamic systems.
The dread should have been that most companies/teams won't care, won't engineer on that level, will just set it up however and make it bigger/faster/more-consumed if it's too slow.
Also a lot of GitHub repos with names like CS550…
Only place I’ve seen it referenced since that class was:
1. HN periodically
2. A YouTube documentary about Disneyland fast pass queues
3. A business optimization consultancy
Ultimately, a pretty niche subject area AFAICT. But implementing a fairly efficient queuing simulation seems to be pretty easy, and the subject seems to be stupid powerful and broadly applicable… so I’ve never understood why it’s not more discussed
Let's say you consume a job, do something, but that something fails. You have to set it to be rescheduled somewhere. Lets say your things writing jobs to the queue fails. That needs to be handled as well.
"They will be processed later [if I use a queue]" is not true by default (without work and thought) and assuming so will get you in trouble, and will also make certain people on your team who do understand that isn't the case upset.
I've seen this come up recently, where our queue reader was(well, still is) over scaled, where it could handle 100x the normal queue throughout. When the input rate gets higher though, instead of the queue filling a bit and controlling the traffic, we hammer our dependencies.
I now see a continuously empty queue as a smell. It should be filling and emptying over and over again. Not all the time, but a queue that never fills is a ticking time bomb
* (Mandatory) Have a back-pressure system.
* (Optional, depends on system requirements) have auto-scaling of queue readers.
* (Mandatory) Make sure the system is regularly load tested.
In the US it was called queue theory, in the USSR it was known as mass service theory.
The thing is, there were no queues in the US and no mass service in the USSR.
If overall customer demand is x widgets per unit of time, if you design your system to process y widgets per unit of time, you will have overprovisioning (and hence impact to profits) if y >> x (y very large relative to x). As a corollary, if you design your system where y < x (y smaller than x), then you will not meet SLAs for some customers.
This analysis can be done at each subsystem level that is a component workflow and in reality some subsystems may be overprovisioned while some are underprovisioned. The goal should be to match the overall rate of customer demand.
You also design the system to handle a certain level of spike in demand, but the system will fail when demand exceeds capacity. When that happens, you want one place where you want to manage flow and invariably that will be a critical queue.
So non-empty queues are a good thing, even for hospital waiting lists.
(To forestall arguments, of course the operating theatre example is not a pure queue example, some patients get put to the head of the queue according to triage - for the purpose of this example it doesn't matter)
If you are using queues for this purpose, then I recommend reading the following:
https://lmax-exchange.github.io/disruptor/disruptor.html#_th...
> When in use, queues are typically always close to full or close to empty due to the differences in pace between consumers and producers. They very rarely operate in a balanced middle ground where the rate of production and consumption is evenly matched.
This has always been my experience. Has anyone ever had a perfect goldilocks scenario between 2 threads?
The reasons for this can be numerous, maybe you produce a lot of jobs at once because of a user interaction, or your consumer is running slower for five minutes because some cron job is stealing your resources.
>This way, the producer doesn’t need to know anything about the consumer and vice versa.
Both components needs to agree on shared data, down to what each field actually means. The producer is just not saying What To Do with the message.
Take an example of processing a purchase order in multiple, asynchronous steps. While the webpage button spins there's a whole dance of agreed messages and expected outcomes that needs to happen in order for it to work. And there is no decoupling in sight.
My take is that Decoupling is a vague word. One should design distributed systems assuming that coupling is manageable, not inexistent.
>you could submit the task and ACK immediately allowing the job to finish in the background.
That's bad practice.
This can be helpful for avoiding resource contention, or hitting an API limit.
This article obviously has a very specific kind of workload/service/target in mind that they fail to define and they overgeneralize. This is bad advice in many situations that they failed to imagine.
This article is saying you basically want an average negative velocity where consumers eat data faster then producers on average.
Any queue that is non empty means the velocity is positive. Producers are making data faster then consumers and you have a limited time window before the queue becomes full and your system starts dropping data.
If you observe that your system never has a full queue but is constantly loaded with partially filled queues it means you built and tuned your system such that it operates on the border of failure.
It means your system is continuously hitting unsustainable velocities all the time, but stops short of total failure (a full queue).
Remember, think about the system in terms of it's derivative. The average derivative must be negative for your system to even function. Any time a queue is filled partially or full it means there was a period of positive velocity. These periods cannot exceed the amount of time the system is in negative velocity.
Come on, it's an ideal; it doesn't have to be realistic.
[Edit] Just read the Slimmons article. His argument is specifically about a queue served by a server that is under stress. His assertion is:
"As you approach maximum throughput, average queue size – and therefore average wait time – approaches infinity."
He's quite right (and yes, it's a bit counter-intuitive). But it's orthogonal to the matter of the ideal queue length.
Thanks for the link.
“Error queues should be empty”