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.
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.
Define publisher/subscriber/demand/backpressure without treating “reactive” as a synonym for “asynchronous.”
Identify which current Neo4j driver paths expose a Reactive API and which do not.
Explain how Bolt result fetching and application demand cooperate to bound client buffering.
Use cancellation/limited consumption deliberately and preserve resource cleanup.
Decide when ordinary async iteration is simpler than a reactive pipeline.
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.
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 |
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
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
mkdir js-reactivecd js-reactivenpm init -ynpm install neo4j-driver@6.2.0 rxjs
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
- Is every async API reactive?
- Does Python driver 6.3 expose a Reactive Streams session?
- What does Python fetch_size=100 mean?
- Can backpressure fix a path-explosion query?
- 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
- Current Neo4j versions — Current server and 5.26 LTS release snapshot.
- Neo4j Python Driver Manual — Official Python driver guide used by the mandatory async lab.
- Python Driver 6.3 API — Current synchronous and asynchronous API contracts and configuration.
- Python async API — Async driver/session/result lifecycle, cancellation and concurrency rules.
- Python concurrency guide — AsyncGraphDatabase and concurrent workflow guidance.
- Python performance recommendations — Driver performance, routing and batching guidance.
- Cypher UNWIND — List-to-row semantics and ordering boundary.
- Transaction management — Server-side transaction limits and timeout concepts.
- JavaScript driver installation — Current full vs lite driver and Reactive API boundary.
- JavaScript driver npm package — Current 6.2.0 package and RxJS dependency.