---
title: "Distributed Compute and Schedulers"
book: "Research, Data and Risk Platforms"
subject: quant
language: en
chapter: 13
exercises: 8
source: https://one-course.com/books/quant/15/en/chapter/13-distributed-compute-and-schedulers
---

# Chapter 13 — Distributed Compute and Schedulers

A [parameter sweep](#def-pl-distributed-compute-and-schedulers-sweep) of ten thousand backtests was submitted at 18:00 on a Friday and finished at 13:47 on the Saturday. From 10:20 that morning, fewer than a tenth of the cluster’s cores had been busy. Forty of the tasks were thirty times longer than the median, and they had been submitted last: the cluster reached them after twelve hours and then waited seven and a half for them. Nothing had failed; the scheduler had done exactly what it was told. This chapter is about what to tell it: how a sweep becomes a [task graph](#def-pl-distributed-compute-and-schedulers-graph), how schedulers order the graph on a cluster, why the last few tasks decide when a sweep ends, and what caching, retries and backups cost and save.

## 13.1 Clusters and what they are for

**Definition 13.1 (Compute cluster, job scheduler).**

A *compute cluster* is a set of machines (nodes), each with cores, memory and sometimes accelerators, managed as one pool of capacity. A *job scheduler* decides which submitted work runs on which part of the pool and when, according to a policy, the resources each piece of work requests and the priorities of the users who submitted it.

A research platform uses its cluster for work that is large, repetitive and mostly independent: [parameter sweeps](#def-pl-distributed-compute-and-schedulers-sweep), walk-forward refits (Book 7, chapter 20), simulations over seeds (chapters 11 and 12), risk scenarios (chapters 18 and 20), model training (Book 12). Most of it is the easiest kind of parallel work there is.

**Definition 13.2 (Embarrassingly parallel workload, parameter sweep).**

A workload is *embarrassingly parallel* when it splits into tasks that do not communicate while they run. A *parameter sweep* is such a workload: the same computation, typically a backtest, run once for each point of a grid of parameters, data slices or seeds.

[Embarrassingly parallel](#def-pl-distributed-compute-and-schedulers-sweep) does not mean embarrassingly easy to finish on time. The tasks of a sweep have unequal durations — a longer lookback, a busier market, a slower data slice — and the sweep ends only when its last task does.

**Definition 13.3 (Makespan, straggler).**

The *makespan* of a set of tasks is the time from the start of the first to the end of the last. A *straggler* is a task that finishes long after most of its peers, because it is intrinsically long, because it runs on a slow or overloaded machine, or because it failed and was run again.

## 13.2 Schedulers and queues

The simplest scheduler keeps one queue and gives the next free core the task at its head: first in, first out. Graham (1969) analysed this *list scheduling* and bounded it: whatever the order, the [makespan](#def-pl-distributed-compute-and-schedulers-makespan) on $m$ identical cores is at most $2 - 1/m$ times the best possible, and taking the tasks longest first — *longest processing time*, LPT — tightens the bound to $4/3 - 1/(3m)$. The bounds rest on a simple observation, which is also the chapter’s measuring stick.

**Proposition 13.4 (Lower bound on the makespan).**

No schedule of tasks with durations $d_i$ on $m$ cores finishes before $\max\bigl(\sum_i d_i / m,\ \max_i d_i,\ L\bigr)$, where $L$ is the longest chain of dependent tasks: the work must be done by $m$ cores, the longest task runs on one core, and a chain runs in sequence.

```python
def lower_bound(tasks: list[Task], total_cores: int) -> float:
    """No schedule beats the work over every core, the longest task, or the longest chain."""
    by_id = {t.tid: t for t in tasks}
    memo: dict = {}

    def chain(tid):
        if tid not in memo:
            t = by_id[tid]
            memo[tid] = t.duration + max((chain(d) for d in t.deps), default=0.0)
        return memo[tid]

    work = sum(t.duration * t.cores for t in tasks)
    longest = max(t.duration for t in tasks)
    return max(work / total_cores, longest, max(chain(t.tid) for t in tasks))
```

***Listing 13.1.** The lower bound of [Proposition 13.4](#prop-pl-distributed-compute-and-schedulers-bound), with the longest chain computed over the task graph. code/firm/jobgraph/firm_jobgraph.py*

LPT needs to know which tasks are long. A scheduler does not know durations, only estimates — from the parameters (a longer window, a longer history), from earlier runs of the same code, or from the user. The chapter’s estimates are right to within a factor $e^{0.3z}$ with $z$ standard normal, a typical spread for a model of run time.

**Definition 13.5 (Fair-share scheduling).**

*Fair-share scheduling* gives each user or team a target share of the cluster and orders the queue so that those using less than their share are served first; a team that submits ten thousand tasks does not lock out a team that submits five hundred an hour later.

```python
class FairShare(FIFO):
    """Per-team queues; next task from the team using fewest cores per unit of target."""

    def __init__(self, targets: dict):
        self.targets, self.qs = targets, {t: deque() for t in targets}

    def push(self, task: Task, front: bool = False) -> None:
        (self.qs[task.team].appendleft if front else self.qs[task.team].append)(task)

    def pop(self, core: int, running: dict):
        use = {t: 0 for t in self.targets}
        for task in running.values():
            use[task.team] += task.cores
        ready = [t for t in self.targets if self.qs[t]]
        if not ready:
            return None
        team = min(ready, key=lambda t: (use[t] / self.targets[t], t))
        return self.qs[team].popleft()

    def __len__(self) -> int:
        return sum(len(q) for q in self.qs.values())
```

***Listing 13.2.** Fair share in its simplest form: the next task comes from the team using the fewest cores per unit of its target. code/firm/jobgraph/firm_jobgraph.py*

**As of September 2026 — Fair share in a production scheduler.**

Slurm, a widely used open-source workload manager, computes a fair-share factor from normalised shares $S$ and normalised usage $U$ (its “Fair Tree” algorithm, the default since Slurm 19.05): an association’s level fair-share is $S/U$, above 1 when it is under-served, and users are ranked so that the children of a better-served account always rank below those of a worse-served one.

**Definition 13.6 (Work stealing).**

*Work stealing* is a decentralised scheduling policy: each worker keeps its own queue of tasks and takes from its head; a worker whose queue is empty takes a task from the far end of another worker’s queue. Blumofe and Leiserson (1999) analysed it for multithreaded computations; it balances load without a central queue.

```python
class WorkStealing(FIFO):
    """Tasks dealt round robin to per-core deques; a core takes its head or steals a tail."""

    def __init__(self, cores: int):
        self.dq = [deque() for _ in range(cores)]
        self.i = 0
        self.steals = 0

    def push(self, task: Task, front: bool = False) -> None:
        self.dq[self.i % len(self.dq)].append(task)
        self.i += 1

    def pop(self, core: int, running: dict):
        if self.dq[core]:
            return self.dq[core].popleft()
        victim = max(range(len(self.dq)), key=lambda c: len(self.dq[c]))
        if self.dq[victim]:
            self.steals += 1
            return self.dq[victim].pop()
        return None

    def __len__(self) -> int:
        return sum(len(d) for d in self.dq)
```

***Listing 13.3.** Work stealing: tasks dealt round robin to per-core deques, an idle core stealing from the fullest. code/firm/jobgraph/firm_jobgraph.py*

## 13.3 Task graphs

**Definition 13.7 (Task graph).**

A *task graph* is a directed acyclic graph whose nodes are tasks, each with a resource request and an estimated duration, and whose edges say which task’s output another task reads; a task is ready when all its predecessors have finished.

A sweep is rarely a flat list. Its backtests read features; the features read cleaned data; the data were built once for many sweeps. Written as a graph, each shared stage is one task that many backtests depend on, and its output can be cached under the content key of Book 7’s research workflow (chapter 29): the hash of its code, parameters and inputs. The next sweep over the same data finds the features in the cache and runs only the backtests.

## 13.4 Parameter sweeps and stragglers

![The chapter’s task graph and cluster. A hundred feature stages of twenty minutes each, each read by a hundred backtests; every backtest’s duration drawn from the measured spread, forty of them thirty times the median in the flat sweep. The feature stages are cached by content key.](https://one-course.com/images/onecourse/chapters/quant-15/pl-distributed-compute-and-schedulers/fig-dc51f863b531.svg)

***Figure 13.1.** The chapter’s [task graph](#def-pl-distributed-compute-and-schedulers-graph) and cluster. A hundred feature stages of twenty minutes each, each read by a hundred backtests; every backtest’s duration drawn from the measured spread, forty of them thirty times the median in the flat sweep. The feature stages are cached by content key.*

The chapter’s cluster is simulated (`firm.jobgraph`); its tasks are not invented. Twenty-four small backtests of chapter 11 — level-3 replays of five-minute `firm.tape` sessions — were run on this laptop through the same scheduler interface, with one worker: their median was 0.089 seconds, the spread of their log durations 0.76, and the longest took 8.0 times the median. The simulated sweep keeps that spread and scales the median to fifteen minutes: 10 000 backtests, forty of them thirty times the median (seven and a half hours, a long lookback on a long history) submitted last, as a sweep over a grid naturally submits its largest parameters. The work adds up to 3 649 core-hours on 384 cores, a lower bound of 9.50 hours.

![Busy cores over time for the same sweep under two schedules, sampled every ten minutes, on a log scale. First in, first out keeps the cluster full for about nine hours, reaches the forty long tasks only after twelve, and runs them on an almost idle cluster until the last ends at 19.8 hours; longest first with speculative backups ends at 10.0 hours. Data: fig_jobgraph.py.](https://one-course.com/images/onecourse/chapters/quant-15/pl-distributed-compute-and-schedulers/fig-2742b48b3ac2.svg)

***Figure 13.2.** Busy cores over time for the same sweep under two schedules, sampled every ten minutes, on a log scale. First in, first out keeps the cluster full for about nine hours, reaches the forty long tasks only after twelve, and runs them on an almost idle cluster until the last ends at 19.8 hours; longest first with speculative backups ends at 10.0 hours. Data: `fig_jobgraph.py`.*

[Table 13.1](#tab-pl-distributed-compute-and-schedulers-policies) is the sweep under each policy. First in, first out reaches the last of the forty long tasks only after twelve hours, and the sweep ends at 19.8 hours — 475.0 node-hours of a cluster held for a job that needed 3 649 core-hours. Longest first starts the long tasks at once and ends at 11.8 hours. [Work stealing](#def-pl-distributed-compute-and-schedulers-stealing), which balances load but does not reorder it, ends at 16.3.

| policy | [makespan](#def-pl-distributed-compute-and-schedulers-makespan) (hours) | node-hours | backups |
| --- | --- | --- | --- |
|  | plain | with backups | plain | with backups |  |
| first in, first out | 19.79 | 20.32 | 475.0 | 487.6 | 152 |
| longest first | 11.80 | 10.02 | 283.3 | 240.4 | 105 |
| [work stealing](#def-pl-distributed-compute-and-schedulers-stealing) | 16.33 | 16.40 | 391.9 | 393.5 | 145 |
| longest first, no slow node | 10.32 | — | 247.6 | — | — |
| lower bound | 9.50 |  |  |  |

***Table 13.1.** The 10 000-task sweep on 384 cores under each policy, without and with speculative backups (a backup after 2.5 times the estimate), and longest first on a cluster without the slow node. Node-hours are the cluster’s 24 nodes held for the [makespan](#def-pl-distributed-compute-and-schedulers-makespan).*

![Makespan of the sweep by policy, without and with speculative backups (after 2.5 times the estimate), against the lower bound. First in, first out (FIFO) and work stealing run the long tasks last; longest first (LPT) runs them first. Data: fig_jobgraph.py.](https://one-course.com/images/onecourse/chapters/quant-15/pl-distributed-compute-and-schedulers/fig-6ce27728b623.svg)

***Figure 13.3.** [Makespan](#def-pl-distributed-compute-and-schedulers-makespan) of the sweep by policy, without and with speculative backups (after 2.5 times the estimate), against the lower bound. First in, first out (FIFO) and [work stealing](#def-pl-distributed-compute-and-schedulers-stealing) run the long tasks last; longest first (LPT) runs them first. Data: `fig_jobgraph.py`.*

**Definition 13.8 (Speculative execution).**

*Speculative execution* launches a second copy of a task that is running much longer than expected, on another machine; the first copy to finish is kept and the other is killed.

Google’s MapReduce paper named the problem and the cure: a [straggler](#def-pl-distributed-compute-and-schedulers-makespan) can be a machine with a bad disk whose read rate falls from 30 to 1 MB/s, and when a job is close to completion the master schedules backup executions of the remaining in-progress tasks; their sort took 44% longer with backups disabled. The threshold matters. The chapter’s backups start when a task has run 2.5 times its estimate: with estimates right to within $e^{0.3z}$, an honest task passes that point with probability 0.11% ($z > 3.05$), so backups go mostly to tasks on the slow node. Longest first gains 1.8 hours from them; first in, first out loses 0.5, because its long tasks are late to start, not slow to run, and the backups only take cores; the wasted work — killed copies, failed and retried attempts — is 91 to 130 core-hours, 2.5 to 3.6% of the sweep.

```python
    late: list = []                       # (core, attempt) of tasks waiting for a backup copy

    def fill():
        free.sort()
        while late and free:                  # a late task's backup goes before the queue
            core, aid = late.pop(0)
            if attempt.get(core, (0, None))[1] != aid or len(copies[running[core].tid]) != 1:
                continue
            node = core // cluster.cores
            other = next((c for c in free if c // cluster.cores != node), None)
            if other is None:
                late.insert(0, (core, aid))
                break
            free.remove(other)
            start(other, running[core], backup=True)
        while free and len(policy):
            core = free[0]
            task = policy.pop(core, running)
            if task is None:
                break
            free.pop(0)
            start(core, task)
```

***Listing 13.4.** Speculative execution in the simulator: a late task’s backup goes to a free core on another node before the queue is served. code/firm/jobgraph/firm_jobgraph.py*

**Example 13.9 (The sweep that finished on Saturday).**

Longest first with backups ends at 10.02 hours, 5.5% above the lower bound of 9.50 hours. Without the slow node, but also without backups, it ends at 10.32: the backups recover more than the slow node costs, because they also rescue the tasks whose estimates were wrong and those that failed. First in, first out on the same cluster ends at 19.79 hours, twice as long, and holds 475 node-hours against 240. The cluster had been almost idle for three and a half hours: from 16.3 hours after submission, fewer than a tenth of its cores were busy.

## 13.5 Fair share between teams

A second team submits 500 backtests an hour after the first team’s sweep. Under first in, first out they wait behind all 10 000: half of them are done 8.28 hours after submission and the last 11.4 hours after. Under fair share with equal targets, every core that frees goes to the second team until it holds half the cluster: half its tasks are done in 0.59 hours and all in 2.64, while the first team’s sweep ends at 16.68 hours instead of 19.79 — earlier, not later, because interleaving the second team’s tasks changes where and when the first team’s long tasks run, and that, not the extra 5% of work, sets its end.

## 13.6 Caching, failure and cost

When every backtest recomputes its features (twenty minutes each), the sweep needs 6 922 core-hours and 19.3 hours; written as a [task graph](#def-pl-distributed-compute-and-schedulers-graph) with a hundred feature stages it needs 3 514 core-hours and 11.4 hours; the next sweep over the same data finds the hundred stages in the cache and needs 3 475 core-hours and 10.3 hours. Failures cost what was lost: a failed attempt is retried from its start, so a failure of a seven-hour task costs up to seven hours, which is why long tasks write checkpoints (Book 12, chapter 23). The cost of a sweep is the capacity it holds: on a cluster rented by the node-hour, the difference between first in, first out and longest first with backups is 235 node-hours for the same answers.

## 13.7 Tutorial: a sweep on a simulated cluster

**Goal.** Calibrate a sweep on real backtests, schedule it under four policies on a simulated cluster, and measure what backups, fair share and caching change. **End state:** [Table 13.1](#tab-pl-distributed-compute-and-schedulers-policies) and [Figure 13.2](#fig-pl-distributed-compute-and-schedulers-profile).

1. **Calibrate** : `bench_jobgraph.py` runs 24 small backtests through `LocalExecutor` with one worker; `measured_sigma()` .
2. **Sweep** : `sweep(sigma)` , `cluster()` , `lower_bound` .
3. **Policies** : `policies(sigma)` ; `profile(sigma)` for the busy-core curves.
4. **Teams** : `fair_share(sigma)` .
5. **Caching** : `caching(sigma)` on the [task graph](#def-pl-distributed-compute-and-schedulers-graph) .

**What to change next.** Give the scheduler perfect estimates and see how close longest first comes to the bound; give it no estimates (every task the median) and see what remains of its advantage.

## 13.8 Build: the job graph

**Purpose.** A task-graph and scheduling library the platform uses to plan and run sweeps, and a simulator to choose policies before spending node-hours.

**Interface.** `Task`, `Cluster`, `simulate(tasks, cluster, policy, spec, cache) -> Schedule`; `FIFO`, `LPT`, `FairShare`, `WorkStealing`; `lower_bound`; `utilisation`; `LocalExecutor`.

**Rules.** The scheduler sees estimates, never durations; a failed or preempted attempt is retried from its start; a backup runs on another node and the loser is killed and counted as waste; a cache hit takes no time; everything is seeded.

**Acceptance tests.** `code/firm/jobgraph/tests/`: Graham’s example ($3,3,2,2,2$ on two cores: 7 against an optimum of 6); first in, first out against longest first; dependencies and the chain bound; fair share splitting cores; [work stealing](#def-pl-distributed-compute-and-schedulers-stealing) on an unbalanced deal; backups rescuing a slow node; retries; cache hits; the local executor’s order.

**Stretch.** Memory classes and multi-core tasks with backfilling; preemptible (spot) capacity at a lower price; usage decaying over time in fair share.

Sources and further reading

- R. L. Graham, “Bounds on multiprocessing timing anomalies”, *SIAM Journal on Applied Mathematics* 17(2), 1969.
- J. Dean and S. Ghemawat, “MapReduce: simplified data processing on large clusters”, OSDI 2004.
- J. Dean and L. A. Barroso, “The tail at scale”, *Communications of the ACM* 56(2), 2013.
- R. D. Blumofe and C. E. Leiserson, “Scheduling multithreaded computations by work stealing”, *Journal of the ACM* 46(5), 1999.
- Slurm documentation, “Fair Tree Fairshare Algorithm”.

## 13.9 Exercises

**Exercise 13.1 ★.**

Five tasks of durations 3, 3, 2, 2, 2 run on two cores. What is the [makespan](#def-pl-distributed-compute-and-schedulers-makespan) under longest first, and what is the best possible?

**Solution of Exercise 13.1.**

Longest first puts 3 and 3 on the two cores, then 2 and 2, then the last 2 on the first core free: 7. The best is 6 (3 and 3 on one core, 2, 2 and 2 on the other). The ratio $7/6$ equals Graham’s bound $4/3 - 1/(3m)$ for $m = 2$: the example is the worst case.

**Exercise 13.2 ★.**

The sweep’s median task takes fifteen minutes. How long do the forty long tasks take on a normal node, and on the slow node?

**Solution of Exercise 13.2.**

Thirty times fifteen minutes, 7.5 hours, on a normal node; four times that, 30 hours, on the slow node.

**Exercise 13.3 ★.**

Why is the lower bound of the chapter’s sweep the total work divided by the cores, and not the longest task?

**Solution of Exercise 13.3.**

The total work is 3 649 core-hours; spread over 384 cores that is 9.50 hours, longer than the longest task (7.5 hours on a normal node) and than any chain (the flat sweep has none). The bound takes the largest of the three.

**Exercise 13.4 ★★.**

Why does the backup threshold of 2.5 times the estimate rarely fire for a healthy task, and what would a threshold of 1.5 do?

**Solution of Exercise 13.4.**

A healthy task runs its true duration, which exceeds 2.5 times its estimate only when the estimate’s error $e^{0.3z}$ is below $1/2.5$, that is $z > 3.05$: a probability of 0.11%. At 1.5 the threshold is passed whenever $z > 1.35$, about 9% of tasks: hundreds of useless backups, each taking a core from the queue and wasting work.

**Exercise 13.5 ★★.**

Why does [work stealing](#def-pl-distributed-compute-and-schedulers-stealing) finish later than longest first, although no core is idle while work remains queued?

**Solution of Exercise 13.5.**

[Work stealing](#def-pl-distributed-compute-and-schedulers-stealing) keeps every core busy while work remains, but it takes tasks in the order they were dealt; the long tasks, submitted last, start late and end late, as under first in, first out. Balance is not order: only a policy that starts the longest tasks first avoids the tail.

**Exercise 13.6 ★★.**

Fair share lets the second team finish sooner. What does it do to the first team’s end, and why?

**Solution of Exercise 13.6.**

In this run it ends earlier, at 16.68 hours instead of 19.79: the second team’s 500 tasks are only about 5% of the work, and the first team’s end is set by where and when its forty long tasks run, which the interleaving changes — here to its advantage, in another run perhaps not.

**Exercise 13.7 ★★★.**

*Coding.* Run `policies` with the long tasks submitted first instead of last. What happens to first in, first out, and to longest first?

**Solution of Exercise 13.7.**

First in, first out improves from 19.79 to 14.41 hours (13.69 with backups): the long tasks start at once. It still trails longest first, which is unchanged at 11.80 (10.02 with backups) because it orders by estimate, not by submission: the other long-tailed tasks are still in random order under first in, first out.

**Exercise 13.8 ★★★.**

*Find the flaw.* “The sweep ran all night on a cluster of 384 cores and used only 3 649 core-hours, so the cluster is too big: we will halve it.”

**Solution of Exercise 13.8.**

Utilisation over a night says nothing about how the capacity was used: the sweep held 475.0 node-hours for 3 649 core-hours of work because of the order of its tasks, and it was almost idle for three and a half hours. On half the cluster the lower bound doubles to 19 hours; the fix is the schedule (longest first, backups), which ends near the bound on the cluster as it is.

## 13.10 Problem: The Sweep That Finished on Saturday

**Problem 13.1.**

Weekend problem — ten thousand backtests and a slow node

The chapter’s sweep, cluster and policies.

**Part I — The workload.**

1. Why is a [parameter sweep](#def-pl-distributed-compute-and-schedulers-sweep) [embarrassingly parallel](#def-pl-distributed-compute-and-schedulers-sweep) ?
2. How was the spread of task durations obtained?
3. What is the total work, and the lower bound on 384 cores?
4. Where are the forty long tasks in the submission order, and why?
5. What makes a task a [straggler](#def-pl-distributed-compute-and-schedulers-makespan) here, in three ways?

**Part II — Policies.**

6. What does Graham’s analysis guarantee for list scheduling and for longest first?
7. What does longest first need that first in, first out does not?
8. Why does [work stealing](#def-pl-distributed-compute-and-schedulers-stealing) not fix the long tasks?
9. What does fair share change between two teams?
10. When does [speculative execution](#def-pl-distributed-compute-and-schedulers-spec) help, and when does it waste?

**Part III — The numbers.**

11. Give the [makespan](#def-pl-distributed-compute-and-schedulers-makespan) of each policy with and without backups.
12. How long was the cluster almost idle under first in, first out?
13. How far from the lower bound is the best schedule, and why?
14. What does fair share do for the second team’s first half of tasks?
15. What do the [task graph](#def-pl-distributed-compute-and-schedulers-graph) and the cache save?

**Part IV — The verdict.**

16. State the *named result* : the [makespan](#def-pl-distributed-compute-and-schedulers-makespan) and node-hours of the sweep under each policy, the gain from [speculative execution](#def-pl-distributed-compute-and-schedulers-spec) , and the distance of the best schedule from the lower bound.
17. What would you change in the platform’s default scheduling of sweeps?
18. What would you ask users to provide with each sweep?
19. How would you decide the size of the cluster?
20. In one sentence: what decides when a sweep ends?

**Solution of Problem 13.1.**

1. Each backtest reads its inputs and writes its result without talking to the others.
2. From 24 small backtests of chapter 11 run on this laptop through the local executor: the spread of log durations, 0.76.
3. 3 649 core-hours; 9.50 hours on 384 cores.
4. At the end, because a grid over parameters submits its largest settings last.
5. Intrinsic length, a slow node, and failure with a retry from the start.
6. At most $2 - 1/m$ times the best [makespan](#def-pl-distributed-compute-and-schedulers-makespan) for any order; at most $4/3 - 1/(3m)$ for longest first.
7. Estimates of duration.
8. It balances load but keeps the dealt order, so the long tasks still start last.
9. It serves the team using less than its share first, so a small late sweep is not queued behind a large one.
10. When a task is slow because of where it runs (a slow node); it wastes when estimates are poor and the threshold fires on healthy tasks.
11. First in, first out 19.79 (20.32 with backups); longest first 11.80 (10.02); [work stealing](#def-pl-distributed-compute-and-schedulers-stealing) 16.33 (16.40); longest first without the slow node 10.32.
12. About three and a half hours: fewer than a tenth of the cores were busy from 16.3 hours to 19.8.
13. 5.5% (10.02 against 9.50), because of the slow node and the estimation error; without the slow node but without backups, 10.32.
14. Half its tasks are done 0.59 hours after submission instead of 8.28.
15. The graph computes the features once (3 514 instead of 6 922 core-hours, 11.4 instead of 19.3 hours); the cache saves the hundred feature stages in the next sweep (3 475 core-hours, 10.3 hours).
16. **Named result.** On 384 cores the 10 000-task sweep takes 19.79 hours and 475.0 node-hours first in, first out, 11.80 and 283.3 longest first, 16.33 and 391.9 with [work stealing](#def-pl-distributed-compute-and-schedulers-stealing) ; backups after 2.5 times the estimate save longest first 1.8 hours and cost first in, first out 0.5; the best schedule, longest first with backups, ends 5.5% above the lower bound of 9.50 hours.
17. Longest first on estimates by default, backups with a threshold set from the estimates’ error, fair share between teams.
18. An estimate of each task’s duration (or the parameters that predict it), and cache keys for shared stages.
19. From the lower bound of the typical sweep and the deadline: cores = work / (deadline $-$ longest task margin), then check with the simulator under the policy actually used.
20. Its last tasks: their length, where they run and when they start.

## 13.11 Interview questions

**Interview question 13.1 ★ developer.**

Your sweep’s last 1% of tasks takes as long as the first 99%. What do you do?

**Solution of Interview question 13.1.**

Find out why they are long: intrinsic (start them first with longest-first ordering on estimates), a slow machine (backups on another node), or failures (checkpoints, retries); then check the [makespan](#def-pl-distributed-compute-and-schedulers-makespan) against the lower bound.

*What the interviewer is looking for: Diagnose the tail before adding capacity.*

**Interview question 13.2 ★★ developer.**

What is [speculative execution](#def-pl-distributed-compute-and-schedulers-spec), and when is it a bad idea?

**Solution of Interview question 13.2.**

Running a second copy of a task that is late, on another machine, keeping the first to finish. It is a bad idea when tasks are late because they are long rather than because of where they run, when estimates are too poor to tell, or when the work has side effects that two copies would duplicate.

*What the interviewer is looking for: Thresholds from estimate error; idempotent tasks.*

**Interview question 13.3 ★★ developer, researcher.**

Two teams share a cluster; one submits huge sweeps. How do you keep the other productive?

**Solution of Interview question 13.3.**

Fair share with targets per team, so that a team under its share is served first; limits on concurrent cores per user; and priority for short interactive work over long sweeps.

*What the interviewer is looking for: Shares, not first come first served.*

**Interview question 13.4 ★★ researcher.**

How would you structure ten thousand backtests that share expensive feature computations?

**Solution of Interview question 13.4.**

As a [task graph](#def-pl-distributed-compute-and-schedulers-graph): the feature computations as shared stages with content-addressed cache keys, the backtests depending on them; then schedule the graph, longest estimated tasks first.

*What the interviewer is looking for: Graph, cache keys, reuse across sweeps.*

**Interview question 13.5 ★★★ developer.**

Design the scheduling layer of a research platform’s [compute cluster](#def-pl-distributed-compute-and-schedulers-cluster).

**Solution of Interview question 13.5.**

[Task graphs](#def-pl-distributed-compute-and-schedulers-graph) with resource requests and estimates; a queue with fair share between teams and longest-first ordering within a sweep; retries with checkpoints; backups with a calibrated threshold; a content-addressed cache; accounting of node-hours per team; and a simulator of the cluster to test policies before they are deployed.

*What the interviewer is looking for: Policy chosen by measurement and simulation, not by default.*
