Skip to content

Running Task Graphs with Dependencies

Build steps, data pipelines and start-up sequences are often graphs: fetch three sources, parse each, join two of them, index the third, then report. Running such a graph with asyncio means starting each job as soon as its dependencies are done — and deciding what a failure does to the rest. Python's standard library already has the scheduling half in graphlib.TopologicalSorter. Measured on Python 3.14 with a ten-node graph of simulated I/O jobs whose critical path was 0.50 s: running it level by level, with gather over each of the four topological levels, took 0.95 s, because each level waited for its slowest job. Starting each node the moment its own dependencies finished — with a task per node, or with a TopologicalSorter-driven loop — took 0.50 s; limited to two jobs at a time, 0.75 s. When one parse step failed, the task-per-node version inside a TaskGroup cancelled an unrelated index job that was already running; the scheduler marked the failure, skipped only its 2 dependents, and completed the other 7 jobs. A cyclic graph was rejected before anything ran, with CycleError: ('nodes are in a cycle', ['a', 'b', 'a']). This guide builds that scheduler.

Prerequisites

1. Describe the graph and reject cycles up front

Represent the graph as a mapping from each job to the jobs it depends on, and let TopologicalSorter validate it before running anything:

import graphlib

GRAPH = {
    "fetch_a": [], "fetch_b": [], "fetch_c": [],
    "parse_a": ["fetch_a"], "parse_b": ["fetch_b"], "parse_c": ["fetch_c"],
    "join_ab": ["parse_a", "parse_b"], "index_c": ["parse_c"],
    "report": ["join_ab", "index_c"], "notify": ["parse_c"],
}

sorter = graphlib.TopologicalSorter(GRAPH)
sorter.prepare()                         # raises CycleError if the graph has a cycle

Measured: a two-node cycle raised CycleError with the cycle itself as the second argument, ['a', 'b', 'a'], which names the offending jobs. Catching that at start-up is far better than discovering a deadlock halfway through a run, where two tasks wait for each other forever. Keep job definitions and dependencies as data like this, separate from the code that runs them, so the graph can be checked, printed and tested on its own.

Verify: a test builds the real graph and calls prepare().

2. Do not run the graph level by level

The simplest correct execution groups jobs into levels — all jobs whose dependencies are in earlier levels — and gathers each level:

async def run_by_levels(sorter, run):
    sorter.prepare()
    while sorter.is_active():
        ready = sorter.get_ready()
        await asyncio.gather(*(run(job) for job in ready))   # waits for the slowest in the level
        sorter.done(*ready)

Measured: 0.95 s for a graph whose critical path is 0.50 s. The first level held fetches of 0.30, 0.05 and 0.10 s, so parse_b — 0.30 s, depending only on the 0.05 s fetch — could not start until the slow fetch_a finished. Every level boundary is a barrier that the graph does not actually require. Levels are easy to reason about and fine when jobs in a level take similar times; with mixed durations, they waste most of the parallelism.

Verify: total run time is compared with the critical path; a large gap means barriers are being added.

Ten-job graph with a 0.50 s critical path 4 horizontal bars comparing level by level (4 gathers) with the others. Ten-job graph with a 0.50 s critical path level by level (4 gathers) 0.95 s task per node, awaits its deps 0.50 s graphlib scheduler, start when ready 0.50 s graphlib scheduler, max 2 at once 0.75 s Job durations 0.05-0.30 s, simulated with asyncio.sleep. Level barriers nearly doubled the run.

3. Start each job when its own dependencies finish

Two designs remove the barriers. The compact one gives every job a task that first awaits its dependencies' tasks:

async def run_task_per_node(graph, run):
    tasks: dict[str, asyncio.Task] = {}

    async def node(name):
        await asyncio.gather(*(tasks[d] for d in graph[name]))
        await run(name)

    order = graphlib.TopologicalSorter(graph).static_order()    # create dependencies first
    async with asyncio.TaskGroup() as tg:
        for name in order:
            tasks[name] = tg.create_task(node(name))

The more controllable one drives the sorter directly, starting jobs as they become ready and feeding completions back:

async def run_graph(graph, run, limit: int | None = None):
    sorter = graphlib.TopologicalSorter(graph)
    sorter.prepare()
    slots = asyncio.Semaphore(limit or len(graph))
    running: dict[asyncio.Task, str] = {}
    failed: set[str] = set()
    skipped: set[str] = set()

    async def guarded(job):
        async with slots:
            await run(job)

    while sorter.is_active():
        for job in sorter.get_ready():
            if any(d in failed or d in skipped for d in graph[job]):
                skipped.add(job)
                sorter.done(job)                       # skip, and let its dependents be skipped too
            else:
                running[asyncio.create_task(guarded(job))] = job
        if not running:
            continue
        done, _ = await asyncio.wait(running, return_when=asyncio.FIRST_COMPLETED)
        for task in done:
            job = running.pop(task)
            if task.exception() is not None:
                failed.add(job)
            sorter.done(job)
    return failed, skipped

