Chapter 14 · Asynchronous, Reactive, Batch, and High-Throughput Application Patterns

Reactive Streams Concepts, Demand, Backpressure, and When Reactive Access Helps

Distinguish asynchronous execution from true reactive demand/backpressure and use each only where it solves an observable flow-control problem.

Advanced160–210 minutesReactive/backpressure comparison labNeo4j 2026.07.1 Community · Cypher 25Python driver 6.3.0 · asyncOptional JS driver 6.2.0 · Reactive APILast reviewed: September 2026

Learning outcomes

An AtlasMart export endpoint may return a very large graph-shaped result to a downstream analytics service. An async API can avoid blocking threads, but it does not automatically express how many records the downstream consumer is prepared to accept. Reactive Streams adds a demand contract: producers should not overwhelm consumers that have requested less work.

01

Define publisher/subscriber/demand/backpressure without treating “reactive” as a synonym for “asynchronous.”

02

Identify which current Neo4j driver paths expose a Reactive API and which do not.

03

Explain how Bolt result fetching and application demand cooperate to bound client buffering.

04

Use cancellation/limited consumption deliberately and preserve resource cleanup.

05

Decide when ordinary async iteration is simpler than a reactive pipeline.

Chapter 14 baseline · reviewed 9 September 2026

The mandatory lab continues Neo4j Community 2026.07.1, database neo4j, explicit CYPHER 25 where language-version behavior matters, container atlasmart-neo4j, loopback Bolt bolt://127.0.0.1:7687, authentication neo4j/atlasmart-course-2026, and no TLS on the disposable loopback-only instance. The primary client is the official Python package neo4j 6.3.0 on Python 3.10–3.14. The optional Reactive Streams exercise uses the official JavaScript driver neo4j-driver 6.2.0; its lite package deliberately omits the Reactive API. Neo4j 5.26.30 remains the LTS comparison line.

Evidence boundary

This generation environment does not connect to your AtlasMart Neo4j container, so no throughput, p95/p99, pool-queue, CPU, page-cache, I/O, cancellation, GC, routing or network measurements are fabricated. The chapter provides deterministic fixtures, runnable harnesses, measurement columns and acceptance invariants. Values such as concurrency, fetch size and batch size are experimental variables—not universal recommendations.

1. Async and reactive answer different questions

Model Question it answers Typical abstraction
synchronous does this call block the calling thread? blocking iterator/call
asynchronous can work suspend while I/O is pending? Future/Promise/coroutine
reactive streams how much data/work has the downstream consumer requested? Publisher → Subscriber demand + cancellation
batching how many logical records share one request/transaction? list parameter / chunk

Reactive access helps most when the result is large, consumers are variably paced, or a broader reactive stack already exists. For a small bounded API response, adding reactive operators can increase complexity without improving the database plan.

2. Current driver support boundary

Driver path Current chapter use Backpressure surface
Python neo4j 6.3.0 mandatory async lab lazy AsyncResult, async iteration and fetch(n); not a Reactive Streams API
JavaScript neo4j-driver 6.2.0 optional concrete reactive example full package includes Reactive API / RxJS dependency
JavaScript neo4j-driver-lite 6.2.0 comparison same main capabilities except reactive sessions/API
other official drivers conceptual only here verify current language-specific manual before adopting a reactive abstraction
Version discipline

The JavaScript full-vs-lite split is current at generation time. Do not infer reactive support from “official driver” generically; choose the actual language/package/version and check its current manual.

3. Python: bounded pull without a reactive API

Python · consume in explicit windows
async def process_in_windows(driver, customer_id, window=100):    async with driver.session(database="neo4j", fetch_size=window) as session:        result = await session.run(            """            CYPHER 25            MATCH (:Customer {customerId:$id})-[:PLACED]->(o:Order)            RETURN o.orderId AS orderId, o.status AS status            ORDER BY o.orderId            """,            id=customer_id,        )        while True:            records = await result.fetch(window)            if not records:                break            await send_downstream([r.data() for r in records])        await result.consume()

