Streaming large results
How em.stream() is implemented for Neo4j: what runs on the server, what the client holds, and which
query shapes stay lazy.
mikro-orm-neo4j ≥ 0.3.0 · @mikro-orm/core 7.x · Neo4j ≥ 5
To process a large amount of entities without loading them all into memory, use em.stream(). It
returns a cursor you can iterate with for await ... of:
const stream = em.stream(Book, {
populate: ['author'],
where: { price: { $gt: 100 } },
orderBy: { id: 'ASC' },
});
for await (const book of stream) {
console.log(book.title);
console.log(book.author.$.name);
}
What it does
The query is run as a Bolt cursor. The driver asks the server for chunkSize records at a time
(PULL n), hands them over as they arrive and asks for more only as the consumer keeps up — the
Bolt driver's own high/low watermarks (70% / 30% of the fetch size) pause and resume the stream. At
most a batch or so is queued client-side, whatever the size of the result.
em.find() | em.stream() | |
|---|---|---|
| Records | Every record is buffered before anything is mapped | Pulled chunkSize at a time, mapped one entity at a time |
| Client memory | The whole result, plus the entity graph | One batch, plus the entity currently being assembled |
| Time to first entity | After the last record | After the first batch |
| Identity map | Entities are managed | Entities come from throwaway forks — nothing is retained |
| To-many relations | A second, select-in query per relation | Expanded in the same query and merged as rows arrive |
| Failure mid-query | Nothing is returned | The entities produced before the failure are already yielded |
That last row is the sharpest observable difference. Cypher evaluates lazily, so a query that fails on its thousandth row still delivers the first nine hundred and ninety nine:
const cursor = em.streamRaw('UNWIND range(1, 10) AS i RETURN CASE WHEN i > 5 THEN 1 / 0 ELSE i END AS value');
// yields 1…5, then throws
The cursor
em.stream(), repo.stream(), em.streamRaw() and qb.stream() all return a Neo4jCursor. It is
an async iterable — a for await loop is the normal way to use it, and ending that loop for any
reason (break, return, a thrown error) cancels the query server-side and returns the session to
the pool. On top of that it carries:
const cursor = em.stream(Book, { orderBy: { id: 'ASC' } });
await cursor.next(); // drive it by hand
await cursor.toArray(); // drain it (defeats the purpose, handy in tests)
cursor.map((b) => b.title); // derive another cursor, same lifecycle
cursor.asStream(); // an object-mode Node Readable, for pipelines
await cursor.close(); // release it explicitly; idempotent
cursor.stats; // { fetchSize, records, yielded, closed, took }
asStream() accepts a transform, so a cursor drops straight into a Node pipeline:
await pipeline(
em.stream(Book, { orderBy: { id: 'ASC' } }).asStream({ transform: (b) => `${b.title}\n` }),
createWriteStream('books.txt'),
);
stats is live, and shared by every layer of one stream: records counts what came off the wire,
yielded counts what the consumer saw. With a to-many relation populated, the ratio between them is
the merge factor — 60 records in, 20 entities out.
Chunk size
chunkSize is the Bolt fetch size, i.e. how many records the server sends per round trip. Lower
values hold less in memory and cost more round trips; higher values do the opposite. The iterator
still yields one entity at a time regardless.
const stream = em.stream(Book, {
orderBy: { id: 'ASC' },
chunkSize: 100, // 1000 (Bolt's own default) unless set here or on the driver
});
Resolution order is chunkSize → driverOptions.fetchSize on the connection → 1000. Passing
FETCH_ALL (-1) asks for everything in one batch, which is what a non-streamed query does. Any
other non-positive value is rejected up front rather than deep inside the Bolt driver.
Inside an explicit transaction the fetch size belongs to the session the transaction was opened
with, so chunkSize is ignored there.
Populating relations
Streaming is the one read path where core runs no second query: em.stream() forces the joined
strategy and never does the select-in pass em.find() relies on for collections. Every populated
relation — to-one, to-many, nested, any depth — therefore has to come out of the one Cypher query.
A to-one relation is an OPTIONAL MATCH that leaves the row count alone. A to-many one multiplies
rows, and those rows are merged back into entities as they arrive:
// em.find(Author, {}, { populate: ['books'] }) — collected server-side
MATCH (this0:Author)
OPTIONAL MATCH (this0)-[this1:WROTE]->(this2:Book)
RETURN this0 AS node, collect(this2) AS rel_0
// em.stream(Author, { populate: ['books'] }) — expanded, merged client-side
MATCH (this0:Author)
OPTIONAL MATCH (this0)-[this1:WROTE]->(this2:Book)
RETURN this0 AS node, this2 AS rel_0
Why not keep collect()? Because an aggregate hands back a whole collection in a single record:
a supernode's thousand relations arrive in one piece however small the fetch size is, which puts the
memory they take outside chunkSize's reach. The planner also only streams an aggregation while it
can prove the input is already grouped (OrderedAggregation); order by anything that is not the
grouping key and it becomes a blocking EagerAggregation that consumes the whole result before
emitting a row. Expansion has neither problem, and it is what the SQL drivers do with a
joined-strategy result set.
Merging is bounded by design: rows for one root arrive together, so only the entity currently being assembled is held. Nested paths work the same way, and the cartesian product two sibling collections produce is de-duplicated per parent:
const cursor = em.stream(Author, { populate: ['books.tags'], orderBy: { id: 'ASC' } });
Streaming row by row
mergeResults: false yields every row as its own entity instead: to-many collections contain at
most one item, and a root with several children comes back several times.
const stream = em.stream(Book, {
populate: ['tags'],
orderBy: { id: 'ASC' },
mergeResults: false,
});
Ordering, and what stays lazy
Without orderBy, a streamed query plans to a label scan and an expand — no blocking operator
anywhere, so the first entity arrives after the first batch:
+ProduceResults
+Projection
+OptionalExpand(All)
+NodeByLabelScan
ORDER BY adds a Sort, which is blocking: the server materializes and sorts the whole result
before sending the first record. Client memory is still bounded — that is what the fetch size
governs — but time to first entity is not. Two ways to keep it cheap:
- Order by an indexed property, so the planner can take the order from the index instead of
sorting (
Ordered byshows up in the plan without aSortoperator). - Order by the primary key, which is backed by the identity constraint the schema generator creates.
Rows of one root have to be adjacent for merging to work. Unsorted, they already are — an expansion
emits all of one root's matches before moving to the next. Sorted, the root's primary key is
appended to the ORDER BY automatically, so two roots with the same sort key cannot interleave and
split an entity in half.
limit and offset window the roots, not the rows an expansion multiplies them into: when a
to-many relation is populated, they are applied in a WITH before the expansion.
MATCH (this0:Author)
WITH this0 ORDER BY this0.id ASC SKIP 1 LIMIT 2
OPTIONAL MATCH (this0)-[this1:WROTE]->(this2:Book)
RETURN this0 AS node, this2 AS rel_0 ORDER BY this0.id ASC
Cancellation
Ending the loop is enough — break, return, throw, or cursor.close() all discard the rest of
the result on the server rather than draining it. An AbortSignal does the same from outside, and is
raced against each pull so a stream parked on a slow query notices immediately:
const controller = new AbortController();
const cursor = em.stream(Book, { signal: controller.signal });
setTimeout(() => controller.abort(), 1000);
The signal an EntityManager fork was created with (em.fork({ signal })) is honoured too.
Transactions
A stream opened inside em.transactional() runs on the transaction's own session, so it sees that
transaction's uncommitted writes:
await em.transactional(async (em) => {
em.create(Book, { id: 'B1', title: 'Draft', author });
await em.flush();
for await (const book of em.stream(Book, { where: { title: 'Draft' } })) {
// the uncommitted book is here
}
});
Keep the cursor inside the transaction: an open cursor holds the transaction open with it.
Raw Cypher
em.streamRaw() is em.run() without the array — same conversion of Neo4j integers and temporals,
one row at a time:
const cursor = em.streamRaw<{ name: string }>(
'MATCH (p:Person) WHERE p.age > $age RETURN p.name AS name',
{ age: 30 },
{ chunkSize: 500 },
);
for await (const row of cursor) {
console.log(row.name);
}
The query builder streams too, mapped the way execute() maps, or raw:
const cursor = em.createQueryBuilder(Movie, 'm')
.match()
.where('released', 1999)
.return(['title'])
.stream({ chunkSize: 200 });
const records = em.createQueryBuilder(Movie, 'm')
.match()
.return(['title'])
.stream({ rawResults: true }); // neo4j-driver Records, mapped by you
Virtual entities
Virtual entities stream as well. The expression callback receives a fourth argument saying whether
a stream was asked for, so one definition can answer both em.find() and em.stream() — and the
streaming branch may hand back a cursor of its own instead of a query to run:
const BookView = defineEntity({
name: 'BookView',
expression: (em: EntityManager, where, options, stream) => {
const cypher = 'MATCH (b:Book) RETURN { title: b.title, price: b.price } AS node';
return stream ? em.streamRaw(cypher, {}, { chunkSize: 500 }) : cypher;
},
properties: {
title: p.string(),
price: p.float(),
},
});
for await (const row of em.stream(BookView)) {
console.log(row.title);
}
When the expression hands back a cursor of its own like that, the stream is driven by that cursor,
so its stats are the ones describing the query — the wrapper em.stream() returns only counts what
it yielded. An expression that returns a string or a query builder is run by the driver as usual, and
its stats are complete.
Coming from Drivine
Drivine's cursor is the same idea, and the shapes map across directly:
| Drivine | Here |
|---|---|
CursorSpecification.batch | chunkSize — sent to the server as Bolt's PULL n, rather than re-running the statement with SKIP/LIMIT per page |
cursor[Symbol.asyncIterator]() | for await ... of cursor |
cursor.asStream({ transform }) | cursor.asStream({ transform }) |
cursor.close() | cursor.close() |
| — | cursor.stats, entity mapping, populate merging, transactions, AbortSignal |
The batching mechanism is the one real difference. Drivine's Neo4j cursor pages with SKIP/LIMIT
— its source calls that "a rudimentary placeholder for a pending implementation that will use the
driver's streaming capabilities" — so page n re-evaluates everything before it, and the result can
shift under a concurrent write. This driver keeps one server-side cursor open for the whole stream:
each batch continues where the last one stopped, and the read stays consistent for the transaction it
runs in.
Constraints
- Streamed entities are not managed. Identity holds within one returned entity graph, not across the stream, and nothing lands in the identity map.
- Relations are loaded through the joined expansion described above; the select-in strategy is not
used for streams, whatever
strategysays. mergeResultsneeds rows of one root to be adjacent. That holds for the queries the driver builds; it will not hold for a hand-written query that sorts the expansion by something unrelated.- Relationship-property payloads on a many-to-many are not populated in a stream — the target
entities are, exactly as with
em.find(). cache, cursor-based pagination (first/last/before/after) andoverfetchare not part ofStreamOptionsin core, and are not accepted here either.
Test coverage
tests/Neo4jStream.test.ts runs against a real Neo4j
(Testcontainers) and pins: unmanaged results, where/orderBy, to-one, to-many, many-to-many and
nested populate with the exact record-to-entity ratios, empty collections, mergeResults: false,
entity-wise limit/offset, chunkSize reaching the session as the fetch size, invalid chunk
sizes, early break releasing the session, AbortSignal, asStream(), streaming inside a
transaction, em.streamRaw(), qb.stream() mapped and raw, virtual entities, partial delivery
before a failing row, and reading the head of a five-million-row query without materializing it.
The unit tests cover the merge algorithm and the cursor's lifecycle without a database.