Most explanations of Kahn's algorithm stop when it prints an order. In real systems nobody wants the order for its own sake. A CI pipeline, a data workflow, a training job made of preprocessing and evaluation stages, or a package installer wants to run the tasks, as many at once as the machines allow, while never starting a task before everything it depends on has finished. Kahn's algorithm is unusually good at this, because its core loop already has the shape of a scheduler: a counter per task of unfinished prerequisites, and a pool of tasks whose counter has reached zero.

This article builds a small but complete DAG executor on Kahn's algorithm: a task state machine, event-driven dispatch, critical-path priority, retries, failure propagation and restarts, measured on a release pipeline with two workers. The static algorithm, its DFS alternative and deterministic orders are covered in Topological Sort in depth, so here they get only a short recap.

Advertisement

From a sort to a scheduler

A directed acyclic graph has an edge from u to v when u must finish before v starts. Kahn's algorithm computes the in-degree of every node, puts every node with in-degree zero into a ready set, and repeatedly removes a node from that set, appends it to the output, and decrements the in-degree of each successor. A successor whose in-degree reaches zero joins the ready set. A short output means a cycle. The cost is O(V + E).

The step that turns this into a scheduler is small. In the textbook version, a node is removed from the ready set and immediately counts as done. In an executor, removing a node from the ready set means starting it, and the decrements happen later, when it finishes. The in-degree counter now means "prerequisites that have not yet succeeded", and the ready set means "tasks that could start now if a worker were free". Everything else in this article follows from keeping those two meanings exact.

The life of a task

An executor needs more states than a sort. A task is PENDING while some prerequisite has not succeeded, READY when all have, RUNNING once a worker picks it up, and then one of three terminal states: SUCCEEDED, FAILED after its retries are exhausted, or SKIPPED because an ancestor failed and it can never become ready.

PENDINGin-degree > 0READYin-degree = 0RUNNINGon a workerSUCCEEDEDdecrement successorsFAILEDretries exhaustedSKIPPEDan ancestor failedlast prereqdispatchokerrorerror, attempts left: back to READYpropagatePENDING descendants of a failure
The life of one task in a Kahn-based executor. Only the coordinator moves tasks between states; workers just run tasks and report back.

Two rules keep the machine honest. Only a success decrements successors; a failed or skipped task must never release its dependants, or a broken build would publish. And the run is over when the count of terminal states equals the task count. Executors that stop when the ready set is empty either hang or report success after a failure, because the failed task's descendants stay PENDING forever. Counting terminal states and skipping descendants explicitly closes both holes.

Advertisement

Why level barriers waste time

A tempting shortcut is to compute Kahn levels and run them one at a time, joining the thread pool after each level. It is correct, but it wastes time.

The barrier makes every task in level k wait for the slowest task in level k minus 1, even when its own prerequisites finished long ago. The event-driven alternative has no levels at all: whenever any task finishes, decrement its successors immediately and start whatever became ready. A task starts the moment its last prerequisite succeeds and a worker is free, which is the earliest any correct schedule could start it on that worker. The worked example below quantifies the difference.

An event-driven executor

Here is a complete executor in Python. One coordinator thread owns all the state; workers in a pool only run tasks. Because the coordinator is single-threaded, plain dictionaries are safe and the decrement needs no lock. The ready set is a heap keyed on a priority, explained in the next section.

import heapq
from concurrent.futures import ThreadPoolExecutor, FIRST_COMPLETED, wait