The session fetch_size controls how many records the driver requests from Neo4j per batch (default 1000 in Python 6.3), but it is not a server-side LIMIT, nor does it cap the database memory required by an earlier Sort or aggregation operator. It controls record transfer/buffering behavior.

4. Optional JavaScript Reactive API example

PowerShell · install the full driver, not lite
mkdir js-reactivecd js-reactivenpm init -ynpm install neo4j-driver@6.2.0 rxjs
JavaScript · reactive session with bounded downstream windows
import neo4j from 'neo4j-driver'import { bufferCount, concatMap } from 'rxjs/operators'const driver = neo4j.driver(  'bolt://127.0.0.1:7687',  neo4j.auth.basic('neo4j', 'atlasmart-course-2026'))const session = driver.rxSession({ database: 'neo4j' })const subscription = session.run(`  CYPHER 25  MATCH (:Customer {customerId:$id})-[:PLACED]->(o:Order)  RETURN o.orderId AS orderId, o.status AS status  ORDER BY o.orderId`, { id: 'C-1001' }).records().pipe(  bufferCount(100),  // concatMap preserves one downstream batch at a time.  concatMap(records => persistBatch(records.map(r => r.toObject())))).subscribe({  error: async err => { console.error(err); await session.close(); await driver.close() },  complete: async () => { await session.close(); await driver.close() }})

This is an optional learning path. Exact RxJS operator behavior is part of the application stack, not a Neo4j server guarantee. Backpressure cannot repair an unbounded Cypher traversal that already generates an enormous server-side intermediate result; query cardinality still matters.

5. Wrong: reactive syntax with unbounded work behind it

Mistake Why reactive does not save it Repair
unbounded variable-length MATCH server may generate huge work before records reach client bound/anchor traversal and inspect plan first
consumer buffers every record to an array downstream discards backpressure benefit process bounded windows / stream to sink
slow subscriber with no queue policy latency/retention can grow indefinitely bounded buffer, timeout/cancellation and load shedding policy
reactive API for a 10-row endpoint extra abstraction with no material flow-control need ordinary async/eager API may be clearer

6. Evidence: demand, bytes and completion

Signal Interpretation
records requested/consumed per window consumer demand shape
time between windows downstream processing pressure
process RSS/heap whether records accumulate client-side
result bytes network/serialization load
query result_available_after vs full consume server availability vs end-to-end consumption
cancel/complete cleanup whether session/driver resources are released

Production judgment

Review area Decision evidence
Graph/workload fit fan-out/degree, result cardinality and traversal shape; avoid making driver concurrency compensate for a poor model
Correctness stable keys, constraints, transaction boundaries and idempotent retry semantics
Latency p50/p95/p99 plus timeout/error/cancellation rate—not average latency alone
Client resources event-loop queue, connection-pool wait, result bytes, process RSS/GC and serialization time
Server resources CPU, page cache, store I/O, active/queued transactions, locks and query memory where observable
Security parameterized Cypher, tenant authorization, TLS/auth policy and bounded user-controlled result sizes
Recovery/rollback batch checkpoint/reconciliation, stable IDs, bounded blast radius and restartable jobs
Edition/topology Community lab is single-instance; Aura/Enterprise routing and managed-service limits require separate evidence

Check your understanding

  1. Is every async API reactive?
  2. Does Python driver 6.3 expose a Reactive Streams session?
  3. What does Python fetch_size=100 mean?
  4. Can backpressure fix a path-explosion query?
  5. When should you prefer ordinary async iteration?
Review the answers

1. No. Async concerns suspension/concurrency; reactive streams add explicit downstream demand/backpressure semantics.

2. No. It exposes async/lazy result APIs; this chapter uses the JavaScript full driver for a concrete Reactive API example.

3. It controls driver record-fetch batches, not Cypher LIMIT or server-side intermediate operator memory.

4. No. It can control result transfer/consumption; server query shape and cardinality must still be fixed.

5. When result volumes and processing are bounded and the reactive abstraction would add more complexity than value.

Summary and next step

Reactive demand is a client flow-control contract, not a substitute for bounded Cypher or capacity planning. Lesson 4 attacks another common overload source: N+1 and chatty graph access.

Authoritative references

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.