Turning PostgreSQL into a queue serving 10k jobs per second (2013)
gist.github.com
gist.github.com
One problem with using PostgreSQL in this way (using either advisory locks or LOCK FOR UPDATE) is that it requires you to keep an open connection to the database whilst the job is being worked on.
For a MySQL database, this would be just fine, but PostgreSQL uses a process-per-connection model which caps the number of active connections to the database to a relatively low number (on the order of 1000x fewer connections than a similarly sized MySQL instance) and tools like PgBouncer do nothing to help with this.
As a result, if your jobs take more than a few milliseconds to execute (let's say you make external HTTP requests as part of your job) this is not a good approach to take.
I use a similar approach which avoids this problem, but it only works because I have relatively low throughput requirements. I essentially implement in-database advisory locks using a separate table - before taking a job, workers create a row in the table, and the primary key of this table is used as a worker ID. Jobs are "taken" by assigning them a worker ID. Each row in the worker table has an expiry date, so if workers die, the corresponding row will be deleted and any linked jobs released back into the queue.
As well as transactional guarantees, using a database as a job queue gives you a lot of power over how jobs are executed: for example, our service for delivering webhooks has a separate queue per customer, and we can ensure that within a single queue jobs are processed strictly in order. Meanwhile, our service for search indexing supports different priority levels, so that newly created records are indexed with a higher priority.
Not necessarily. You lock the row for an instant update to a field, for example called "status" into "running" and then disconnect from the database within milliseconds. Finish your job taking as much time as you want. And then connect again to change the row's status to "finished".
This is how it's always designed, as I have seen. The locking problem is for querying rows where status="waiting" and then instantly change it.
It's not to "keep the record locked, and DB connection up, until I finish my batch job". That would be a bad design.
Maybe simply using a timeout per job type is a better way. (That of course trades off simplicity.)
* If the entire worker process dies, then it will lose its Postgres connection which is holding an advisory lock on the jobs being worked. This releases those jobs to be worked by another worker. I don't recall how the built-in retry & back-off mechanisms work in this scenario. This advisory lock is indeed held for the entire time the jobs are being worked on, but only from a single supervisor connection (rather than one connection per job).
* If the job thread crashes, the worker supervisor catches this and the job is marked for a retry at a later time.
[1]: https://github.com/que-rb/que/blob/master/lib/que/migrations...
But come to think of it, these two classes of problems also occur with systems that hold a db lock for the duration of the job. If the worker loses its connection it needs to somehow cancel itself unless the work is idempotent and computation waste doesn't matter. And if the job crashes you need to make sure you release the lock.
Btw to add another related point, databases do have a lock timeout that you have to worry about if you hold a lock for the duration of the job. Your job execution time cannot exceed the lock timeout.
Otherwise, I completely agree on the benefits of transactional job enqueueing. This lets you push off so much complexity until you actually need it (when you’ve scaled such that a single database doesn’t handle your needs well). At that point, you get to solve the same challenges you will have been solving all along to deal with jobs that might run which depend on transactions that may not have committed (yet or ever). I believe this model is almost always the right starting point for a web application, barring some unusual job requirements or massive initial scale.
For example: if a single worker thread crashes but the process doesn't realise it, it may be possible for jobs to become stuck, because other workers rely on the process releasing the job back to the queue. It's definitely solvable but something to be aware of.
Would you run this job queue on the same postgres database as the rest of the application or rather use a different one specific for workers?
You only get this benefit if you’re doing everything in a single Postgres database. By the time you outgrow that setup, you may as well move to a more dedicated job queue or something built on Redis, because you’ve already lost the benefits of having a single system.
I suppose you could maintain the queue in a separate Postgres database, but without sharing ACID with the app I’m not sure What you’d gain vs another system.
FYI, we've been using QueueClassic at Rainforest for a while. Not sure how many actual jobs we do, as we've reset at least when moving hosts, but we're currently at 849202793 jobs (~850m) though QC.
[edit; looking at the last failed job from ~24h ago, we're doing 2.5-3m jobs per day]
It’s changed significantly since this post as the 1.x betas use a very different structure which should actually be more efficient, use fewer Postgres connections, cause less lock contention, and cause less table bloat.
Not sure if the benchmarks have been run recently or not but I’m definitely curious how things stack up to this post from 6 years ago :)
If it is possible to set up a smaller test environment with one or a few instances of your DB and Redis, and you have time to do it, you could do some testing where you purposely shut down each of them at various points in time and inspect what happens to your application state compared to how it behaves when all is well.
Perhaps you could even outfit the testing version of your application with two proxies, one that will proxy the DB connection and one that will proxy the Redis connection, and have these proxies randomly decide within some threshold whether to pass the data along to their upstreams, or to simulate a broken connection. Then test the application with some different threshold values, for example, 1% probability of failure, 5% probability of failure, 50% probability of failure and 100% probability of failure. Repeat the tests some number of times for each threshold value and observe the behavior of the application each time. Also, make sure that you log the decisions made by the proxies each run, so that you are able to look at these afterwards.
queue_classic jobs per second: avg = 1879.1, max = 2072.1, min = 1779.4, stddev = 128.3
que jobs per second: avg = 1500.5, max = 1550.8, min = 1405.2, stddev = 58.1
Full deets: https://github.com/QueueClassic/queue_classic/pull/303#issue...TDLR the lock time grows exponentially depending on the number of dead tuples in the table, which naturally grows as you use long running transactions.
When you open a transaction, you need to have a guarantee that you can touch rows that existed at the moment when transaction has started. You job queue is chugging along and processes let's say a 1000 jobs per minute. Processing a job involves deleting the row from the queue, but since you have a transaction running, Postgres only marks the row as deleted, and keeps it around in case a transaction would want to access it at some point. Each time you need to process a job, Postgres needs to lock the row. The way this mechanism works involves iterating over the rows until you find one you can lock on. If your transaction is running 30 minutes, each job would have to iterate through 30k dead rows (deleted, but still around for the sake of the transaction). Slowing down lock time leads to overal degreaded performance of the job queue, which leads to jobs being added faster than they're being processed, which further exacerbates the problem
Guess you lose job atomicity that way.
In particular, Que’s new design locks jobs in a single connection per worker process, not a connection per-job. Job assignments are handled in an in-memory/unlogged table. Jobs are also automatically assigned to a free worker upon being enqueued via LISTEN/NOTIFY and an unlogged table of available workers. And jobs are locked in batches by each worker process, not one-at-a-time. So polling becomes much less frequent and is far more efficient when it happens.
This does not mean that it’s suddenly ok for jobs to take out transactions which run a long time (that is usually a bad idea in a production database) but it does substantially minimize the scenarios where these problems might occur.
I'm glad the state of things has improved, since I really love database-backed job queues and all transactional guarantees it gives you.
Percona has a good article on this https://www.percona.com/blog/2018/08/10/tuning-autovacuum-in...
That said, using a common persistence store initially makes sense, but trying to compartmentalize the jobqueue/batch-processing stuff never hurts.
One place this has come up for me is when the next job that's picked depends on currently running jobs, e.g., each job is associated with a user, and if a single user already has N tasks running you may want to prioritize another user's tasks for the N+1 slot.
I have seen benchmarks reaching millions of messages per second.
I had no difficulty with long-running jobs because servicing jobs out of PG was simply a matter of pushing them onto the Kafka queue for immediate uptake there.
> So, many developers have started going straight to Redis-backed queues (Resque, Sidekiq) or dedicated queues (beanstalkd, ZeroMQ...), but I see these as suboptimal solutions - they're each another moving part that can fail, and the jobs that you queue with them aren't protected by the same transactions and atomic backups that are keeping your precious relational data safe and consistent. Inevitably, a machine is going to fail, and you or someone else is going to be manually picking through your data, trying to figure out what it should look like.
I disagree though. BLPOP is far easier to grok than any Postgres solution and Redis is rock-solid. Either using the new Redis streams or a good queueing library is going to guarantee you don’t miss any jobs and have a need for transactions across both systems either.
I would also be very hesitant to add work to my database. That’s often a sensitive part of systems.
Edit: oh, someone mentioned this being from 2013. No Redis Streams back then.
I've lost gigabytes of data with Redis; nothing with Postgres. I run both in production.
Redis has a lot of obscure failure modes, lacks transactions (with rollback), and has no query language. For simple stuff where it's okay to lose data regularly, it's fantastic (that's when we use it).
That said, we've been doing more and more with Postgres over time, and less and less with Redis. The benefits of having all of your data in the same storage, transactionally consistent, with a query planner for ad hoc visualizations and reporting is just too great.
It's well designed so that when you do venture into territory where you need a specialized solution it's not difficult to replace. But you'll definitely miss ACID and relational data if you're leveraging them.
There is a great operational economy to making Postgres your default solution for these sorts of problems.
> many developers have started going straight to Redis-backed queues (Resque, Sidekiq) or dedicated queues (beanstalkd, ZeroMQ...), but I see these as suboptimal solutions - they're each another moving part that can fail, and the jobs that you queue with them aren't protected by the same transactions and atomic backups that are keeping your precious relational data safe and consistent
Edit: I'm rate limited so to elaborate a bit.
Queueing technologies are not "new interesting bit of technology", they're tailored solutions to solve a specific problem.
Using a RDBMS for queueing is not what it was built to do, and you will run into issues doing so (I have).
By trying to use one tool for everything, you lose out on all kinds of optimizations, features, and performance enhancements that are specific to the problem you're trying to solve.
It's bad engineering to try and force Postgres into the role of queue when vastly superior technologies exist. You're actively hurting the engineers you work with, and the company you work for, if you force the same tech into use cases it isn't optimal for.
I cannot stress this enough; fear of learning is anathema to software, and trying to hide it behind a veneer of caution is not only disingenuous but potentially malicious as well, maximizing exclusively for the benefit of the individual against the interests of the group.
Fighting the tendency of everyone wanting to draw in each new interesting bit of technology is important.
Every piece of complexity you incur should be proven as required, and you should build things to just somewhat larger than near term projected scale... and not expect that you are a temporarily embarrassed unicorn with loads of unexpected demand showing up.
If you run out of performance, you have some headroom from tuning; from that point, you can then choose whether it makes sense to retrofit to a higher scale technology or throw larger hardware at the problem (or both).
And then you have a component that is designed for the job instead of trying to use a database as a poor man's queue.
It's still another moving part. Thing should be as simple as they can be, but no simpler.
> trying to use a database as a poor man's queue.
Of course, it's not really a "poor man's" queue-- it's got some superior capabilities. It just loses on top-end performance. (Of course, using those capabilities is dangerous, because it creates some degree of lock-in, so go into it open-eyed).
For as much as you accuse others of looking down their nose / not willing to seriously consider other technologies... you seem to be inclined that way yourself.
Redis is great. But if you have Postgres already, and modest to moderate queuing requirements, why add another piece to your stack? Postgres-by-default is not a bad technology sourcing strategy.
It comes from a place of ignorance, and you're promoting ignorance. Learn why technologies exist and make an informed decision about tradeoffs, instead of being lazy and incompetent by blindly choosing technology based on what's a very locally maxima for you personally. It's selfish and damaging.
Edit since I'm rate limited: The folks landing on Postgres are not landing there after due consideration, they're landing on a technology they're familiar with because they don't know how to learn.
It is lazy, ignorant and selfish, and while of course people don't like having truth spoken to them, that's what it is; truth.
Edit 2: Making these poor technology choices will kill your business because you won't have the agility that using a specific message broker technology gives you. You won't have built-in solutions to common queueing problems, you won't have dedicated specific logging/metrics to monitor, you won't have the libraries in your preferred language to directly tackle your problem, you won't have the support community available to you (it will be much smaller), you won't be able to pivot onto related patterns as your needs change, your scaling will always be more complicated because fewer people are doing it and the tool you use isn't tailored for your use case; the list goes on. Your competitors will swallow you because they move faster than you do, and your business will die.
This is annoying to me because I've been in situations where I've had to maintain and write features against technology that was a poor fit for its use, but people like those in this comment sections bitched and moaned about a "new thing" existing in our stack. The reality was they didn't want to learn anything, and were fine pushing the hard work off onto the developers, so they could safely continue to do as little as they possibly could get away with.
Do your job, learn technology that actually fits your use case, and stop trying to push work off onto other people.
It's not lazy, blind or selfish as you claim - and frankly when you disagree with authority & condescended tone - it's a turn off.
I can't stress this enough; claiming anything is "vastly superior", especially with no examples, isn't useful - things aren't this black and white in reality. Most tech-choices like this are somewhere between trade-offs, preferences, ignorance or wrongly held opinions about something being vastly superior > something else.
You're applying generic heuristics against a problem where folks have specific domain knowledge that contradicts those heuristics.
I cannot stress this enough; you are flat wrong if you think Postgres is appropriate to use for queueing.
Nothing is likely be correct about your arguments if you don't know what you are trying to fix or the sitation, which should also be blatantly obvious. You're applying some unstated subjective situation in your head, then dumping out something you've heard / done before. LOL.
Also "if you've ever used both" - I have, which is why I know. It's gray; sometimes you should, sometimes you shouldn't. You're assuming I haven't as I have differing point of view to you.
If you already have tooling and knowledge for managing Postgres, and Postgres is sufficient for your queueing needs, then why would you incur the cost of introducing a new service? Most situations don't require the best tool for the job, and in that case, you are well served to use a tool you already have.
If you are using an abstraction layer over your queues (you almost certainly are) then you should be able to easily change to a different queue implementation later, should the need arise.
This has nothing to do with getting shit done, and everything to do with people who don't want to learn new things because they're afraid of not being smart enough to make the new thing work.
If a PG-based solution meets the needs then this reduces complexity in the large. And this has nothing to do with learning or not learning.
Sorry but your argument just does not hold up.
The only argument that doesn't hold up is "use one tool for everything". This ideology is for lazy people who want to anchor themselves in a period of time, and refuse to learn new things.
Agree on devops. Of course.
There’s still the security surface area argument and the overall complexity argument—complexity of the system of systems is NOT reduced by adding technologies.
By way of analogy, using military aircraft, take the F35 vs. the A-10. The F35 is insanely more complex than the nearly indestructible A10 in large part due to the massive number of subsystems, software, sensors, etc. Sure the F35 way more capable (this is NOT a head to head argument about the military capabilities or said aircraft) than the A10, but it also has orders of magnitude more failure modes. More parts. More complexity. More opportunities for failure. This is engineering fact.
What's frustrating about this conversation is that I'm literally, right now, supporting two different queueing systems based on PG and Redis, so I get on a very real level, the tradeoffs. I know in great detail the problems that come up, but HN is not conducive to talking at that level of detail. At this point every comment I make is flagged and downvoted, so why would I pay time into a system that has clearly decided my opinion isn't relevant?
You're stuck in PG because your DevOps team is a bunch of incompetent dolts, so yeah, make the best out of the bad situation and go ahead and use something like Que.
Just don't develop Stockholm Syndrome while doing it.
People in this thread are protecting their own egos by expressing how great PG is as a message broker, but it's harmful to the industry to let that kind of attitude perpetuate.
So if my language doesn't have a library like Que ( no thanks using Ruby ) I can't do anything with PG whereas using Redis or any other solution you can use "subscribe" / "publish" ect ...
Like many things in this space there are a number of good ways - each with minor trade offs
Others posters have detailed those trades very clearly.
So now we can all be more educated about which one of the six good ways we'd choose.
This is wrong.
You are ultimately getting paid to solve the company's problems. If those problems actually require deployment of a new technology to satisfy the requirements, then it is fine. Note that I said satisfy requirements, not exceed them.
But more often than not, * that's not the case *. If your queue requirements can be served with PostgreSQL, and you already have it, then why not? To do otherwise would be over engineering.
Because you may know a new and fancy technology that would be objectively better. Cool. Now, do you have enough know-how in the company to use it effectively? Do you know what the best practices are? Can you deploy, monitor, audit and patch this in production? If you have answered NO to any of those questions, you'll either hire people or you shouldn't deploy. Period.
And if you have only one person answering YES to the previous questions, you still shouldn't deploy it. Because one day that person will leave, and you are now left the company in a worse situation than it was before. For no reason other than satisfying your engineering itch.
THAT is anathema to working, reliable production software. Toys, you can do whatever.
If you want to introduce something new, you can. But you need to do it responsibly. It needs a reason to exist - a valid reason, not "it's better". How does it being "better" help the company? Will it reduce costs? Maintenance? Future development work will happen faster? Will it make a faster user experience? Scale to the projected company growth?
It will have a cost that will have to be accounted for – including opportunity costs.
> By trying to use one tool for everything, you lose out on all kinds of optimizations, features, and performance enhancements that are specific to the problem you're trying to solve.
You only care about this if you have a reason to care about this. Otherwise it's all irrelevant.
The reason you are completely wrong and a toxic member of your team is that your logic can be used to justify all kinds of terrible, short term wins that sacrifice any kind of long term/sustainable design. You'd rather make a bad technology play so your boss gets off your back, than think through your decision and make a long term investment into a sustainable, maintainable, improvable, and agile technology that's actually suited to solve the real problem you have.
OF COURSE you're there to accomplish the tasks your business needs to succeed. That's a given. What's not a given is that you need to slap together the shittiest piece of software that'll do that job right then and there.
Your attitude is what keeps companies from solving problems, and forces them to pay dollar after dollar to fix the same issue over and over again. You're selfish, you're greedy, and you cost your company insane amounts of money.
You should not be employed in this industry if what you wrote here is how you think, period.
I put out what is given.
Would you please pick one account to post with? Using multiple accounts to work around moderation restrictions is obviously abusive.
As for the other comment you made (can't reply directly, rate limit), it isn't a low information post, I'm showing a screenshot of the message I get when I hit the rate limit.
You're literally gas lighting right now, it's insane.
You also know if you actually do anything else beyond rate limiting I'll just go dark, which is even "worse" for you, so this so-called "leniency" is nonsense. Just stop.
Edit: You can't use multiple accounts to get around rate limits, as you and I have discussed multiple times, so I'm not "doing" anything that constitutes a ban, but feel free to ban at your leisure. Over the years you've established yourself firmly as an unreasonable person, so there's no point trying to treat you like one.
Rate limits are meaningless if people can just use multiple accounts to get around them. If you keep doing that, we're going to have to ban at least one of your accounts—but more likely all of them. After several years of trying to persuade you to use HN as intended, I'm beginning to lose patience.
That seems pretty personal. Not that I'm hopping on dang's bandwagon, but you can't post that sort of attack and say you're not doing anything like that.
I've wound up seeing a lot of your recent posts and I can understand where a lot of your responders are coming from... you are very strongly opinionated and have no problem crossing the line from "I hold this opinion" to "you're wrong for not agreeing and therefore you <insert consequence>".
THAT's the behavior that's getting you the pushback. Be more open to discussing things and comparing experiences and maybe you'll see more cordial response back.
I am aware, generally, how things are. My complaining here is not confusion, it's frustration. I know I can change my behavior to elicit a better response, but I'm frustrated that it's relevant. People should have thicker skin, and when people dish it out (as what was going on here), I should be allowed to give it back. Dang's not commenting on anyone else's posts here, though there were many other rule violators.
It's an unfair application of the rules.
Edit: This post is 0 minutes old, how does it have a downvote already? This is the shit I'm talking about...
> A queue implemented in the RDBMS will never match the performance of a fast dedicated queueing system, even one that makes the same atomicity and durability guarantees as PostgreSQL. Using SKIP LOCKED is better than existing in-database approaches, but you’ll still go faster using a dedicated and highly optimised external queueing engine.
People only realize these warnings when they have built a toy system that is rolled out to a production/product that takes off. By that time you will realize "ohhh I don't need all the ACID guarantees for every job in my system", and I want to run my workers as lambdas/elastic scaling workers (you need more connections). That is where companies will be spending efforts from their best engineers to move away from Postgres as queue. Which comes down to question; why do it in first place? What makes it cheaper (other than local dev) to deploy a system that is bound to fail in future? Don't get me wrong I love Postgres, I just don't believe that it's the right tool for doing this job.
The great thing is that it is a mature tool and you don't need to bring another beast to the zoo.
There are different performance constraints to this approach, but data integrity and robustness of it are unmatched, really.
The only case where this would be problematic is if the worker crashed after doing some side effect external to the queue server that rendered the job non-idempotent, but before committing an acknowledgement. But that's not really a “what you use as a queue server” issue as a “are you also using a distributed transaction system to coordinate all side effects including queue updates” issue. (If the same Postres DB is the operational DB and the queue server, you may be able to avoid distributed transactions by having all the side effects in the DB, though that creates its own issues.)