def run_dag(cost, deps, run, workers=4, max_retries=1):
    """cost: {task: estimated seconds}; deps: {task: [prerequisites]};
    run(task) raises on failure. Returns {task: terminal state}."""
    succ = {t: [] for t in cost}
    indeg = {t: 0 for t in cost}
    for t, pres in deps.items():
        for pre in pres:
            succ[pre].append(t)
            indeg[t] += 1
    prio = bottom_levels(cost, succ, indeg)   # raises if the graph has a cycle
    state = {t: "PENDING" for t in cost}
    tries = {t: 0 for t in cost}
    ready = []

    def make_ready(t):
        state[t] = "READY"
        heapq.heappush(ready, (-prio[t], t))  # largest bottom level first

    for t in cost:
        if indeg[t] == 0:
            make_ready(t)
    running, terminal = {}, 0
    with ThreadPoolExecutor(workers) as pool:
        while terminal < len(cost):
            while ready and len(running) < workers:
                _, t = heapq.heappop(ready)
                state[t] = "RUNNING"
                tries[t] += 1
                running[pool.submit(run, t)] = t
            if not running:                       # nothing runs, nothing is ready
                stuck = sorted(t for t in cost if state[t] == "PENDING")
                raise RuntimeError(f"no progress possible; pending: {stuck}")
            finished, _ = wait(running, return_when=FIRST_COMPLETED)
            for fut in finished:
                t = running.pop(fut)
                if fut.exception() is None:
                    state[t] = "SUCCEEDED"
                    terminal += 1
                    for s in succ[t]:
                        indeg[s] -= 1
                        if indeg[s] == 0:
                            make_ready(s)
                elif tries[t] <= max_retries:
                    make_ready(t)                 # retry: same priority, back in the heap
                else:
                    state[t] = "FAILED"
                    terminal += 1 + skip_descendants(t, succ, state)
    return state

def skip_descendants(t, succ, state):
    skipped, stack = 0, list(succ[t])
    while stack:
        s = stack.pop()
        if state[s] == "PENDING":                 # already SKIPPED via another path: ignore
            state[s] = "SKIPPED"
            skipped += 1
            stack.extend(succ[s])
    return skipped

Three details matter. The no-progress check cannot fire on a validated acyclic graph, so it is a defensive assertion against bugs or a graph mutated mid-run. Retries go back into the heap instead of running inline, so a flaky task does not hold a worker. And the skip walk touches only PENDING nodes, which keeps it safe when two failures share descendants.

Choosing what runs next: critical-path priority

When more tasks are ready than workers are free, the executor must choose, and the choice decides how long the run takes. A FIFO queue starts tasks in the order they became ready, which depends on how edges happen to be listed. A better rule, used by classic list scheduling, is to prioritise by bottom level: the length of the longest path from the task to the end of the graph, counting its own cost. A task with a large bottom level heads a long chain, and delaying it delays everything behind it.

Bottom levels come from one pass over a topological order in reverse, which is also where the executor validates the graph. If the static Kahn pass cannot place every node, there is a cycle and nothing should start.

def bottom_levels(cost, succ, indeg):
    indeg = dict(indeg)
    order = [t for t in cost if indeg[t] == 0]
    for t in order:                               # Kahn: the list grows while we walk it
        for s in succ[t]:
            indeg[s] -= 1
            if indeg[s] == 0:
                order.append(s)
    if len(order) != len(cost):
        raise ValueError(f"cycle among {sorted(set(cost) - set(order))}")
    bl = {}
    for t in reversed(order):
        bl[t] = cost[t] + max((bl[s] for s in succ[t]), default=0)
    return bl

Be clear about what this buys. Minimising the makespan of precedence-constrained tasks on m identical machines is NP-hard, so no fast rule is optimal in general. Graham showed that any list schedule, one that never idles a worker while a task is ready, finishes within a factor of 2 minus 1/m of the optimum. Critical-path priority does not improve that worst-case bound, but in practice it usually lands much closer to optimal than FIFO and tolerates rough cost estimates. Use historical median durations as costs.

Worked example: a release pipeline on two workers

Take a release pipeline with eight tasks and estimated costs in minutes: A checkout (1), B compile (4), C lint (1), D unit tests (2), E build docs (1), F package (1), G integration tests (3) and H publish (1). The edges are A to C, A to E and A to B, listed in that order; B to D and B to F; F to G; and D, G, C and E all to H. Total work is 14 minutes, and with two workers no schedule can beat the larger of 14 divided by 2 and the longest chain.

Bottom levels, computed backwards: H is 1, G is 3 + 1 = 4, D is 2 + 1 = 3, F is 1 + 4 = 5, C and E are 1 + 1 = 2, B is 4 + max(3, 5) = 9, and A is 1 + 9 = 10. The critical path is A, B, F, G, H, with length 10, so 10 minutes is a lower bound here.

PolicyKey decisionMakespan
Event-driven, bottom levelAt t=1, B (9) starts before C and E (2)10 (optimal)
Event-driven, FIFOC and E, listed first, take both workers; B waits to t=211
Level barriersG waits for D, not just F, until t=711
Event-driven, bottom-level priority (makespan 10)W1AB compileDHW2CEFG testsEvent-driven, FIFO in edge order (makespan 11)W1ACB compileDHW2EFG testsLevel barriers, any priority (makespan 11)W1AB compileDG testsHW2CEF01234567891011time units
Three policies on the same eight-task release pipeline and two workers. The barrier version wastes time because G waits for D, which it does not depend on.

