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.

Advertisement

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.

PropertyMeaningExamples
StatelessProcessing one element needs no information about othersfilter, map, flatMap, mapMulti, peek
StatefulThe operation must remember elements already seen, possibly all of themdistinct, sorted, limit, skip, takeWhile, dropWhile
Short-circuitingCan finish without consuming all input, so it can work on infinite streamslimit, 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 read

Insert .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.

Advertisement

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.

How a pipeline executes: stages are linked at build time, elements flow one at a time at run timeSourcespliteratorfilterstatelessmapstatelesssortedstateful barrierlimit + toListterminalBuild time: each call returns a new stage that records its op and flags; nothing runsRun time: the terminal op wraps sinks back to front, then pushes elements front to backforEachRemainingor tryAdvancefilter sinkaccept(e)map sinkaccept(f(e))sorted sinkbuffers alllimit sinkcancellationeeat end()Before the barriereach element goes through filter and map, one by oneAfter the barriersorted emits only when upstream ends, then limit stopsSink chaining describes the OpenJDK implementation; the specification only promises laziness and the documented semantics
Stages are linked when the pipeline is built; sinks are wrapped back to front and fed front to back when the terminal operation runs.

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 a TreeSet source.
  • 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. The Stream.count javadoc illustrates this with a list source followed by peek and count(): the peek action 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 peek or map. They may be skipped (count), reordered (parallel) or run on another thread. Treat peek as a debugging aid only.
  • A barrier before short-circuiting. sorted().findFirst() or sorted().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.lines holds a file handle until closed; use try-with-resources.
  • Parallel limit on ordered input. Often slower than sequential; use unordered() 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

  1. 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.
  2. Replace sorted().findFirst() with min or max, and sorted().limit(k) on large inputs with a bounded heap.
  3. Remove side effects from peek and map; move required effects to forEach or outside the stream.
  4. Swap flatMap over tiny streams for mapMulti where it simplifies the code.
  5. On JDK 24 or later, replace collect-and-restart workarounds with Gatherers.windowFixed, windowSliding, scan or a custom gatherer, and test it with empty, single-element and short-circuited inputs.
  6. Use mapConcurrent for bounded concurrent I/O instead of parallel streams.
  7. Measure parallel pipelines against sequential ones before keeping them.
Key takeaway: Stream operations are defined by three properties: whether they keep state, whether that state is a barrier, and whether they can stop early. Elements travel vertically through the pipeline until a barrier, the terminal operation drives everything, and flags let the library skip work you thought was guaranteed. With those rules you can predict any pipeline, and with Gatherers you can add the operation the library is missing instead of breaking the pipeline apart.