System design hack: Postgres is a great pub/sub and job server
layerci.com
layerci.com
The architecture is to stick a list of your input shards in a Postgres table, have a state flag that goes PENDING->WORKING->FINISHED->(ERROR?), and then spin up a bunch of worker processes as EC2 spot instances that check for the next PENDING task, mark it as WORKING, pull it, process it, mark it as FINISHED, and repeat. They write their output back to the DB in a transaction; there's an assumption that aggregation can happen in-process and then get merged in a relatively cheap transaction. If the worker fails or gets pre-empted, it retries (or marks as ERROR) any shards it was previously working on.
Postgres basically functions as the MapReduce Master & Reducer, the worker functions as the Mapper and Combiner, and there's no need for a shuffle phase because output <<< input. Almost all the actual complexity in MapReduce/Hadoop is in the shuffle, so if you don't need that, the remaining stuff takes < 1 hour to implement and can be done without any frameworks.
Sounds like you'd end up with a bunch of dangling shards orphaned in WORKING state. And now you need timeouts and health checks and something to coordinate all that.
For exceptions, crashes, and other unexpected errors, you just let it crash. In the shell script that starts the worker process, you wrap it in "while (1) { ... }", so it's just permanently restarting. Then the aforementioned retry check in main() handles the previous shard just as if the worker had been pre-empted or finished normally.
There are a variety of low-tech ways to detect when the job as a whole has finished, ranging from manually shutting down the cluster to writing a small cronjob that counts PENDING/WORKING shards and shuts down the cluster if there are none to having your workers return a different exit status if there's no more work to do and shutting down the box if so. A neat side effect of this is that you get an easy status report of exactly which machine is working on which shard and how long you have to go through "SELECT * FROM work_queue". Another neat side effect is that you can view partial results in the DB before the job as a whole completes, something you can't do with MapReduce, and potentially adjust your code and restart the job if you're getting bad data.
Example: https://gist.github.com/risicle/f4807bd706c9862f69aa
Does that mean workers poll the database?
Was "the database" a single node, or multiple systems?
This architecture originally evolved from an Amazon-SQS based queue system, where I'd stick the jobs in SQS and pop them off. Then I realized that I needed the DB anyway for my reducer-substitute, and if I'm writing SQL/ORM code anyway it was simpler to use the DB to manage the queue than to use SQS, and it gave me some other nifty abilities like being able to record timing & size stats with each shard, being able to query progress while the job is going, and being able to manually restart shards after correcting the errors that led to their failure. If you're running dozens-to-hundreds of workers, the cost of a t2 RDS instance (that's all you need for 1 QPS) is a rounding error, and a lot cheaper than EMR would be.
I'm not going to say "I don't know why anybody is surprised," but I'll tell you that I'm surprised! This is good, you're all one of 10,000![1]
I'd call parent's a "basic" workflow, but Wikipedia seems to tag it as "state-based."[2] I'm glad someone pointed out that failed steps do not update the db unless you have good traps in the process itself. I'd probably suggest a housekeeping process that resets any errors after some amount of time so that they're run through again from the beginning. As long as the error is surfaced in the meantime, that is.
And this is definitely not a new pattern - if you wanna go back in time, this is how things often worked in the mainframe era. My point - and IIUIC the point of the article - is that people often reach for flavor-of-the-month heavyweight frameworks when they can easily adapt tools they're already using to solve the problem within hours.
(Another great example of this is the single-process in-RAM architecture used by PlentyOfFish, StackOverflow, Mailinator, and Hacker News when they were solo-developer projects, where you store all state in RAM and run the whole site on a single box, periodically dumping out checkpoints to disk in case the process crashes. This can scale to thousands of requests per second - that's millions of queries per day - and is a lot simpler to program than 3-tier architectures where everything runs through a data abstraction layer and you spend 90+% of the code shuffling data around.)
So what happened to a guy who had the bright idea to do CPU mining on Google's infra?
Absolutely, which is why I think it's important to call concepts the same thing they called it in the olden times if you know it's an old thing!
At the core of our implementation is this library to schedule the jobs to different workers: https://github.com/kagkarlsson/db-scheduler
Would highly recommend that library!
This is pretty good advice in general.
I was pleased to see they are using `SELECT FOR UPDATE SKIP LOCKED`. That is what this 2nd Quadrant article recommends, which I think is required reading if you want to implement this yourself:
https://www.2ndquadrant.com/en/blog/what-is-select-skip-lock...
It goes into more detail about wrong ways to implement a queue and what the downsides are for its preferred approach.
I had seen it also for soft that sends massmail (our case around 100k/day).. it's state is a postgres queue.
We also use Pg for transactional mail. We insert it on a table. (There is a process that sends the row mails).. the so nice part is that the mail is joining the dB transaction for free.. (all or nothing)
It's also the easiest one to start with, excelling in all the mainstream categories for relational databases.
I really like the look of this approach, but don't want to build it myself if I can avoid it.
We had an internal library that used advisory locks which had all sorts of strange behavior we couldn't figure out until we just moved to SKIP LOCKED
https://github.com/malthe/pq - postgres queue based on rq and ruby queue_classic
https://github.com/coleifer/huey - sqlite, redis and in memory by coleifer (peewee creator). Possible to implement postgres storage layer simply.
https://github.com/closeio/tasktiger - flexible redis-based python task queue alternative to celery
https://dramatiq.io/ - actor based python job queue on redis/rabbitMQ
https://github.com/GoogleCloudPlatform/psq - gcp pub/sub based task queue
https://github.com/celery/celery/blob/master/celery/backends...
I wish all these libraries agreed on a common schema >_<
The code is open-sourced here in case anyone wants to re-use
https://github.com/rudderlabs/rudder-server/blob/master/jobs...
We had to built additional logic to clean-up old jobs (similar to Level merges in similar queing systems)
The feature of postgres that makes this viable in comparison to most other databases is the "channel"
INSERT INTO JOBS
versus INSERT INTO JOBS
redis-cli LPUSH job-queue job-id
It's also easier to configure, redis has BRPOPLPUSH for acked-queues, but it's much easier to just query jobs with a select statement from the database.1.) *They want to transactionally commit work along with the change that caused it
2.) They are already using Postgresql not using Redis
3.) Requiring users install yet another service(Redis) just for this one item isn't worth the costs
Though I agree on points 2 and 3, especially 3 since it adds complexity.
That's only true if they do not error—there is no "rollback" feature in Redis, and if you do error, anything you've done up to that point remains.
That's not transactional, it's more "you can, sometimes, execute a bunch of separate statement atomically".
Typically you would implement visibility timeouts and other such stuff. Depending on the use case you could make specific optimizations or keep it generic and have SQS like semantics or something.
If so if worker crashes you get a connection timeout on postgresql and it rolls back the transaction and releases row locks - then next select SKIP picks it up.
If you want fast response with safety listen and on restart do a poll always to pick up anything that's ready that you might have missed a pub for. If you want lots of coverage do a poll every 10s - postgresql can handle it, you still get immediate response by listening to notify.
If you need fancier approaches you can do a last_updated time on the job and or watch waiting on pg_stat_activity or similar etc - but all totally overkill I think - the simple approach get's pretty far (not an expert though at all).
There is a solution to every perceived problem.
But we still see issues.
No thanks, I'll stick with GCP PubSub or AWS SQS, which are explicitly designed for this use case, and for which I have to setup no infrastructure.
- A lot of use cases involve not dropping messages after they are processed (like CI jobs, in this example), so you don't have to vacuum the rows
- If you're comfortable with SQS there's no really big reason to switch, but it makes it so that your project can only run on amazon cloud servers, which is annoying.
Why? Can't you use it as a standalone product that integrates into whatever service you want? The last time I used it was a few years ago for batch processing. I had Heroku instances that dispatched jobs directly to SQS. I'm doubtful but it's been a while since I've used the service so it might have changed.
> There is a queue that holds notifications that have been sent but not yet processed by all listening sessions. If this queue becomes full, transactions calling NOTIFY will fail at commit. The queue is quite large (8GB in a standard installation) and should be sufficiently sized for almost every use case.
My understanding of MVCC (correct me if I'm wrong), is every time you do an update, dead tuples are left behind that need to eventually be vacuumed. If vacuum isn't running or can't keep up, you'll run out of space or have a tx id wrap around.
Work queues are a simple paradigm and swapping out one queue for another is a trivial exercise that won't change the architecture of a system, and SQS or PubSub could be easily used, regardless of where the rest of the project is hosted. I'd rather pay that one-time cost than the ongoing cost of maintaining my own job queue service and infrastructure.
This is an often untold superpower of PostgreSQL when arguing about how it supposedly "don't scale well". You often have a large array of optimizations available before your project really scale out of scope.
My point isn’t that Postgres can’t do it. It can. But it just requires so much more effort to get it right. And this article glosses over so many critically important details about doing these things at scale in production environments. Most of them are things I don’t have to worry about when using a real queue service.
hell, even its clustering system got a lot better!
> In the list above, I skipped things similar to pub/sub servers called "job queues" - they only let one "subscriber" watch for new "events" at a time, and keep a queue of unprocessed events:
If your job queue only allows one single worker (even per named queue), I'd argue it's a shit job queue.
I'm not here to defend celery generally, but it has no such limitation.
The use case I have in mind has a lot of logically separate queues. Is it better for each queue to have its own channel (so subscribers can listen to only the queue they need) or have all queues notify a global channel (and have subscribers filter for messages relevant to them). I am mainly confused about whether I need a dedicated db connection per LISTEN query and also how many channels is too much.
I'm not sure what the maximum scale is, but I've never hit it.
After a few weeks of running on a multi-TB table, you'll find the dead tuples massively outnumber the live tuples, and database performance is dropping faster than the Vacuum can keep up. Vacuum is inherently single-threaded, while your processes making dead tuples as part of queries are multi-threaded, so it's obviously the vacuum that fails first if most queries are long running and touch most rows. Your statistics will get old because they're also generated by the vacuum process, making everything yet slower.
Even if you can live with gradually dropping performance, eventually the whole thing will fail when your txids wrap around and the whole database goes read-only.
At that point, people usually start looking at the connection pooling tools. Depending on how much work you need from the DB, connections pools can be a win. Anyone know how connection pooling works with listeners?
A
They don't. LISTEN is per connection and pgbouncer multiplexes many sessions onto a single one. The poller holds the connection and has no way to propagate back the notification while still maintaining multiplexed sessions as isolated.
So while postgres makes for a pretty awesome database, and a pretty good queueing system, some folks may be seriously impacted by the max number of connections
It also eats a lot more cpu cycles:
In our case we had about ~200 connections to database. After we placed pgbouncer in front and reduced number of connections to 60. Our commit throughput doubled.
And I'm more than perfectly fine with that.
Disclaimer: I work on Debezium
[1] https://debezium.io/documentation/reference/0.10/connectors/...
Could you please elaborate?
Instead you use a cloud like google cloud where you can add NVMe SSDs to whatever instance type you need and configure custom RAM and CPU instead of picking from the super expensive AWS instances with no configurable options and almost always the wrong resource allocations for your workload.
Source: testing my infrastructure that requires 60,000 iops on both google cloud and AWS and it being 1/4 the cost and higher performance on Google. Of note: this was a very high throughput streaming data application. YMMV for other applications.
Most advices I have seen say that I'll probably want to code it myself, but I was wondering about the latency of that solution? I'll likely have a SQL store and that would be a good argument to use postgres...
We have switched to an Phoenix (elixir) server which listens to the database, then clients subscribe via websockets. We're opensourcing it here: https://github.com/supabase/realtime
It's still in very early stages (although I am using it in production). Basically the Phoenix server listens to PostgreSQL's replication functionality and converts the byte stream into JSON, which it then broadcasts over websockets. This is great since you can scale the Phoenix servers without any additional load on your DB. Also it doesn't require wal2json (a Postgres extension). The beauty of listening to the replication functionality is that you can make changes to your database from anywhere - your api, directly in the DB, via a console etc - and you will still receive the changes via Phoenix.
It's similar to Apache Kafka but very light on resources and easy to selfhost.
Ow yeah, the throughput is around 150 k msg/s. Without much configuration.
I.e. is this good advice because Postgres in particular is a great implementation of sql, or because sql in general is good enough to solve this problem, or a mix of the two?
- has strong performance vs, say, sqlite
- has "channel" and "trigger" support so you can avoid polling (which either slows down your jobs or limits your number of workers)
- is actually OSS (versus, say, mysql)
I believe Citus DB(?) is something similar.
"Citus Unforks From PostgreSQL, Goes Open Source"
https://www.citusdata.com/blog/2016/03/24/citus-unforks-goes...
https://news.ycombinator.com/item?id=11353322
As for mysql - I don't know why anyone wouldn't prefer mariadb - but I suppose it's conceivable mysql has features that makes it worth dealing with oracle licensing. But I doubt it.
I agree that PostgreSQL has some very nice extensions and if you know you are going to have multiple users (or multiple processes using it), SQLite can't compete.
But if you only have one process using the DB, then SQLite is immune to N+1 query problem since everything is in-process and there are no round-trips.
INSERT INTO ci_job_status(repositor ...
BTW, it looks Go becomes so popular that many tutorials are using Go for examples. ;DCan this hack not be achieved by a mariadb table too?
The notification system comes with a cost. You can't scale up the number of connection to a PG database instance.
Therefore you need to go pooling.
And for SKIP LOCKED, I just checked and mysql/maria do have it.
However, my design was more bare bones. I was picking jobs by chaining CTE's that did the status update as they returned the first element of the queue.
Key-value stores could do a lot of things theoretically.
If you needed the existence of completed jobs to compute future jobs, use a supplementary datastructure that stores spans or other heuristics, updated in a transaction with job completion
It's a great solution.
https://github.com/davidbanham/kewpie_go
I’ve been using it for years across multiple products.
I also built this for a client lately. It lets them interop their Haskell services with kewpie a bit easier without needing to write a native wrapper.
On my first go I was taking batches, on the reasoning that it optimized for query performance. It was a complete disaster and I quickly realized that I was engaging in premature optimization.
Batching is out, stream-processing is in. If you design your job table correctly, PG will perform very well as an advanced stream processor.