The idea: divide and conquer on a few threads, each with its own deque

A ForkJoinPool runs tasks that split themselves into smaller tasks. A big task splits in two, each half splits again, and so on until a piece is small enough to compute directly. Then the results are combined on the way back up. The number of threads stays fixed (the parallelism, by default the number of CPU cores). The task tree can have thousands of nodes.

Task tree for SumTask on int[32] with THRESHOLD 4: [0,32) splits into [0,16) and [16,32), those into ranges of 8, and those into 8 leaves of 4 elements, [0,4) to [28,32), which are summed directly; each parent returns l + r. 15 tasks in total.
A task splits in halves until a range is at most THRESHOLD long, then the sums are added back up the tree.

Each worker thread has its own deque (double-ended queue) of tasks. The page sums an int[32] with a SumTask extends RecursiveTask<Long>. On the canvas you see the array and the task tree (top left), the code with a marker for each worker's current line (top right), and each worker's deque and call stack (middle). The timeline at the bottom has one cell per tick for each worker.

Reading the canvas: an array cell and a tree node take the colour of the worker that worked on them. Worker W1 is the thread ForkJoinPool-1-worker-1, and so on; W5 is a spare that only starts for a ManagedBlocker. A numbered box left of the code shows which line worker n is on in this tick. In each worker's column the deque and the call stack grow upwards: the owner pushes and pops at the top (LIFO), idle workers steal from the base at the bottom (FIFO).

The model counts in ticks: adding one element is 1 tick, a split is 1 tick, return l + r is 1 tick and stealing a task from another thread is 1 tick. Popping your own task and calling a method are free.

fork(), compute() and join()

class SumTask extends RecursiveTask<Long> {
    protected Long compute() {
        if (hi - lo <= THRESHOLD) return sumDirectly();
        int mid = (lo + hi) >>> 1;
        SumTask left = new SumTask(a, lo, mid), right = new SumTask(a, mid, hi);
        left.fork();                   // push on my own deque
        long r = right.compute();      // keep working in this thread
        long l = left.join();          // pop it back, help, or wait
        return l + r;
    }
}
long sum = pool.invoke(new SumTask(a, 0, a.length));
  • fork() does not start a thread. It pushes the task onto the top of the current worker's deque and returns at once.
  • right.compute() is a plain method call. The current thread keeps working on the other half, so it never sits idle while there is work.
  • join() returns the result, and gets it in one of the ways described below.
  • RecursiveAction is the same without a result. invokeAll(left, right) does the fork/compute/join for you.

(Demo: invoke() — fork, compute, join on 4 workers: 21 ticks instead of 47 on one worker. Demo: one worker: with nobody to steal, every fork() is popped right back. The deque then behaves like a second call stack and the tree is walked depth-first.)

Work stealing: LIFO for the owner, FIFO for thieves

An owner pushes and pops at the top of its deque (last in, first out). A worker with nothing to do steals from the base of another worker's deque (first in, first out). There are three reasons for this split:

  • The owner and the thieves work at different ends, so they rarely touch the same slot. In the real implementation the owner's push and pop need no lock, only a CAS when the deque is nearly empty.
  • The base holds the oldest task, which was forked highest in the tree, so it is the biggest one. One steal gives the thief a lot of work, and it will split that work further on its own deque. Few steals are needed.
  • The top holds the newest, smallest task, whose data is probably still in the owner's CPU cache.
W1's deque holds, from base to top, [0,16), [16,24) and [24,28). W1, the owner, pushes and pops at the top (LIFO). Idle W2 steals from the base (FIFO) and gets [0,16), the oldest and biggest task: half of all the work in one steal.
The owner works at the top of its deque; a thief takes from the base, where the oldest and biggest task waits.

(Demo: two workers: W2 makes a single steal, takes [0,16), half of everything, and both workers stay busy until the end.) The root task comes from a thread outside the pool, so it waits in an external submission queue until a worker scanning for work picks it up.

What join() really does

situation when left.join() runswhat the worker doesname in the JDK source
the task is donetake the resultstatus < 0
still on top of my own deque (nobody stole it)pop it back and run it in this threadtryUnpush / tryRemoveAndExec
a thief is running ithelp: steal a task from the thief's deque (the thief's own subtasks), and if the thief is joining too, follow the chainhelpStealer (JDK 8), helpJoin (JDK 19+)
nothing to help withwait; the pool may start a spare thread so the parallelism does not dropawaitJoin, tryCompensate

