Command Palette

Search for a command to run...

PHASE 10Intermediate Java 8+ ~34 min· topic 8 of 8

Topic 10.8

Parallel Streams and Their Pitfalls

In one line

Calling parallelStream() or .parallel() splits a stream's work across the threads of the shared fork/join pool. It gives the same answer as a sequential stream only if your operations are stateless, non-interfering and associative, and it's faster only for large, CPU-heavy work on sources that split well.

Think of it like this

Counting votes after an election. Instead of one person counting a million ballots, you split the boxes among many counters, each counts their own pile, and at the end you add up the totals. That works because each pile is independent and adding is order-free. It goes wrong if all the counters scribble on one shared tally sheet at the same time, and it's pointless for ten ballots: organising the team takes longer than counting alone. Parallel streams are exactly this, with CPU cores as the counters.

Words you'll meet

New words in this topic, in plain English. Come back here whenever one feels fuzzy.

Parallel stream
A stream whose work is split into pieces processed at the same time on several threads, then combined.
Thread
An independent path of execution inside a program. Several threads can run at once on different CPU cores (Phase 13).
ForkJoinPool
A pool of worker threads built for splitting a task into smaller tasks (fork) and combining their results (join). Parallel streams use its shared common pool.
Spliterator
The source-side object that hands out elements and can split itself into two halves for parallel processing. Short for "splittable iterator".
Shared mutable state
Data that several threads can change at the same time, such as one ArrayList that every thread adds to.
Race condition
A bug where the result depends on the exact timing of threads, so it changes from run to run (Topic 13.3).
Overhead
Extra work that isn't the real job, like splitting data and coordinating threads.
Encounter order
The order in which the source presents elements, such as list index order.

Step by step

01Turning it on

Two ways: start parallel with collection.parallelStream(), or switch an existing pipeline with .parallel(). .sequential() switches it back. The setting applies to the whole pipeline, and the last call before the terminal operation wins, so stream.parallel().map(f).sequential().forEach(g) runs entirely sequentially.

isParallel() tells you which mode a stream is in. Everything else (the operations, the lambdas, the result type) stays exactly the same, which is both the appeal and the danger: nothing in the code warns you that the lambdas now run on several threads.

Main.javawhole filejava
long a = list.parallelStream().filter(this::isValid).count();      // parallel from the start
long b = list.stream().parallel().filter(this::isValid).count();   // same thing
long c = list.parallelStream()
        .map(this::expensive)
        .sequential()                                               // last call wins:
        .count();                                                   // the whole pipeline is sequential

02What happens under the hood: split, compute, combine

The terminal operation wraps the pipeline into a fork/join task. The task asks the source's Spliterator to trySplit() itself; each half becomes a new task, which splits again, until the pieces are small (roughly the total size divided by four times the pool's parallelism).

Each leaf task runs the pipeline sequentially over its piece, producing a partial result (a partial sum, a partial list). Then results are combined up the tree: sums added, lists concatenated in order, maps merged. Idle threads steal waiting tasks from busy threads' queues, which keeps all cores working.

What happens under the hood: split, compute, combinediagram
Rendering diagram…

03Sources that split well, and sources that don't

An ArrayList spliterator splits by index: half the range each time, instantly, and both halves know their exact size. That's ideal. A LinkedList can't jump to the middle; its spliterator copies a batch of elements into an array to split off, so for a small list the "split" takes everything and leaves nothing for other threads.

Infinite or sequential-by-nature sources (Stream.iterate, Stream.generate, lines read from a file or socket) report no size and split in batches, if at all. IntStream.range, arrays, ArrayList, HashMap key sets and ConcurrentHashMap split well.

04Rule 1: never mutate shared state from the lambdas

forEach(list::add) into an ArrayList works sequentially, but in parallel several threads call add at once. ArrayList isn't thread-safe: two threads write into the same slot, or one resizes the array while another writes. The result is a list with missing elements, or an ArrayIndexOutOfBoundsException, and it's different on every run. A shared int[] counter loses increments the same way.

The fix isn't a synchronized list (that makes the threads queue up, killing the speed-up). The fix is to let the stream build the result: collect(...), toList(), count(), sum(), reduce(...). Each thread fills its own container and the library merges them safely.

terminal
$ java Main.java # output varies on every run
── expected output ──
run 1: size 23339
run 2: size 16864
run 3: size 14245
counter: 38309

05Rule 2: associative operations and a true identity

In parallel, reduce(identity, op) starts every chunk from the identity. So the identity must change nothing: reduce(0, Integer::sum) is fine, reduce(10, Integer::sum) adds 10 once per chunk, and the error grows with the number of chunks, which depends on the machine.

