The naive version reads the whole file first:
List<String> lines = Files.readAllLines(file); // the entire dataset, resident in heapThat line is fine at ten thousand rows and fatal at a hundred million: readAllLines builds
the reservoir before you process a drop. The streaming alternative changes one method and
the entire memory story — the canal in the photo above: water moves through it continuously,
and at no point does the canal contain the river.
long countRejected(Path file) throws IOException { try (Stream<String> lines = Files.lines(file)) { // lazy: reads as the pipeline pulls return lines .filter(line -> !line.isBlank()) .map(Event::parse) .filter(Event::isRejected) .count(); // aggregates without retaining }}Files.lines does not read the file when called — it hands the pipeline a lazy source, and
each line is read, examined, and becomes garbage as the terminal operation pulls the next.
Heap usage is bounded by the line currently in flight (a pathological ten-megabyte line
costs ten megabytes) plus whatever the pipeline’s stages allocate per element — never by the
dataset, whether the file is a megabyte or a terabyte.
This post — item by item, with the sharp edges included — is about that discipline:
processing data whose size is none of your heap’s business.
The lazy sources the JDK actually ships
Streaming starts at the source; everything downstream inherits its laziness:
Files.lines(path)/BufferedReader.lines()— line-at-a-time file and reader streams. Two sharp edges: the stream holds the file handle, so it must be closed — hence the try-with-resources above, easily forgotten because most streams need no closing (the library’sResource.fromAutoCloseablecomposes the same acquire-use-release discipline as a value, if you prefer it over the block) — and read failures after opening surface asUncheckedIOExceptionfrom inside the pipeline, not from theFiles.linescall itself.Stream.iterate(seed, next)/Stream.generate— computed sources: unbounded sequences of pages, IDs, retries-with-timestamps. Infinite by construction (the three-argiterate(seed, hasNext, next)overload is the self-terminating variant), usable because short-circuiting (limit,takeWhile,findFirst) stops the pull.- Your own sources — anything that can produce elements on demand (a paginated API, a
cursor, a message poll) becomes a stream via an
IteratororSpliteratorhanded toStreamSupport.stream(...)— with the lifecycle wired up: registeronCloseon the stream to close the underlying cursor, and keep theStatement/Connectionin an enclosing try-with-resources, so early termination (limit, an exception) releases everything. The database counterpart is driver-level: a forward-onlyResultSetwith a boundedfetchSize, wrapped the same way, streams rows a batch at a time — when the driver honors it:fetchSizeis a JDBC hint, and the common drivers need convincing (PostgreSQL streams only with autocommit off; MySQL needs its streaming mode). Verify with your driver before betting the heap on it.
The unifying property: the source answers “give me the next one,” never “give me everything.”
The pipeline is only as streaming as its greediest stage
A lazy source buys nothing if a downstream stage rebuilds the reservoir. The stages divide cleanly:
Flow-preserving — hold one element (or one bounded window) at a time: filter, map,
takeWhile, limit, and bounded
gatherers. Aggregating
terminals whose accumulator has a fixed size — count, sum, a max, a running statistics
object — also keep the flow. Two need a second look: flatMap preserves the flow only as
far as its mapper does (a mapper that materializes a large collection per element pays that
cost per element), and reduce is flow-preserving only while its accumulator stays bounded —
reducing into a growing List is toList() wearing a disguise.
Reservoir-building — sorted() buffers the entire upstream before emitting anything
(sorting a terabyte through a stream is still sorting a terabyte), and distinct() carries
a seen-set that in the worst case is the dataset again. toList() materializes the output
outright. Collectors.groupingBy deserves precision: the two-arg form with a folding
downstream — groupingBy(Event::category, counting()) — produces a small summary per
category, but plain groupingBy(classifier) builds a Map holding every element,
regrouped: the reservoir with extra keys.
The practical test before running a pipeline over big data: for each stage, ask what it must remember. The canonical large-data shape remembers only a window — batching a huge import into bounded inserts:
try (Stream<String> lines = Files.lines(file)) { lines.map(Event::parse) .gather(Gatherers.windowFixed(500)) // one 500-element batch in memory at a time .forEach(batch -> repository.insertAll(batch));}Memory holds one batch, not one dataset — the gatherers post
covers windowFixed and friends in depth.
Failures per element, without stopping the river
At a hundred million rows, some rows are garbage — and one bad line must not abort hour six of the import. The typed-outcome move applies per element: capture the failure as a value, keep flowing.
var rejected = new LongAdder(); // contained effect: count what we drop
try (Stream<String> lines = Files.lines(file)) { lines.map(line -> Try.of(() -> Event.parse(line))) // failure becomes a value, not an abort .filter(t -> { if (t.isFailure()) { rejected.increment(); return false; } return true; }) .map(Try::get) .gather(Gatherers.windowFixed(500)) .forEach(repository::insertAll);}log.info("import done, {} malformed lines rejected", rejected.sum());Each line’s outcome is data, handled and accounted for without holding more than the
current element. (When you genuinely do not need the count, the terse keep-successes form is
.flatMap(Try::stream) — the library’s one-method bridge from a Try to a zero-or-one
element stream.) A caution worth appending to the
streams post’s
Result.toList example: collectors that gather all outcomes are aggregation tools for
bounded data — on an unbounded stream they are the reservoir again. Per-element handling
stays per-element.
One honest boundary on the abort-proofing itself: Try.of wraps your parse, not the
source. The UncheckedIOException edge from the sources section lives upstream of it — a
single malformed byte in a UTF-8 file aborts the pull from inside Files.lines, before any
Try sees it. If the input’s encoding is untrusted, read bytes and decode per line with a
CharsetDecoder configured with CodingErrorAction.REPORT — the String constructor
silently replaces malformed bytes rather than failing — so each decode failure surfaces
as one more per-element outcome instead of corrupted-but-green data.
The honest limits
- A stream is one pass. Consumed is consumed. If two computations need the data, either fuse them into one pass (a combined fold, a gatherer) or accept reading the source twice.
parallelStreamis the wrong tool for IO-bound streaming — not because the source splits poorly (Files.lineshas split efficiently since JDK 9) but because per-element blocking IO runs on the sharedForkJoinPool.commonPool, starving every other parallel stream in the JVM while threads sit in waits. For concurrent per-element IO inside a streaming pipeline,Gatherers.mapConcurrent(n, ...)bounds the fan-out on virtual threads while preserving the flow — wrap the mapper inTry.ofthere too, since an exception from one call cancels the remaining in-flight work, while aTrykeeps it a per-element outcome — the gatherers post has the fullermapConcurrentstory, and the concurrency post the fan-out problem it solves.- Backpressure is the async version of this discipline. When the producer is a network
peer rather than a file you pull from, reactive streams (
Flux, withfun-reactorfor typed outcomes) make “don’t send more than I can hold” an explicit protocol. A pull-basedStreamsidesteps the problem differently: it consumes the source at its own pace, in pipeline order — nothing pushes — but bounded demand or prefetch against an async producer exists only where an adapter explicitly enforces it. - Laziness here is about sequences. Its sibling — deferring and memoizing a single
expensive value with
Lazy<T>— is the lazy-evaluation post’s territory. (Do not wrap a stream itself in aLazy: memoization would hand every caller the same one-shot, handle-holding stream — defer the inputs to a pipeline, not the pipeline.)
The through-line with the rest of this blog: streaming is the pipeline shape under a memory constraint. Pure per-element functions, outcomes as values, one materialization — at the sink, sized to the answer rather than the input. Build the canal, not the reservoir.
Further reading
- Streams, Immutable Collections, and Efficient Data Processing — fusion and short-circuiting, the machinery this post runs at scale.
- Stream Gatherers: Custom Intermediate Operations
—
windowFixed,mapConcurrent, and building your own bounded stages. - Lazy Evaluation: When It Helps — the
single-value side of laziness, with
Lazy<T>. - Modeling Data Transformation Pipelines — the pipeline shape, before the memory constraint.
- Functional Concurrency: Parallel Work Without Shared State — bounded fan-out when per-element work is slow.
- Designing a Good Error Type — the typed failures that let bad rows become data instead of aborts.
Found a bug or have a suggestion? Open an issue on GitHub.
