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.
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.
Create one application-lifetime AsyncDriver and
one short-lived AsyncSession per concurrent
unit of work.
Bound coroutine fan-out explicitly and relate the limit to pool capacity, server concurrency and latency SLOs.
Use server-side transaction timeouts and client-side deadlines without confusing timeout with confirmed rollback.
Handle asyncio.CancelledError with
driver-supported cleanup and recognize ambiguous commit
outcomes.
Measure queue time, pool acquisition, database result timing and end-to-end latency separately.
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.
Reproducible AtlasMart setup
$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 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 |
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
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.
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
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()
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 |
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 |
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
- Why is an async driver not permission to create unbounded tasks?
- Can one AsyncSession be shared by many concurrent coroutines?
- What should happen when CancelledError interrupts driver I/O?
- Why record queue time separately from service time?
- 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
- 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.
- Python async cancellation — Documented cancellation and connection invalidation behavior.
- Python driver configuration — Pool, connection and retry configuration reference.