skip to content

How does ForkJoinPool's work-stealing algorithm work, and what is the role of the per-worker deques?

level: seniorimportance: must knowfreq 68%

answer

  1. Per-worker deque, not one shared queue
  2. Owner: push/pop HEAD, LIFO (cache-hot, depth-first)
  3. Thief: steal from TAIL, oldest = biggest chunk
  4. Opposite ends → low contention, CAS only on conflict
  5. Pull model: idle threads steal work to themselves

basics

~20 s

Each worker thread has its own double-ended queue (deque) of tasks. A worker pushes and pops its own tasks from one end. When a worker runs out of work, it steals a task from the opposite end of another busy worker's deque, so idle threads stay busy and load stays balanced.

solid answer

~50 s

ForkJoinPool gives every worker thread its own work-stealing deque (double-ended queue). When a task forks subtasks, they are pushed onto the current worker's deque. That worker processes its own deque in LIFO order, popping from the top (head) — this keeps it on recently forked, cache-hot subtasks and naturally does depth-first execution. When a worker's own deque is empty, it becomes a thief: it picks another worker (the victim) and steals from the tail (the opposite end), taking the oldest, largest-grained task. Pushing and popping at your own end is mostly contention-free; stealing from the other end via CAS means a thief and the owner rarely touch the same node. The payoff is automatic load balancing without a central queue bottleneck: instead of distributing work centrally, idle threads pull work to themselves. It assumes plentiful, splittable, non-blocking tasks; if subtasks block, the stealing can't compensate and cores idle.

code

java · 27 lines
java
// Sum a large array via divide-and-conquer on a ForkJoinPool.
import java.util.concurrent.RecursiveTask;
import java.util.concurrent.ForkJoinPool;

final class SumTask extends RecursiveTask<Long> {
    private static final int THRESHOLD = 10_000;
    private final long[] data; private final int lo, hi;

    SumTask(long[] data, int lo, int hi) { this.data = data; this.lo = lo; this.hi = hi; }

    @Override protected Long compute() {
        if (hi - lo <= THRESHOLD) {            // base case: small enough, do it directly
            long sum = 0;
            for (int i = lo; i < hi; i++) sum += data[i];
            return sum;
        }
        int mid = (lo + hi) >>> 1;
        SumTask left = new SumTask(data, lo, mid);
        SumTask right = new SumTask(data, mid, hi);
        left.fork();                            // push left onto THIS worker's deque (stealable)
        long rightResult = right.compute();     // compute right inline (no extra task overhead)
        long leftResult  = left.join();         // wait for left (maybe stolen by another worker)
        return leftResult + rightResult;
    }
}

// Usage: ForkJoinPool.commonPool().invoke(new SumTask(data, 0, data.length));

go deeper

for a junior

Can state that each worker has its own queue and idle workers steal tasks from others so no thread sits idle while work remains.

for a middle

Knows the queue is a double-ended deque, the owner uses one end and thieves the other, and stealing balances load without a central queue.

for a senior

Explains LIFO-own-end vs FIFO-steal-end and why (cache locality, depth-first bounding, stealing the biggest chunk), and that opposite-end access plus CAS keeps contention low.

for a principal

Reasons about the push-vs-pull scheduling inversion, scalability vs a single-queue pool, randomised victim selection, and the failure mode when tasks block; relates it to sizing and to alternatives for blocking workloads.

## The data structure: a per-worker deque A **deque** ("deck", double-ended queue) is a queue you can add to and remove from at *both* ends. In ForkJoinPool, **every worker thread owns one deque** holding the tasks it has forked. There is no single shared queue that all threads fight over — that decentralisation is the whole point. Two ends matter: - **Head / top** — the owner's end. The owner **pushes** newly forked subtasks here and **pops** from here (LIFO, last-in-first-out). - **Tail / base** — the thieves' end. Idle workers **steal** from here. ## How a worker runs its own work (LIFO, head) When `compute()` calls `fork()`, the subtask is pushed onto the **head** of the current worker's deque. When the worker needs more work, it **pops from the head** — the most recently forked task. Why LIFO for yourself? - **Cache locality:** the most recent subtask shares data with what you just touched, so it's likely still in CPU cache. - **Depth-first traversal:** popping the newest task drives the recursion down to base cases quickly, which keeps the deque small and memory bounded rather than fanning out breadth-first. ## How stealing works (FIFO, tail) When a worker's own deque is **empty**, it doesn't sleep — it tries to **steal**. It picks another worker (the **victim**, chosen with some randomisation to spread contention) and takes a task from the victim's **tail** — the *oldest* task the victim pushed. Why steal the oldest/biggest task? The oldest task sits highest in the recursion tree, so it represents the **largest remaining chunk** of work — not yet split into tiny pieces. Stealing one big task gives the thief plenty to do (and more subtasks to later be stolen from *it*), minimising how often threads have to steal again. ## Why opposite ends = low contention The owner works the **head**; thieves work the **tail**. Because they operate on opposite ends, they almost never touch the same slot, so the owner's push/pop can be nearly lock-free and only the rarer steal (and the edge case where head meets tail on a near-empty deque) needs an atomic **compare-and-swap (CAS)** to resolve who gets the last task. CAS is a hardware instruction that atomically updates a value only if it still holds an expected value — it lets threads coordinate without a lock. ## Push vs pull: the key inversion A classic thread pool **pushes** work from a central queue out to workers — that central queue becomes a contention hotspot as core counts rise. Work stealing **inverts** this: work stays local to the worker that created it, and **idle threads pull work to themselves**. Balancing is demand-driven and happens only when a thread actually goes idle, so there's no coordinator and no global bottleneck. This is what lets Fork/Join scale across many cores on irregular, recursively-generated workloads where you can't predict task counts up front. ## Assumptions and failure modes Work stealing balances load **only if there are stealable tasks and workers don't block**: - If subtasks **block** (I/O, locks, sleeps), the blocked worker holds a core but generates no stealable work, and parallelism (sized to cores by default) collapses. - If you **don't split finely enough**, there's nothing to steal and some cores idle; split **too** finely and per-task overhead dominates. Choosing a sensible threshold (or letting parallel streams choose) matters. ## Mental model Picture each worker with a personal stack of work. You always grab from the top of *your own* stack (fresh, cache-hot). When your stack is empty, you reach over and grab from the *bottom* of a neighbour's stack — the big, old job they haven't gotten to. Everyone stays busy; nobody fights over one shared pile.

  • Why does the owner use LIFO on its own end but thieves take FIFO from the other end?
    LIFO keeps the owner on freshly forked, cache-hot subtasks and drives depth-first recursion (bounding deque size). The thief takes the oldest task from the far end because it's the largest, least-split chunk — one steal yields lots of work and minimises future steals, and using the opposite end avoids contending with the owner.
  • What happens to throughput if your compute() tasks frequently block on I/O?
    Blocked workers occupy cores without producing stealable tasks, so other workers find nothing to steal and cores idle. Since parallelism defaults to the core count, a handful of blocked tasks can starve the pool. You'd use a different executor, ManagedBlocker, or virtual threads for blocking work.

saying these in an interview costs you the question

  • Saying there is one shared task queue — each worker has its own deque
  • Claiming workers steal the newest task — they steal the oldest/biggest from the tail
  • Saying a worker processes its own deque in FIFO — it's LIFO from the head
  • Thinking stealing requires a global lock — opposite-end access + CAS keeps it mostly contention-free
  • Assuming work stealing fixes load imbalance even when tasks block — it can't reclaim a blocked core

context