Split a recursive computation into tasks and hand them to P threads, and you hit the load-balancing problem immediately: you do not know in advance how big each subproblem is, so any fixed division leaves some threads idle while others finish. Work stealing solves it without a central queue. Every worker keeps its own double-ended queue of ready tasks, pushes and pops at one end like a stack, and only when it runs dry does it pick another worker at random and steal a task from the other end.

This algorithm sits under Cilk, Java's ForkJoinPool, Intel TBB, Rust's Rayon and Tokio, and the Go scheduler, and it comes with one of the cleanest performance guarantees in parallel computing. This article covers the model that makes the guarantee precise, measures it with a small simulator, explains why thieves must take the oldest task, and walks through the Chase-Lev deque that makes the fast path cheap. Runtime-specific details of particular pools are left to the linked articles.

Work, span and the guarantee

Model a parallel program as a directed acyclic graph of unit-time instructions. An edge says one instruction must finish before another starts; a spawn node has two successors, and a sync node waits for several predecessors. Two numbers describe the graph. The work T1 is the total number of nodes, the time on one processor. The span T∞ is the length of the longest path, the time on infinitely many processors. Their ratio T1/T∞ is the parallelism: the largest speedup any scheduler could hope for.

Two lower bounds hold for every scheduler on P processors: TP ≥ T1/P, because only P nodes run per step, and TP ≥ T∞, because the critical path is sequential. Any greedy scheduler that never idles a processor while work is ready achieves TP ≤ T1/P + T∞, within a factor of two of optimal. The catch is that a greedy scheduler needs global knowledge of all ready work, which means a shared queue and contention on every operation.

Blumofe and Leiserson (JACM 1999) proved that randomised work stealing, which uses no global knowledge at all, achieves expected time T1/P + O(T∞) on fully strict fork-join computations, with an expected O(P T∞) steal attempts, and space at most P times the serial stack space. The practical reading: when parallelism greatly exceeds P, the T1/P term dominates, speedup is nearly linear, and steals are rare because they scale with the span, not the work.

Randomised work stealing: owners work at the bottom, thieves take from the topworker 0runs one task at a timetop (oldest)bottom (newest)worker 1runs one task at a timetop (oldest)bottom (newest)worker 2 (idle)runs one task at a timetop (oldest)bottom (newest)fib(18) contfib(17) contfib(16) contfib(15) contfib(17) contfib(16) contpush and pop at bottom: no contentionown deque after a stealsteal oldest (largest) taskpicks a victimuniformly at randomBound: expected T_P = T_1 / P + O(T_inf); expected steals O(P T_inf)Old tasks sit near the root of the computation, so one steal moves a large share of the remaining work.The owner pays for an atomic only when it races a thief for the last task.
Figure 1. Each worker owns a deque. The owner pushes and pops newest tasks at the bottom; an idle worker steals the oldest task from a random victim's top.

The scheduling loop and the work-first principle

The scheduling loop for one worker is short. The version below is the model the simulator implements: when a node enables two successors, the worker pushes one and keeps executing the other; when its deque is empty it makes one steal attempt per step against a uniformly random victim.

def worker_loop(me, deques, rng):
    task = None
    while not computation_done():
        if task is None:
            task = deques[me].pop_bottom()          # newest first: depth-first, cache-warm
        if task is None:
            victim = rng.choice(others(me))
            task = deques[victim].steal_top()       # oldest first: the biggest chunk
            continue                                 # a failed attempt costs a step
        enabled = execute(task)                      # returns successors now ready
        if len(enabled) == 2:
            deques[me].push_bottom(enabled[1])       # the continuation waits
            task = enabled[0]                        # work-first: run the child now
        elif len(enabled) == 1:
            task = enabled[0]
        else:
            task = None

Executing the spawned child immediately and leaving the parent's continuation in the deque is the work-first principle from Cilk-5 (Frigo, Leiserson and Randall, PLDI 1998). The single-worker execution then looks exactly like the serial program, the deque holds at most one continuation per level of recursion, and all scheduling overhead is pushed onto the rare steal path. The alternative, help-first or child stealing, pushes the child and keeps running the parent. That is simpler to implement in libraries without compiler support, which is why ForkJoinPool and Rayon push the forked task, but a loop that spawns n tasks can then fill the deque with n entries before any of them runs.

Because the owner always works at the bottom, it executes the most recently created, smallest, cache-warm task, giving depth-first order. The top of the deque holds the oldest task, which is the continuation closest to the root and therefore represents the largest remaining piece of the computation.