Measured: both took 0.50 s, the critical path. With limit=2, the scheduler took 0.75 s — the limit is a single semaphore around each job, which is how to keep a graph from overwhelming a shared resource.

Verify: run time approaches the critical path when unlimited, and never more than limit jobs run at once.

4. Decide what a failure does to the rest of the graph

The two designs differ most when a job fails. In the task-per-node version the TaskGroup cancels everything still running as soon as any job raises:

task per node, parse_b fails:
  completed: fetch_a, fetch_b, fetch_c, parse_a, parse_c, notify
  cancelled while running: index_c          (does not depend on parse_b)

Measured: index_c, which shares nothing with the failed parse_b, was cancelled mid-run. The scheduler instead recorded the failure and skipped only the jobs downstream of it: join_ab and report were skipped, and the other seven jobs completed, including index_c and notify. Neither behaviour is wrong; they answer different questions. Fail-fast suits graphs where a partial result is useless — a deployment, a transaction-like batch. Skip-dependents suits graphs whose branches are independent outputs — reports, indexes, notifications — where finishing what can finish is the point. Report both sets at the end, since a "successful" run that skipped half the graph should not look like a clean one, as discussed in handling per-item errors in async pipelines.

Verify: a test that fails one job checks exactly which jobs completed, were skipped, and were cancelled.

When parse_b fails A grid of 2 rows by 5 columns. When parse_b fails design parse_b its dependents (join_ab, report) unrelated running job (index_c) outcome task per node in a TaskGroup failed never started cancelled mid-run ExceptionGroup raised graphlib scheduler, skip dependents failed skipped completed 7 completed, 1 failed, 2 skipped Fail-fast or skip-dependents: choose per graph.

5. Make reruns cheap

Graphs that run regularly benefit from remembering which jobs already succeeded for a given input, so a rerun after a failure only repeats what failed and what depends on it. Record each job's success with a key derived from its inputs, and skip jobs whose key is recorded:

async def run_cached(job: str, inputs_key: str, store, run):
    marker = f"{job}:{inputs_key}"
    if await store.exists(marker):
        return                                      # already done for these inputs
    await run(job)
    await store.set(marker, "ok")

With the scheduler from step 3, a rerun after parse_b failed would repeat parse_b, join_ab and report only. Durable checkpoints of this kind are covered in checkpointing progress in long-running async jobs. For graphs that outgrow a single process — many machines, retries across days, a UI — this is the point where a workflow engine earns its complexity; for anything that fits in one process and one run, graphlib and a short loop are enough.

Verify: rerunning after a single failure re-executes only the failed job and its dependents.

How should this graph run? A decision on What should the run optimise for with 4 outcomes. How should this graph run? What should the run optimise for? any graph TopologicalSorter.prepare() first CycleError names the cycle all-or-nothing result task per node in a TaskGroup first failure cancels all independent outputs scheduler skipping dependents 7 of 10 still completed shared resource semaphore around each job 0.75 s at limit 2 Level-by-level gather only when jobs in a level take similar time.

Verification

A task graph runs correctly when:

  • The graph is data, validated with TopologicalSorter.prepare() before any job starts.
  • Jobs start when their own dependencies finish, not at level boundaries.
  • Failure behaviour is chosen per graph — fail-fast or skip-dependents — and tested.
  • Concurrency is limited where jobs share resources, and reruns skip completed work.

Diagnostic Hook: log each job's start, end and the dependency that finished last before it started. The chain of "last dependency" links from the final job back to the start is the graph's actual critical path in that run — the jobs worth making faster — and a large gap between it and total run time means something is adding barriers.

Pitfalls & edge cases

  • Gathering level by level. Measured: 0.95 s against a 0.50 s critical path.
  • A TaskGroup when branches are independent. Measured: an unrelated job was cancelled.
  • Cycles discovered at run time. They deadlock; prepare() catches them first.
  • Unlimited parallelism on a shared resource. Put a semaphore around each job.

Frequently Asked Questions

How do I run tasks with dependencies in asyncio?

Describe the graph as a mapping of job to dependencies, validate it with graphlib.TopologicalSorter, and start each job as soon as its dependencies finish; this ran a ten-job graph in 0.50 s, its critical path.

Why is running a dependency graph level by level slow?

Each level waits for its slowest job before the next starts, even for jobs that do not depend on it: 0.95 s against 0.50 s in testing.

What happens when one task in a dependency graph fails?

In a TaskGroup, all running tasks are cancelled, including unrelated ones. A graphlib-driven scheduler can skip only the failed job's dependents: 7 of 10 jobs still completed.

How do I detect cycles in a task graph?

graphlib.TopologicalSorter(graph).prepare() raises CycleError naming the cycle, such as ['a', 'b', 'a'].