I fail to understand the following aspects though, maybe you can clarify them a bit:
* How does a director recover from a failure? From my understanding it would require fetching all job IDs from jobs that are in an active state via the job transactions table (which sounds expensive) and then loading the associated meta-data from the jobs table? Is that correct?
* Do you assume that the director will have archived all non-completed jobs when deleting the database? Do you try to gracefully shut down the director first then? From my understanding it seems you perform a "drop table" statement on a given database and then regenerate the tables, but this would require being sure that all the jobs have been processed or archived.
On your second point, archiving actually happens in the “drainers”, not the director. It’s the component that picks up unused databases and flush their data back into the set of running directors. We initially built archiving into the directors but it turned out to steal too much resources away from job processing so we moved it out into the drainers, which actually gives us opportunities to do it more efficiently (for example creating larger archives, which helps with compression as well).
Sorry if we cut some details off of the post, there is plenty more to tell about this system but it’s a lot for a single blog post ;)
> How does a director recover from a failure? From my understanding it would require fetching all job IDs from jobs that are in an active state via the job transactions table (which sounds expensive) and then loading the associated meta-data from the jobs table? Is that correct?
This is correct, if a director crashes for any reason, it needs to scan the database on boot. It is a more expensive operation, which is part of the reason that we try and cap the number of total entries in a given database.
> Do you assume that the director will have archived all non-completed jobs when deleting the database? Do you try to gracefully shut down the director first then? From my understanding it seems you perform a "drop table" statement on a given database and then regenerate the tables, but this would require being sure that all the jobs have been processed or archived.
Great question, we glossed over this aspect this a bit in the post itself.
Before a given database is transitioned to the 'spare' state, and its tables dropped, a single Drainer process is responsible for moving any non-completed jobs from that database to another active Director. The Drainer will not successfully exit and transition the database to 'spare' until it is certain it has processed all the non-completed jobs. We never drop any tables which have non-terminal jobs. Similar to the Directors, the Drainer will acquire a lock in consul to ensure only a single process is draining at a time.
We're hoping to go into a bit more depth on how the drainers work and these jobs move around in an upcoming post on Centrifuge's two-phase commit semantics. Ensuring that your data has moved to another system does require fairly complex transactional semantics, so we're hoping to go into depth about how this works.
I would imagine you're taking the database version from Consul and using that as a fencing key for writes within a transaction in the database, but I didn't see that mentioned in the Go code snippet.
I have been on the lookout for a similar system to SideKiq.
Unfortunately for Go, there isn't anything that matches it 100%.
I have looked into:
https://github.com/celrenheit/sandglass
https://github.com/RichardKnop/machinery
https://github.com/contribsys/faktory
In the end I went with machinery due to supporting:
- Batches (or Groups, tasks executed at the same time)
- Chains (tasks executed one by one)
- Callbacks (a task executed, on completion of a Batch)
- Cron Jobs
- Rescheduling failed jobs
- Long-runnng Jobs
- Distributed/Fault Tolerant DB (for Jobs)
- Distributed/Fault Tolerant Workers
+ other things I cant remember now.
It would be very interesting to know if this is going to be opened sourced and whether it supports the list above.
From their own ReadMe:
> Fairway isn't meant to be a robust system for processing queued messages/jobs. To more reliably process queued messages, we've integrated with Sidekiq.