Chapter 14 · Asynchronous, Reactive, Batch, and High-Throughput Application Patterns
Batch Writes with UNWIND, Transaction Chunking, Payload Size, and Throughput/Latency Tradeoffs
Batch AtlasMart writes with explicit payload, transaction and concurrency boundaries, then prove restartability and invariants.
Learning outcomes
AtlasMart receives 20,000 product updates. Sending 20,000 separate transactions wastes protocol/transaction overhead; sending one 20,000-row payload in one giant transaction may create a large network message, long lock lifetime, high transaction memory and expensive retries. This lesson treats batch size and transaction chunking as measured control variables.
Use parameterized list-of-map payloads with
UNWIND for repeatable batch writes.
Distinguish rows-per-query, rows-per-transaction and concurrent-batch count.
Measure serialized payload bytes, transaction duration and p95/p99 rather than maximizing batch size blindly.
Keep batch writes idempotent with stable keys, uniqueness constraints and deterministic mutation semantics.
Choose application chunking versus
CALL { … } IN TRANSACTIONS based on the
ingestion path and transaction ownership.
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. Three batching dimensions are independent
| Dimension | Example | Primary pressure |
|---|---|---|
| rows in one Bolt parameter payload | 100 product maps in $rows |
network/serialization/client memory |
| rows in one database transaction | same 100 rows under one managed transaction | locks, transaction state/memory, retry blast radius |
| concurrent batches | 4 transactions each processing 100 rows | pool, CPU, lock contention, page cache/I/O |
“Batch size = 500” is incomplete unless all three dimensions are stated. Keep them in benchmark output so a later change is reproducible.
2. Idempotent UNWIND upsert
CYPHER 25UNWIND $rows AS rowMERGE (p:Product {productId: row.productId})ON CREATE SET p.createdAt = datetime()SET p.name = row.name, p.category = row.category, p.price = toFloat(row.price), p.updatedAt = datetime()RETURN count(*) AS inputRows;
from neo4j import QueryUPSERT_PRODUCTS = Query("""CYPHER 25UNWIND $rows AS rowMERGE (p:Product {productId:row.productId})ON CREATE SET p.createdAt = datetime()SET p.name=row.name, p.category=row.category, p.price=toFloat(row.price), p.updatedAt=datetime()RETURN count(*) AS inputRows""", timeout=10.0, metadata={"app":"atlasmart", "op":"batch-products"})async def upsert_batch(tx, rows): result = await tx.run(UPSERT_PRODUCTS, rows=rows) record = await result.single() summary = await result.consume() return record["inputRows"], summary.countersasync def write_batch(driver, rows): async with driver.session(database="neo4j") as session: return await session.execute_write(upsert_batch, rows)
Because a managed transaction callback may retry, the mutation
must remain safe when repeated. The stable
productId uniqueness contract plus
MERGE identity keeps the node count stable;
timestamp values may intentionally change on each successful
retry/execution, so do not call every property bit-for-bit
idempotent unless that is your actual contract.
3. Chunk in the application and bound concurrent batches
import asyncio, jsonfrom time import perf_counterdef chunks(rows, size): for i in range(0, len(rows), size): yield rows[i:i+size]async def ingest(driver, rows, batch_size=100, concurrency=4): gate = asyncio.Semaphore(concurrency) async def one(batch_no, batch): payload_bytes = len(json.dumps(batch, separators=(",", ":")).encode()) async with gate: t0 = perf_counter() input_rows, counters = await write_batch(driver, batch) ms = (perf_counter() - t0) * 1000 return { "batch": batch_no, "rows": input_rows, "payload_bytes": payload_bytes, "elapsed_ms": ms, "nodes_created": counters.nodes_created, } jobs = [one(i, b) for i, b in enumerate(chunks(rows, batch_size), 1)] return await asyncio.gather(*jobs)
A safe experiment holds total input and concurrency constant while sweeping batch size, then separately holds batch size constant while sweeping concurrency. Otherwise two independent causes change at once.
4. Wrong extremes: row-at-a-time and giant transaction
| Pattern | Why it looks attractive | Failure mode | Repair |
|---|---|---|---|
| one transaction per row | simple code | protocol/transaction overhead; poor throughput | parameterized UNWIND batches |
one enormous $rows list |
few round trips | large serialization/memory; long locks; expensive retry | bounded chunks with measured payload/transaction duration |
| hundreds of concurrent batches | high offered load | pool queue, contention, server overload, p99 collapse | bounded concurrency + saturation test |
CREATE in retryable batch |
fast first run | duplicates on replay/retry |
stable keys + constraints +
MERGE/idempotent design
|
5. Application chunking vs CALL IN TRANSACTIONS
| Mechanism | Who owns chunking? | Good fit | Important boundary |
|---|---|---|---|
| application chunks + managed tx | client application | parameter/list ingestion, service jobs, explicit retry/checkpoint control | each managed callback may retry; keep side effects safe |
LOAD CSV ... CALL { … } IN TRANSACTIONS
|
Cypher/server statement | CSV import and server-managed transaction batching | must be run as an auto-commit statement because the clause manages inner transactions |
| single huge transaction | nobody | rarely appropriate for bulk migration | large failure/lock/memory blast radius |
Do not cargo-cult the import chapter into application batching.
If data is already parsed in your service, list parameters plus
UNWIND usually give the application clearer
checkpointing and payload control.
6. Benchmark matrix and invariants
| Run | Batch size | Batch concurrency | Capture |
|---|---|---|---|
| A | 1 | 1 | baseline transaction overhead |
| B | 10 | 1 | payload/transaction improvement |
| C | 100 | 1 | larger-batch efficiency and p95/p99 |
| D | 100 | 4 | throughput vs contention |
| E | 500 | 4 | whether payload/transaction duration now hurts tails |
CYPHER 25MATCH (p:Product)WHERE p.productId STARTS WITH 'P-BENCH-'WITH count(p) AS nodes, collect(p.productId) AS idsRETURN nodes, size(ids) AS idsObserved, size(ids) = size(reduce(s=[], x IN ids | CASE WHEN x IN s THEN s ELSE s + x END)) AS idsUnique;
Expected invariant: after rerunning the exact same generated
dataset, the count of P-BENCH-* product nodes stays
equal to the number of distinct input keys. Command success
alone is not reconciliation.
7. 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
- What does UNWIND change?
- Why can a larger batch reduce throughput?
- What makes the example retry-safe?
- Why record payload bytes?
- When is CALL IN TRANSACTIONS different from application chunking?
Review the answers
1. It turns each list element into a row; it does not by itself define transaction chunking or guarantee row order.
2. Serialization, transaction memory, lock duration, retry cost or server saturation can outweigh fewer round trips.
3. Stable product IDs, a uniqueness constraint and a match-or-create mutation keyed only by the stable identity.
4. Rows alone hide the fact that one row can be much larger than another and network/client memory costs scale with bytes.
5. The Cypher clause owns inner transaction batching inside an auto-commit statement; application chunking owns payloads/checkpoints and can use managed transactions.
Summary and next step
Batching is a three-dimensional control problem: payload size, transaction size and concurrency. Lesson 3 generalizes bounded demand into a true Reactive Streams mental model.
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 advanced queries — LOAD CSV / CALL IN TRANSACTIONS ownership and transaction configuration.
- Python driver performance — Batch data creation and driver performance guidance.