SQL Maxis: Why We Ditched RabbitMQ and Replaced It with a Postgres Queue
prequel.co
prequel.co
It ultimately doesn't matter because of the low volume they're dealing with, but gang, "just slap a queue on it" gets you the same results as "just slap a cache on it" if you don't understand the tool you're working with. If they knew that some jobs would take hours and some jobs would take seconds, why would you not immediately spin up four queues. Two for the short jobs (one acting as a DLQ), and two for the long jobs (again, one acting as a DLQ). Your DLQ queues have a low TTL, and on expiration those messages get placed back onto the tail of the original queues. Any failure by your consumer, and that message gets dropped onto the DLQ and your overall throughput is determined by the number * velocity of your consumers, and not on your queue architecture.
This pg queue will last a very long time for them. Great! They're willing to give up the easy fanout architecture for simplicity, which again at their volume, sure, that's a valid trade. At higher volumes, they should go back to the drawing board.
In my experience, the benefits of a SQL table for a problem like this are real big - easier to see what's in the queue, manipulate the queue, resolve head-of-queue blocking problems, etc.
Your multiple queue solution might work but it is most efficient to have just one queue with a pool of workers where where each worker doesn't pop a job unless it's ready to process it immediately. In my experience, this is the optimal solution.
I can probably bolt some of these properties onto a queue that doesn't support all the features I need.
Did you do long running jobs like they did? It's a stereotype, but I don't think they used the technology correctly here -- you're not supposed to hold onto messages for hours before acknowledging. They should have used RabbitMQ just to kick off the job, immediately ACKing the request, and job tracking/completion handled inside... a database.
It did take some configuring to get it working. Between acking appropriately and the prefetch (qos perhaps? Can’t remember, don’t have it in front of me). We were able to make it work. It was pretty straightforward it never even crossed my mind that this isn’t a correct use case for RMQ.
(Used the Java client.)
Assume that your producers will be able to spike and generate messages faster than your consumers can process them. This is normal! This is why you have a queue in the first place! If your jobs take 5 seconds or 5 hours, your strategy is influenced by the answers to those three questions. For example -- if you're willing to drop a message if a consumer gets power-cycled, then yeah, you'd immediately ack the request and put it back onto a dead letter queue if your consumer runs into an exception. Alternatively, if you're unwilling to block and you want to be very tolerant of consumer failure, you'd fan out your queues and have your consumers checking multiple queues in parallel. Etc etc etc, you get the drift.
Keep in mind also that this isn't specific to RabbitMQ! You'd want to answer the same questions if you were using SQS, or if you were using Kafka, or if you were using 0mq, or if you were using Redis queues, or if you were using pg queues.
I’d like to hear why people chose Kafka over some RDBMS tables.
Kafka has specific use cases but it seems a lot of people just go "ok use Kafka here" and wait for the load that rarely comes
If your load is, say, a few hundred writes/second, stick with the database only, and it will be much simpler.
If the answer is a call to a shared database, you might as well not have RabbitMQ.
You can build a queuing system with a database, but you have to do that. Some of the features and constraints of the database might even make your life harder than it has to be.
Instead, view it like that: there is a need for a queuing system and a job system. Either or both can be implemented using a database for certain concersn, but it can also be a custom implementation. It's not a great idea to mix the two things unless the operational and infrastructure costs and complexity outweigh the benefits of a clear separation.
I'm not saying that libraries like pg-boss and co. cannot sometimes replace a full queue implementation. But the tradeoffs need to be clear.
RabbitMQ is completely happy to do this. On the other hand, literally every python client library would fall over and/or wedge forever periodically. Celery in particular was a flaming pile of garbage (at least circa 2018).
I wound up vendoring the best of the bunch (amqpstorm, fyi), and adding a bunch of additional logic to handle all the various corner cases. That, plus a bunch of careful thought about how the queues are structured, and I've been pushing millions of messages a day without issue.
I wasn't able to get a lot of my changes upstream due to the extent I had to alter the assumptions the library made about control flow.
That's what I think too. Use the message queue for communication and potentially work distribution and the database for the rest like internal state.
In my experience, RabbitMQ isn’t a good fit for long running tasks. This was 10 years ago. But honestly, if you have a short number of long running tasks, Postgres is probably a better fit. You get transactional control and you remove a load of complexity from the system.
yeah I've seen 3 different workplaces run into this exact issue, usually when they started off with a default Django-Celery-Redis approach
all of those cases were actually easily fixed with Postgres SELECT FOR UPDATE as a job queue
Our problem was that when a consumer hung on a poison pilled message, the prefetched messages would not be released. We fixed the hanging, but hit a similar issue, and then we fixed that, etc.
We moved to SQS for other reasons (the primary being that we sometimes saturated a single erlang process per rabbit queue), but moving to the SQS visibility timeout model has in general been easier to reason about and has been a better operations experience.
However, we've found that all the jobs are in postgres anyway, and being able to index into our job queue and remove jobs is really useful. We started storing job metadata (including "don't process this job") in postgres and checking it at the start of all our queue workers and we've decided that our lives would be simpler if it was all in postgres.
It's still an experiment on our part, but we've seen a lot of strong stories around it and think it's worth trying out.
You do want to make sure to set the VM high water mark below 50% of RAM as the GC phase can double the used memory. If high water mark is set too high, the box will thrash swap really badly during GC and hang entirely.
Also, sharing queues for multiple types of jobs just feels like frustration waiting to happen.
I'd estimate more like dozens to hundreds per second should be pretty doable, depending on payload, even on a small DB. More if you can logically partition the events. Have implemented such a queue and haven't gotten close to bottlenecking on it.
Rabbit is a phenomenal tool but you need to know how to use it.
Here is why I would not recommend that.
Do that and you have to rewrite your system around predictions how long each job will take, deal with 4 sources of failure, and have more complicated code. All this to maintain the complication of an ultimately unneeded queueing system. You call it an "easy fanout architecture for simplicity." I call it, "an entirely unnecessary complication that they have no business having to deal with."
If they get to a volume where they should go back to the drawing board, then they can worry about it then. This would be a good problem to have. But there is no need to worry about it now.
What do you mean by "broken"? Are you implying that the behavior they're describing is not the way the consumer library is supposed to work? They linked to RabbitMQ's documentation basically saying that's exactly how it works. Also, where do you get the sense that they've misconfigured it? You've made these statements, but did not exactly enlightened us as to how one should set things up to have consumers handle exactly one job at a time. That was their only problem (Edit: an answer by @whakim https://news.ycombinator.com/item?id=35530108 is providing more light on this).
The rest of your answer sanctimoniously presumes that they don't know how to use the tool, but your own proposed solution is moot, as it seems to address a different problem, not the one that they have (1 job per consumer max).
RabbitMQ jobs is to handle message transfert, if you tie business logics (jobs completion) to its state, you have a problem.
So basically, producer and consumer MUST have their own storage for tracking processes, typically by ID, RabbitMQ jobs being the message handler for synchin' states.
how do you deal with worker getting killed in that scenario ? if you can't rely on the queue job state then you need another whole set of code somewhere to handle timeouts & retries, don't you ?
A worker, or at least, the consumer receiving the messages and spawning sub workers, when getting killed, should restart at some point, inspecting jobs ids unfinished, the one stored just before ack'ing reception, and either resume work or notify failure.
Answering is the worker sole responsibility, and no viable implementation of an app can reliably try to substitute to that.
As for timeouts, they are necessary when waiting for answers from services that are not under your control, eg, you cannot assume they will answer you in all cases. Under your control, you have to keep yourself in a situation where worker response is guaranteed, success or failure.
At this point, why not throw the whole job queue away and simply use that system ?
Having the jobs processing itself being independent of the messaging solution seems like good way to go as the actual message implementation, might have to change. I really think that you don't want ossification of your messaging implementation detail on your business logics.
I did not pretend that it was the only solution, merely just something I know would be reliable, an implicit «one way to it».
I really was focusing on low volume, long jobs, where RabbitMQ is overkill, and workers respawn does not need proper dequeuing.
One aspect to keep in mind is also the kind of use case, when the job is just a subtask, and the final result is stored, on completion it produces as message to be handled for continuation, you have direct id tracking without further additional logics, and a simple ack is not enough.
For eg, sending emails, the jobs does not create further message, and no shared mutable state occurs. The scenario with ack on end of process covers all cases.
You’re right that if the business-processing is a long heavy-weight operation then that should be decoupled from the message queue and failure handled internally, in which case your definition does make sense as an intermediary hand-off step. Just in the general case, I think most people wouldn’t consider hand-off alone as processing.
At some point using what you understand is easier
Monitor by having some health check evaluating that no row stays in PROCESSING for too long, and that no row stays NOT_STARTED for too long, etc. Introspect by making a nice little HTML screen that shows this work queue and its states.
As I wrote in another comment, this is somewhat similar to a "work queue pattern" I've described here: https://mats3.io/patterns/work-queues/
If you your needs aren't that expensive, and you don't anticipate growing a ton, then it's probably a smart technical decision to minimize your operational stack. Assuming 10k/jobs a day, thats roughly 7 jobs per minute. Even the most unoptimized database should be able to handle this.
Scaled much, much further than I would’ve guessed at the time when I called it a short-term solution :) — now I have much more confidence in Postgres ;)
That's very refreshing to hear. In a previous role I was in a similar situation than yours, but I pushed for RabbitMQ instead of postgres due to scaling concerns, with hypothetical seilings smaller than the ones you faced. My team had to make a call without having hard numbers to support any decision and no time to put together a proof of concept. The design pressures were the simplicity of postgres vs paying for the assurance of getting a working message broker with complexity. In the end I pushed for the most conservative approach and we went with RabbitMQ, because I didn't wanted to be the one having to explain why we had problems getting a RDBMS to act as a message broker when we get a real message broker for free with a docker pull.
I was always left wondering if that was the right call, and apparently it wasn't, because RabbitMQ also put up a fight.
If there were articles out there showcasing case studies of real world applications of implementing message brokers over RDBMS then people like me would have an easier time pushing for saner choices.
You mean "industrial scale RDBMS" that you can license for thousands of dollars? No, you can't really implement message brokers on those.
You will never see those showcase articles. Nobody paying wants them.
Almost everything you see on how to use a DBMS is an amateur blog or one of those studies. One of those is usually dismissed on any organization with more than one layer of management.
Your comment reads like a strawman. I didn't needed "studies". It was good enough if there was a guy with a blog saying "I used postgres as a message broker like this and I got these numbers", and they had a gitlab project page providing the public with the setup and benchmark code.
I'm interested in hearing more about this (making a similar decision right now!). What pains did RabbitMQ give you?
Really depends on the needs but this can unlock some very impressive and sustainable throughputs.
Can you elaborate on how do you do this?
https://github.com/bensheldon/good_job/blob/10e9d9b714a668dc...
That's heavily inspired by Rail's Action Cable (websockets) Adapter for Postgres, which is a bit simpler and easier to understand:
https://github.com/rails/rails/blob/be287ac0d5000e667510faba...
Briefly, it spins up a background thread with a dedicated database connection and makes a blocking Postgres LISTEN query until results are returned, and then it forwards the result to other subscribing objects.
https://www.postgresql.org/docs/15/sql-listen.html
https://www.postgresql.org/docs/15/sql-notify.html
I think they have been around for ages, but handling the LISTEN responses may need special client library support.
Then, whenever you write a new job to the queue, you also do a NOTIFY on the same channel.
This lets you keep latency low while still polling relatively infrequently.
NOTIFY is actually transactional which makes this approach even better (the LISTENer won't be notified until the NOTIFY transaction commits)
No problem with millions of enqueue+dequeue per day.
A table for a queue is also going to be so tiny that postgresql might even outdo my own expectations.
Anybody had any success running a queue on top of... sqlite?
With the way the sqlite file locking mechanisms work, are you basically guaranteed really low concurrency? You can have lots of readers but not really a lot of writers, and in order to pop a job off of the queue you need to have a process spinning waiting for work, move its status from "to do" to "in progress" and then "done" or "error", which is sort of "write" heavy?
> An EXCLUSIVE lock is needed in order to write to the database file. Only one EXCLUSIVE lock is allowed on the file and no other locks of any kind are allowed to coexist with an EXCLUSIVE lock. In order to maximize concurrency, SQLite works to minimize the amount of time that EXCLUSIVE locks are held.
> You can avoid locks when reading, if you set database journal mode to Write-Ahead Logging (see: http://www.sqlite.org/wal.html).
But that was ... gremlins. More than one person tried to figure out wtf was going one and eventually it was better business wise to declare 'gremlins' and everybody involved in the incident has been annoyed about not figuring it out since.
SQLite also supports in-memory databases with shared cache.
https://www.sqlite.org/inmemorydb.html
If a message queue does not require any form of persistence, writes don't sound like an issue.
Sometimes this is fine.
On modern hardware you should be able to trivially handle more than that.
I don’t know how many messages per second it does but for a podcast crawling side project I have processed hundreds of millions of messages through this little Python wrapper around SQLite. Zero problems. It just keeps running happily.
When you’re writing billions of messages per day, I don’t see how a file system scales.
But also keep in mind every executable has at minimum its own executable as an open file. Picking a random python process I currently have running, lsof reports it has 41 open *.so files.
There's already a spec for transferring SOAP over SMTP so the request/response part is already built for you.
(couldn't resist, sorry)
That is, without the custom mail server you describe. You can feed every incoming message to a custom command with ".qmail", or forward with ".forward", so we used that for "push" delivery of messages and plain old POP3 for pull. E-mail really does have everything you need to build complex routing topologies.
We ran a mail provider, and so had a heavily customised Qmail setup, and when we needed a queue we figured we might as well use that. Meant we could trivially do things like debugging by cc:'ing messages to a mailbox someone connected a mail client to, for example.
However, I'd hope the SOAP part made it clear I was exaggerating a little for effect :)
... also because the idea of the custom server meaning you could have the main system's qmail handle spaced retries based on SMTP status codes amused me.
Our main CGI (yes...) written in C++ (I can hear you squirm from here) was also loosely inspired by CORBA (stop that screaming) in terms of "sort-of" exposing objects via router that automatically mapped URLs to objects, and which used XML for persistence and an XML based rendering pipeline (not that different from React components, actually, except all C++ and server side).
Hit all the 90's buzz words.
Plus lots of NFS for the shared filestore that the qmail SMTP nodes and the courier/qmail-pop3d mail receipt nodes mounted.
Plus ... yeah, you can see why I thought we'd not find each others' setups -too- surprising.
So, no, not going to squirm, because I mean, yes, I know, but it all (mostly) worked and the customers weren't unusually unhappier with us than they are with any provider ;)
Interestingly, journaled rollback mode in SQLite is probably worse durability than the file/directory queue method. (https://www.sqlite.org/lockingv3.html#how_to_corrupt) (https://www.sqlite.org/atomiccommit.html#sect_9_0) It depends on the same system-level guarantees about the filesystem and hardware, except it also depends on POSIX advisory locks, which are fallible and implementation-specific; while the file/directory queue solely depends on rename() being atomic, which it always is.
No, it doesn't. Once COMMIT returns, your data is durable.
> You can use exactly the same method that SQLite uses to flush data to disk.
Good luck with that! It took them about ten years to get it right, and the SQLite people are world class. (You can't just copy their code wholesale, since they have one file to worry about and your scheme has tons.)
> With a journaled filesystem, fdatasync() after each atomic operation, and battery-backed RAID, I can't imagine it getting more durable [in comparison to SQLite].
The problem is lack of imagination. :-) See https://danluu.com/deconstruct-files/ (and the linked papers) for an overview of all the things that can go wrong. In particular, if you only only ever fdatasync(), your files are not even guaranteed to show up after a power loss.
I've not run it myself in production but I would definitely have heard if it went wrong, I think OpenSuSe's OpenQA thing does build worker boxes that way (though of course they'll have fairly slow jobs so it may just be the write throughput doesn't matter).
This being HN, I'd like to point out you can steal the SQL from the source for your own uses if you want to, this just happens to be the example full implementation I'm familiar with.
I've seen systems at that scale where that's roughly true. But I've also seen systems where those jobs come in a daily batch, at a point in time of day, and then nothing until the next day's batch.
Postgres could still handle that though, IMHO :)
A million requests a day sounds really impressive, but it’s 12req/s which is not a lot. I had a project that needed 100 req/s ages ago. That was considered a reasonably complex problem but not world class, and only because C10k was an open problem. Now you could do that with a single 8xlarge. You don’t even need a cluster.
10k tasks a day is 7 per minute. You could do that with Jenkins.
Brandur wrote a great piece about a related pattern here: https://brandur.org/job-drain
He recommends using a transactional "staging" queue in your database which is then written out to your actual queue by a separate process.
I built a system ages ago that had modest queue needs.. maybe 100 jobs a day. It involved syncing changes in the local database with external devices. Many changes would ultimately update the same device, and making the fewest number of updates was important.
The system used an extremely simple schema: A table with something like [job_id, device, start_after, time_started, time_finished]
When queueing a job for $device, do an upsert to either insert a new record, or bump up the start_after of a not yet started job to now+5 minutes. When looking for a job to run, ignore anything with a start_after in the future.
As edits were made, it would create a single job for each device that would run 5 minutes after the last change was made.
I know a lot of queueing systems have the concept of a delayed job, but I haven't come across any that had the concept of delayed jobs+dedup/coalescence.
We utilise a decorator for our job addition to external queues, such that the function that does the addition gets attached to Django's "on transaction commit" signal and thus don't actually get run until the outer database transaction for that request has been committed.
A possible solution to this is to use a "transactional outbox" pattern, but that has many of the same drawbacks of using postgres as a queue.
That it's beneficial to use Postgres messaging (or Oracle AQ or whatever) for its transactional semantics is kind of accidental and a consequence of folks not wanting to bother with dtx. Even though databases are accessed via networks, truly scalable work distribution can't be achieved using SQL, much less with SQLite. Or in other words, if you're using messaging queues in databases, you could use tables and row locks directly just as well.
Keeps things simple and reliable.
The only drawback we saw was that there will be few second of latency as system scale; if that doesn't matter
- You probably want FOR NO KEY UPDATE instead of FOR UPDATE so you don't block inserts into tables that have a foreign key relationship with the job table. [1]
- If you need to process messages in order, you don't want SKIP LOCKED. Also, make sure you have an ORDER BY clause.
My main use-case for queues is syncing resources in our database to QuickBooks. The overall structure looks like:
BEGIN; -- start a transaction
SELECT job.job_id, rm.data
FROM qbo.transmit_job job
JOIN resource_mutation rm USING (tenant_id, resource_mutation_id)
WHERE job.state = 'pending'
ORDER BY job.create_time
LIMIT 1 FOR NO KEY UPDATE OF job NOWAIT;
-- External API call to QuickBooks.
-- If successsful:
UPDATE qbo.transmit_job
SET state = 'transmitted'
WHERE job_id = $1;
COMMIT;
This code will serialize access to the transmit_job table. A more clever approach
would be to serialize access by tenant_id. I haven't figured out how to do that
yet (probably lock on a tenant ID first, then lock on the job ID).Somewhat annoyingly, Postgres will log an error if another worker holds the row lock (since we're not using SKIP LOCKED). It won't block because of NOWAIT.
CrunchyData also has a good overview of Postgres queues: [2]
[1]: https://www.migops.com/blog/2021/10/05/select-for-update-and...
[2]: https://blog.crunchydata.com/blog/message-queuing-using-nati...
Skip locked main utility is to provide work queue capability to PG. Even documentation of PG is referring to skip locked as a main driver for work queues
SELECT job.job_id, rm.data
FROM qbo.transmit_job job
JOIN resource_mutation rm USING (tenant_id, resource_mutation_id)
WHERE job.state = 'pending'
AND job.job_id = (
SELECT qj2.job_id
FROM qbo.transmit_job qj2
WHERE job.tenant_id = qj2.tenant_id
AND qj2.state = 'pending'
ORDER BY qj2.create_time
LIMIT 1
)
ORDER BY job.create_time
LIMIT 1 FOR NO KEY UPDATE OF job SKIP LOCKED
> you should just use TemporalDo you mean the Airflow-esque DAG runner? I prefer Postgres because it's one less moving part, I like transactional guarantees, my volume is tiny, and I can tweak the queue logic like above with a simple predicate change.
You don’t need a queue, a database, a blob store, and a cache. You just need Postgres for all of these use cases. Once your project scales past what Postgres can handle along one of these dimensions, replace it (but most of the time this will never happen)
It also does wonders for your uptime and SLO.
It becomes impossible to ever lose in-flight data. The moment you persist to your queue you can ack back to the client.
If anything, I prefer to use ZeroMQ and make sure everything can recover from an outage and settle eventually.
To ingest large inputs, I would just use short append only files and maybe send them over to the other node over ZeroMQ to get a little bit more reliability, but rarely are such high volume data that critical.
There is nothing like free lunch when talking distributed fault tolerant systems and simplicity usually fares rather well.
Relational databases also have this feature.
You mean like... ACID transactions?
https://www.postgresql.org/docs/current/tutorial-transaction...
Use whatever you like to actually implement the queue, Postgres wouldn't be my first or second choice but it's fine, I've used it for small one-off projects.
Are y'all not running your stuff highly available?
The answer to what happens if the queue goes down is that it runs highly available.
For way more than 99% of the people, all of that should be irrelevant.
If you push a change to your frontend server that does all the work itself and it breaks you have a service interruption. That same system based on queues it's a queue backup and you have breathing room.
If you're doing a longer operation and your app server crashes that work is lost and the client has sit there and wait, in a queue system you can return almost immediately and if there's a crash it'll get picked up by another worker.
Having a return queue and websockets let you give feedback to the client js in a way that is impervious to network interruptions and refreshes, once you reconnect it can catch up.
These aren't theoretical, this architecture has saved my ass on more occasions than I can count.
I'd say it's more like:
- 95.0%: SQLite
- 4.9%: Postgres
- 0.1%: Other
I'm a bit late to the party (also, why is everyone standing in a circle with their pants down) but, does no one care about high availability with zero data loss failover, and zero downtime deployments anymore?
(and litestream exists for sqlite if needed)
I want to use Postgres for JSON, I know it has specific functionality for that.
But still, what do you mean by that and does it apply to JSON why or why not?
But do you think postgres latency is actually going to be fine for many things people use redis or memcached for?
Fixing that really needs to be at the protocol level, like a way to re-establish a previous session, or rollback the session state, or something. It's definitely hard mode for library authors to fix this in any kind of transparent way.
The qpid client libraries supported automatic transparent reconnection attempts, but in the end I usually had to disable them in order to add logic for what to do after reconnecting. IE, I needed to know the connection was lost in order to handle it anyways.
We've had RabbitMQ as part of our stack at my day job since time began, I think it's great software overall but boy are the client libraries a challenge.
We've built a generalised abstraction around first Pika and then pyamqp (because Pika had some odd issues, I forget the details of which) and while pyamqp seems better, it's still not without its odd warts.
We ended up needing to develop a watchdog to wrap invocation of amqp.Connection.drain_events(timeout: int) because despite using the timeout, that call would very occasionally inexplicably block forever (with the only way to break it free being to call amqp.Connector.collect()).
My other data point was a time I built something to slice off a copy of production data for testing purposes (from instances of the system above) using Benthos (pretty cool software tbh, Go underneath), but it would inexplicably just stop consuming messages and I had no idea why (so I just went back to our gross but proven Python abstraction to achieve the same).
Historically people like to point out the common locking issues etc... with SQL but in modern datbases you have a good number of tools to deal with that ("select for update nowait").
If you think about it a queue is just a performance optimisation (it helps you get the 'next' item in a cheap way, that's it).
So you can get away with "just a db" for a long time and just query the DB to get the next job (with some 'reservations' to avoid duplicate processing).
At some point you may overload the DB if you have too many workers asking the DB for the next job. At that point you can add a queue to relieve that pressure.
This way you can keep a super dynamic process by periodically selecting 'next 50 things to do' and injecting those job IDs in the queue.
This gives you the best of both worlds because you can maintain granular control of the process by not having large queues (you drip feed from DB to queue in small batches) and the DB is not overly burdened.
DB queues allow easy prioritization, blocking, canceling, and other dynamic queued job controls that basic FIFO queues do not. These are all things that add contention to the queue operations. Keep your queues as dumb as you possibly can and go with FIFO if you can get away with it, but DB queues aren't the worst design choice.
love it
Obviously it makes sense to not use complex tech when simple tech works, especially at companies with a lower traffic volume. That is just practical engineering.
The inverse, however, can also be true. At super high volumes you run into issues really quickly. Just got off a 3 hour site-wide outage due to the database unable to keep up with the unprecedented queue load, and the db system basically ground to a halt. The proposed solution is actually to move off of a dedicated db queue for SQS.
This was a system running that has run well for about 10 years. Granted there was an unprecedented queue volume for this system, but sometimes a scaling ceiling is hit, and it is hit faster than you might expect from all these comments saying to always use a db always, even with all the proper indexing and optimizations.
Db queues are simple to implement and so given the volume it's one way to approach working around an mq client issue.
Personally, and I mean personally, I have found messaging platforms to be full of complexity, fluff, and non-standard "standards", it's just alot of baggage and in the case of messaging alot of bugs.
I have seen Kafka deployed and ripped out a year later, and countless bugs in client implementations due to developer misunderstanding, poor documentation, and unnecessary complexity.
For this reason, I refer to event driven systems as "expert systems" to be avoided. But in your life "there will be queues"
It’s meant to be a simpler reimagining of Amazon SQS.
It has an HTTP API and behaves mostly like SQS.
I wrote it to support Postgres, Microsoft’s SQL server and also MySQL because they all support SKIP LOCKED.
At some point I turned it into a hosted service and only maintained the Postgres implementation though the MySQL and SQL server code is still in there.
It’s not an active project but the code is at https://github.com/starqueue/starqueue/
After that I wanted to write the worlds fastest message queue so I implemented an HTTP message queue in Rust. It maxed out the disk at about 50,000 messages a second I vaguely recall, so I switched to purely memory only and in the biggest EC2 instance I could run it on it did about 7 million messages a second. That was just a crappy prototype so I never released the code.
After that I wanted to make the simplest possible message queue so I discovered that Linux atomic moves are the basis of a perfectly acceptable message queue that is simply file system based. I didn’t put it into a message queue, but close enough to be the same I wrote an SMTP buffer called Arnie. It’s only about 100 lines of Python. https://github.com/bootrino/arniesmtpbufferserver
What is the prefetch value for RabbitMQ mean? > The value defines the max number of unacknowledged deliveries that are permitted on a channel.
From the Article: > Turns out each RabbitMQ consumer was prefetching the next message (job) when it picked up the current one.
that's a prefetch count of 2.
The first message is unacknowledged, and if you have a prefetch count of 1, you'll only get 1 message because you've set the maximum number of unacknowledged messages to 1.
So, I'm curious what the actual issue is. I'm sure someone checked things, and I'm sure they saw something, but this isn't right.
tl;dr: prefetch count of 1 only gets one message, it doesn't get one message, and then a second.
Note: I didn't test this, so there could be some weird issue, or the documentation is wrong, but I've never seen this as an issue in all the years I've used RabbitMQ.
If prefetch was the issue; they could have even used AMQP's basic.get - https://www.rabbitmq.com/amqp-0-9-1-quickref.html#basic.get
1.) You (hopefully) know a bit about how your DB works, what the workload is, what your tasks are. You also (hopefully) know a bit about SQL and Postgres. So you learn a bit more and build upon that knowledge, and implement the queue there (which comes with other benefits).
2.) You learn about a new package, deal with how to set that up, and how to integrate it with your existing database (including how tasks get translated between the queue and your existing DB). This also increases your maintenance and deployment burden, and now developers need to know not only about your DB, but the queueing package as well.
There are certainly cases where #2 makes sense. But writing off #1 as NIH often leads to horrifically over-engineered software stacks, when 10s/few hundred lines of code will suffice.
To their defense they openly admit this
> and didn't read the documentation.
They clearly did, they even linked to it.
But, as a sibling comment mentioned, they could probably got away with a basic get instead.
As documentation writers (and often that should be read "as developers") it is our task to make our users fall into the success pit - and stay there.
Unfortunately, by not considering this, thousands of hours get lost every year and lots of potential users leave because they were somehow (Google, random blogs, our documentation) led into the hard path.
My favourite example is how, for years, if you tried to find out how to use an image in a Java application you would end up with the documentation for an abstract class 2dgraphics or something, while the easy way was to use an Icon.
Channel prefetch:
https://www.rabbitmq.com/confirms.html
"Once the number reaches the configured count, RabbitMQ will stop delivering more messages on the channel unless at least one of the outstanding ones is acknowledged"
consumer prefetch:
https://www.rabbitmq.com/consumer-prefetch.html
So a prefetch count of 1 = 1 un-ACKed message -> what they want
btw, publisher confirms used in conjunction with prefetch setting can allow for flow control within a very well behaved band.
People run into issues with Rabbit for two reasons. You noted one (they are documentation averse), and number is two is mistaking a message queue for a distributed log. Rabbit does -not- like holding on to messages. Performance will crater if you treat it as a lake. It is a river
> took half a day to implement + test
so seems like there are maybe 2 or 3 services using rabbitmq
Your link results in a 404.
The scale isn't large enough for this to at all be a worry. The biggest worry here I imagine is ensuring that a job isn't processed by multiple workers, which they solve with features built into Postgres.
Usually I caution against using a database as a queue, but in this case it removes a piece of the architecture that they have to manage and they're clearly more comfortable with SQL than RabbitMQ so it sounds like a good call.
There are a lot of ways to implement a queue in an RDBMS and a lot of those ways are naive to locking behavior. That said, with PostgreSQL specifically, there are some techniques that result in an efficient queue without locking problems. The article doesn't really talk about their implementation so we can't know what they did, but one open source example is Que[1]. Que uses a combination of advisory locking rather than row-level locks and notification channels to great effect, as you can read in the README.
They also say that some jobs take hours and that they use 'SELECT ... FOR UPDATE' row locks for the duration of the job being processed. That strongly implies a small volume to me, as you'd otherwise need many active connections (which are expensive in Postgres!) or some co-ordinator process that handles the locking for multiple rows using a single connection (but it sounds like their workers have direct connections).
I've also used in-database queuing, and it worked well enough for some use cases.
However, most importantly: calling yourself a maxi multiple times is cringey and you should stop immediately :)
I am afraid I've moved to a default three-way architecture:
- backend autoscaling stateless server
- postgres database for small data
- blobstore for large data
it's not that other systems are bad. its just that those 3 components get you off the ground flying, and if you're struggling to scale past that, you're already doing enormous volume or have some really interesting data patterns (geospatial or timeseries, perhaps).
I don't remember anymore what that general kind of datatype is called, sadly.
I'm curious to hear what the team thinks the pros/cons of a row vs advisory lock are and if there really are any performance implications. I'm also curious what they do with job/task records once they're complete (e.g., do they leave them in that table? Is there some table where they get archived? Do they just get deleted?)
The memory space reserved for locks is finite, so if you were to have workers claim too many queue items simultaneously, you might get "out of memory for locks" errors all over the place.
> Both advisory locks and regular locks are stored in a shared memory pool whose size is defined by the configuration variables max_locks_per_transaction and max_connections. Care must be taken not to exhaust this memory or the server will be unable to grant any locks at all. This imposes an upper limit on the number of advisory locks grantable by the server, typically in the tens to hundreds of thousands depending on how the server is configured.
– https://www.postgresql.org/docs/current/explicit-locking.htm...
Without more information about your setup, the advisory locking sounds like dead weight.
I avoided implicit locking by manually handling transactions. The query that acquired the lock was a separate transaction from the query that figured out which jobs were eligible.
> Without more information about your setup, the advisory locking sounds like dead weight.
Can you expand on this? Implementation-wise, my understanding is that both solutions require a query to acquire the lock or fast-fail, so the advisory lock acquisition query is almost identical SQL to the row-lock solution. I'm not sure where the dead weight is.
I'm imagining this is the sorta setup we're comparing:
Row Lock - https://pastebin.com/sgm45gF2 Advisory Lock - https://pastebin.com/73bqfBfV
And if that's accurate, I'm failing to see how an advisory lock would leave the table unblocked for any amount of time greater than row-level locks would.
The point of the explicit row-level locking is to allow a worker to query for fresh rows without fetching any records that are already in-progress (i.e. it avoids what would otherwise be a race condition between the procedural SELECT and UPDATE components of parallel workers), so if you've already queried a list of eligible jobs, and then have your workers take advisory locks, what are those locks synchronizing access to?
In my solution, the id I was passing to pg_try_advisory_lock was the id of the record that was being processed, which would allow several threads to acquire jobs in parallel.
The second difference is that my solution filters the table containing jobs with the pg_locks table and excluds records where the the lock ids overlapped and the lock type was an advisory lock. Something like:
SELECT j.* FROM jobs j WHERE j.id NOT IN ( SELECT pg_locks l ON j.id = (l.classid::bigint << 32) | l.objid::bigint WHERE l.locktype = 'advisory' ) LIMIT 1;
The weird expression in the middle comes from the fact that Postgres takes the id you pass to get an advisory lock and splits it across two columns in pg_locks, forcing the user to put them back together if they want the original id. See https://www.postgresql.org/docs/current/view-pg-locks.html.
Thanks for sharing your design! :)
LISTEN and NOTIFY are the common approaches to avoiding polling, but I've not used them myself (yet).
This is what you get when you hire using leet code exercises and dump all your design thinking to "Senior" developers that have 1-3 years under their belt.
"a million jobs a minute"
https://getoban.pro/articles/one-million-jobs-a-minute-with-...
Uses both listen notify and advisory locks so it is using all the right features. And you can enqueue a job from sql and plpgsql triggers. Nice!
Worker is in Node js.
Another problem is that INSERT followed by SELECT FOR UPDATE followed by UPDATE and DELETE results in a lot of garbage pages that need to be vacuumed. And managing vacuuming well is also an annoying issue...
This assumes the processing is idempotent in the rest of the system and is only committed transactionally when it's done. Some workers might do wasted work, but you can tune the expiration future time for throughput or latency.
One feature I like is change event stream which you can subscribe to. It is pretty fast and reliable and for good reason -- the same mechanism is used to replicate MongoDB nodes.
I found you can use it as a handy notification / queueing mechanism (more like Kafka topics than RabbitMQ). I would not recommend it as any kind of interface between components but within an application, for its internal workings, I think it is pretty viable option.
One caveat to this is that you can only start from wherever the beginning of your oplog window is. So for large deployments and/or situations where your oplog ondisk size simply isn't tuned properly, you're SOL unless you build a separate mechanism for catching up.
Putting the queue in its own deployment is a good insulation against this (assuming you don't need to use aggregate() with the queue across collections obviously).
Also, I am firm believer that you should not put actual data through notifications. Notifications are meant to wake other systems up, not carry gigabytes of data. You can pack your data into another storage and notify "Hey, here is data of 10k new clients that needs to be processed. Cheers!"
The message is meant to ensure correct processing flow (message has been received, processed, if it fails somebody else will process it, etc.), but it does not have to carry all the data.
I have fixed at least one platform that "reached limits of Kafka" (their words not mine) and "was looking for expert help" to manage the problem.
My solution? I got the component that publishes upload the data to compressed JSON to S3 and post the notification with some metadata and link to the JSON. And the client to parse the JSON. Bam, suddenly everything works fine, no bottlenecks anymore. For the cost of maybe three pages of code.
There is few situation where you absolutely need to track so many individual objects that you have to start caring if they make hard drives large enough. And I managed some pretty large systems.
We're in agreement, I think we may be talking past each other. I use mongo for the exact use case you're describing (messages as signals, not payloads of data).
I'm just sharing a footgun for others that may be reading that bit me fairly recently in a 13TB replica set dealing with 40mm docs/min ingress.
(Its a high resolution RF telemetry service, but the queue mechanism is only a minor portion of it which never gets larger than maybe 50-100 MB. Its oplog window got starved because of the unrelated ingress.)
> I dont think I've ever seen a benchmark for any DB that's gotten above >30k writes/sec
Mongo's own published benchmarks note that a balanced YCSB workload of 50/50 read/write can hit 160k ops/sec on dual 12-core Xeon-Westmere w/ 96GB RAM [1].
Notably that figure was optimized for throughput and the journal is not flushed to disk regularly (all data would be lost from last wiredtiger checkpoint in the event of a failure). Even in the durability optimized scenario though, mongo still hit 31k ops/sec.
Moving beyond just MongoDB though, Cockroach has seen 118k inserts/sec OLTP workload [2].
[1] https://www.mongodb.com/scale/mongodb-benchmark [2] https://www.cockroachlabs.com/docs/stable/performance.html#t...
https://www.postgresql.org/docs/15/logicaldecoding.html
Quite a number of bits of software use that to do things with these events. For example:
https://debezium.io/documentation/reference/stable/connector...
Depends on the use-case, but the original article smells like FUD. This is because the connection C so lib allows you to select how the envelopes are bound/ack'ed on the queue/dead-letter-route in the AMQP client-consumer (you don't usually camp on the connection). Also, the expected runtime constraint should always be included when designing a job-queue regardless of the underlying method (again, expiry default routing is built into rabbitMQ)...
Cheers =)
The dude's being a bit too self-deprecating with that (sarcastic quip).
But there's valuable something buried there - if it is possible to efficiently solve a problem without introducing novel mechanisms and concepts, it is highly desirable to do so.
Don't reinvent the wheel unless you need to.
> You could set the prefetch count to 1, which meant every worker will prefetch at most 1 message. Or you could set it to 0, which meant they will each prefetch an infinite number of messages.
O what the actual fuck?
I'm really hoping he's right about holding it wrong, because otherwise ???
I’d venture to guess that the median RabbitMQ-using app in production could not easily be replaced with postgres though. The main reasons this one could are very low volume and not really using any of RMQ’s features.
I love postgres! But RMQ and co fulfill a different need.
Because the chances of a cloud instance randomly going down is not insignificant.
If you're measuring throughput in jobs/s, use a real work queue.
I've written about a somewhat similar problem in context of the Mats3 library I've made: https://mats3.io/patterns/work-queues/
And yes, the point is then to use a database to hold the work, dispatching from that. In the Mats3 context I describe, the reason is to pace the dispatching to a queue (!), but for them, it should be to just run the work from that table. Also, the introspection/monitoring argument presented should be relevant for them.
That a message queue library fetches more messages than the one it is working on is totally normal: ActiveMQ per default uses a prefetch of 1000 for queues, and Short.MAX_VALUE-1 (!) for topics. The broker will backfill you when you've depleted half of that. This is obviously to gain speed, so that once you've finished with one message, you already have another available, not needing to go back to the broker with both network and processing latencies: https://activemq.apache.org/what-is-the-prefetch-limit-for
In summary, I feel that the use case they have, "thousands of jobs per day", which is extremely little for a message queue, where many of these jobs are hours-long, is .. well .. not optimal use case for a MQ. It is the wrong tool for the job, and just adds complexity.
When programming my side hustle I also had the requirement of a SIMPLE queue. I didn't want to introduce AWS SQS or RabbitMQ, so I wrote a few C# classes which where dequeueing from a MongoDB collection. It works pretty well. It basically leverages the atomic MongoDB operation "findAndModify", so you can ensure that dequeueing will find only messages in status "enqueued" and in the same operation sets the status to "processing" so you can ensure that only one reader processes the message. (https://www.mongodb.com/docs/manual/reference/method/db.coll...)
I created a small NuGet package, which you can find here: https://allquiet.app/open-source/mongo-queueing.
- clear case of misconfigured instance (prefetch) and a bug somewhere in the stack (reconnection issue)
- prefetch behavior sounds like they receive a message, ack it then process it
- i wouldn't recommend ack before processing, because you become responsible for tracking and verifying if the worker ran to completion or not.
- work then ack is the way. the other way around ignores key job processing benefits like rabbitmq automagically requeueing messages when a worker crashes and failure related logic like deadletter queues.
- the trick i've started leaning on with rabbitmq is giving each worker their own instance queue (i call it their mailbox).
- when a worker starts a job, it writes the job id, start time and the worker's mailbox to a db. any system can now look up the "running job" in the database, know how long it has been running and can even talk to the worker using its mailbox to inquire if it is running and if that job state in the db is accurate.
- happy the writer and team found what works for them. ultimately, what you understand best would serve you better, so they made a good choice to lean on their strengths (postgre).
Just occurred to me there's a different possibility:
The extra connection bug mentioned could be from a rogue thread/process/code-path connecting to the queue and fetching an item, but somehow not doing the work.
Verifying this likely requires access to the codebase, so treat this as little more than idle speculation.
Wouldn't this lead to contention issue when a lot of multiple workers are involved?
[0] https://www.2ndquadrant.com/en/blog/what-is-select-skip-lock...
With SKIP LOCKED you can resolve this easily, as long as you know about the feature. But almost every tutorial and description of implementing a job queue in Postgres mentions this now.
Like everything else in life, it's always a tradeoff. Know your workload, the tradeoffs your tools are making, and make sure to mix and match appropriately.
In the case of Prequel, it seems they possibly have a low throughput situation at hands, i.e. in the case of periodic syncs the time spent queuing the instruction <<< the time needed to execute it. Postgres is great in this case.
https://webapp.io/blog/postgres-is-the-answer/
I think its a reasonable option and the webapp.io people scaled it out pretty high. At Zapier we utilize RabbitMQ heavily and I cannot imagine scaling the amount of tasks we handle each day on postgres.
This article mostly is lesson learns and misconfigurations of RabbitMQ though. Which is probably a good reason to simplify if you don't need its power. No reason to have to learn how to configure it if postgres is good enough.
I'm always surprised, even with the older version of ActiveMQ, what kind of throughput you can get, on modest hardware. A 1gb kvm with 1 cpu easily pushes 5000 msgs/second across a couple hundred topics and queues. Quite impressive and more than we need for our use case. ActiveMQ Artemis is supposed to scale even farther out.
We’ve also had similar issues as OP, except fixing it just came down to configuring the Java client to have 0 prefetch so that long jobs don’t block other msgs from being processed by other clients. Also using separate queues wide different workloads.
Nowadays for small stuff I just create a simple script that uses asyncio, I can async request stuff and run_forever etc and handle it all as a docker service.
So now you have to deal with deciding how tight your polling loop is, and with reads that are happening regardless of whether you have messages waiting to be processed or not, expending both CPU and requests, which may matter if you are billed accordingly.
Not in any way knocking it, just pointing out some trade-offs.
Moved to Python-RQ [1] + Redis and been rock solid for years now. Redis is also great for locking to ensure only one instance of a job/task can run at a time.
All the stuff on Agenda works perfectly all the time. (which is basically using mongo's find and update)
Actually, using Postgres stored procedures they can do anything in Postgres. I am quite sure they can rewrite their entire product using only stored procedures. Doesn't mean they really want to do that, of course.
https://streamnative.io/blog/comparison-of-messaging-platfor...
Google search for [sql maxis] just returns this article.
SQL server won't callback clients [or application-server] informing that new data is available. So how do you poll it? A query in periodic loop? Some other way? How much that scales, or how much load does that create on the server?
Also, problem described seems like a logical error. Worker shouldn't ack the job before finishing it.
If you ever get into a case where "we don't think we're using it right", then you didn't understand it when you implemented it. That is a much bigger problem to understand and prevent in the future than the problem of picking a tool.
If an external process is responsible for marking a job as done, you could add a timestamp column that will act as a timeout. The column will be updated before the job is given to the worker.
SELECT ... WHERE ts > NOW()
UPDATE ... SET ts = NOW() + INTERVAL '1 HOUR'
My question would be why you're in the business of writing a job manager in the first place.
Most often, because it's easier than configuring a ready made one.
In Mysql you can directly call
Update queue set working_id=xxxx where working_id='' limit 1
to atomically pop one task from queue.
anyone else have trouble completely disregarding the whole of the article when they see things like this?
always go back to fundamentals. the rdbms is giving you replication , queries , locking but at what cost ?