RQ – Simple Job Queues for Python
github.com
github.com
After replicating most of Sidekiq’s pro and enterprise behavior using older data structures I attempted to migrate to streams. What I discovered is that all the features I really wanted were available in SQL (specifically PostgreSQL). I’m not the first person to discover this, but it was such a refreshing change.
That led me to develop a Postgres based job professor in Elixir: https://github.com/sorentwo/oban
All the goodies only possible by gluing Redis structures together through lua scripts were much more straightforward in an RDBMS. Who knows, maybe the recent port of disque to a plug-in will change things.
We're currently using ecto_job which works really well for us, so we have no reason to switch. Plus we like it's in different tables.
PG is an amazing tool that can handle more than most people think (or at least more than I thought).
What exactly is fragile about this approach? We've supported scheduled tasks in TaskTiger (https://github.com/closeio/tasktiger) via Redis's sorted sets and haven't had any issues.
Using sorted sets for scheduled tasks isn't the fragile part, gluing it all together to prevent losing jobs by shifting into backup lists (or hashes) is the fragile part in my experience.
That definitely does the job. My point is that it is much more complex than a select/update clause in SQL.
[0] https://github.com/closeio/tasktiger/blob/master/tasktiger/r...
Stream does solve the backup queue issue, but not being able to filter messages from the consumer's pending list makes implementing retry mechanism and exponential backoff hard (without reaching to sorted set).
Though, I can understand the Redis team's decision to keep it simple.
I use Postgres SKIP LOCKED as a queue. Postgres gives me everything I want. I can also do priority queueing and sorting.
All the other queueing mechanisms I investigated were dramatically more complex and heavyweight than Postgres SKIP LOCKED.
Here is a complete implementation - nothing needed but Postgres, Python and psycopg2 driver:
import psycopg2
import psycopg2.extras
import random
db_params = {
'database': 'jobs',
'user': 'jobsuser',
'password': 'superSecret',
'host': '127.0.0.1',
'port': '5432',
}
conn = psycopg2.connect(**db_params)
cur = conn.cursor(cursor_factory=psycopg2.extras.DictCursor)
def do_some_work(job_data):
if random.choice([True, False]):
print('do_some_work FAILED')
raise Exception
else:
print('do_some_work SUCCESS')
def process_job():
sql = """DELETE FROM message_queue
WHERE id = (
SELECT id
FROM message_queue
WHERE status = 'new'
ORDER BY created ASC
FOR UPDATE SKIP LOCKED
LIMIT 1
)
RETURNING *;
"""
cur.execute(sql)
queue_item = cur.fetchone()
print('message_queue says to process job id: ', queue_item['target_id'])
sql = """SELECT * FROM jobs WHERE id =%s AND status='new_waiting' AND attempts <= 3 FOR UPDATE;"""
cur.execute(sql, (queue_item['target_id'],))
job_data = cur.fetchone()
if job_data:
try:
do_some_work(job_data)
sql = """UPDATE jobs SET status = 'complete' WHERE id =%s;"""
cur.execute(sql, (queue_item['target_id'],))
except Exception as e:
sql = """UPDATE jobs SET status = 'failed', attempts = attempts + 1 WHERE id =%s;"""
# if we want the job to run again, insert a new item to the message queue with this job id
cur.execute(sql, (queue_item['target_id'],))
else:
print('no job found, did not get job id: ', queue_item['target_id'])
conn.commit()
process_job()
cur.close()
conn.close()The parent post's model also tracks "attempts", but this only captures known or acknowledged failures -- unacknowledged failures (i.e. crashes) will not be recorded, so a task which explodes the process would be re-run ad infinitum.
An alternative method could be to record an attempt in a separate transaction, so that the next scheduler execution can detect that the job/message was serviced before, even if the message itself appears fresh.
Seriously - don't add complex dependencies to your stack unless you need them. The database makes a great task queue and the filesystem makes a great cache. You really might not need anything more.
Otherwise, I agree, I would love to be able to use Postgres for queues, and I think there is a Dramatiq Postgres backend that will do queueing properly (EDIT: Yes, Dramatiq-PG).
1. Message compaction — notifications are deduplicated, so if two jobs are inserted in the same queue the system may only get one notification. 2. Message saturation — with high levels of activity you’ll need to start discarding messages, essentially denouncing, otherwise the database gets throttled. 3. Dedicated connections — pubsub requires a dedicated connection for listening, which requires a single connection with custom dispatching or one connection per queue.
Relying on PG for everything is awesome regardless!
I looked for a solution that would use Postgres+Celery but there wasn't anything obvious. There is a way to just use the filesystem for Celery, but that doesn't seem ideal.
So that's how I ended up with a Redis dependency.
You can roll your own background task in a few dozen lines of pure Python.
There was no easy way to do complex background tasks - but Django is just Python. Spawning a new process is easy as is scheduling via cron or similar. I've used both to implement simple background tasks.
Half the time people seem to reach for Celery etc all they want to do is run a long task without blocking the response. That's a helleva lot of code to add to your project to just have that functionality.
I need a queue to handle long running jobs. I looked around the Python ecosystem (bc Flask app) and found RQ. So now I add code to run RQ. Then I add Redis to act as the Queue. Then I realized I needed to track jobs in the queue, so I put them into the DB to track their state. Now I effectively have a circular dependency between my app, Redis, and my Postgres DB. If Redis goes down, I’m not really sure what happens. If the DB goes down I’m not really sure what’s going on in Redis. This added undue complexity to my small app. Since I’m trying to keep it simple, I recently found that you can use Postgres as a pub/sub queue which would have completely solved my needs while making the app much easier to reason about. Using Postgres will have plenty of room to grow and buy you time to figure out a more durable solution.
Say you have an object, the object load some config from a settings module which in turn fetch from env. No matter how many times I restarted RQ, the config won't changed. Due to the code changes, the old config cause the job to crash and keep retrying.
Until I got frustrated and get into Redis, pop the job, and yet all the setting was in there. In other words, RQ serialize the whole object together with all properties.
RQ isn't that good IMHO. You will have to add monitoring, health check, node ping, scheduler, retrying, middleware like Celery eventually it grow into a home grown job queue system that make it harder to on board new devs.
Just use Celery. Celery isn't that bloated. It has many UI to support it and backend and very flexible. Celery beat schedule is great as well.
Instead of relying on trusted exposed endpoints and just invoking them by URL, it does a bytecode dump of task functions and stores those in Redis before restoring them from bytecode at execution time.
This has a few drawbacks:
- Payloads in the queue are potentially a fair bit larger for complex jobs - Serialization for stuff that has decorators (and especially stateful one, like `lru_cache`) is not really possible, even with `dill` instead of `pickle` - It's not trivial, but this exposes a different set of security risks compared to the alternative
I don't want to say this is a bad piece of software, it's super easy to set up and way more lightweight than Celery for example, but it's not my tool of choice having worked with the alternatives.
https://github.com/Bogdanp/dramatiq
Stay away from Celery if you can, or stick to v3.2 because 4.x has a ton of bugs.
It's just a fairly higher amount of boilerplate and setup compared to most of the alternatives either way.
Haven't had the opportunity to work with them, but I've heard good things about SQS and other cloud queues, but I think those are fundamentally message queues, so you'd be responsible for writing a job queue layer over it, which can be a huge hassle.
The pickle logic lives here: https://github.com/rq/rq/blob/80c82f731f57186f8c97239b38100f... so if you've got a minute you can skim that and hopefully get an answer to your question.
I don't want to import "count_words_at_url" just to make it a "job" because that couples my web app runtime to whatever the job module needs to import even though the web app runtime doesn't care how the job is handled.
I want to send a message "count-words" with the URL in the body of the message and let a worker pick that up off the queue and handle it however it decides without the web app needing any knowledge of the implementation. The web app and worker app can have completely different runtime environments that evolve/scale independently.
The web framework is a way to handle http, rest, or graphql - deserializing and serializing those protocols, not a way to handle my business logic.
Decoupling these things lets you write one-off scripts or have task queues that don't need to load the context of a large web framework - they can just be simple python.
After opening GitHub issue about it and asking if there were plans to implement sub-minute intervals, library's author (coleifer on github) had dismissive/arrogant attitude, and his reply was something along the lines of "too bad you cannot fork it and do it by yourself" and deleted my thread.
This threw me off from using this library, and I went back to Celery.
* how did you settle on 10 seconds? Keeping two databases in sync is a complex process. I'd suggest that the difference between running 1x/min or 6x/min is negligible -- and if it's not then probably you need something more sophisticated than a simple cronjob.
* I am providing free software. You are not contracting with me to provide developer support, so in my book I'm under no obligation to be courteous and polite all the time. I try to be most of the time, but nobody is perfect. Luckily the source is available if you don't want to talk to me. That was my point and I'm surprised that is so triggering to some people.
* Closing an issue is not deleting your thread.
It doesn't matter how I chose 10 seconds as a periodic task interval, problem was that to a simple question about feature support I got dismissive/borderline rude answer which made me lose confidence in Huey as a project and it's long-term maintainability. And it's not like I asked for something exotic, but a support for subminute periodic tasks which almost every other task queue has.
That issue doesn't seem to be anywhere in Github, so it's deleted: https://github.com/coleifer/huey/issues?utf8=%E2%9C%93&q=is%...
As you said, it's a free software and you can do whatever you want with it, but to me this is a huge red flag and I'm happy to not be part of it. ;)
Celery<4.0
I'm thinking that after we go to prod, I'll move to celery and MQ overtime.
Quick tutorial of Flask + RQ: https://testdriven.io/blog/asynchronous-tasks-with-flask-and...
More advanced features like queue throttling and complex job workflows are available.
Please don't take my question negatively. I acknowledge that not everything is about ROI.
For example I often let my team scratch their own itches by developing tools just for the sake of having fun and learning new toys. Perhaps this is the case with TaskTiger?
It is true that we've spent quite some time working on TaskTiger where we could've used an off-the-shelf task queue system. However, since we've created it, it allowed us to develop new features a lot faster and spend a lot less time debugging issues.
We were using Celery before and it was a big pain with memory leaks, periodic tasks not working as expected, some tasks being lost for no apparent reason, and poor visibility into both aggregate statistics and individual tasks. We spent too much time debugging Celery's issues instead of working on our core business. Here's some more context: https://github.com/closeio/tasktiger/issues/141#issuecomment...
Advanced features like scheduled tasks, unique queues, and task locks have helped us move faster, too.
Especially never had problems with zombie workers. That wouldn't be a problem with your choice of queue library, more an infra problem.
> and some small but annoying bugs while testing
Do you recall any details?
Appreciate you trying TaskTiger, even if you've moved on since!
It came down to I like the features of RabbitMQ:
* RabbitMQ scales messages without needing tons of RAM * I don't have to decide between not persisting messages to disk with Redis vs only using half the machine's RAM [1] * Queue and task visibility is better in RabbitMQ * Support for purging all tasks in a queue * Tasktiger had lower throughput than Celery and Dramatiq, maybe needs lazy-forking?
Things I did like about Tasktiger: * no feature bloat and it's possible to actually read it's source code[2].
[1]: To enable persisting to disk (Redis fork + save snapshotting) you must limit Redis to only using half the available RAM on a machine.
[2]: Celery is split up into multiple convoluted, bloated, difficult to read repos.
I really like the design of beanstalkd and I used that at one place, but using rq + redis was one less thing to deploy and/or fail.
[1]: https://lightbus.org
For less simple use cases, just use Celery.
Unless your use case is absurdly simple, and will always + forever be absurdly simple, then celery will fit better. Else you find yourself adding more and more things on top of RQ until you have a shoddy version of Celery. You can also pick and choose any broker with Celery, which is fantastic for when you realize you need RabbitMQ,