Most Java developers can write a stream pipeline. Fewer can predict what it will do: in which order side effects run, why sorted() makes findFirst() read the entire input, why a peek sometimes never prints, or why a parallel pipeline with limit is slower than the sequential one. Those answers all come from the semantics of individual operations, and from how a pipeline is executed.
This article is about operations, not the API tour. It classifies operations by the properties that matter, traces execution element by element, looks at how OpenJDK implements a pipeline as a chain of sinks, shows the optimisations that flags permit, and then teaches you to write your own intermediate operations with Gatherer, final since JDK 24 (JEP 485). For sources, collectors, primitive streams and the wider API, read the Java Streams API guide first; for the lambdas operations take, see lambda expressions.
Three properties that classify every operation
Every stream operation is either intermediate, returning a new stream and doing no work yet, or terminal, which triggers execution and consumes the stream. Within those, three properties predict behaviour.
| Property | Meaning | Examples |
|---|---|---|
| Stateless | Processing one element needs no information about others | filter, map, flatMap, mapMulti, peek |
| Stateful | The operation must remember elements already seen, possibly all of them | distinct, sorted, limit, skip, takeWhile, dropWhile |
| Short-circuiting | Can finish without consuming all input, so it can work on infinite streams | limit, takeWhile; findFirst, findAny, anyMatch, allMatch, noneMatch |
Statefulness comes in degrees. distinct() remembers every element seen but can still emit each new one immediately. sorted() cannot emit anything until it has seen the last element, because the smallest element might arrive last. That second kind is a barrier: everything after it waits for everything before it. limit(n) is stateful only in the sense of keeping a counter, and its cost depends on encounter order, which matters a great deal in parallel.
Vertical execution: one element at a time
A common mental model is that each operation processes the whole collection and hands a new collection to the next. That model is wrong. Without a barrier, each element travels through the whole pipeline before the next element is read. The terminal operation pulls the work.
List<String> out = Stream.of("kiwi", "fig", "banana", "plum")
.filter(s -> { System.out.println("filter " + s); return s.length() > 3; })
.map(s -> { System.out.println("map " + s); return s.toUpperCase(); })
.limit(2)
.toList();
// filter kiwi
// map kiwi
// filter fig
// filter banana
// map banana <- limit is satisfied here; "plum" is never readInsert .sorted() after filter and the trace changes: all four filters print first, then the two maps. Short-circuiting downstream of a barrier can only stop work after the barrier, never before it. That is why sorted().findFirst() reads everything, and why min(comparator), which keeps one element, is the right way to find the smallest item.
Laziness has a second consequence: nothing happens until the terminal operation. A stream built and never consumed runs no lambdas, and a stream can be consumed only once; a second terminal call throws IllegalStateException.
How OpenJDK runs a pipeline: stages and sinks
The specification describes what operations do. The OpenJDK implementation, which is what you run, makes the behaviour above easy to see. This section describes the implementation in the java.util.stream package; it is not a contract, but it explains the traces.
Each intermediate call creates a pipeline stage that links to its predecessor and records the operation and its flags. When the terminal operation runs, it walks the stages from last to first, asking each to wrap the downstream Sink in its own sink. A sink has four methods: begin(size), accept(element), end() and cancellationRequested(). The source spliterator then pushes each element into the first sink.
A filter sink's accept calls its downstream only if the predicate passes. A map sink calls downstream with the mapped value. A sorted sink appends to a buffer and does nothing more until end(), when it sorts and replays the buffer downstream. A limit sink counts and, once full, reports cancellationRequested() as true. When any sink in the chain is short-circuiting, the source switches from bulk forEachRemaining to an element-at-a-time tryAdvance loop that checks for cancellation between elements.
Stream flags and the work the library may skip
Stages carry characteristics such as SIZED, ORDERED, DISTINCT and SORTED, starting from the source spliterator's characteristics and changing as operations are added: filter clears SIZED, map clears DISTINCT and SORTED, sorted() sets SORTED. The library uses them to skip work.
sorted()on a stream already known to be sorted by natural order can do nothing, for example after aTreeSetsource.distinct()on a stream already flagged DISTINCT can pass elements through.count()may not execute the pipeline at all if the count can be computed from the source. TheStream.countjavadoc illustrates this with a list source followed bypeekandcount(): thepeekaction may never run.
The lesson is the one the specification keeps repeating: behavioural parameters should be free of side effects you depend on. If logging or metrics must run for every element, do it in a forEach terminal operation or outside the stream.
Stateful operations in parallel
Parallel streams split the source into chunks and run the stateless stretch of the pipeline independently on each. A barrier forces the chunks to join. In ordered parallel pipelines, limit(n) and skip(n) must respect encounter order, so the library buffers chunks to know which elements are truly the first n; that can cost more than running sequentially. If you do not need encounter order, call unordered() so limit can take any n elements. Likewise findAny is cheaper than findFirst in parallel. Thread pools and splitting are covered in the fork/join pool article.
mapMulti: one-to-many without allocating streams
flatMap needs a stream for every input element, which is wasteful when most elements map to zero or one outputs. mapMulti, added in Java 16, gives you a consumer instead, and you push as many results as you like.
// Expand order lines into shipments, skipping cancelled lines, without a Stream per element
List<Shipment> shipments = orders.stream()
.<Shipment>mapMulti((order, sink) -> {
for (Line line : order.lines()) {
if (!line.cancelled()) sink.accept(new Shipment(order.id(), line.sku()));
}
})
.toList();The explicit type witness <Shipment>mapMulti is usually needed because the compiler cannot infer the output type from the consumer. Prefer flatMap when you already have a stream per element, and mapMulti when you would otherwise build tiny streams.
Gatherers: writing your own intermediate operation
Before JDK 24 the set of intermediate operations was closed. If you needed fixed-size batches, a running total or deduplication of consecutive values, you had to collect and restart, or write a spliterator. Stream.gather(Gatherer) opens the set. A Gatherer<T, A, R> consumes elements of type T, keeps private state A and emits R. It is defined by four functions: an initializer that creates state, an integrator called per element, an optional combiner that merges states for parallel evaluation, and an optional finisher called at the end, which can emit buffered output.
The integrator receives the state, the element and a Downstream. Downstream.push returns false once downstream wants no more, and the integrator returns false to stop consuming input. That is how short-circuiting is expressed. Integrator.ofGreedy marks an integrator that never initiates short-circuiting, which lets the implementation optimise. Gatherer.ofSequential(...) builds a gatherer that runs sequentially even in a parallel stream; Gatherer.of(...) with a combiner can run in parallel.
import java.util.*;
import java.util.function.Predicate;
import java.util.stream.*;
final class MoreGatherers {
// Drop consecutive duplicates: 1,1,2,2,2,1 -> 1,2,1. Sequential, greedy.
static <T> Gatherer<T, ?, T> distinctConsecutive() {
class State { T last; boolean seen; }
return Gatherer.ofSequential(
State::new,
Gatherer.Integrator.ofGreedy((state, element, downstream) -> {
if (state.seen && Objects.equals(state.last, element)) return true;
state.seen = true;
state.last = element;
return downstream.push(element);
}));
}
// Like takeWhile, but also emits the element that ends the run. Short-circuiting.
static <T> Gatherer<T, ?, T> takeUntilInclusive(Predicate<? super T> stop) {
return Gatherer.<T, T>ofSequential(
Gatherer.Integrator.of((unused, element, downstream) ->
downstream.push(element) && !stop.test(element)));
}
}The Gatherers class supplies built-ins. windowFixed(n) groups elements into lists of n, with a shorter last window. windowSliding(n) emits overlapping windows. fold is an ordered reduction that emits one value, for cases where no combiner exists. scan emits each running accumulation. mapConcurrent(maxConcurrency, mapper) runs the mapper on virtual threads with at most that many in flight, preserves stream order, and on failure rethrows the exception and cancels remaining tasks. Gatherers compose with andThen.
Worked example: a log-processing pipeline
Suppose you read a service's access log, need the per-minute request counts, want to flag minutes whose count is more than double the trailing five-minute average, and need to fetch recent deployments around each flagged minute from an HTTP service. Here is the pipeline with every operation's role annotated.
record Hit(long minute, String path, int status) {
static Hit parse(String line) { throw new UnsupportedOperationException("parsing omitted"); }
}
record Alert(long minute, long count, double baseline) {}
// Emits {minute, count} each time the minute changes. Constant memory; input must be time-ordered.
static Gatherer<Long, ?, long[]> countRuns() {
class Run { long minute; long count; }
return Gatherer.ofSequential(
Run::new,
Gatherer.Integrator.of((run, minute, downstream) -> {
if (run.count > 0 && minute != run.minute) {
boolean more = downstream.push(new long[] {run.minute, run.count});
run.count = 0;
if (!more) return false; // downstream is done: stop reading
}
run.minute = minute;
run.count++;
return true;
}),
(run, downstream) -> { // finisher flushes the last run
if (run.count > 0) downstream.push(new long[] {run.minute, run.count});
});
}
try (Stream<String> lines = Files.lines(Path.of("access.log"))) { // source; must be closed
List<Alert> alerts = lines
.map(Hit::parse) // stateless
.filter(h -> h.status() < 500) // stateless
.map(Hit::minute) // stateless; input is time-ordered
.gather(countRuns()) // custom: emits [minute, count] per run
.gather(Gatherers.windowSliding(6)) // stateful, bounded memory: 6 lists
.filter(w -> w.size() == 6) // short input yields one short window
.map(w -> {
double base = w.subList(0, 5).stream().mapToLong(m -> m[1]).average().orElse(0);
long[] now = w.get(5);
return new Alert(now[0], now[1], base);
})
.filter(a -> a.count() > 2 * a.baseline())
.toList();
List<String> context = alerts.stream()
.gather(Gatherers.mapConcurrent(16, a -> deployLog.changesAround(a.minute()))) // HTTP, virtual threads
.toList();
}Notice what is absent: no sorted(), no collect(groupingBy(...)) over the whole file. Because the log is already in time order, a sequential gatherer that counts runs of equal minutes uses constant memory, and the sliding window holds six elements. (Minutes with no traffic produce no run; a production version would emit zero counts for gaps.) A groupingBy version would hold one entry per minute for the whole file and lose the ordering needed for a trailing window. The network calls run sixteen at a time without a thread pool to manage, as virtual threads make that cheap, and results come back in order.
Failure modes
- Side effects in
peekormap. They may be skipped (count), reordered (parallel) or run on another thread. Treatpeekas a debugging aid only. - A barrier before short-circuiting.
sorted().findFirst()orsorted().limit(10)reads everything; on an infinite stream it never terminates. - Unbounded state.
distinct()on a long stream of mostly unique values keeps all of them in a hash set. - Stateful lambdas. A lambda that mutates an external list or counter breaks in parallel and is fragile sequentially. Put state inside a gatherer or collector, where its lifecycle is defined.
- Unclosed I/O sources.
Files.linesholds a file handle until closed; use try-with-resources. - Parallel
limiton ordered input. Often slower than sequential; useunordered()or stay sequential. - Gatherer without a combiner in a parallel stream. It is evaluated sequentially for that stage, which can quietly serialise a parallel pipeline.
What to do next
- For each pipeline on a hot path, mark every operation as stateless, stateful or barrier, and check that no short-circuit sits behind a barrier.
- Replace
sorted().findFirst()withminormax, andsorted().limit(k)on large inputs with a bounded heap. - Remove side effects from
peekandmap; move required effects toforEachor outside the stream. - Swap
flatMapover tiny streams formapMultiwhere it simplifies the code. - On JDK 24 or later, replace collect-and-restart workarounds with
Gatherers.windowFixed,windowSliding,scanor a custom gatherer, and test it with empty, single-element and short-circuited inputs. - Use
mapConcurrentfor bounded concurrent I/O instead of parallel streams. - Measure parallel pipelines against sequential ones before keeping them.