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
| Release | Status | Spec |
|---|---|---|
| Java 22 | First preview | JEP 461 |
| Java 23 | Second preview | JEP 473 |
| Java 24 | Standard feature | JEP 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:
| Piece | Role |
|---|---|
| Initializer | Creates per-pipeline (or per-partition) mutable state |
| Integrator | Consumes one upstream element; may push zero or more downstream |
| Finisher | Runs at end-of-stream; may flush remaining state |
| Combiner | Merges 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 anotherStream.- 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::newholds keys already emitted. - Integrator — push only when
addreturns true; returnfalsefrompush/ respectisRejecting()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.ofSequentialfor 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
| Need | Prefer |
|---|---|
| Sliding / fixed windows, running scan, mid-pipeline fold | Built-in Gatherers |
| Custom cross-element state still inside the pipeline | Custom Gatherer |
| One-in, one-out pure transform | map |
| Drop or keep by predicate | filter |
| One-in, many-out expansion | flatMap |
| Whole-stream result with no further intermediate ops | collect / reduce |
| Heavy custom graph of joins across multiple sources | Leave 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
Gatherersbefore 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
mapalready 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
Downstreamand 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.