Retry and Idempotency for Graph Writes
Retrying a failed write is only safe if two things are true, and most ingestion code checks neither. The write has to be idempotent, or a retry after a commit the client never saw about doubles the data; and the error has to be transient, or the retry is a guaranteed second failure that consumes a connection and delays the batches behind it. Get the first wrong and a deadlock retry silently duplicates relationships. Get the second wrong and a syntax error becomes five identical syntax errors with exponential backoff between them. This page separates the two decisions, and adds the third that stops a recovering database from being knocked over by its own clients.
Prerequisites & Versions
Standard driver exceptions; the idempotency is a property of the Cypher.
| Requirement | Minimum version | Install |
|---|---|---|
| Python | 3.11 | — |
| neo4j (async driver) | 5.20 | pip install "neo4j>=5.20" |
| Neo4j Server | 5.15 | uniqueness constraints |
Implementation
import asyncio
import random
from dataclasses import dataclass
from neo4j import AsyncGraphDatabase
from neo4j.exceptions import (
ClientError,
DatabaseError,
Neo4jError,
ServiceUnavailable,
SessionExpired,
TransientError,
)
# Idempotent by construction: MERGE on a key that identifies the element, and
# SET only properties derived from the payload. Re-running this with the same
# batch converges to the same graph, so a retry after an unseen commit is a
# no-op rather than a duplicate.
UPSERT = """
UNWIND $batch AS row
MERGE (src:Junction {id: row.src_id})
ON CREATE SET src.location = point({latitude: row.src_lat, longitude: row.src_lon})
MERGE (tgt:Junction {id: row.tgt_id})
ON CREATE SET tgt.location = point({latitude: row.tgt_lat, longitude: row.tgt_lon})
MERGE (src)-[s:SEGMENT {id: row.edge_id}]->(tgt)
SET s.length_m = row.length_m, s.drive_s = row.drive_s
RETURN count(s) AS written
"""
@dataclass(frozen=True)
class RetryPolicy:
max_attempts: int = 5
base_delay_s: float = 0.2
max_delay_s: float = 8.0
def delay_for(self, attempt: int) -> float:
"""Exponential backoff with full jitter.
The jitter is not decoration. Without it, every client that failed
during the same blip retries at the same instant, and the database
recovers into a synchronised thundering herd that knocks it over again.
"""
ceiling = min(self.max_delay_s, self.base_delay_s * 2 ** attempt)
return random.uniform(0, ceiling)
def is_retryable(exc: BaseException) -> bool:
"""Transient means 'the same request might succeed later'.
Deadlocks, leader elections and acquisition timeouts qualify. A syntax
error, a constraint violation or a type mismatch will fail identically
every time, and retrying them wastes a connection slot that a batch which
could succeed is waiting for.
"""
if isinstance(exc, (TransientError, ServiceUnavailable, SessionExpired)):
return True
if isinstance(exc, ClientError):
return False # our bug — schema, syntax, constraint
if isinstance(exc, DatabaseError):
return False # server-side fault a retry will not fix
return False
class ResilientWriter:
def __init__(self, uri: str, auth: tuple[str, str],
policy: RetryPolicy | None = None, concurrency: int = 16) -> None:
self._driver = AsyncGraphDatabase.driver(
uri, auth=auth, max_connection_pool_size=concurrency * 2
)
self._policy = policy or RetryPolicy()
self._semaphore = asyncio.Semaphore(concurrency)
async def close(self) -> None:
await self._driver.close()
async def write(self, batch: list[dict]) -> int:
last: BaseException | None = None
for attempt in range(self._policy.max_attempts):
try:
async with self._semaphore:
async with self._driver.session() as session:
result = await session.run(UPSERT, batch=batch)
record = await result.single()
return int(record["written"])
except Neo4jError as exc:
last = exc
if not is_retryable(exc):
# Fail fast and loudly. A ClientError retried five times is
# five identical failures and eight seconds of delay.
raise
if attempt == self._policy.max_attempts - 1:
break
await asyncio.sleep(self._policy.delay_for(attempt))
raise RuntimeError(
f"batch of {len(batch)} failed after {self._policy.max_attempts} "
f"attempts: {last}"
) from last
How It Works
Idempotency comes from MERGE on a stable key, not from the retry logic. The retry can only be safe if repeating the write is safe, and that is a property of the Cypher. MERGE on edge_id converges; CREATE does not, and no amount of care in the client makes it. The ON CREATE SET for the coordinate is deliberate too: it writes the location when the node is first seen and leaves it alone afterwards, so a retry cannot overwrite a corrected coordinate with the original one.
Classification decides whether to retry at all. TransientError covers the cases where the same request genuinely might succeed later — a deadlock between two batches touching the same high-degree junction, a leader election, a momentarily exhausted pool. ClientError covers the cases where it will not: a typo in the Cypher, a constraint violation, a parameter of the wrong type. Retrying the second category is worse than useless, because it occupies a connection and delays batches that would have succeeded.
Full jitter is what prevents the second outage. A blip that fails a hundred concurrent batches gives a hundred clients the same backoff schedule, and without jitter they all return at the same instant to a database that has just started recovering. Randomising each delay across the whole window spreads the return, and it is the difference between a recovery and a sawtooth of repeated collapses.
Common Failure Patterns
1. Catching Exception and retrying. The broadest possible net, and it turns a schema bug into a slow schema bug. A Cypher syntax error retried with backoff takes eight seconds to report something that was knowable immediately, and in a batch loop it does that for every batch. Catch the driver’s own exception hierarchy and let everything else propagate.
# WRONG: a typo in the Cypher is now a five-attempt, eight-second failure.
except Exception:
await asyncio.sleep(backoff)
# RIGHT: only the errors where a later attempt could genuinely differ.
except (TransientError, ServiceUnavailable, SessionExpired):
await asyncio.sleep(policy.delay_for(attempt))
2. Assuming a failed write did not commit. A connection dropped between commit and acknowledgement leaves the client believing the write failed and the database holding it. That is precisely the scenario idempotency exists for, and it is why CREATE in an ingestion path is a defect rather than a style choice — the duplicate it produces is invisible until something counts relationships.
3. Retrying the batch rather than the transaction. If a batch is split across multiple statements without a transaction boundary, a retry re-runs the whole batch including the parts that already committed. With idempotent statements that is harmless; with any non-idempotent step it compounds. Keep the retryable unit and the transactional unit the same thing.
Performance Notes
The cost of getting classification wrong is easy to quantify. A batch that fails permanently and is retried $n$ times with exponential backoff occupies a connection for
$$T_{\text{wasted}} \approx \sum_{i=0}^{n-1} \frac{b \cdot 2^{i}}{2}$$
which for five attempts at a 200 ms base is roughly three seconds of a pool slot, per doomed batch. On an ingestion run where a schema problem affects every batch, that is the entire run spent waiting to fail.
The concurrency interaction is worth stating explicitly: retries consume the same semaphore permits as first attempts, so a burst of transient failures reduces effective throughput exactly when the system is already struggling. That is the correct behaviour — it is backpressure — but it means the retry budget and the concurrency limit have to be chosen together. A large retry budget with a tight semaphore converts a brief blip into a long stall, because the retrying batches hold the permits that new work needs.
Deadlocks deserve a specific note because they are the most common transient error in graph ingestion and they are partly self-inflicted. Two batches that both MERGE the same high-degree junction contend for the same lock, and the probability rises with concurrency and with batch size. Partitioning batches so that co-located writes land in the same batch — the same discipline the POI enrichment path uses for its H3 buckets — removes most of them at source, which is better than retrying them well.
Two things are worth instrumenting rather than inferring. The first is the retry rate as a proportion of attempts, which is a leading indicator of contention long before it becomes a latency problem — a run that quietly climbs from one per cent to fifteen is telling you the concurrency and the batch partitioning have drifted out of balance, usually because the data got denser rather than because anything was changed. The second is the classification breakdown: counting retryable against non-retryable failures separately means a schema regression shows up as a spike in the second series rather than as a vague slowdown in the first.
It is also worth being explicit about what happens after the retry budget is exhausted. Raising is correct, but the batch is then lost unless something catches it, and an ingestion run that loses batches silently is worse than one that stops. The pattern that holds up is a dead-letter list: on final failure, record the batch and its last error, continue with the rest of the run, and report the count at the end. That turns “the import failed” into “the import wrote 4.2 million rows and could not write these 340, for this reason”, which is the difference between a run you can act on and one you have to repeat.
Related
- Async Batch Processing for Graphs — the write path this policy wraps.
- Scaling Async Graph Ingestion with Python asyncio — the semaphore and pool sizing retries compete for.
- Backpressure with Bounded asyncio Queues — why a stalled retry loop must not let the producer run ahead.
- Enriching POI Data with Real-Time Demographics — partitioning writes so the deadlocks never happen.
This guide is part of Async Batch Processing for Graphs, within Spatial Graph Construction & OSM Ingestion.