The gap grows with graph width and with the spread of durations, which is why ML and data pipelines, where one stage takes hours and its siblings seconds, feel it most.

Failures, retries and skipped tasks

Make the failure policy explicit. The executor above retries up to a limit, then marks the task FAILED and its PENDING descendants SKIPPED while unrelated branches keep running. That "continue" policy reports every independent problem in one run. "Fail fast" stops dispatching and cancels running tasks after the first failure, which saves compute when the rest of the run is worthless.

Retries need two guards. Retry only plausibly transient errors such as timeouts, preemption and rate limits, never assertion failures. And make tasks safe to run twice: write to a temporary path and rename atomically on success. Wrap each run in a deadline too, so a hung task becomes an observable failure.

Cycles and graphs that change while running

Validate the graph before starting anything; the bottom-level pass does this for free and names the nodes it could not place. Report the actual cycle, not the leftover set, which includes innocent downstream nodes; a three-colour DFS over the leftovers finds it, as shown in cycle detection in directed graphs.

Graphs that change during a run are the harder case. When a task spawns new tasks, such as one shard per input file, add them through the coordinator, count in-degrees only against prerequisites that have not succeeded, and re-check for cycles. For graphs edited often while running, topological sort, batch versus online covers incremental ordering. A cycle that slips through looks exactly like the no-progress condition, which is why the executor raises instead of waiting.

Concurrency, scale and restarts

A single coordinator does O(1) work per edge, so it scales far. Beyond that, workers decrement counters themselves. Then the decrement must be atomic, and exactly one decrement must observe the transition to zero, or a task runs twice or never. In Java that is if (pending.get(s).decrementAndGet() == 0) readyQueue.add(s);; in a database it is a conditional update such as UPDATE task SET pending = pending - 1 WHERE id = ? RETURNING pending inside the same transaction that records the parent's success. Never read the counter and then decrement it in two steps.

For restarts, persist each state transition before acting on it. On restart, treat SUCCEEDED as final, reset RUNNING to READY (hence idempotent tasks), and recompute in-degrees from prerequisites that have not succeeded rather than persisting counters, which can drift. For very large graphs, use compressed adjacency arrays; the traversal patterns are in BFS and DFS in depth.

Failure modes and trade-offs

  • Release on failure. Decrementing successors in a finally block, so they run after a failed parent. Only success may decrement.
  • Hang after failure. Ending the loop when the ready set is empty, or never ending it at all, because descendants of a failure stay PENDING. Count terminal states and skip descendants explicitly.
  • Double dispatch. A non-atomic read-then-decrement lets two finishing parents both see zero and enqueue the child twice. Use an atomic decrement, or a single coordinator.
  • Retry storms. Retrying deterministic failures, or retrying instantly against a rate-limited service. Classify errors, back off, cap attempts.
  • Non-idempotent tasks. A restart or retry reruns a task that appends rather than overwrites, and outputs double. Write to a temporary path and rename.

Trade-offs: a single coordinator is simplest; distributed decrements remove its bottleneck but add atomicity concerns. Fail fast saves compute; continue gives more information per run.

What to do next

  1. Write down your executor's state machine and check that only success decrements successors.
  2. Change the termination test to count SUCCEEDED, FAILED and SKIPPED tasks against the total.
  3. Replace level barriers with event-driven dispatch, and measure makespan before and after on a real pipeline.
  4. Record per-task durations and use their medians to compute bottom levels for priority.
  5. Validate the graph before starting, and report a real cycle, not just the leftover nodes.
  6. Make every task idempotent, persist state transitions, and test a crash and restart in the middle of a run.
  7. If workers update counters directly, make the decrement atomic and test it under concurrency.
Key takeaway: Kahn's algorithm becomes a scheduler when removing a task from the ready set means starting it and decrementing successors means it succeeded. Dispatch on every completion instead of by level, prioritise ready tasks by bottom level, release dependants only on success, skip the descendants of failures, and end the run when terminal states equal the task count. Validate for cycles first, keep decrements atomic, and make tasks idempotent so retries and restarts are safe.