Processing millions of records in Java: chunking, streams, and memory-efficient pipelines
Process huge datasets in Java without OOM — combine keyset-paginated batches or true database streams with Stream pipelines, lambdas, and checkpoints.
A nightly job that recomputes a field on every customer has run fine for a year. The table now has five million rows, and the job either dies with OutOfMemoryError or stalls in GC long enough that the orchestrator kills it. Someone “fixes” it by adding .stream() — the code looks functional now, and it still holds the entire table in the heap.
Environment assumed below: Java 21, Spring Boot 4.1, Spring Data JPA (Hibernate 7), PostgreSQL 18. Where behaviour is driver- or database-specific, that is stated rather than implied.
The common mistake: stream() is not a data source
List<Customer> all = customerRepository.findAll();all.stream() .map(this::toCommand) .filter(this::isEligible) .forEach(this::process);map and filter are genuinely lazy — nothing happens until forEach runs. But laziness describes how the operations behave, not where the data lives. findAll() already materialised every row — entity objects, column values, Hibernate’s persistence-context bookkeeping — before stream() was called. The pipeline traverses a fully-loaded collection; memory stays O(table size) no matter how elegant the chain looks.
The same trap applies to Lists.partition(all, 500) and List.subList(): those are views over the loaded list. They divide the work, not the memory. To bound memory, the bound must happen where the data enters the pipeline — the query itself.
The imperative baseline: keyset pagination
Fetch a bounded batch ordered by a stable, unique, indexed key, carry the last value forward as a cursor:
SELECT id, email, total_spendFROM customersWHERE id > :lastSeenId AND id <= :jobMaxIdORDER BY idLIMIT :batchSize;As Spring Data JPA, a projection keeps the persistence context out entirely:
package in.o612.eng.billing;
public record CustomerRow(long id, String email, java.math.BigDecimal totalSpend) { }public interface CustomerRepository extends JpaRepository<Customer, Long> {
@Query(""" select new in.o612.eng.billing.CustomerRow(c.id, c.email, c.totalSpend) from Customer c where c.id > :lastSeenId and c.id <= :maxId order by c.id """) List<CustomerRow> findBatch(@Param("lastSeenId") long lastSeenId, @Param("maxId") long maxId, Pageable pageable);
@Query("select max(c.id) from Customer c") long findMaxId();}PageRequest.of(0, batchSize) supplies the LIMIT — JPQL has none. Returning List (not Page) means no count query fires. The projection must be a top-level class: Hibernate resolves constructor expressions by canonical name and cannot resolve in.o612.eng.Outer.Inner.
public void run(int batchSize) { long maxId = customerRepository.findMaxId(); // captured once, at job start long cursor = 0; while (true) { List<CustomerRow> batch = customerRepository.findBatch(cursor, maxId, PageRequest.of(0, batchSize)); if (batch.isEmpty()) break; batchProcessor.process(batch); // own transaction — see below cursor = batch.getLast().id(); // List.getLast(): Java 21 }}This baseline is boring on purpose. Each iteration is a query, a bounded list, a transaction — memory is O(batchSize), progress is inspectable, and a failure after batch k restarts at cursor k. For most jobs this alone is sufficient; everything below is refinement.
Streams and lambdas — inside the batch
With the source bounded, a Stream is the right tool for the per-batch transform:
batch.stream() .map(this::toCommand) // lazy intermediate .filter(this::isEligible) // lazy intermediate .forEach(processor::handle); // terminal — the only step that runs workOnly forEach triggers traversal; map/filter fuse into a single pass over the batch, so each row visits toCommand then isEligible then handle in sequence — no intermediate collections are created for the intermediate stages. Two disciplines keep this honest:
- Prefer method references and named functions over anonymous chains.
this::isEligiblecan be unit-tested, reused, and read in a stack trace; a six-line lambda cannot. - Never
collect(toList())a result you then iterate — it materialises what the pipeline just avoided, and at dataset scale it’s thefindAll()bug again. Terminal operations that reduce (forEach,count,collectinto a bounded aggregate) are fine; terminal operations that re-materialise the data are not.
Each batch is held in memory whole — that’s the design. Total memory stays O(batchSize + per-batch working set), flat from row 1 to row 5,000,000.
A genuinely streaming source
When one pass over everything is the actual job, Spring Data can return a Stream whose source is the result set itself:
@Query(""" select new in.o612.eng.billing.CustomerRow(c.id, c.email, c.totalSpend) from Customer c where c.id <= :maxId order by c.id """)Stream<CustomerRow> streamAllOrdered(@Param("maxId") long maxId);@Transactional(readOnly = true)public void run() { try (Stream<CustomerRow> rows = repo.streamAllOrdered(maxId)) { rows.forEach(processor::handle); }}The try-with-resources is not decorative — the stream owns a JDBC ResultSet and its close() releases it; forgetting it leaks a cursor until connection close. And the caveat that matters: a Stream<T> return type does not guarantee streaming. On PostgreSQL the driver buffers the whole ResultSet unless autocommit is off and fetchSize is positive — inside @Transactional you get both by default, but a misconfigured connection or a different driver changes the behaviour silently. Verified on this stack, it streamed a million rows in one call; verify on yours with a heap dump or -verbose:gc, not by reading the signature.
The trade-off is the transaction’s lifespan: one transaction spans the whole scan. There are no per-batch commits, so nothing is durable until the end — a crash at row 4,999,999 replays everything. That’s why the streaming variant suits read-mostly work (transform-and-publish, export), while jobs that write per row usually pair the stream with explicit chunking — buffer batchSize rows, flush, clear:
List<CustomerRow> chunk = new ArrayList<>(batchSize);rows.forEach(row -> { chunk.add(row); if (chunk.size() == batchSize) { writer.writeBatch(chunk); // own transaction boundary if needed chunk.clear(); }});if (!chunk.isEmpty()) writer.writeBatch(chunk);Functional concepts, stated plainly
map/filter/ terminal op are a pipeline specification: declarative shape, one traversal, no intermediate collections.- Pure vs effectful functions.
toCommandandisEligibleare pure — same input, same output, no writes — so they’re trivially testable and reorderable.handle, which writes to the database or publishes an event, is a side effect. Keep the two visibly separate; a pipeline of pure transforms feeding one effectful sink is far easier to reason about than effects scattered throughmapstages. - Exceptions change shape inside
forEach. A throw mid-stream aborts the terminal operation with no built-in resume point — you cannot checkpoint “row 12,341 of the stream” the way you checkpoint “batch 24 complete.” Retries therefore belong outside the pipeline: retry the batch, or the job, not the lambda. parallelStream()is rarely the answer here. It fans work onto the common ForkJoinPool, makes side effects unordered, and multiplies load against the real bottleneck — the database. For CPU-bound per-row work with no writes, it can help; for anything effectful or DB-bound, sequential streams plus deliberately-sized batches are faster in practice and far easier to debug.
Ways to form chunks — compared
- Batch at the source (keyset), stream per batch. Memory O(batch), a transaction per batch, natural checkpoint points, resumable. The default choice.
- Stream the whole result set, buffer into chunks for batched writes. Memory O(chunk), one long transaction, checkpoint only if you add it yourself — good for read-heavy exports, risky for long write jobs.
List.subListon a loaded collection. Memory O(table). Not a strategy — the bug this article removes.
Production concerns
- Persistence context. If you load entities instead of projections, every managed entity accumulates in the
EntityManageruntil the transaction ends — the memory problem returns inside Hibernate. Per batch:em.flush()thenem.clear()(in that order — clearing first silently drops pending changes). And call the@Transactionalbatch method through a separate injected bean: invoking it viathis.bypasses the proxy and silently collapses everything into one giant transaction — the most reliable way to write this bug. - Checkpoints belong after commit, inside the same transaction as the batch’s own writes, so progress and results commit or roll back atomically. Persist
last_idin ajob_checkpointtable; resume from it on restart. If results leave the database (Kafka, email), a committed transaction is not delivery — use an outbox table written in the same transaction. - Idempotency before retries. Any replay after a partial failure re-runs work, so processing must tolerate it: upserts, dedupe keys, not blind inserts.
- Concurrent writes. Keyset traversal sees inserts above the cursor and misses updates to already-passed rows; deletes vanish silently. Capturing
maxIdat job start bounds the window. A true consistent snapshot means a held-openREPEATABLE READtransaction orpg_export_snapshot()— possible, but you’re trading away the per-batch-commit design for it. - Batch size is a measurement, not a constant. Bigger batches amortise round trips; smaller ones shorten transactions and the replay window. Start around a few hundred rows, then tune against query time and lock lifetime — there is no universal number.
- Observe it. Log cursor and rows/second at a fixed cadence; a six-hour job with no progress signal is indistinguishable from a hung one.
The worked example, twice
Recompute loyalty_tier for 4M customers nightly.
A — keyset + per-batch stream (default). findMaxId() at start; loop findBatch; per batch, batch.stream().map(this::toCommand).filter(this::isEligible).forEach(writer::apply) inside a @Transactional method on a separate bean; write the checkpoint in that same transaction; advance the cursor. Flat memory, a commit every 500 rows, resumable after any crash.
B — database stream + chunk buffer. One @Transactional(readOnly = true) method; try (Stream<CustomerRow> rows = repo.streamAllOrdered(maxId)); buffer into chunks, hand each chunk to a writer bean with its own REQUIRES_NEW transaction for writes and checkpoints; close the stream. Fewer round trips and no pagination queries — at the price of a long-lived transaction and cursor, plus care that the JDBC fetch actually streams.
A is the better default: bounded transactions, real checkpoints, simpler failure semantics. Reach for B when the job is genuinely a single read pass and the write path batches well.
Decision table
| Approach | Memory | Resource lifetime | Complexity | Best use |
|---|---|---|---|---|
findAll().stream() / subList chunks | O(table) | tx per call only | trivial | Never for large data — the anti-pattern |
| Keyset batch + per-batch stream | O(batch) | short tx per batch | low | Updateable jobs, restarts — the default |
DB-backed Stream | O(1) rows | one long tx, open cursor | medium | Read-only passes, exports |
JDBC cursor + fetchSize | O(fetchSize) | connection held | medium-high | No JPA, maximum control |
| Spring Batch chunks | O(chunk) | managed | framework | Scheduling, skip/retry, restart semantics |
Checklist
- No
findAll()/subListas the chunking mechanism; the bound happens in SQL - Cursor unique and indexed;
maxIdupper bound captured at start - Streams used per batch, or backed by a verified streaming fetch (
fetchSize, autocommit off, open tx) -
try-with-resourceson every database-backed stream - Pure functions (
map/filter) separate from effectful sinks; retries outside the pipeline -
forEach/count/boundedcollectas terminals — nevertoList()at dataset scale - Checkpoint committed atomically with each batch’s writes; processing idempotent for replay
- Progress logged; batch size measured against query time, not guessed
- No
parallelStream()over DB-backed or side-effecting pipelines