Sometimes a stream job needs memory of earlier elements: groups of N items, overlapping pairs, or a total that updates after each value. map, filter, and flatMap work on one element at a time, so they do not cover those cases well. People often stop the stream early with collect, hide a list inside a lambda, or switch to a for loop.

Stream Gatherers fill that gap. They let you add a step in the middle of a stream that can remember earlier values, still pass results along, and keep working with the rest of the pipeline.

New to pipelines? Start with Java Streams basics, then Streams advanced for collectors, parallel pitfalls, and production habits — then come back here.

A Gatherer is a custom intermediate stream operation — you decide how each upstream element updates state and which (if any) elements go downstream. Reach for gatherers when built-in ops force awkward intermediate collections or sneaky mutable captures — not when map / filter / flatMap already say the same thing clearly.

The problem gatherers solve

Classic intermediate ops are powerful but fixed. Need every pair of consecutive readings?

List<Integer> temps = List.of(18, 19, 21, 20, 22);

// Awkward: leave the stream, build windows by hand
List<List<Integer>> windows = new ArrayList<>();
for (int i = 0; i < temps.size() - 1; i++) {
    windows.add(List.of(temps.get(i), temps.get(i + 1)));
}

Or you smuggle state into a lambda and hope nobody runs the stream in parallel:

List<Integer> buffer = new ArrayList<>();
temps.stream()
        .peek(buffer::add) // mutable capture — brittle and parallel-hostile
        .filter(t -> buffer.size() >= 2)
        // ... still does not emit windows cleanly
        .toList();

Gatherers keep the work inside the pipeline. You stay lazy until a terminal op, and the JDK ships common shapes so you do not invent them twice.

import java.util.stream.Gatherers;

List<List<Integer>> windows = temps.stream()
        .gather(Gatherers.windowSliding(2))
        .toList();
// [[18, 19], [19, 21], [21, 20], [20, 22]]

Same intent. No early collect. No mutable peek hacks.

When gatherers shipped

ReleaseStatusSpec
Java 22First previewJEP 461
Java 23Second previewJEP 473
Java 24Standard featureJEP 485

Use Java 24+ for gatherers with no --enable-preview. The API lives on Stream.gather(...) and java.util.stream.Gatherer / Gatherers.

Mental model

Think of a gatherer as four optional pieces that turn upstream elements into downstream elements:

PieceRole
InitializerCreates per-pipeline (or per-partition) mutable state
IntegratorConsumes one upstream element; may push zero or more downstream
FinisherRuns at end-of-stream; may flush remaining state
CombinerMerges parallel partitions (only when you need parallel gather)
upstream ──► integrator(state, element, downstream) ──► downstream
                    │
                    └── finisher(state, downstream) at the end
  • Stream.gather(gatherer) is an intermediate op — it returns another Stream.
  • Built-ins live in java.util.stream.Gatherers.
  • Custom gatherers usually start with Gatherer.ofSequential(...) unless you deliberately support parallel.
  • Prefer records for emitted DTOs and window payloads so equals / toString stay honest.

You still prefer map / filter when the transform is one-in, zero-or-one-out with no cross-element memory. Gatherers earn their keep when output depends on more than the current element.

Built-in gatherers

1. Fixed and sliding windows

Batch into fixed-size lists, or emit overlapping windows:

import java.util.stream.Gatherers;

List<String> events = List.of("a", "b", "c", "d", "e");

List<List<String>> batches = events.stream()
        .gather(Gatherers.windowFixed(2))
        .toList();
// [[a, b], [c, d], [e]]  — last window may be short

List<List<String>> sliding = events.stream()
        .gather(Gatherers.windowSliding(3))
        .toList();
// [[a, b, c], [b, c, d], [c, d, e]]

Use fixed windows for “flush every N messages.” Use sliding windows for moving averages, consecutive-pair diffs, and n-gram style views.

2. Fold: collapse to one value mid-pipeline

fold is like a running reduce that emits a single result when the stream ends — useful when a later stage still wants a Stream of that one value (or you want to keep chaining):

Optional<Integer> sum = Stream.of(1, 2, 3, 4)
        .gather(Gatherers.fold(() -> 0, Integer::sum))
        .findFirst();
// Optional[10]

Contrast with scan, which emits the intermediate accumulation after every element.

3. Scan: running totals that keep flowing

List<Integer> running = Stream.of(1, 2, 3, 4)
        .gather(Gatherers.scan(() -> 0, Integer::sum))
        .toList();
// [1, 3, 6, 10]

Good for cumulative metrics, balance-after-each-txn style ledgers, and progressive progress values without leaving the stream.

4. Concurrent map with a concurrency limit

mapConcurrent applies a mapper with at most maxConcurrency in-flight tasks — handy for I/O-bound transforms while preserving encounter order of results:

List<String> bodies = urls.stream()
        .gather(Gatherers.mapConcurrent(8, url -> httpClient.get(url)))
        .toList();

Prefer this over a hand-rolled executor + queue when you already live in a stream pipeline. Keep the mapper side-effect-aware and timeout-bounded like any concurrent I/O code.

One custom gatherer

Built-ins cover windows and folds. Custom gatherers shine for domain rules — for example, keep the first occurrence of each key while streaming (dedupe-by-key without collecting the whole stream first):

import java.util.HashSet;
import java.util.Set;
import java.util.function.Function;
import java.util.stream.Gatherer;

