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
ArrayListthat 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.
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 sequential02What 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.
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.
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
Remove forEachOrdered
In the first example, change
forEachOrderedtoforEachand run it a few times. The order changes between runs (and on a single-core machine it may not). Then explain why theStringBuilderis now unsafe, even if the output looks fine. - 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)). Predictgoodandbadbefore running. Does the correct identity still give 36? - 3
Measure honestly
Time
LongStream.rangeClosed(1, 1_000_000).sum()against its parallel version withSystem.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
forEachOrdered guarantees each action happens-before the next, so appending to one StringBuilder is safe there. With plain forEach it would not be.
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: trueCompare with forEach(evens::add) into a plain ArrayList, which loses elements or throws, differently on every run (see the walkthrough).
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}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.
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: 46LinkedList 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.
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: falseBreak 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.
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".
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.
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
AbstractTasksubclasses (ReduceTask,ForEachTask,FindTask...), which areCountedCompleters. Each splits while the estimated size exceedssizeEstimate / (commonPoolParallelism * 4)(AbstractTask.LEAF_TARGET), forks one half and continues with the other.CountedCompleterlets parents complete without blocking threads injoin(). - ▸
Short-circuiting in parallel (
findFirst,anyMatch,limit) uses a shared cancellation flag that tasks poll.findFirstmust still wait for all tasks to the left of a found element, which is whyfindAnyis 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.forkuses 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. BoxedStream<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 (LongAdderexists partly to avoid it).
Remember this
- 1
list.parallelStream()orstream.parallel()marks the whole pipeline as parallel (the lastparallel()/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 **commonForkJoinPool**, and the partial results are combined. The calling thread works too. - 2
Splitting is done by the source's **
Spliterator** (trySplit). Arrays,ArrayListandIntStream.rangesplit 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
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, anint[]counter, aHashMap), 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
Order:
toList(),collect,forEachOrdered,findFirst,limitandskiprespect encounter order even in parallel, at a cost (order-preservinglimit,distinctandsortedneed extra buffering).forEachandfindAnydon't, so they can be faster. If order doesn't matter,.unordered()lets the library drop that cost. - 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
The common pool is shared by the whole JVM: every parallel stream (and
CompletableFuturewithout 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
How does a parallel stream split and execute work?
What conditions must a pipeline meet to give correct results in parallel?
When does a parallel stream actually make things faster, and when slower?
Why is blocking I/O inside a parallel stream a problem?
How do ordering guarantees differ between forEach, forEachOrdered, findFirst, findAny and toList in parallel?
Practice
Compute the number of primes below 200,000 with a parallel IntStream and a simple isPrime method, and check it equals the sequential count.
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.
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,distinctandsorted.unordered()andfindAnytrade 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
Spliteratorand commonForkJoinPool.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.