Worked example: a simulator

The simulator builds the fork-join graph of fib(20) with unit-cost nodes: a spawn node, a continuation node that spawns the second call, and a sync node per call, with fib(0) and fib(1) as single leaf nodes. That gives T1 = 43,781 nodes and T∞ = 40, so the parallelism is about 1,094. Each step every worker either executes one node or makes one steal attempt. Results, averaged over three random seeds:

PT_PT_1/P + T_infSpeedupStealsAttempts
143,78143,8211.0000
221,90421,9302.00727
410,96910,9853.993195
85,4995,5137.9655208
162,7752,77615.78175614
321,4111,40831.033111,371
6473172459.896263,003

Speedup is nearly linear up to 64 workers, and the measured time tracks T1/P + T∞ closely; at 32 and 64 workers it slightly exceeds it, which is why the theorem has a constant inside the O and you should not quote the simple sum as a hard ceiling. The steal counts tell the more important story: 626 successful steals to execute 43,781 nodes on 64 workers. Steals per worker per unit of span stayed around 0.2 to 0.3 from 16 to 64 workers, exactly the O(P T∞) scaling the theory predicts.

Now change one line so thieves take the newest task from the bottom of the victim's deque instead of the oldest. At P = 16 successful steals rose from 175 to 1,391, about eight times as many, and TP rose from 2,775 to 2,901. At P = 64, steals rose from 626 to 2,998 and TP from 731 to 830. The newest task is a tiny subtree near the leaves, so a thief that takes it runs out almost immediately and has to steal again. Taking the oldest task is not a detail; it is what makes steals rare.

The Chase-Lev deque

The fast path is the owner's push and pop, which happen on every spawn, so they must avoid locks and, ideally, any atomic read-modify-write. The Chase-Lev deque (SPAA 2005) achieves this with a circular array and two indices: top, advanced only by successful steals via compare-and-swap, and bottom, written only by the owner. The owner needs a CAS only when it races a thief for the last element. Lê, Pop, Cohen and Zappa Nardelli (PPoPP 2013) gave a version for the C11 memory model with proofs; the core, following their fence placement, with signed indices so an empty deque cannot wrap:

typedef struct { _Atomic int64_t top, bottom; _Atomic(Array *) array; } Deque;

void push(Deque *q, Task *x) {
    int64_t b = atomic_load_explicit(&q->bottom, memory_order_relaxed);
    int64_t t = atomic_load_explicit(&q->top, memory_order_acquire);
    Array *a = atomic_load_explicit(&q->array, memory_order_relaxed);
    if (b - t > a->size - 1) a = grow(q, a, t, b);       /* copy into a bigger ring */
    atomic_store_explicit(&a->buf[b % a->size], x, memory_order_relaxed);
    atomic_thread_fence(memory_order_release);           /* publish the task first */
    atomic_store_explicit(&q->bottom, b + 1, memory_order_relaxed);
}

Task *take(Deque *q) {                                   /* owner only */
    int64_t b = atomic_load_explicit(&q->bottom, memory_order_relaxed) - 1;
    Array *a = atomic_load_explicit(&q->array, memory_order_relaxed);
    atomic_store_explicit(&q->bottom, b, memory_order_relaxed);
    atomic_thread_fence(memory_order_seq_cst);           /* the fence people forget */
    int64_t t = atomic_load_explicit(&q->top, memory_order_relaxed);
    Task *x = NULL;
    if (t <= b) {
        x = atomic_load_explicit(&a->buf[b % a->size], memory_order_relaxed);
        if (t == b) {                                    /* last element: race thieves */
            if (!atomic_compare_exchange_strong_explicit(&q->top, &t, t + 1,
                    memory_order_seq_cst, memory_order_relaxed))
                x = NULL;                                /* a thief won */
            atomic_store_explicit(&q->bottom, b + 1, memory_order_relaxed);
        }
    } else {
        atomic_store_explicit(&q->bottom, b + 1, memory_order_relaxed);  /* was empty */
    }
    return x;
}

Task *steal(Deque *q) {                                  /* any other thread */
    int64_t t = atomic_load_explicit(&q->top, memory_order_acquire);
    atomic_thread_fence(memory_order_seq_cst);
    int64_t b = atomic_load_explicit(&q->bottom, memory_order_acquire);
    if (t >= b) return NULL;                             /* empty */
    Array *a = atomic_load_explicit(&q->array, memory_order_acquire);
    Task *x = atomic_load_explicit(&a->buf[t % a->size], memory_order_relaxed);
    if (!atomic_compare_exchange_strong_explicit(&q->top, &t, t + 1,
            memory_order_seq_cst, memory_order_relaxed))
        return ABORT;                                    /* lost a race: retry elsewhere */
    return x;
}

