5 ms·
Hey HN, author of the post here. A number of Segment engineers are hanging around today, and we are happy to answer questions in the comments. Thanks in advance
by calvinfo 8y ago
Hey HN, author of the post here. A number of Segment engineers are hanging around today, and we are happy to answer questions in the comments. Thanks in advance for any feedback and thoughtful discussion!
- ryanworl 8y agoHow are you blocking multiple director processes from writing to the JobDB instances? I.e. node A gets partitioned away, a new node B picks up in its place, but then node A comes back and still thinks it owns the lock for the JobDB instance. 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.
- rbranson 8y agoIn short... atypically long timeouts on these locks. We can afford to inject a few minutes of delay on a small number of in-flight messages to preserve this safety property. I’d agree that we should probably have a second layer of fencing.
- ThePhysicist 8y agoThis is a very interesting blog post and some great engineering, congrats! 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.
- calvinfo 8y agoThank you! > 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.
- achille-roussel 8y agoOn your first point, you’re correct. Directors scan their databases on start to rebuild their cache and reschedule the jobs that need to get retried. The scan is usually quick for data that was recently inserted since they’re likely in cache, it make take a couple of minutes to scan everything. Because we keep the database to a small size we can cap how long this operation will take. We also scan the database in reverse order (using the primary key, thanks to the rough ordering of KSUIDs), to help reschedule he most recent jobs first, which is possible because the scan can happen concurrently with tye jobs processing. 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 ;)
- ronreiter 8y agoIs Centrifuge going to be open source?
- achille-roussel 8y agoWe do want to open source the code. In order to get to production we took some shortcuts which tightly integrated some components of this system with Segment’s infrastructure and would make it hard for anyone to deploy. We discussed open sourcing the code along with this blog post but we felt like it would have more value once we’ve put a bit more work into it to make it easier to work with, and add a couple of features we have in the pipeline.
- no1youknowz 8y agoWould like to know this as well. 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/celrenheit/sandglass https://github.com/RichardKnop/machinery https://github.com/RichardKnop/machinery https://github.com/contribsys/faktory 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.
- mnutt 8y agoI don't know if it exactly fits, but maybe check out some of the stuff Customer.io has been doing, such as https://github.com/customerio/fairway https://github.com/customerio/fairway? It seems like they're solving similar problems and contribute a lot of Go projects.
- no1youknowz 8y agoUnfortunately it's not. 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.