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

Async Driver APIs, Concurrency Limits, Cancellation, Timeouts, and Resource Cleanup

Scale AtlasMart I/O concurrency without turning the connection pool or database into an accidental unbounded queue.

Advanced180–230 minutesAsync concurrency and cancellation labNeo4j 2026.07.1 Community · Cypher 25Python driver 6.3.0 · asyncOptional JS driver 6.2.0 · Reactive APILast reviewed: September 2026

Learning outcomes

AtlasMart launches an asynchronous order API. The first load test looks excellent at 20 concurrent requests, so an engineer replaces the limit with asyncio.gather() over every incoming item. Soon hundreds of coroutines wait for the same finite connection pool, cancellation tears down connections, and p99 latency rises while average throughput appears healthy. The mechanism is not “async is slow”; it is unbounded demand meeting bounded database and client resources.

01

Create one application-lifetime AsyncDriver and one short-lived AsyncSession per concurrent unit of work.

02

Bound coroutine fan-out explicitly and relate the limit to pool capacity, server concurrency and latency SLOs.

03

Use server-side transaction timeouts and client-side deadlines without confusing timeout with confirmed rollback.

04

Handle asyncio.CancelledError with driver-supported cleanup and recognize ambiguous commit outcomes.

05

Measure queue time, pool acquisition, database result timing and end-to-end latency separately.

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.

Reproducible AtlasMart setup

PowerShell · environment
$env:NEO4J_URI='bolt://127.0.0.1:7687'$env:NEO4J_USER='neo4j'$env:NEO4J_PASSWORD='atlasmart-course-2026'$env:NEO4J_DATABASE='neo4j'py -3 -m venv .venv.\.venv\Scripts\Activate.ps1python -m pip install --upgrade pippython -m pip install neo4j==6.3.0
Cypher · stable fixture and constraints
CYPHER 25CREATE CONSTRAINT customer_id IF NOT EXISTSFOR (c:Customer) REQUIRE c.customerId IS UNIQUE;CREATE CONSTRAINT product_id IF NOT EXISTSFOR (p:Product) REQUIRE p.productId IS UNIQUE;CREATE CONSTRAINT order_id IF NOT EXISTSFOR (o:Order) REQUIRE o.orderId IS UNIQUE;MERGE (c:Customer {customerId:'C-1001'})SET c.name='Mina Rahimi', c.tier='GOLD'MERGE (p1:Product {productId:'P-1001'})SET p1.name='Trail Camera', p1.category='Cameras', p1.price=129.90MERGE (p2:Product {productId:'P-2001'})SET p2.name='Smart Shelf Sensor', p2.category='Store IoT', p2.price=79.50MERGE (o:Order {orderId:'O-5001'})SET o.status='PAID', o.orderedAt=datetime('2026-09-08T16:30:00Z')MERGE (c)-[:PLACED]->(o)MERGE (o)-[r1:CONTAINS]->(p1)SET r1.quantity=1, r1.unitPrice=129.90MERGE (o)-[r2:CONTAINS]->(p2)SET r2.quantity=2, r2.unitPrice=79.50;

The fixture is intentionally tiny. Performance conclusions come from the synthetic workload generated by the harness, not from pretending two products are representative of production. The stable keys and uniqueness constraints make repeated batch-write experiments safe to reconcile.

1. Async changes waiting behavior, not database capacity

Concept Precise meaning Design consequence
Coroutine Python task that can suspend while waiting for I/O many coroutines can exist; that does not create infinite DB connections
AsyncDriver application-level concurrency-safe driver for coroutines; owns connection pool create once; do not close while work is still using it
AsyncSession lightweight transactional context; not safe for concurrent Tasks one task/workflow owns a session at a time
Connection pool finite reusable Bolt connections per host excess demand waits or hits acquisition timeout
Backpressure mechanism that slows/rejects upstream demand when downstream is saturated implement queue/concurrency limits instead of spawning forever
Timeout/deadline upper bound on waiting/execution classify where timeout occurred and reconcile ambiguous writes
Python · one maintained async driver
from __future__ import annotationsimport osfrom neo4j import AsyncGraphDatabase, RoutingControlclass AsyncAtlasMart:    def __init__(self) -> None:        self.database = os.getenv("NEO4J_DATABASE", "neo4j")        self.driver = AsyncGraphDatabase.driver(            os.getenv("NEO4J_URI", "bolt://127.0.0.1:7687"),            auth=(os.getenv("NEO4J_USER", "neo4j"),                  os.getenv("NEO4J_PASSWORD", "atlasmart-course-2026")),            max_connection_pool_size=20,            connection_acquisition_timeout=5.0,            connection_timeout=5.0,            max_transaction_retry_time=15.0,        )    async def verify(self) -> None:        await self.driver.verify_connectivity()    async def close(self) -> None:        await self.driver.close()    async def get_order(self, order_id: str) -> dict | None:        records, summary, _ = await self.driver.execute_query(            """            CYPHER 25            MATCH (o:Order {orderId:$orderId})            OPTIONAL MATCH (o)-[r:CONTAINS]->(p:Product)            WITH o, collect({productId:p.productId, quantity:r.quantity}) AS lines            RETURN {orderId:o.orderId, status:o.status, lines:lines} AS order            """,            orderId=order_id,            database_=self.database,            routing_=RoutingControl.READ,        )        return records[0]["order"] if records else None

