samphilipd
How to best handle tasks when pulling from a job queue
I’m building a postgres-backed job queue for Elixir. The way I have it setup, I have a JobDispatcher that polls the DB every N milliseconds for available jobs. It queries the number of jobs that it has available capacity for and then distributes these to worker processes.
I can imagine two ways to implement this:
One would be to have a separate WorkerSupervisor that spins up a certain number of Worker GenServers (let’s say 25). Each Worker registers itself with the JobDispatcher to say that it is available to receive jobs. Then the JobDispatcher sends the job via message to the existing GenServer and monitors it for either a “job completed” message (meaning it became available for work again) or an exit (if it crashed).
The other way is to not worry about keeping the worker processes running at all, and instead to have the JobDispatcher spawn a one-shot tasks for each new job (using Task.Supervisor.async_nolink) which it keeps a monitor on, and handles both normal and dirty exits. This would appear to have the advantage of simplicity, plus garbage collection is kept localized to each Task so it might also be more efficient.
I haven’t worked enough with OTP patterns to know which one of these is the better approach, does anybody have any ideas?
Trending in Questions
Other Trending Topics
Categories:
Sub Categories:
Forums
Popular Tags
- #ecto
- #liveview
- #troubleshooting
- #learning-elixir
- #deployment
- #library
- #erlang
- #testing
- #genserver
- #mix
- #absinthe
- #remote-other
- #otp
- #plug
- #how-to-question
- #macros
- #postgres
- #channels
- #elixirconf
- #exunit
- #discussion
- #code-sync
- #javascript
- #podcasts
- #onsite
- #dialyzer
- #docker
- #authentication
- #umbrella
- #full-time-contract
- #podcasts-by-brainlid
- #ecto-query
- #elixir-ls
- #phoenix_html
- #iex
- #blog-post
- #graphql
- #genstage
- #ai
- #websockets
- #elixirconf-us
- #supervisor
- #advent-of-code
- #distillery
- #processes
- #api
- #forms
- #metaprogramming
- #performance
- #security










First 10 of 33 Posts!
Qqwy
What is the reason you are building a postgres-backed queue?
Because:
I believe there already are quite a few packages doing either of these:
Taskmodule or ex_job,samphilipd
Our requirements are:
I have been looking at existing solutions for some time, and none of them suit our needs, so I opted to build one.
Qqwy
All right, then let me give you a slightly more in-depth answer to your original question
.
Both approaches (working with a worker pool vs a new task for each job) are valid: The main idea behind a worker pool is that frequently you don’t want too many background jobs running at the same time, because some OS resources (like the outgoing internet connection used to fetch data from elsewhere, or like the email handler) can be easily flooded if not rate-limited in some way.
You should not care about (garbage collection) efficiency at this time, because this would be a premature optimization (and the BEAM handles its garbage quite well). I believe using one or more worker pools is the cleaner solution.
As for your requirements:
I gather you already are using an SQL database for something else in your application, so it is not considered ‘additional’?
I’d like to note that Mnesia is able to write all changes to disk just like an SQL database does, and you can obviously have one (or multiple) nodes that just run Mnesia just like you’d run the SQL database in a separate location, besides replicating its data across all nodes. Obviously Mnesia is going to be faster because it is in-memory and lives right next to your application, using the same datatypes that the application uses.
Could you explain the reasoning behind this one? Ecto is not that heavyweight of a dependency, and it abstracts the database-specific details away from you, allowing you to swith between a whole slew of DBs.
Either use something that can handle erlang-terms natively like Mnesia, or use Erlang’s external term format to convert to/from a binary format that you could store on disk (or in a binary column in a database).
minhajuddin
Genstage would be ideal for this use case, You’d have a producer hooked up via a postgres pubsub mechanism which knows how many jobs are available to be processed at any given time (no polling). And a bunch of consumers which would try to get 1 job at a time and process it.
sasajuric
This is how I usually end up doing it. Based on your description, I’d have a
GenServerdispatcher process, which is responsible for the lifecycle of jobs. I’d implement polling in a separate process (possibly also aTask). That way, a crash when reading from the database would not affect the dispatcher, and therefore would not cause currently running jobs to fail.Instead of having a separate
Tasksupervisor, I tend to use the dispatcher as the parent of the job processes. I setup the dispatcher to trap exits, and start each new task withTask.async. With such setup, the dispatcher process receives all the necessary events (task result, task failure) in form of messages.Using a separate
Task.Supervisorshould also work. I usually don’t do that, since the supervisor wouldn’t really lift a lot of responsibilities from the dispatcher. The dispatcher would have to setup a monitor to each task, and respond to:DOWNmessages, and keep monitors and running tasks in its state. So it doesn’t seem that a dedicated supervisor would bring anything useful to the table.The supervision subtree could look something like:
outlog
take a look at
All though it’s called ectojob - I don’t see that ecto should be a requirement for writing the job to the DB or that anything prevents you from wrapping it in an ecto-less api..
(although jobs are often tied to DB insert/update and the ecto multi job hook looks pretty darn awesome)
your “no ecto” requirement is the only one that sticks out to me.. what is the reason for “no ecto” ?
OvermindDL1
Sounds costly, why not use a PostgreSQL Stream and get results streamed back to you in real-time whenever a table gets a new entry?
You could then just toss it into GenStage/Flow and let it all handle the concurrency as well.
For note, the PostgreSQL Library that Ecto uses (and thus you’d already have it if you use Ecto) has the Stream calls wrapped up in a nice API for you to use.
/me hopes Ecto gets such interfaces too so we get the auto-conversions and such
Yeah Ecto is quite lightweight and does a lot of caching for you and so forth, really useful to have.
Yep! This is it!
brightball
Only other thing I’d suggest looking into would be how you plan to lock jobs. There’s several options from using a lock column to advisory locks and a couple of others that can have a significant performance impact.
samphilipd
Correct. And I would guess that a large majority of Elixir deployments in the wild are also backed by a Postgres DB.
Sure. But mnesia requires your nodes to have state, and it is a PITA to manage (I have used it in production for a while).
It’s not that heavyweight, but it doesn’t make sense for this application for several reasons:
Sure, I looked into Postgres pubsub and it seems interesting, but adds a little extra complexity. It’s not clear-cut whether it would be more performant than a poll-based approach. My current polling implementation is simple and fast enough, once that is stable in production I’ll look into pubsub.
That supervision tree looks interesting… but why separate the Dispatcher and the Poller? Couldn’t you have the dispatcher set a timer to poll itself and fetch/dispatch the jobs synchronously itself inside the handle_info function?
This library aims to solve a slightly different problem than ecto_job. Note that ecto_job uses SELECT FOR UPDATE and holds transactions open for the duration of the job, this limits your concurrency to one job per connection. For some use-cases this is fine, but I want this to be more of a general purpose job queue.
Polling is surprisingly very performant. The problem with streaming events is lock contention since every Dispatcher will try to get a lock at the same time. While I think you could make this work with advisory locks which are very lightweight, you’d always need a polling system as a failsafe so I’d prefer to get that really rock solid first then consider Postgres pubsub as a potential optimisation.
In addition, you don’t need GenStage at all with a polling system since it’s purely pull-based - even simpler.
I am using a combination of advisory locking and SELECT FOR UPDATE SKIP LOCKED.
Initial benchmarks show a throughput of about 1.5k jobs per second on my local machine which is plenty fast enough for our needs.
sasajuric
This separation allows the dispatcher to do its work without any significant delays or interrupts if database polling becomes extremely slow and/or it crashes.
Last Post!
outlog
latest Oban supports sqlite oban | Hex