# Downstream

> Nothing runs until what it needs has run.

Canonical URL: <https://datadriven.io/problems/downstream>

Domain: Python · Difficulty: hard · Seniority: L6 · Asked in: Meta

## Problem

Dependencies that circle back on themselves leave no valid run order, and that case has to raise rather than hang or return a partial result. Otherwise `tasks` maps each task name to its `deps` and an `op` such as `"add 3"` or `"mul 2"`, applied to the combined output of those deps once they have all run; return every task's output keyed by name.

## Example 1

Input:

```
tasks = {"load":{"op":"add 0","deps":["clean","enrich"]},"clean":{"op":"add 1","deps":["extract"]},"enrich":{"op":"mul 2","deps":["extract"]},"extract":{"op":"add 10","deps":[]}}
```

Output:

```
{"load":31,"clean":11,"enrich":20,"extract":10}
```

## Example 2

Input:

```
tasks = {"score":{"op":"mul 3","deps":["ingest"]},"ingest":{"op":"add 2","deps":[]},"publish":{"op":"add 5","deps":["score"]}}
```

Output:

```
{"score":6,"ingest":2,"publish":11}
```

## Worked solution and explanation

### What this problem really is

Take off the orchestration costume and this is dependency-ordered graph execution. A task needs its prerequisites' outputs in hand before it can run, and those outputs have to flow forward into it. Everyone can parse 'mul 2'. What separates candidates is noticing that dict order is not run order: in the first example `load` is listed before `extract`, so a single pass reads `outputs['clean']` before it exists and dies with a KeyError. And when the graph circles back on itself, a naive retry-until-done version spins forever waiting on tasks that can never become ready.

> **Schedule by readiness, not by position**
>
> Count how many unfinished dependencies each task still has. Run any task whose count is zero, then decrement every task that depends on it. A task is only ever touched after all of its inputs exist, and the leftovers at the end are exactly the cycle.

---

### Break down the requirements

#### Step 1: Build the waiting counts and the reverse edges

`waiting_on` holds each task's count of unfinished dependencies; `dependents` inverts the edges so that finishing a task tells you whom to notify. Building the reverse map up front is what keeps each edge to one visit.

#### Step 2: Seed the queue with the roots

Every task with zero dependencies goes into a `deque`. These are the roots; their combined input is 0, which is why 'add 10' on `extract` yields 10 and a root 'mul' would yield 0.

#### Step 3: Run a ready task, then release its dependents

Split `op` into a verb and an amount, convert the amount with `int`, sum the dependencies' stored outputs, and apply the verb. Then decrement each dependent and queue any that reach zero. `load` in the first example only becomes ready after both `clean` (11) and `enrich` (20) finish, giving 31.

#### Step 4: Count what ran to catch the cycle

When the queue drains, either every task has an output or some tasks were never released. The second case can only mean a cycle, so compare `len(outputs)` to `len(tasks)` and raise.

---

### The solution

**Schedule by readiness, forward outputs, catch the cycle**

```python
import operator
from collections import deque

OPERATIONS = {'add': operator.add, 'mul': operator.mul}


def execute_dag(tasks):
    if not tasks:
        return {}

    waiting_on = {name: len(spec['deps']) for name, spec in tasks.items()}
    dependents = {name: [] for name in tasks}
    for name, spec in tasks.items():
        for dep in spec['deps']:
            dependents[dep].append(name)

    ready = deque(name for name, count in waiting_on.items() if count == 0)
    outputs = {}
    while ready:
        name = ready.popleft()
        spec = tasks[name]

        verb, amount_text = spec['op'].split()
        try:
            apply_op = OPERATIONS[verb]
        except KeyError:
            raise ValueError("task '" + name + "' has unknown op '" + verb + "'")

        combined = 0
        for dep in spec['deps']:
            combined += outputs[dep]
        outputs[name] = apply_op(combined, int(amount_text))

        for child in dependents[name]:
            waiting_on[child] -= 1
            if waiting_on[child] == 0:
                ready.append(child)

    if len(outputs) != len(tasks):
        raise ValueError('cycle detected: some tasks can never run')
    return outputs
```

> **One visit per task, one per edge**
>
> O(V + E) time and space: each task enters the `deque` once and each dependency edge decrements one counter once. Re-scanning the task dict until nothing changes looks simpler but is O(V^2) on a long chain listed backwards, which is exactly the shape of the second example.

> **Fail with the task's name attached**
>
> A dispatch table like `OPERATIONS` keeps the op handling to one lookup, and the `try` around it turns an unknown verb into a named failure instead of a bare KeyError from deep inside the loop. Interviewers notice when a failure message tells the on-call engineer which task broke.

> **Retry-until-done never terminates on a cycle**
>
> Retrying unready tasks in a `while` until every task has output. On a cycle that condition never becomes true, so the function hangs instead of raising. The readiness count makes the stopping point explicit: an empty queue ends the run, every time.

---

## Common follow-up questions

- A task fails halfway through the run. What should happen to the tasks downstream of it? _(Tests wrapping each task's execution and deciding whether downstream tasks are skipped or the run aborts.)_
- How would you run independent tasks concurrently instead of one at a time? _(Tests dispatching every zero-count task in a wave to a worker pool.)_
- How would you resume a partially completed run without redoing finished tasks? _(Tests persisting `outputs` and re-seeding the counts from what already finished.)_
- How would you tell the user which tasks form the cycle? _(Tests reporting the actual cycle path rather than just the leftover set.)_

## Related

- [All practice problems](https://datadriven.io/problems)
- [Mock interview mode](https://datadriven.io/interview/downstream)
- [Python Interview Questions](https://datadriven.io/python-interview-questions)
- [Data Engineering Interview Prep Guide](https://datadriven.io/data-engineer-interview-prep)
- [Daily Challenge](https://datadriven.io/daily)

---

Source: DataDriven (https://datadriven.io). DataDriven is the data engineering interview community. Live code execution in SQL, Python, and Spark sandboxes. Every feature is open to every member.