4 ms·
I recently started working on my own DAG execution framework, after failing to get some patches into Airflow to make the scheduling easier to reason about. My
by iroddis 5y ago
I recently started working on my own DAG execution framework, after failing to get some patches into Airflow to make the scheduling easier to reason about.
My typical use case was orchestrating DAGs with thousands of vertices, and airflow would silently wedge itself and fail to report errors.
Daggy is just starting out, but I’m hoping it’ll become more robust and scalable as time goes on.
https://gitlab.com/iroddis/daggy https://gitlab.com/iroddis/daggy
- slotrans 5y ago(workflow/DAG systems are an area of special interest for me...) Happy to see you're supporting runtime-dynamic DAG shape. That's one of the critical things Airflow is missing. But, if I'm inferring correctly, it looks like you're recording task state in a similar way to Airflow, i.e. in some kind of data store. I suggest you look at how Luigi models task state: using "targets" which are a task's materialized outputs (e.g. a file, or anything whose existence can be tested for). Luigi's model isn't perfect but it's very valuable for tasks to be able to recognize the condition "my output exists, therefore I don't need to run". ...which leads me to the most important feature you need to check off, that most competitors are getting wrong: resuming a partially failed build. It's often the case that one node of a DAG will fail with a non-recoverable error that can't be solved with retries, and so the DAG as a whole will fail. But if we can fix that error in a way that's compatible with the work that was already done (e.g. fix a network ACL, IAM policy, etc), we should be able to start the DAG again and have it resume at the point of failure, and not do any work over again. Luigi does this very well, Airflow does it adequately, and I think all the other tools in this space (Prefect, Dagster, Nextflow... tell me about others!) don't do it at all.
- thundergolfer 5y agoI know this thing you want by the name “incremental” processing. GrailBio’s pipeline tech is incremental, as is Lyft’s Flyte I think. Build systems such as Bazel are build on the same incremental DAG processing foundation. I agree that it’s the most important feature, as it dramatically improves correctness and performance (when at scale). IAM policy trouble is a good example of a bug that crashes a pipeline but has no influence on precomputed outputs of succeeded nodes.
- iroddis 5y agoThanks for the feedback. I'll take a look at how Luigi models task state. Right now each TaskExecutor type is responsible for running and reporting on tasks (e.g. the Slurm executor submits jobs and monitors them for completion). I was considering adding a companion "verify" stage for every vertex, which would be a command that ran and verified output. It might be a way to do what I think you're describing above without having to build in a variety of expected outputs into the daggy core. I'll check what Luigi is doing, though. > resuming a partially failed build Daggy does this! Right now it will continue running the DAG until every path is completed or all vertices in a processing state (queued, running, retry, error) are in the error state, then the DAG goes to an error state. It's possible to explicitly set task/vertex states (e.g. mark it complete if the step was manually completed), then change the DAG state to QUEUED, at which point the DAG will resume execution from where it left off. [1] is a unit test that walks through that functionality. [1] https://gitlab.com/iroddis/daggy/-/blob/master/tests/unit_server.cpp#L275 https://gitlab.com/iroddis/daggy/-/blob/master/tests/unit_se...