How Marin's StepRunner Handles DAG Execution: A Deep Dive into the Core Scheduler
Marin's StepRunner executes directed acyclic graphs (DAGs) by eagerly pulling steps from an iterable, expanding transitive dependencies via post-order traversal, and launching steps as soon as all parent dependencies complete, while supporting parallel execution through a thread pool and intelligent caching.
Marin is an open-source pipeline orchestration framework maintained by marin-community/marin. At its heart lies the StepRunner, a deterministic DAG scheduler that orchestrates StepSpec objects into a coherent execution flow. Understanding how this component handles dependency resolution, concurrency, and caching is essential for building efficient, reproducible data pipelines with Marin.
What is StepSpec? The Building Block of Marin's DAG
Before execution begins, Marin represents each unit of work as a StepSpec—an immutable data object defined in lib/marin/src/marin/execution/step_spec.py (lines 25–102). Unlike imperative task definitions, a StepSpec describes what to run, not how to run it.
Key fields include:
name– A human-readable identifier for the step.hash_attrs– JSON-serializable configuration used for content hashing and cache invalidation.deps– A list of otherStepSpecobjects representing upstream dependencies.fn– The callable (orRemoteCallable) that performs the actual computation.resources– An optionalResourceConfigthat triggers remote execution via Fray.writes_record– A boolean indicating if the step writes its own artifact record for lazy artifacts.
The output path of a step—used for status files, cache checks, and as a unique execution key—is derived from a prefix plus the step's name_with_hash. This path serves as the canonical identity for the step throughout the DAG execution process.
The Execution Loop: How StepRunner Schedules Steps
The StepRunner.run() method in lib/marin/src/marin/execution/step_runner.py (lines 27–44) implements the core scheduling logic. It maintains several bookkeeping sets to track execution state:
completed– Paths of steps that finished successfully.failed– Paths of steps that errored or had failed dependencies.running– Mapping ofoutput_path → JobHandlefor active executions.waiting– List of steps whose dependencies are not yet satisfied.pruned– Paths of cache-only leaves skipped due to existing outputs.scheduled– Global set of all paths already seen for deduplication.
The runner follows a four-phase cycle for each incoming step: expand the dependency graph, harvest completed jobs, flush waiting steps with satisfied dependencies, and determine the fate of the current step based on its state and cache status.
Expanding the DAG with _expand_unseen
When StepRunner.run() receives a step, it invokes _expand_unseen(step, seen, is_built, pruned) (lines 39–56) to traverse the dependency graph post-order. This function walks the step's transitive dependencies and returns a list of steps that have not yet been scheduled.
If _expand_unseen encounters a cycle, it raises ValueError, though cycles should never appear because the graph is explicitly built as a DAG. When is_built(node) returns True—indicating the step is already cached—the node's output path is added to pruned, and its children are not traversed. This optimization prunes entire sub-trees of already-successful work from the schedule.
The Core Scheduling Algorithm
After expansion, the runner processes each step through a decision tree:
- Cache confirmation – If the path exists in
pruned, call_confirm_pruned()to verify the cached artifact. - Failure propagation – If any dependency path exists in
failed, mark the step as failed. - Immediate launch – If all dependencies are in
completed, submit the step via_do_launch(). - Defer execution – Otherwise, append the step to
waitingfor later processing.
The loop continues until both running and waiting collections are empty. If any steps failed during execution, the runner raises a combined exception containing all failure details.
Parallel Execution and Concurrency Control
Marin utilizes a ThreadPoolExecutor to enable parallel step execution. The executor is initialized with max_concurrent workers (defaulting to 8) and persists for the duration of the run() call.
Each ready step is submitted to the pool via _launch_step, which wraps the step's run_step call in a worker function. This wrapper also propagates the current Fray client context, ensuring that remote jobs inherit the same authentication and configuration as the parent process. The executor guarantees that no more than max_concurrent steps execute simultaneously, regardless of how many steps are ready in the waiting queue.
Caching, Drift Detection, and Pruning
Before submitting a step to the thread pool, _launch_step (lines 65–78) performs several validation checks to avoid redundant work:
- Dry-run guard – When
dry_run=True, the runner skips all I/O operations and simulates execution. - Drift check – The
check_driftfunction determines if mutable (development) artifacts require rebuilding regardless of cache status. - Cache validation – The runner reads the step's status file; if it contains
STATUS_SUCCESSand the step is immutable, the execution is skipped entirely.
These mechanisms ensure that Marin's DAG execution is idempotent—re-running a pipeline will automatically skip steps that completed successfully in previous runs, while still detecting configuration changes that necessitate rebuilding.
Remote Execution via Fray Integration
When a StepSpec defines a resources field, the StepRunner automatically dispatches the step as a Fray job rather than executing it locally. This remote execution flow (lines 155–215) branches based on the callable type:
- Standard functions with resources trigger
_run_iris_job(). RemoteCallableinstances invoke_run_remote_step().
Both paths ultimately call _submit_iris_job(), which constructs a JobRequest, submits it to the current Fray client, and blocks until completion. This allows Marin pipelines to seamlessly scale from local development to distributed computing without changing the DAG structure.
Practical Examples
Building a Simple Dependency DAG
from marin.execution.step_spec import StepSpec
from marin.execution.step_runner import StepRunner
# Leaf step – tokenizes some data
tokenize = StepSpec(
name="tokenize",
fn=lambda out: open(out, "w").write("tokens"),
hash_attrs={"vocab": "v1"},
)
# Dependent step – builds a vocab from the tokenized output
build_vocab = StepSpec(
name="build_vocab",
deps=[tokenize],
fn=lambda out: open(out, "w").write("vocab"),
hash_attrs={"size": 1000},
)
# Top step – trains a model that depends on both previous steps
train = StepSpec(
name="train",
deps=[tokenize, build_vocab],
fn=lambda out: open(out, "w").write("model"),
hash_attrs={"layers": 12},
)
# Run the DAG
StepRunner().run([train])
The runner expands the transitive dependencies (tokenize, build_vocab) and executes them in topological order before launching the train step.
Parallel Execution with Concurrency Limits
# Create many independent steps
steps = [
StepSpec(name=f"step_{i}", fn=lambda out, i=i: open(out, "w").write(str(i)))
for i in range(20)
]
# Run up to 5 steps concurrently
StepRunner().run(steps, max_concurrent=5)
This configuration processes all 20 steps using only 5 concurrent workers, respecting resource constraints while maximizing throughput.
Remote Execution with Resource Requirements
from fray.types import ResourceConfig
remote = StepSpec(
name="remote_job",
fn=lambda out: open(out, "w").write("remote result"),
resources=ResourceConfig(cpu=4, memory_gb=16), # forces Fray submission
)
StepRunner().run([remote])
When resources is specified, the StepRunner automatically submits the step to the Fray cluster rather than executing it in the local process.
Summary
- StepSpec objects in
lib/marin/src/marin/execution/step_spec.pyserve as immutable DAG nodes containing dependency references, hashable configuration, and execution logic. - Post-order expansion via
_expand_unseenprunes cached sub-trees and detects cycles before scheduling begins. - State-machine scheduling tracks completed, failed, running, and waiting steps to ensure dependencies complete before downstream execution.
- ThreadPoolExecutor enables parallel execution with configurable concurrency limits (
max_concurrent). - Intelligent caching checks status files and drift conditions to skip redundant work, ensuring idempotent pipeline runs.
- Fray integration automatically promotes steps with
ResourceConfigto remote jobs without altering the DAG structure.
Frequently Asked Questions
How does StepRunner handle circular dependencies in the DAG?
The _expand_unseen method raises a ValueError if it encounters a cycle during the post-order traversal of the dependency graph. However, cycles should never appear in practice because Marin's StepSpec graphs are explicitly constructed as directed acyclic graphs by design.
What happens when a step fails during DAG execution?
When a step fails, its output path is added to the failed set. The StepRunner then marks all downstream steps as failed during the scheduling phase, preventing wasted computation on dependent tasks. Once the execution loop completes, the runner raises a combined exception containing details of all failures encountered during the run.
Can StepRunner resume a partially completed pipeline?
Yes. The runner automatically prunes steps where is_built(node) returns True by adding them to the pruned set and skipping their children. When re-running a pipeline, successfully completed steps with valid status files and matching hashes are skipped entirely, allowing the pipeline to resume from the point of failure or interruption.
What is the difference between local and remote step execution?
Local execution runs the step's fn callable directly within the ThreadPoolExecutor. Remote execution occurs when a StepSpec includes a resources field or contains a RemoteCallable; in these cases, the StepRunner submits the work to a Fray cluster via _submit_iris_job(), blocking until the remote job completes while still respecting the local concurrency limits for job submission.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →