Chapter 14 · Asynchronous, Reactive, Batch, and High-Throughput Application Patterns
Benchmark Driver-Side Throughput and Separate Database Time from Network/Serialization/Client Time
Build an evidence-driven AtlasMart throughput benchmark that separates admission, driver, database, transfer and serialization costs and selects a safe operating envelope.
Learning outcomes
The AtlasMart team reports “Neo4j takes 180 ms” because the HTTP endpoint takes 180 ms. But 25 ms is spent waiting for a concurrency permit, 8 ms acquiring a connection, the server reports results available after 22 ms, 70 ms is spent consuming/serializing 2 MB, and the rest is network/framework time. Without a stage model, optimization becomes guesswork.
Build a benchmark that records requests/sec and p50/p95/p99 with errors, not only average throughput.
Separate offered-load queue time, driver/pool time, server summary timing, result consumption, serialization and total request time.
Sweep concurrency and batch size independently and identify the first overload/saturation signals.
Correlate client metadata with server transactions/queries where the deployment permits.
Choose a safe operating envelope from correctness, tail latency, resource and fairness evidence.
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. Define the timing model before collecting numbers
| Stage | Clock boundary | Representative signal |
|---|---|---|
| queue/admission | request arrival → semaphore acquired | queue_ms |
| driver/pool/network start | permit → query submitted/header available | driver logs + application clock |
| server availability |
driver summary result_available_after
|
server-reported result availability; not whole request |
| record transfer/consume | first/header → all required records consumed | application clock, result count/bytes |
| serialization | DTO → HTTP/JSON bytes | application clock + byte count |
| end-to-end | arrival → response complete | API/request metric |
Some stages overlap and driver/server summaries are not a distributed tracing system. The purpose is a useful causal decomposition, not pretending every millisecond can be perfectly attributed without clock-synchronized tracing.
2. Generate deterministic synthetic data
def make_products(n=5000): return [ { "productId": f"P-BENCH-{i:06d}", "name": f"Benchmark Product {i}", "category": "Benchmark", "price": 10.0 + (i % 100) / 10.0, } for i in range(n) ]
Keep the random/dataset seed, record count and graph density fixed between variants. If a benchmark changes the data distribution while changing concurrency, it cannot isolate causality.
3. Benchmark harness with percentiles and stage metrics
import asyncio, json, math, statistics, timefrom neo4j import AsyncGraphDatabase, QueryREAD = Query("""CYPHER 25MATCH (p:Product)WHERE p.category=$categoryRETURN p.productId AS productId, p.name AS name, p.price AS priceORDER BY p.productIdLIMIT $limit""", timeout=5.0, metadata={"app":"atlasmart-bench", "op":"product-feed"})def percentile(values, q): xs = sorted(values) if not xs: return None i = min(len(xs)-1, max(0, math.ceil(q * len(xs)) - 1)) return xs[i]async def run_once(driver, gate, request_id, limit=100): arrival = time.perf_counter() async with gate: admitted = time.perf_counter() records, summary, _ = await driver.execute_query( READ, category="Benchmark", limit=limit, database_="neo4j") consumed = time.perf_counter() payload = json.dumps([r.data() for r in records], separators=(",", ":")).encode() done = time.perf_counter() return { "id": request_id, "queue_ms": (admitted-arrival)*1000, "driver_to_consumed_ms": (consumed-admitted)*1000, "serialize_ms": (done-consumed)*1000, "total_ms": (done-arrival)*1000, "bytes": len(payload), "server_available_ms": summary.result_available_after, "server_consumed_ms": summary.result_consumed_after, }async def benchmark(driver, requests=200, concurrency=8): gate = asyncio.Semaphore(concurrency) t0 = time.perf_counter() rows = await asyncio.gather(*(run_once(driver, gate, i) for i in range(requests))) seconds = time.perf_counter() - t0 totals = [x["total_ms"] for x in rows] return { "requests": requests, "concurrency": concurrency, "rps": requests / seconds, "p50_ms": percentile(totals, .50), "p95_ms": percentile(totals, .95), "p99_ms": percentile(totals, .99), "mean_queue_ms": statistics.fmean(x["queue_ms"] for x in rows), "total_bytes": sum(x["bytes"] for x in rows), }, rows
result_available_after and
result_consumed_after are useful server/driver
summary timings, but they are not the same thing as HTTP
latency, connection-pool queue time, Python serialization time
or an end-to-end distributed trace. Keep them as separate
columns.
4. Sweep one variable at a time
| Experiment | Fixed | Sweep | Stop/flag when |
|---|---|---|---|
| read concurrency | dataset, query, page size | 1,2,4,8,16,32… | p95/p99 or error/queue rises disproportionately |
| page/result size | concurrency, dataset | 10,100,500,1000… | bytes/memory/consume time dominate |
| write batch size | concurrency, input set | 1,10,100,500… | transaction duration/p99/retries worsen |
| batch concurrency | batch size, input set | 1,2,4,8… | pool/server contention or fairness degrades |
The safe operating envelope is normally below the cliff, with margin for noisy neighbors, GC, backups/checkpoints, failover, larger tenants and temporary graph-density changes. “Maximum RPS before failure” is not a production target.
5. Observe server state without overfitting to one edition
SHOW TRANSACTIONSYIELD transactionId, currentQuery, elapsedTime, allocatedBytes, metaDataWHERE metaData.app = 'atlasmart-bench'RETURN transactionId, elapsedTime, allocatedBytes, metaData;
Available columns and privileges can vary by version/edition/security configuration. On Community, capture what is available plus host/container CPU, process memory, I/O and page-cache-relevant metrics. On Aura, use the observability exposed by the service tier rather than assuming access to self-managed files or OS counters.
6. Diagnose common benchmark lies
| Misleading claim | What is missing | Repair |
|---|---|---|
| “async doubled performance” | maybe offered concurrency simply increased | compare same offered load and tails/resource use |
| “batch 1000 is best” | one dataset/payload shape | sweep payload bytes and graph/write contention |
| “database took 180 ms” | client queue/network/serialization not separated | record stage timings and server summary separately |
| “p95 is fine” | p99/error/cancellation/fairness absent | report distribution + failures + tenant mix |
| “no errors at 500 RPS” | queue may grow without bound | record queue depth/wait and fixed-duration/closed-loop behavior |
| “warm benchmark proves cold-start SLO” | cache state ignored | run explicit warmup and separately document cold/cache-changing scenarios |
7. Select the operating envelope
| Acceptance area | Example criterion style—not a universal value |
|---|---|
| correctness | 0 invariant violations; stable unique IDs after retries/replays |
| tail latency | p95/p99 remain inside your service SLO across representative tenant/result shapes |
| overload | bounded queue; explicit rejection/timeout rather than unbounded memory growth |
| pool | acquisition wait/errors stay controlled with headroom |
| server | CPU/I/O/transaction/memory signals have recovery margin |
| fairness | large tenant/query does not starve small requests beyond policy |
| failure | cancellation/retry/ambiguous-outcome tests reconcile correctly |
8. Cleanup/reset
CYPHER 25MATCH (p:Product)WHERE p.productId STARTS WITH 'P-BENCH-'DETACH DELETE p;
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 report RPS and p99 together?
- Is result_available_after the database query time seen by the user?
- What is saturation?
- Why benchmark sparse and dense graph shapes?
- What makes a safe operating envelope production-worthy?
Review the answers
1. Throughput can improve while the slowest requests become unacceptable; both capacity and tail quality matter.
2. No. It is one summary timing and excludes/overlaps other client/network/serialization stages.
3. The region where more offered concurrency stops giving proportional useful throughput and causes queue/tail/errors/resource pressure to rise.
4. The same query text can expand very differently with degree/cardinality skew.
5. Correctness plus latency/error/fairness/resource headroom under representative workload and failure conditions—not the maximum observed RPS.
Summary and next step
High-throughput Neo4j application design is a closed-loop capacity discipline: bound demand, shape results, batch deliberately, measure every stage, find saturation, and operate with headroom. Chapter 15 can now move from single-instance application throughput into cluster architecture and availability without confusing more client concurrency with high availability.
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 result summary — ResultSummary timing/counter API.
- Connection management — Server-side connection visibility and client user-agent evidence.
- Query planning/profile — Server operator evidence that complements driver benchmarks.