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.

Advanced210–270 minutesEnd-to-end throughput benchmark labNeo4j 2026.07.1 Community · Cypher 25Python driver 6.3.0 · asyncOptional JS driver 6.2.0 · Reactive APILast reviewed: September 2026

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.

01

Build a benchmark that records requests/sec and p50/p95/p99 with errors, not only average throughput.

02

Separate offered-load queue time, driver/pool time, server summary timing, result consumption, serialization and total request time.

03

Sweep concurrency and batch size independently and identify the first overload/saturation signals.

04

Correlate client metadata with server transactions/queries where the deployment permits.

05

Choose a safe operating envelope from correctness, tail latency, resource and fairness evidence.

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. 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

Python · benchmark input generator
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

Python · compact async harness
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
Summary timing caveat

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

Cypher · correlated transaction view where privileges allow
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 · remove benchmark-only products
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

  1. Why report RPS and p99 together?
  2. Is result_available_after the database query time seen by the user?
  3. What is saturation?
  4. Why benchmark sparse and dense graph shapes?
  5. 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

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.