static <T, K> Gatherer<T, ?, T> distinctBy(Function<? super T, ? extends K> keyFn) {
    return Gatherer.ofSequential(
            HashSet<K>::new,
            (Set<K> seen, T element, Gatherer.Downstream<? super T> downstream) -> {
                if (seen.add(keyFn.apply(element))) {
                    return downstream.push(element);
                }
                return !downstream.isRejecting();
            }
    );
}

Use it:

record Event(String id, String payload) {}

List<Event> unique = events.stream()
        .gather(distinctBy(Event::id))
        .toList();

Anatomy of that gatherer:

  • Initializer — HashSet::new holds keys already emitted.
  • Integrator — push only when add returns true; return false from push / respect isRejecting() so short-circuiting terminals (findFirst, limit) can stop early.
  • No finisher — nothing left to flush.
  • Sequential-only — a parallel distinct-by-key would need a combiner (and usually a different design).

Another common pattern is “buffer until a flush signal,” with a finisher that empties the leftover buffer when the upstream ends. Same four pieces; the finisher is what makes end-of-stream correct.

Parallel and the combiner

Gatherer.of(...) (as opposed to ofSequential) expects a combiner when state must merge across parallel partitions. Windowing and distinct-by-key are often sequential by nature — forced parallelism can change semantics (duplicate keys across partitions, windows that span split boundaries).

Rules of thumb:

  • Start with Gatherer.ofSequential for domain gatherers.
  • Add a combiner only when parallel splits are meaningful and mergeable (e.g. independent numeric aggregates).
  • Do not assume parallel() makes every gatherer faster — coordination and ordering costs dominate for small streams and I/O-heavy mappers.

When they win vs when they don’t

NeedPrefer
Sliding / fixed windows, running scan, mid-pipeline foldBuilt-in Gatherers
Custom cross-element state still inside the pipelineCustom Gatherer
One-in, one-out pure transformmap
Drop or keep by predicatefilter
One-in, many-out expansionflatMap
Whole-stream result with no further intermediate opscollect / reduce
Heavy custom graph of joins across multiple sourcesLeave Streams; use a clearer imperative or reactive design

Gatherers are not a replacement for collectors. Collectors terminate (or feed collect). Gatherers sit in the middle so you can keep mapping, filtering, and gathering after them.

Gotchas

Short-circuiting and Downstream

Integrators should honor downstream.push(...)’s boolean return and/or downstream.isRejecting(). Ignoring rejection wastes work after findFirst or limit already stopped caring.

Mutable state is local, not shared

Initializer state is owned by the gatherer invocation. Do not close over a shared ArrayList from outside the gatherer the way peek-hacks do — that reintroduces the bug gatherers were meant to remove.

Last incomplete window

windowFixed may emit a shorter final window. If you need only full batches, filter on size after the gather (or write a custom gatherer that withholds a partial flush).

Parallel + ordered I/O

mapConcurrent helps with concurrency limits, but the surrounding pipeline’s parallelism and ordering still matter. Measure; do not assume “gather = faster.”

Preview vs final

On Java 22–23 you needed --enable-preview. On Java 24+ (JEP 485) you do not. Pin your examples and CI to 24+ so readers are not stuck on preview flags.

Cheat sheet

stream.gather(gatherer)          // intermediate op
Gatherers.windowFixed(n)
Gatherers.windowSliding(n)
Gatherers.fold(initializer, folder)
Gatherers.scan(initializer, scanner)
Gatherers.mapConcurrent(max, mapper)

Gatherer.ofSequential(initializer, integrator)
Gatherer.ofSequential(initializer, integrator, finisher)
Gatherer.of(initializer, integrator, combiner, finisher)

Java 24+ (JEP 485); no preview flag
Good: windows, running totals, dedupe-by-key, bounded concurrent map
Avoid: replacing simple map/filter; parallel when semantics are sequential
Watch: Downstream rejection; short final fixed windows; external mutable captures

Do:

  • Reach for built-in Gatherers before writing a custom integrator.
  • Keep gatherer state inside the initializer; emit records when the downstream type is a data carrier.
  • Respect short-circuiting via Downstream.

Don’t:

  • Use gatherers to hide what a plain map already expresses.
  • Capture and mutate outer lists from the integrator.
  • Flip on parallel() for sequential-by-nature gatherers and hope the results match.

Pros and cons

Pros

  • True intermediate extension point — stay lazy and keep chaining
  • Built-ins cover windows, fold, scan, and bounded concurrent map
  • Custom gatherers encode domain streaming rules without leaving Streams
  • Final in Java 24 — no preview tax for new codebases on current JDKs

Cons

  • Easy to over-abstract — a small loop can still be clearer than a clever gatherer
  • Parallel support needs an honest combiner; many useful gatherers are sequential-only
  • Requires Java 24+ for the standard API
  • Learning curve around Downstream and finisher flush semantics

Wrap-up

Stream Gatherers close the gap between “Streams are elegant” and “I need windows / running state / custom flush rules.” They became a standard feature in Java 24 (JEP 485) after two preview rounds.

Start with Gatherers.windowFixed, windowSliding, scan, and fold. When the domain rule is not one of those shapes, write a sequential gatherer with a tight initializer and integrator — and keep simple transforms on map / filter. For immutable payloads leaving the gatherer, lean on records; when those payloads form a closed set of variants, sealed classes + pattern switch stay the right companion tools downstream.