2. Wrong: unbounded task fan-out

Python · deliberately unsafe fan-out
async def wrong(driver, order_ids):    # Every item becomes an immediately scheduled Task.    # A finite pool/server now becomes the queue by accident.    return await asyncio.gather(*[        driver.execute_query(            "CYPHER 25 MATCH (o:Order {orderId:$id}) RETURN o.status",            id=oid, database_="neo4j"        )        for oid in order_ids    ])

This can create a large client-side working set and move queueing into driver pool acquisition. A larger pool is not automatically a repair: it can merely shift overload into the database and worsen lock, CPU or memory pressure.

Python · bounded concurrency with one operation per permit
import asynciofrom time import perf_counterasync def bounded_read(store, order_ids, concurrency=8):    gate = asyncio.Semaphore(concurrency)    async def one(order_id):        queued_at = perf_counter()        async with gate:            started = perf_counter()            order = await store.get_order(order_id)            finished = perf_counter()        return {            "orderId": order_id,            "queue_ms": (started - queued_at) * 1000,            "service_ms": (finished - started) * 1000,            "found": order is not None,        }    return await asyncio.gather(*(one(x) for x in order_ids))

Run a sweep such as 1, 2, 4, 8, 16 and 32—not because those values are universally correct, but because the saturation point is empirical. Keep request mix and fixture scale constant while changing one variable.

3. Session ownership and cancellation

Python · cancellation-safe session ownership
import asyncioasync def read_lines(driver, order_id):    session = driver.session(database="neo4j", fetch_size=100)    try:        result = await session.run(            """            CYPHER 25            MATCH (:Order {orderId:$id})-[r:CONTAINS]->(p:Product)            RETURN p.productId AS productId, r.quantity AS quantity            ORDER BY productId            """,            id=order_id,        )        return [record.data() async for record in result]    except asyncio.CancelledError:        # Documented purpose: forcefully close the held connection/work.        session.cancel()        raise    finally:        await session.close()
Cancellation is not a commit oracle

The Python driver documentation states that cancellation may close the connection and there is no guarantee that a server-side commit did or did not finish before the cancellation became visible to the client. Never turn “client saw CancelledError” into “write definitely rolled back.” Stable operation IDs and reconciliation are the safer contract.

4. Timeouts have layers

Layer Example What it bounds What it does not prove
Pool acquisition connection_acquisition_timeout waiting for/creating a pooled connection database query duration after acquisition
Socket connect connection_timeout network connection setup query execution
Transaction/query neo4j.Query(..., timeout=...) or transaction timeout server-side transaction execution budget HTTP/API request total time
Application deadline asyncio.timeout() around the full operation client business-request budget definite server rollback after an ambiguous disconnect
Python · server query timeout + outer request deadline
import asynciofrom neo4j import Queryasync def timed_lookup(driver, order_id):    query = Query(        "CYPHER 25 MATCH (o:Order {orderId:$id}) RETURN o.status AS status",        timeout=2.0,        metadata={"app":"atlasmart", "op":"order-status"},    )    async with asyncio.timeout(3.0):        async with driver.session(database="neo4j") as session:            result = await session.run(query, id=order_id)            record = await result.single()            await result.consume()            return None if record is None else record["status"]

5. Observable experiment

Metric Where to measure Why it matters
arrival → permit wait application harness shows overload before DB call
permit → result available application + result summary service/DB/network path
records consumed / result bytes application large responses can dominate serialization/memory
pool acquisition failures driver exception/logging finite connection budget is saturated
active transactions SHOW TRANSACTIONS where supported server concurrency, duration and metadata
CPU/page cache/I/O server monitoring available to your edition/deployment distinguishes client queueing from server saturation
p50/p95/p99 + errors benchmark harness tail behavior and overload quality
Cypher · transaction metadata inspection
SHOW TRANSACTIONSYIELD transactionId, currentQuery, elapsedTime, metaDataWHERE metaData.app = 'atlasmart'RETURN transactionId, elapsedTime, metaData;

Transaction-listing privileges and available columns are deployment/edition/security dependent. If your Community user cannot view other transactions, keep the client-side correlation metadata and use the evidence available to that deployment rather than weakening security solely for a lab.

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. Why is an async driver not permission to create unbounded tasks?
  2. Can one AsyncSession be shared by many concurrent coroutines?
  3. What should happen when CancelledError interrupts driver I/O?
  4. Why record queue time separately from service time?
  5. What determines the concurrency limit?
Review the answers

1. The driver/database still have finite pools, transaction capacity, CPU, memory and I/O; unbounded Tasks merely create uncontrolled queueing.

2. No. Sessions are not concurrency-safe; separate concurrent units of work need separate sessions.

3. Use the driver-supported cancel/close path and treat any write outcome as potentially ambiguous until reconciled.

4. Otherwise rising overload may be hidden inside total latency and falsely blamed on query execution.

5. Measured saturation/latency/error behavior plus pool/server/tenant budgets; not a universal constant.

Summary and next step

Async improves utilization while work waits on I/O, but reliability comes from ownership, bounded concurrency, explicit deadlines and cleanup. Lesson 2 changes the unit of work from one row/request to a measured batch.

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.