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.

Advanced190–240 minutesUNWIND batch-size 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

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.

01

Use parameterized list-of-map payloads with UNWIND for repeatable batch writes.

02

Distinguish rows-per-query, rows-per-transaction and concurrent-batch count.

03

Measure serialized payload bytes, transaction duration and p95/p99 rather than maximizing batch size blindly.

04

Keep batch writes idempotent with stable keys, uniqueness constraints and deterministic mutation semantics.

05

Choose application chunking versus CALL { … } IN TRANSACTIONS based on the ingestion path and transaction ownership.

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. 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 · one parameterized batch
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;
Python · managed transaction batch
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

Python · chunk + semaphore + payload evidence
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 · post-run reconciliation
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 · scoped benchmark cleanup
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. What does UNWIND change?
  2. Why can a larger batch reduce throughput?
  3. What makes the example retry-safe?
  4. Why record payload bytes?
  5. 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

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.