So a worker blocked in join() keeps doing useful work when it can. Look for brown "help" cells in the timeline (Demo: invoke() at t = 15, where W2 helps W4; Demo: THRESHOLD too small, where W1 helps twice, the second time by following a chain of thieves to W3).

Choosing THRESHOLD

Each task costs something to create, push, steal and join. If the leaves are too small, that overhead outweighs the work. If they are too big, there are not enough tasks to keep every worker busy or to even out the load. The ForkJoinTask Javadoc suggests a leaf should do between about 100 and 10 000 basic computational steps.

THRESHOLD (4 workers)taskstickswhat limits it
2312415 splits + 15 returns of overhead, more steals and join waits
41521a good balance for this array
8715exactly 4 leaves for 4 workers: good here, but leaves no slack if the leaves were uneven
16320only 2 leaves, so 2 workers stay idle

Common mistakes

codewhat happens
left.fork(); l = left.join(); r = right.compute();The task is joined before this thread has done anything else. It is still on top of the deque, so the same thread pops it back and runs it. Everything runs on one thread (Demo: mistake: 47 ticks, 0 steals, 3 idle workers).
left.fork(); right.fork(); left.join(); right.join();Works, but the current thread forks work it could have done itself, and it joins in the wrong order (right is on top). Prefer fork / compute / join or invokeAll. Joins should go innermost first, in the reverse order of the forks.
left.invoke() or compute() both halvesNo fork at all, so there is nothing to steal: plain sequential recursion.
blocking I/O, Thread.sleep, synchronized waits inside a taskThe worker sleeps and the pool does not know, so it runs with fewer threads. See the next section.
tasks with side effects on shared dataTasks run in parallel; give each one its own slice and combine the results, as SumTask does.

Blocking inside a pool, and ManagedBlocker

A ForkJoinPool assumes its tasks compute. If a task makes a blocking call (an HTTP request, a JDBC query, a lock), its thread sleeps in the kernel. The pool now has one worker fewer, and the tasks waiting in that worker's deque wait too, until another worker is free to steal them (Demo: a blocking call: 34 ticks with 2 workers).

Wrapping the call in a ForkJoinPool.ManagedBlocker and running it through ForkJoinPool.managedBlock(blocker) tells the pool that the thread is about to block. The pool starts or wakes a spare (compensation) thread, which keeps the number of running workers at the parallelism. The spare parks again afterwards (Demo: ManagedBlocker: spare W5 steals the work from the blocked worker's deque, and the run takes 26 ticks).

Spare threads are limited (maximumPoolSize, 256 more than the parallelism by default). For lots of blocking I/O, use a normal ThreadPoolExecutor or virtual threads instead.

The common pool

ForkJoinPool.commonPool() is a shared, static pool. Parallel streams, CompletableFuture.supplyAsync(...) without an executor, and Arrays.parallelSort all run in it. Its parallelism is availableProcessors() − 1, because the thread that calls invoke() or starts a terminal stream operation also works on the tasks. You can change it with -Djava.util.concurrent.ForkJoinPool.common.parallelism=N.

Because every part of the program shares the common pool, one library that blocks inside it slows all the parallel streams in the JVM.

ForkJoinPool vs ThreadPoolExecutor

ForkJoinPoolThreadPoolExecutor
queuesone deque per worker + submission queuesone shared BlockingQueue
idle workersteals from another worker's dequetakes the next task from the shared queue
a task waiting on a subtaskjoin() runs or helps instead of blockingfuture.get() blocks the thread (can deadlock a full pool)
good forrecursive, CPU-bound divide and conquer; many small tasksindependent tasks, I/O-bound request handling

See also Concurrency vs Parallelism (Amdahl's law; why only more cores speed up CPU-bound work), Bounded Buffer (the shared queue a classic thread pool is built on), Context Switch, The Node.js Event Loop (a different answer: one thread and non-blocking I/O), How Tomcat handles an HTTP request (a ThreadPoolExecutor at work), Latency vs Throughput and Stack vs Heap in Java.

What the page leaves out

  • Real thieves pick a random victim and then scan; here they scan in a fixed order so that the demos repeat.
  • Idle workers park and are woken by signalWork when tasks are pushed; the page has no wake-up delay.
  • The deque is a circular array with base and top indexes, and the memory-ordering rules that make push and pop lock-free.
  • tryCompensate during a plain join() wait, CountedCompleter, asyncMode (FIFO local queues, used for event-style tasks), and the ForkJoinPool that carries virtual threads.
  • Real costs: here a steal costs only one tick, but in reality a steal involves cache misses and CAS contention.