The operation must be associative, because chunks are combined in a tree, not left to right. Addition, multiplication, max, min and string concatenation are associative; subtraction, division and averaging pairs ((a + b) / 2) are not. If you need an average, use average() or averagingInt, which carry a sum and a count.

06Order: what's kept and what it costs

A parallel toList() or collect(toList()) still returns elements in encounter order: each chunk's partial list is concatenated in position order. findFirst() returns the true first match, and forEachOrdered runs the action in order (one at a time, losing most of the parallelism for that step).

forEach runs actions on whatever thread processes the element, in any order. findAny returns whichever match is found first. Ordered limit(n), skip(n), distinct() and sorted() must coordinate across chunks; on a parallel stream they can be slower than sequential. .unordered() removes the obligation when you don't need it.

07When parallel is actually faster

Ask four questions. Is there enough work (many elements, or expensive work per element)? Does the source split well? Is the work CPU-bound, not waiting on I/O? Is combining cheap? Summing a million numbers from an array: maybe. Squaring ten numbers: never. Calling a web service for each element: no, use an executor or virtual threads (Topic 13.10) instead.

Boxing hurts parallel streams twice: Stream<Integer> chases pointers to scattered objects, which wastes the memory bandwidth the cores share. IntStream.range(...).parallel() over primitives scales far better. And measure with JMH, not System.currentTimeMillis around one run: JIT warm-up dominates single runs.

08The common pool is shared

ForkJoinPool.commonPool() has availableProcessors() - 1 worker threads by default (you can change it with -Djava.util.concurrent.ForkJoinPool.common.parallelism=N). It's shared by every parallel stream in the JVM. In a web server, ten requests each running a parallel stream don't get ten times the cores; they compete for the same few threads.

A well-known workaround runs the stream inside your own pool: myPool.submit(() -> list.parallelStream()...).get(). The tasks then fork into myPool. It works in current JDKs but is an implementation detail, not a documented guarantee. For blocking work, use a proper executor (Topic 13.6) instead.

Try it yourself

  1. 1

    Remove forEachOrdered

    In the first example, change forEachOrdered to forEach and run it a few times. The order changes between runs (and on a single-core machine it may not). Then explain why the StringBuilder is now unsafe, even if the output looks fine.

  2. 2

    Break and fix the identity

    In the identity example, change the chunks to three uneven pieces (subList(0, 3), subList(3, 5), subList(5, 8)). Predict good and bad before running. Does the correct identity still give 36?

  3. 3

    Measure honestly

    Time LongStream.rangeClosed(1, 1_000_000).sum() against its parallel version with System.nanoTime(), running each 20 times in a loop and printing the last result's time. Then repeat with 1,000 elements. On your machine, at what size does parallel start winning? (Single runs are misleading; that's why serious measurements use JMH.)

Code & diagrams

Correct pipelines give the same answer in parallel Java 16+ New tab

forEachOrdered guarantees each action happens-before the next, so appending to one StringBuilder is safe there. With plain forEach it would not be.

Sign in to run this example in your browser.

Expected output

sequential sum: 500000500000
parallel sum:   500000500000
toList keeps encounter order: [1, 4, 9, 16, 25, 36, 49, 64, 81, 100]
forEachOrdered: 1 2 3 4 5 6 7 8 9 10
joining: 1,2,3,4,5,6,7,8,9,10
findFirst: 37
isParallel: true
Accumulate results the safe way Java 16+ New tab

Compare with forEach(evens::add) into a plain ArrayList, which loses elements or throws, differently on every run (see the walkthrough).

Sign in to run this example in your browser.

Expected output

collect size: 50000, first 0, last 99998
count: 33334
LongAdder: 100000
groupingByConcurrent: {0=10000, 1=10000, 2=10000, 3=10000, 4=10000, 5=10000, 6=10000, 7=10000, 8=10000, 9=10000}
Why the identity must be neutral Java 16+ New tab

The exact wrong value depends on how many chunks the machine splits into, which is why this bug passes tests on a laptop and fails on a 64-core server.

Sign in to run this example in your browser.

Expected output

4 chunks, identity 0:  36
4 chunks, identity 10: 76 (10 added once per chunk)
sequential, identity 10: 46
identity 0, parallel equals sequential: true
identity 10, parallel equals 46: false
fix, add the start value once: 46
How sources split: ArrayList versus LinkedList Java 16+ New tab

LinkedList splits by copying a batch (starting at 1024 elements) into an array, so a small list goes to one thread entirely: no parallelism at all.

Sign in to run this example in your browser.

Expected output