The seq_cst fence in take sits between the owner's store to bottom and its load of top. Without it, on x86 as well as on ARM, the store can be delayed past the load, so owner and thief both believe they hold the last task and execute it twice. Two other details are easy to miss: the old array cannot be freed when the deque grows, because a thief may still be reading it, so it needs epoch-based or hazard-pointer reclamation, the same problem as in a lock-free stack; and indices only grow, so a 64-bit counter is required. Unless you are writing a runtime, use a vetted implementation such as crossbeam-deque, which is what Rayon builds on.

How production schedulers adapt it

Real schedulers adjust the textbook algorithm in a few recurring ways:

  • Steal half. Go and Tokio give each worker a fixed-size local run queue and let a thief take about half of the victim's tasks in one operation, which suits many small independent tasks better than stealing one at a time.
  • A LIFO slot. Both also keep a single next-task slot so a task that wakes another runs it immediately on a warm cache; the slot is protected from immediate stealing so the woken task stays local. The Rust async runtime article covers Tokio's version.
  • Global injection queue. Tasks submitted from outside the pool go to a shared queue that idle workers poll alongside stealing.
  • Parking idle workers. Spinning thieves waste CPU; runtimes cap the number of searching workers and park the rest, waking one when work appears.
  • Locality-aware victims. On multi-socket machines, trying victims on the same NUMA node first keeps data local, at some cost to the randomness the proof relies on.
  • Granularity cutoffs. Below some problem size, run serially; a spawn that creates less than a few microseconds of work costs more than it saves.

Failure modes

  • Not enough parallelism. If T1/T∞ is only a few times P, the span term dominates and steals multiply. Measure work and span, for example with Cilkscale, before blaming the scheduler.
  • Blocking inside a task. A task that waits on I/O or a lock pins its worker and its deque; the rest of the pool cannot help. ForkJoinPool's compensation and Tokio's blocking pool exist for this, described in the ForkJoinPool internals article.
  • Grain too fine. Millions of tiny tasks turn the fast path into the main cost; raise the serial cutoff until per-task overhead is a small fraction.
  • Grain too coarse. A handful of large tasks leaves nothing to steal at the end; the last task determines the finish time.
  • Help-first loops. Spawning every iteration of a large loop before running any fills the deque and memory; split the range recursively instead, as in divide and conquer.
  • False sharing. Putting top and bottom on the same cache line makes every steal attempt slow down the owner; pad them apart.
  • Memory ordering bugs. A missing fence in a hand-written deque shows up as tasks run twice or lost, rarely and only on some hardware.

Trade-offs

Compared with a single shared queue, work stealing removes contention from the common path and keeps execution depth-first and cache-friendly, at the price of no global priority order: a worker may run its own small tasks while an older, more urgent task sits in another deque. Compared with static partitioning, it adapts to irregular work but adds per-spawn overhead and nondeterministic execution order. Work-first execution gives the best space bound and serial efficiency but needs compiler or language support for continuations; help-first is a library-friendly approximation. For latency-sensitive services, fairness mechanisms such as global queues and steal-half matter more than the asymptotic bound, which is about throughput of a single computation.

What to do next

  1. Write the simulator: build the fib(20) graph, confirm T1 = 43,781 and T∞ = 40, then reproduce the speedup and steal columns above.
  2. Change thieves to take from the bottom and confirm the jump in steals; keep both numbers as a test that your real deque steals from the right end.
  3. Measure the work and span of one of your own parallel jobs and compute its parallelism; if it is under about ten times your core count, fix the algorithm first.
  4. Tune the serial cutoff by timing several grain sizes; plot time per element against grain and pick the knee.
  5. Find every blocking call reachable from pool tasks and move it to a blocking pool or an async API.
  6. Read how a production runtime adapts the algorithm in goroutines at scale and compare its choices with the textbook loop.
Key takeaway: Work stealing gives each worker its own deque, runs newest tasks locally and lets idle workers steal the oldest task from a random victim. With enough parallelism it achieves expected time T_1/P plus O(T_inf) with steals proportional to the span, as the simulator showed: near-linear speedup to 64 workers with 626 steals for 43,781 nodes, and about eight times more steals at 16 workers when thieves took the newest task instead. Keep tasks coarse enough, never block in them, and use a proven Chase-Lev implementation rather than hand-placing fences.