5 ms·
Job queues are deceptively tricky
- serhii777_123 3mo ago[flagged]
- zmj 3mo agoI haven't modeled it, but I wonder how far you'd get on randomizing the policy choice for concurrency limit 1. Maybe weighted by past results, but bounded to allow it to shift instead of falling permanently into a basin.
- wewewedxfgdf 3mo agoYou solve this simply with two cron jobs, one for weekend and one weekdays.
- ktimespi 3mo agoI find it very annoying when queue problems break into queue-of-queue patterns like in the `wait`scenario
- irjustin 3mo agoI remember learning about CSV parsing and how it's conceptually simple, yet beyond the simple , and quotes: the corner cases bloat your parser 10-15x.
- edoceo 3mo agoOr, you could be a government agency that implemented a method for citizens to upload compliance data - so they have to upload broken CSVs into a broken queue system. Get the worst of both worlds, and deployed as a green-field project in 2021.
- groundzeros2015 3mo agoIt’s pretty simple to make one that’s RFC compliant. The rules aren’t much more than what you said. Are you talking about trying to interpret malformed data?
- nosefrog 3mo agoMajor lesson from when I worked on Google Search indexing is that queues have a lot of hidden complexity and can make your outages much longer than they need to be. We had a big project to get rid of a bunch of queues by just scaling up our synchronous backends and making them faster.
- esafak 3mo agoCare to share more about the issues?
- nostrademons 3mo agoNot OP but also worked on Google Search once upon a time. I'm not sure if I'm remembering the same issues as OP, but basically the two biggest issues are: 1.) What they do to your 95th percentile latency. Users are often very sensitive to tail latency: a service that responds in 150ms 19 times and then takes 2s on the 20th is still perceived as annoyingly slow. With job queues, the reason for that slowness could be as simple as "it was the 20th request to arrive during a period of high demand, and backends couldn't keep up". The whole point of queues is so you can gracefully handle this case without overprovisioning your backends by a factor of 20x, but if the user is still going to consider this a miss anyway, you have to overprovision the backends anyway. There isn't really another way to handle this other than having spare backend capacity. Also note that in many cases the user hitting "refresh" doesn't cancel the existing queued job, it just adds another one to the queue. Which brings us to... 2.) They can turn simple failures into cascading failures. There were several postmortems that went something like "Service X became overloaded because of an unexpected flood of requests, leading to several individual replicas shutting down. This led to more requests being routed to the remaining replicas, which overloaded them too and led to all of Service X going down. When SRE attempted restart Service X, requests queued in the job queue were all retried en masse, which led to an overload of the partially-restarted service and a subsequent failure. SRE had to limit requests upstream and manually drain all job queues and bring Service X back cluster by cluster to restore service health." The root principle here is that any distributed system needs a concept of backpressure. When critical downstream dependencies are overloaded, they need to pass this information back up the stack to the entry point, which needs to start denying requests from the user or do a simpler fallback that doesn't put load on the overloaded service. Naive queueing does not work, because the requests are still sitting there in the queue waiting to overload the downstream service once it becomes available again. You can bolt backpressure onto a job queue system (by eg. rate-limiting requests to a service that has just come back up, or rate-limiting based on response time, and/or falling back to simpler algorithms), but at that point, it's a backpressure system, not a job queue. The semantics are very different from a system that guarantees eventual delivery, just not sure when. You need to be able to handle partial failures and adapt with different algorithms at multiple points within the system.
- theamk 3mo agoI think the the second part was kinda obvious? The moment I read this: > If you’re anything like me, you would probably have said Parallel Spawn, Prefer New, and Wait are perhaps defensible, whereas Prefer Old feels weird/backward. it was pretty obvious I was not anything like him. My intuitive answers are pretty different. - Parallel Spawn is useful, but it's orthogonal to the rest. Even if you have 4 workers, you'll still have to worry about hitting concurrency limit once you have enough jobs. What is it even doing in this list? - Wait is very useful for non-scheduled tasks: if user uploaded 100 files to process, you better process them all. Sometimes you need limit, sometimes you do not (let them queue for a while until devops notices and either allocate more workers or clear them and has some harsh words with consumer). For scheduled tasks, "Wait" seems much less useful. I can come up with a reason but they are all somewhat convoluted - perhaps you are hitting 3rd-party service, and it has a quota, so you've decided to use scheduler to avoid hitting ratelimits? - "Prefer Old" is normally the best way. You repack takes 3-8 hours, so you set your timer to "every 1 hours, skip if running already" and you can be sure that your job finishes and the next one will start. - "Prefer New" seems almost useless. You've already spent all this effort doing the job, why are you cancelling it and throwing away the results? If you want to add a timeout, add a timeout, preferably to specific operation. For example, if there job starts by fetching data, and this fetch can be super slow, use 'Prefer Old' and put a timeout on the fetch part. This way your job won't be interrupted if the fetch just finally succeed minutes before next scheduled interval hit. Oh, and re "If the Prefer Old semantics are not offered, you can’t really emulate them using the two primitives of regular scheduling and limiting concurrency." - you totally can. Set concurrency to 2, and add as a first thing: "fetch the list of jobs running; if there is anyone except me, exit right away". Really, not that tricky at all.
- eru 3mo agoAt Google we actually had 'prefer new' (ie a stack instead of a queue) for certain jobs, that were likely no longer useful after some time had passed; and where less and less useful the longer you waited until you started. One example was running certain ad auctions when rendering websites, or something like that. You don't want to delay serving the side, if the ads are delayed. So you have a certain wall clock time budget until the rest of the page is assembled to be sent to the user, and if you can fit your ad-serving in there, that's good. If you have more work than you can currently handle, then it makes sense to continue with the newest open request after you handled the previous request. The actual system was a lot more complicated, and combined a short, bounded queue on the inside with a large stack on the outside or something like that. Similar considerations can apply, when you are one of many competing market makers for some financial assets on an exchange. Basically, if you have a situation where serving quick is a lot more important than the distinction between late and very late (or even dropping the request).
- zdc1 3mo agoTangential, but when dealing with queues, the first thing you want to do is have a basic grounding of queuing theory, and know whether you're optimising for throughput or worker utilisation (i.e. what are your SLAs and efficiency targets?). IME each goal involves fairly different metrics and scaling rules, so you'll want to know what you're prioritising.
- rienbdj 3mo agoAnyone know a good into to queuing theory?
- mianos 3mo agoThere are lots of good books and some great vids on youtube, but I'd start with this statement and work backward, because this is the non-obvious thing most bootcamp trained, promoted to CTO don't know: The single most important lesson from queuing theory for software systems is the non-linear relationship between utilisation and latency. As system utilisation approaches 1.0 (100% capacity), the average waiting time does not scale linearly, it scales hyperbolically. A system running at 95% utilisation is vastly more fragile and slow than one running at 80%, even though the load difference is minor.
- inigyou 3mo agoTo explain this: if the system is 95% utilized, and a new request comes in, there's a 95% chance the system is already busy and the request has to wait. But after the currently in-progress request finishes, there's still a 95% chance the system is busy (with another request that was already queued behind it) and request X has to wait. After that one finishes, same thing. On average, request X has to wait for about 20 other requests at 95% utilisation - or 5 requests at 80% utilisation - or 10000 requests at 99.99% utilisation. And that's just the mean, not percentiles.
- thaumasiotes 3mo agoYour math doesn't make any sense.
- teleforce 3mo agoThe most popular resource manager for job submission and queueing system is Slurm. It's being used in majority of TOP500 supercomputers, and overwhelming majority of the world HPCs [1]. SchedMD the leading developer of Slurm has recently being acquired by Nvidia, while Slurm remaining free and open source, but somehow it's Wikipedia entry is not yet updated accordingly. [1] Slurm Workload Manager: https://en.wikipedia.org/wiki/Slurm_Workload_Manager https://en.wikipedia.org/wiki/Slurm_Workload_Manager [2] Nvidia Acquires SchedMD (7 comments): https://news.ycombinator.com/item?id=46277190 https://news.ycombinator.com/item?id=46277190
- Unified-Mentor 3mo ago[flagged]
- Unified-Mentor 3mo ago[flagged]
- forrestthewoods 3mo agoNice post. Thanks for sharing.
- dkdbejwi383 3mo agoSeems like OP has conflated the job queue and the job scheduler here.
- groundzeros2015 3mo agoThe UNIX pipe really is an incredible concurrency invention which is not well understood, and attempts to work around its features turn into bugs. A buffer to accumulate data that blocks when it’s full allows you to handle bursty loads. It solves the back pressure problem of readers and writers operating at different speeds. It doesn’t over consume resources. It also solves the architecture problem of when to trigger work. Both consumer and producer act on the pipe imperatively, rather than one end being imperative and the other being a declarative graph of callbacks (all those “reactive” libraries). Even using a term like “back pressure” is a tell to me that someone is confused snd doing something architecturally wrong.
- nh2 3mo agoThe UNIX pipe has the (for many systems) undesirable negative property of losing data in the pipe when the receiving process terminates: Anything in the kernel buffer of the pipe gets lost. Since the pipe is generally unidirectional, the sending process has no way of knowing whether the receiving process has received, or even more successfully processed, anything sent. For that, one needs to make a pipe in the opposite direction and that is not as easy anymore; it also requires building your own protocol to identify and acknowledge sent work items.
- zbentley 3mo ago> the pipe is generally unidirectional So are most message queues. If the producer needs to know when the consumer has finished work, that's not a queue; that's an RPC.
- groundzeros2015 3mo agoAnd to have a hope of solving that problem you need to restore producer and consumer thread states without encountering that problem.
- 10000truths 3mo ago> The UNIX pipe has the (for many systems) undesirable negative property of losing data in the pipe when the receiving process terminates: Anything in the kernel buffer of the pipe gets lost. There are solutions for that, though: * You can observe and/or persist intermediate data with `tee` * You can use a named pipe, whose lifetime is bound to an inode instead of the processes * You can dup the read end of the pipe and pass the extra file descriptor to a crash handler process None of these are esoteric. You do need a basic understanding of Unix primitives, but any half-decent computer science course will cover them.
- Ellis_dev 3mo agoThe schedule interval vs. hard timeout distinction is the useful bit here. Treating them as one setting makes backlog behavior much harder to reason about.
- contentpulse 3mo ago[flagged]
- Izkata 3mo agoThere is 4th option that looks like a combination of how all three are described here: If J1 is running and J2 gets queued, then when J3 gets queued you cancel J2 but not J1. Prefer the oldest running (well, any running, just don't kill them) and newest queued, which waits on the running one. This is how we had long-running test suites configured on Jenkins/Hudson/buildbot ages ago (though it wasn't exactly a queue, more just a flag that the job needed to be run and the job pulled in the latest state).
- rmarshallATD 3mo ago[flagged]
- BedVibe_Studios 3mo ago[flagged]