ByteScrollGet the app
☰ Topics
forkjoin-work-stealing11 / 200‹›
JAVA / CONCURRENCY3 minute read

ForkJoinPool and work-stealing

Hard

ForkJoinPool runs divide-and-conquer tasks. Each worker has its own deque: it pushes and pops its subtasks at one end, and idle workers steal the oldest, biggest tasks from the other end. Little contention, busy cores. Parallel streams use its common pool.

How it works

  1. Split work into tasks. Extend RecursiveTask<V> (returns a value) or RecursiveAction (no value) and implement compute(): if the piece is small, do it directly; otherwise split it.
  2. fork() pushes a subtask onto the current worker's own deque. No shared queue, no lock.
  3. LIFO for yourself. A worker pops its newest task first. That's the smallest piece, and its data is likely still in cache.
  4. FIFO for thieves. An idle worker picks another worker's deque and takes from the opposite end: the oldest task, usually a big chunk that will keep it busy for a while.
  5. join() doesn't just block. While waiting for a result, the worker runs other pending tasks (its own or stolen), so threads rarely sit idle.
  6. The common pool. ForkJoinPool.commonPool() has availableProcessors() - 1 workers by default. Parallel streams and the *Async methods of CompletableFuture (when no executor is given) use it.
Worker 2 dequeWorker 1 dequepop newest (LIFO)steal oldest (FIFO)sum 0..1M(oldest)sum 0..500ksum 0..250k(newest)(empty)Worker 1 runsitWorker 2 runsit
Worker 2 dequeWorker 1 dequepop newest (LIFO)steal oldest (FIFO)sum 0..1M(oldest)sum 0..500ksum 0..250k(newest)(empty)Worker 1 runsitWorker 2 runsit

Example

Example.javaJava
class SumTask extends RecursiveTask<Long> {
    private static final int CUTOFF = 10_000;
    private final long[] data;
    private final int from, to;

    SumTask(long[] data, int from, int to) { this.data = data; this.from = from; this.to = to; }

    @Override protected Long compute() {
        if (to - from <= CUTOFF) {
            long s = 0;
            for (int i = from; i < to; i++) s += data[i];
            return s;
        }
        int mid = (from + to) >>> 1;
        SumTask left = new SumTask(data, from, mid);
        left.fork();                                   // let someone steal it
        long right = new SumTask(data, mid, to).compute(); // do the other half here
        return right + left.join();
    }
}

long[] readings = LongStream.rangeClosed(1, 50_000_000).toArray();
long total = ForkJoinPool.commonPool().invoke(new SumTask(readings, 0, readings.length));

Edge cases

  • An unchecked exception thrown in compute() is rethrown by join()/invoke() in the waiting thread.
  • If the common pool's parallelism is below 2, CompletableFuture.supplyAsync without an executor starts a new thread per task instead of using it.
  • Blocking calls inside tasks starve the pool. If you must block, wrap it in ForkJoinPool.ManagedBlocker so the pool can add a spare thread.
  • Common pool threads are daemon threads, so they won't keep the JVM alive.
  • The default scheduler for virtual threads is also a ForkJoinPool (a separate one, run in FIFO mode).

Common mistakes

  • left.fork(); right.fork(); left.join(); right.join(); works, but wastes the current thread. Fork one half, compute the other directly.
  • A tiny cutoff: millions of microscopic tasks cost more to schedule than to run.
  • Doing I/O or Thread.sleep in the common pool and slowing every parallel stream in the JVM.

Likely follow-up

"Why is work-stealing better than one shared queue?" With a shared queue every worker contends on the same lock or CAS for every task. With per-worker deques, a worker only touches another's deque when it has run out of work, and it steals large tasks, so steals are rare.

Get every deep dive in the app

Coming soon to the App StoreComing soon to Google Play