Skip to main content

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.

Applies to

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()
RecordsEvery record is buffered before anything is mappedPulled chunkSize at a time, mapped one entity at a time
Client memoryThe whole result, plus the entity graphOne batch, plus the entity currently being assembled
Time to first entityAfter the last recordAfter the first batch
Identity mapEntities are managedEntities come from throwaway forks — nothing is retained
To-many relationsA second, select-in query per relationExpanded in the same query and merged as rows arrive
Failure mid-queryNothing is returnedThe 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 by shows up in the plan without a Sort operator).
  • 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:

DrivineHere
CursorSpecification.batchchunkSize — 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 strategy says.
  • mergeResults needs 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) and overfetch are not part of StreamOptions in 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.