ArrayList, first split: 8 + 8
second split: 4 + 4 + 8
first quarter holds: 1 2 3 4
ArrayList knows exact sizes: true
LinkedList split: 16 + 0
iterate + limit knows its size: false

Break it on purpose

Errors are the best teachers. Make each change, read the error, guess what went wrong, then reveal the answer.

Break #1

Add to a plain ArrayList from a parallel forEach

Write List<Integer> evens = new ArrayList<>(); IntStream.range(0, 100_000).parallel().filter(n -> n % 2 == 0).forEach(evens::add); and print the size three times.

terminal
$ java Main.java # output varies on every run
── what you'll see ──
run 1: size 23339
run 2: size 16864
run 3: size 14245

Break #2

Use a non-neutral identity in a parallel reduce

Write nums.parallelStream().reduce(10, Integer::sum) hoping to get "the sum plus 10".

terminal
$ java Main.java
── what you'll see ──
identity 10, parallel equals 46: false

Break #3

Block inside a parallel stream

In a web service, call a remote API from urls.parallelStream().map(this::fetch).toList(), where fetch waits about a second per call.

terminal
$ jstack <pid> | grep -A3 ForkJoinPool.commonPool
── what you'll see ──
"ForkJoinPool.commonPool-worker-1" #31 daemon prio=5 os_prio=0 ... waiting on condition
java.lang.Thread.State: TIMED_WAITING (parking)
at jdk.internal.misc.Unsafe.park(java.base@21/Native Method)
"ForkJoinPool.commonPool-worker-2" #32 daemon prio=5 os_prio=0 ... waiting on condition
java.lang.Thread.State: TIMED_WAITING (parking)

Myth vs fact

Myth

Adding .parallel() makes any stream faster.

Fact

Only large, CPU-bound, well-splitting workloads speed up. Small inputs, boxed elements, LinkedList or I/O sources, ordered limit/sorted, and expensive merges often make parallel slower than sequential.

Myth

Parallel streams use as many threads as there are elements.

Fact

They use the common ForkJoinPool (cores minus one workers, plus the calling thread). The data is split into a limited number of chunks, roughly four per worker.

Myth

A parallel toList() returns elements in random order.

Fact

Collecting preserves encounter order for ordered sources. It's forEach and findAny that ignore order.

Myth

Wrapping the shared list in Collections.synchronizedList fixes parallel forEach.

Fact

It fixes the corruption but serialises every add behind one lock, so threads mostly wait, and element order is still random. collect is both correct and fast.

When it breaks

Parallel streams inside request handlers starve the common pool

What you see

A REST service uses items.parallelStream().map(this::callPricingService) on every request. Under load, p99 latency climbs from 200 ms to many seconds; thread dumps show every ForkJoinPool.commonPool-worker parked in socket reads, and unrelated features that use CompletableFuture.supplyAsync without an executor slow down too.

Fix & prevent

Move I/O fan-out to a dedicated, bounded executor or virtual threads, with timeouts. Keep parallel streams for CPU-bound batch work. Alert on common-pool queue size (ForkJoinPool.commonPool().getQueuedSubmissionCount()).

A report total is slightly different on the production server

What you see

A nightly report uses amounts.parallelStream().reduce(openingBalance, Double::sum). On a 4-core laptop the total is off by 4 opening balances (nobody noticed); on a 48-core server it's off by many more, and floating-point rounding also varies with chunking. Finance flags mismatches.

Fix & prevent

Use a neutral identity and add constants once: openingBalance + amounts.parallelStream().reduce(0.0, Double::sum). For money, use long cents or BigDecimal with an associative sum, and test with .parallel() on inputs larger than a few thousand elements.

Logging context and security context missing in parallel lambdas

What you see

Logs written from inside a parallel stream lack the request ID (MDC), and code that reads a ThreadLocal-based security context throws or treats the call as anonymous, but only for elements processed on worker threads, so the bug looks random.

Fix & prevent

Don't rely on ThreadLocal state inside parallel lambdas. Capture the values you need before the pipeline and pass them in explicitly, or keep that work sequential. Scoped values (Topic 13.11) are the modern, structured alternative for passing context to child tasks.

Pro corner

