Skip to main content
PRISM

Work stealing

intermediate · occasionally asked

One shared queue every worker must touch, against per-worker deques where only an idle worker coordinates.

Computing the curve…

The problem it solves

Sixteen workers, one queue. Every worker takes every task from the same place, so every task costs one contended operation on one shared structure. The queue becomes the bottleneck long before the work does — the pool flattens out and adding workers stops helping, for exactly the reasons on the lock granularity page.

The alternative is a queue per worker. Now the common case touches nothing shared: a worker pushes and pops its own end, uncontended, at cache speed. The obvious objection is that the work will not be evenly distributed and some workers will finish early while others are still loaded — and the answer to that objection is the whole idea. An idle worker steals from somebody else’s queue.

The elegance is in who pays. A busy worker never coordinates with anyone. Only a worker that has run out does anything expensive, and it does it once per steal rather than once per task. The coordination cost is charged precisely to the situation that needs coordinating.

The mechanism

Each worker owns a double-ended queue. It pushes and pops at one end — the bottom, by convention — which is a purely local operation needing no synchronisation in the common case. Thieves take from the other end, the top, so the owner and the thief are working at opposite ends of the structure and only collide when it is nearly empty.

That choice of ends does two things. It minimises contention, since owner and thief rarely touch the same slot. And it steals the oldest task, which for recursive divide-and-conquer is the largest remaining chunk — so one steal transfers a big subtree rather than a leaf, and the number of steals stays small. Meanwhile the owner works on the newest task, which is the one whose data is still warm in cache.

That is why the pattern fits fork-join especially well: the same discipline that minimises contention also maximises locality and minimises the number of times anybody has to coordinate.

What the enumeration shows

A model page: throughput against worker count for the two designs.

The shared queue line flattens. Beyond the point where the queue’s own throughput is the limit, extra workers add nothing — the line is horizontal, which is the picture of one structure that everybody must touch.

The work-stealing line keeps climbing, because most operations are local pops costing a few nanoseconds rather than contended dequeues costing tens.

Then drag the imbalance parameter. At zero imbalance nobody ever runs dry, no steal happens, and every operation is a cheap local one — stealing costs nothing when it is not needed. As imbalance rises, the fraction of operations that are steals rises with it and the advantage narrows, because you are paying the coordination cost more often.

That is the honest shape of the trade: work stealing is not free, it is charged only when needed, and how often it is needed is a property of your workload.

The numbers worth carrying

  • A local deque pop is roughly 5ns; a contended shared-queue operation is roughly 50ns. Ten times, and the shared one gets worse with more workers while the local one does not.
  • At 16 workers and moderate imbalance, per-worker deques do around 30× the work of one shared queue in this model.
  • Steals should be rare — typically well under 1% of operations in a well-balanced fork-join computation. If your profiler shows many steals, the work is badly partitioned, not the scheduler misbehaving.
  • Steal from the opposite end to the owner. Same-end stealing collides constantly and gives up most of the benefit.

Where it breaks down

It cannot fix a single long task. If one task takes ten seconds and everything else takes a millisecond, no amount of stealing helps — nothing can be taken from a task already running. Work stealing balances queues, not tasks, so the granularity of your decomposition sets the floor.

Locality can suffer. A stolen task runs on a core whose cache knows nothing about it. For memory-bound work, stealing too eagerly can cost more in cache misses than it gains in balance, which is why some schedulers steal reluctantly or prefer to steal from a nearby core in a NUMA topology.

Blocking tasks break the model. A worker blocked on I/O or a lock is not doing work and not available to steal, so its queue sits idle. This is the thread pool deadlock hazard in a different form, and it is why fork-join pools want CPU-bound tasks. Java’s ManagedBlocker exists to tell the pool “I am about to block, please compensate”, which is an admission that the model does not otherwise handle it.

It is more complex than it looks. A correct lock-free deque with the owner and thieves at opposite ends is genuinely subtle — the classic Chase-Lev implementation is a well-studied piece of concurrent algorithm design with real ABA considerations. This is a “use the library” situation.

What people get wrong

“Work stealing balances the load.” It balances queues. A single oversized task cannot be split by any scheduler, which is why decomposition granularity matters more than the scheduler does.

“Stealing is expensive so minimise it.” Stealing is what keeps workers busy; the thing to minimise is the need for it, by partitioning better. A pool with zero steals and idle workers is worse than one with a few steals and none.

“A shared queue is simpler and fine.” Fine up to a handful of workers. The whole point of the chart is where it stops being fine, and it is a lower number than most people expect.

“More workers is always better.” Past the point where the queue or the memory bus saturates, more workers add contention and nothing else. Amdahl sets the other ceiling.

In production

ForkJoinPool in Java is the reference implementation and is what parallel streams and CompletableFuture’s default executor run on. Go’s scheduler gives each P a local run queue with stealing from others and from a global queue. Rust’s Rayon and Tokio both use work-stealing schedulers. .NET’s TPL does the same. libdispatch, TBB and Cilk all descend from the same lineage.

The practical guidance is mostly about how you decompose, since the scheduler is somebody else’s code:

Make tasks small enough to balance and large enough to be worth scheduling. A common rule of thumb is a few tens of microseconds per task — small enough that no worker is stuck holding an outsized one, large enough that scheduling overhead is a small fraction.

Recurse rather than pre-partitioning. Splitting into exactly n chunks for n workers assumes uniform cost and cannot recover when that is wrong. Recursive halving until a threshold lets the scheduler adapt.

Keep blocking work off the fork-join pool. Use a separate pool for it, for the reasons on the thread-pool page.

The follow-up questions

“Why is a shared queue a bottleneck?” — One contended operation per task, so the queue’s own throughput caps the pool regardless of worker count.

“Why steal from the opposite end?” — Less contention with the owner, and it takes the oldest and usually largest task, so one steal moves a lot of work.

“When does work stealing not help?” — One long indivisible task. The scheduler balances queues, not tasks.

“How big should a task be?” — Tens of microseconds as a rule of thumb, and recurse rather than pre-partition so the scheduler can adapt to uneven cost.

In an interview

What ForkJoinPool, Go’s scheduler and Rayon are all doing, and a good test of whether someone can name where the contention went.

  • work stealing
  • deque
  • load balancing
  • contention

Run these next

The rest of patterns in practice