Extra depth for experienced readers. New to this? Skip it for now and come back later.

  • ▸

    Execution: terminal ops become AbstractTask subclasses (ReduceTask, ForEachTask, FindTask...), which are CountedCompleters. Each splits while the estimated size exceeds sizeEstimate / (commonPoolParallelism * 4) (AbstractTask.LEAF_TARGET), forks one half and continues with the other. CountedCompleter lets parents complete without blocking threads in join().

  • ▸

    Short-circuiting in parallel (findFirst, anyMatch, limit) uses a shared cancellation flag that tasks poll. findFirst must still wait for all tasks to the left of a found element, which is why findAny is faster when you don't care which element you get.

  • ▸

    The ForkJoinPool.submit(() -> stream.parallel()...) trick works because fork/join tasks fork into the pool of the thread that runs them (ForkJoinTask.fork uses the current worker's pool). It's undocumented behaviour for streams and has changed subtly between releases, so treat it as a workaround, not an API.

  • ▸

    Memory bandwidth, not core count, is often the real limit: a parallel sum over a large int[] saturates memory bandwidth with a few cores, so 32 cores rarely give 32 times the speed. Boxed Stream<Integer> makes it worse (each element is a pointer to a separate heap object), and false sharing between threads updating nearby memory can erase the gains entirely (LongAdder exists partly to avoid it).

Remember this

  1. 1

    list.parallelStream() or stream.parallel() marks the whole pipeline as parallel (the last parallel()/sequential() call wins; you can't make only one stage parallel). When the terminal operation runs, the source is split into chunks, the chunks are processed as tasks on the **common ForkJoinPool**, and the partial results are combined. The calling thread works too.

  2. 2

    Splitting is done by the source's **Spliterator** (trySplit). Arrays, ArrayList and IntStream.range split perfectly in half with known sizes. LinkedList, Stream.iterate, BufferedReader.lines() and most I/O sources split badly, because finding the middle means walking the elements. A badly splitting source gives you overhead without parallelism.

  3. 3

    For the result to be correct, the rules from earlier topics become strict: lambdas must be stateless and non-interfering, must not modify shared mutable state (an ArrayList, an int[] counter, a HashMap), and reductions need an associative function with a true identity (Topic 10.5). Break them and you get wrong answers that change from run to run, not exceptions you can catch reliably.

  4. 4

    Order: toList(), collect, forEachOrdered, findFirst, limit and skip respect encounter order even in parallel, at a cost (order-preserving limit, distinct and sorted need extra buffering). forEach and findAny don't, so they can be faster. If order doesn't matter, .unordered() lets the library drop that cost.

  5. 5

    Parallel is faster only when the work is big enough to pay for splitting, scheduling and merging. A rough guide (from Doug Lea and Brian Goetz): number of elements × cost per element should be well above 10,000 simple operations, the source must split well, the work must be CPU-bound, and combining results must be cheap (adding numbers yes, merging big maps less so). Measure with a proper benchmark (JMH) before and after.

  6. 6

    The common pool is shared by the whole JVM: every parallel stream (and CompletableFuture without an executor, Topic 13.7) uses it. Its size is the number of CPU cores minus one. A parallel stream that blocks (HTTP calls, JDBC, Thread.sleep) holds those threads hostage and slows every other parallel task in the application. Parallel streams are for computation, not I/O.

Explain it without notes

01

How does a parallel stream split and execute work?

02

What conditions must a pipeline meet to give correct results in parallel?

03

When does a parallel stream actually make things faster, and when slower?

04

Why is blocking I/O inside a parallel stream a problem?

05

How do ordering guarantees differ between forEach, forEachOrdered, findFirst, findAny and toList in parallel?

Practice

01

Compute the number of primes below 200,000 with a parallel IntStream and a simple isPrime method, and check it equals the sequential count.

02

Fix this racy code without using any locks: int[] total = {0}; IntStream.rangeClosed(1, 1000).parallel().forEach(n -> total[0] += n);. Print the correct total.

03

Count words by length in parallel with groupingByConcurrent and print the result as a sorted map, for "a bb cc ddd e ff ggg hhhh".

Trade-offs

  • ↔

    Parallel streams give multi-core speed-ups with almost no code change, but they hide threads behind familiar syntax, so thread-safety bugs slip in easily and only show up under load or on bigger machines.

  • ↔

    Using the shared common pool means no thread management, but also no isolation: one slow or blocking parallel stream affects the whole JVM. Explicit executors cost more code and give control over size, naming and isolation.

  • ↔

    Preserving encounter order keeps results predictable but costs buffering and coordination for limit, skip, distinct and sorted. unordered() and findAny trade predictability for speed when order genuinely doesn't matter.

Done when you can

  • Done when you can explain split, compute and combine, and the role of the Spliterator and common ForkJoinPool.

  • Done when you can spot shared mutable state in a parallel pipeline and replace it with a collector or reduction.

  • Done when you can explain why the identity must be neutral and the operation associative.

  • Done when you can say which terminal operations keep encounter order in parallel.

  • Done when you can judge whether a workload is worth parallelising and know to measure with JMH.

  • Done when you never put blocking I/O inside a parallel stream.