Chapter 20 · Change Data Capture, Event Integration, Kafka/Connect Patterns, and Near-Real-Time Graph Synchronization

At-Least-Once Delivery, Ordering, Idempotent Consumers, Replay, and Schema/Model Evolution

Design an at-least-once AtlasMart consumer that survives duplicates, crashes, replay, retention gaps, and model changes by using durable business keys, idempotent target writes, receipts, and explicit schema-version policies.

Advanced240–330 minutesIdempotent consumer/replay labNeo4j 2026.07.1 · Community simulator mandatory · Enterprise/Aura CDC optionalCypher 25 · db.cdc.* · txLogEnrichment · Kafka Connect 5.5.3Java 21/25 · Python driver 6.3Last reviewed: September 2026

AtlasMart deploys a CDC consumer that updates a search read model. During a network interruption the target transaction commits, but the consumer process dies before its checkpoint file is written. On restart the source legitimately sends the same change again. If the handler performs a blind CREATE, duplicate products or duplicate notifications appear. This is the canonical at-least-once failure: the source did nothing wrong; the consumer lacked an idempotency boundary.

Correctness rule

Assume replay. Make the downstream business mutation and its durable event receipt atomic when possible. Advance the external checkpoint only after that atomic target transaction commits.

Learning outcomes

01

Explain at-least-once delivery and distinguish replay from data corruption.

02

Use event receipts, constrained business keys and transactional target writes to make duplicate application safe.

03

Reason about ordering with txId/seq without assuming a global order across independent streams.

04

Detect retention/cursor gaps and transition explicitly to a backfill/reinitialization state.

05

Version event/model transformations so additive and breaking schema changes can be rolled out safely.

Chapter 20 baseline · reviewed 9 September 2026

Current Neo4j Database is 2026.07.1; current 5.26 LTS patch is 5.26.30. Version-sensitive examples use explicit CYPHER 25. Application examples pin the official Python driver to neo4j==6.3.0; Kafka material pins Neo4j Connector for Kafka 5.5.3. The existing AtlasMart local baseline remains database neo4j, user neo4j, disposable password atlasmart-course-2026, loopback Bolt 7687 and HTTP 7474, no TLS only because the mandatory lab is loopback/local. Java 21 or 25 is the 2026.x server baseline. APOC and GDS are not required.

CDC edition/tier boundary

Neo4j CDC is not available in Community Edition. Current self-managed CDC documentation labels it Enterprise Edition; current Aura CDC documentation labels AuraDB Business Critical and AuraDB Virtual Dedicated Cloud. Therefore the mandatory free path is a deterministic event-log simulator plus an optional Community target projection. Optional real db.cdc.* commands are clearly marked and must not be presented as Community output.

Lab contract and assumptions

Dimension Pinned Chapter 20 assumption
server Neo4j Community 2026.07.1 for the mandatory downstream target; optional real CDC requires matching Enterprise 2026.07.1 or supported Aura tier
Cypher CYPHER 25 for version-sensitive commands
Java 21 or 25 for Neo4j 2026.07.x
driver Python 3.10+ with neo4j==6.3.0 for optional target-graph application examples
database/auth neo4j / neo4j / atlasmart-course-2026
transport bolt://localhost:7687 and http://localhost:7474 only for disposable loopback lab; production uses verified TLS such as neo4j+s
plugins none required; APOC/GDS not part of mandatory path
graph size small deterministic AtlasMart projection: P-2001/P-2002, C-2001, O-2001 plus SyncReceipt records
source semantics free simulator emits ordered synthetic change records with sourceEpoch + ordinal; real CDC uses opaque database-local change identifiers
measurement expected simulator output is deterministic and locally executable; server timing/lag must be measured by the learner and is not fabricated
Term Mechanism-first meaning
CDC Change Data Capture: a transaction-log-derived stream of graph data changes for create/update/delete events; it is not a backup or byte-for-byte replica.
txLogEnrichment Per-database setting that enriches transaction-log records so CDC can reconstruct changes; modes are OFF, DIFF, FULL.
change identifier Opaque cursor/pointer returned by db.cdc.*. It is local to one database history, treated as exclusive when querying, and must not be parsed or reused across restore/copy/import histories.
event id The id returned for a change record; it can become the next cursor. It is not a portable global business identifier.
txId / seq Transaction identifier plus within-transaction sequence. seq orders changes in a transaction. Transaction IDs can legitimately have gaps, so numeric continuity is not a loss detector.
business key Stable domain identity such as productId, customerId or orderId, preferred over elementId for cross-system synchronization.
selector Server-side db.cdc.query filter for entity kind, operation, labels/type, logical keys, changed properties or transaction metadata.
checkpoint Durable consumer state indicating the last safely processed source cursor/offset. Advance it only after the downstream effect is durable.
at-least-once Delivery model in which a change may be seen more than once after retry/restart; consumers must make repeated application safe.
idempotency Applying the same logical event repeatedly produces the same final business state and does not duplicate side effects.
replay Re-reading previously available events from a checkpoint/cursor for recovery or rebuild.
retention gap Consumer checkpoint points to change history already pruned or invalidated; incremental replay cannot continue safely and requires backfill/reinitialization.
backfill Bulk reconstruction of downstream state from a current source snapshot, followed by a new incremental checkpoint boundary.
Kafka Connect offset Connector-managed source progress persisted by Kafka Connect; losing it can cause source replay from configured start behavior.
outbox Rows/events written atomically with relational business changes so an asynchronous publisher can emit them without a dual-write race.
feedback loop A target write is captured again by the source and re-emitted indefinitely because source and sink paths are not separated or origin-filtered.

1. Exactly once is an end-to-end claim, not a source checkbox

Failure point What can happen Consumer requirement
after source read, before target write event will be read again safe retry
after target commit, before checkpoint same event is replayed although target already changed receipt/idempotency
after checkpoint, before side effect event is skipped forever after restart forbidden ordering: never checkpoint first
target timeout/connection loss commit outcome can be ambiguous read/receipt verification before non-idempotent retry
connector offset loss source may honor start-from again and replay history durable Kafka Connect offsets + idempotent sink

2. Make receipt + target mutation one atomic Neo4j transaction

On the target graph, a unique SyncReceipt.eventId constraint turns duplicate detection into an executable invariant. The application checks the receipt and applies the business-key upsert/deletion in the same transaction that creates the receipt. A separate file/Kafka offset may still lag behind, but a replay then sees the receipt and becomes a no-op.

Cypher 25 · target constraints
CYPHER 25
CREATE CONSTRAINT ch20_product_replica_id IF NOT EXISTS
FOR (p:ProductReplica) REQUIRE p.productId IS UNIQUE;
CREATE CONSTRAINT ch20_customer_replica_id IF NOT EXISTS
FOR (c:CustomerReplica) REQUIRE c.customerId IS UNIQUE;
CREATE CONSTRAINT ch20_order_replica_id IF NOT EXISTS
FOR (o:OrderReplica) REQUIRE o.orderId IS UNIQUE;
CREATE CONSTRAINT ch20_receipt_id IF NOT EXISTS
FOR (r:SyncReceipt) REQUIRE r.eventId IS UNIQUE;
Python driver 6.3 · target-transaction pattern
from neo4j import GraphDatabase

URI = "bolt://localhost:7687"
AUTH = ("neo4j", "atlasmart-course-2026")

def apply_event(tx, ev):
    # Receipt and graph mutation share one Neo4j transaction.
    already = tx.run(
        "MATCH (r:SyncReceipt {eventId:$id}) RETURN count(r) AS n",
        id=ev["id"],
    ).single()["n"]
    if already:
        return "duplicate"

    if ev["entity"] == "Product":
        if ev["op"] == "d":
            tx.run("MATCH (p:ProductReplica {productId:$key}) DETACH DELETE p", key=ev["key"])
        else:
            tx.run("""
                MERGE (p:ProductReplica {productId:$key})
                SET p += $after, p.labTag='ch20'
            """, key=ev["key"], after=ev["after"])
    # CustomerReplica / OrderReplica follow the same constrained-business-key pattern.
    tx.run("CREATE (:SyncReceipt {eventId:$id, labTag:'ch20', appliedAt:datetime()})", id=ev["id"])
    return "applied"

with GraphDatabase.driver(URI, auth=AUTH) as driver:
    driver.verify_connectivity()
    # Use session.execute_write(apply_event, event) in the consumer loop.

3. Controlled crash proves the replay path

The simulator writes target state before it writes the checkpoint. --crash-after exits in that gap. Restarting therefore replays sim-2004; the durable receipt converts the replay into DUPLICATE and the consumer continues.

Terminal · crash/restart expected output
$ python cdc_sim.py reset
$ python cdc_sim.py consume --crash-after sim-2004
...
APPLIED id=sim-2004 txId=505 seq=0 key=P-2001
CRASH after durable target apply, before checkpoint advance

$ python cdc_sim.py consume
DUPLICATE id=sim-2004 txId=505 seq=0 key=P-2001
APPLIED id=sim-2005 txId=506 seq=0 key=P-2002
CHECKPOINT ordinal=5 id=sim-2005

# The replay is expected. Receipt-based idempotency turns the duplicate into a no-op.
What this proves

The test proves that the specific target-state + receipt mechanism survives one crash window. It does not prove Kafka, Neo4j CDC, network partitions or external APIs are exactly-once; those require their own failure tests.

4. Ordering: scope the guarantee to the business invariant

Within a Neo4j transaction, seq establishes event order. Across transactions, the CDC result is chronological, but downstream systems may introduce partitions, retries or parallelism. Decide the smallest ordering domain your business needs. Inventory for one SKU may require per-productId order; unrelated SKUs should not be serialized behind each other.

Ordering requirement Practical strategy Risk if overclaimed
same transaction process txId + seq in returned order relationship event applied before required node state if reordered carelessly
same aggregate key Kafka partition key/business-key queue same product versions applied out of order
all AtlasMart events single global lane only if truly required unnecessary throughput bottleneck and large failure blast radius
cross-database domains explicit application/event contract no single Neo4j txId or cursor spans independent database histories

5. Gap/replay recovery uses cursor evidence, not txId arithmetic

Real CDC may raise Neo.DatabaseError.ChangeDataCapture.ScanFailure when the requested log entry has been pruned, or Neo.ClientError.ChangeDataCapture.InvalidIdentifier for a wrong/obsolete cursor. Restores also invalidate old cursors. At that point “retry harder” is unsafe: stop incremental application, establish a new source snapshot/backfill, reconcile target state, then anchor a new current/earliest cursor according to the runbook.

Terminal · deterministic missing-event teaching exercise
$ python cdc_sim.py reset
$ python cdc_sim.py consume --drop sim-2003
APPLIED id=sim-2000 txId=500 seq=0 key=P-2001
APPLIED id=sim-2001 txId=500 seq=1 key=P-2002
APPLIED id=sim-2002 txId=502 seq=0 key=C-2001
GAP expectedOrdinal=3 got=4 -> BACKFILL_REQUIRED

$ python cdc_sim.py backfill
BACKFILL entities=3 products=1 customers=1 orders=1
CHECKPOINT ordinal=5 id=sim-2005

6. Schema and graph-model evolution are part of the event contract

Change Safe rollout pattern Failure to avoid
add optional property consumer tolerates unknown fields; target adds field when understood strict decoder rejects event
rename property dual-read/dual-write or versioned mapping during migration new producer silently removes field old consumer needs
change DIFF↔FULL consumer recognizes captureMode/shape before server switch handler assumes full before/after state but gets diff
split Product into Product + Offer versioned transformer and backfill to new graph model old replay applied directly into incompatible target model
change business key migration mapping + overlap constraints + reconciliation receipts dedupe but identity forks into two target nodes
JSON · versioned integration envelope
{
  "contract": "atlasmart.product-change",
  "schemaVersion": 2,
  "sourceDatabase": "neo4j",
  "sourceEpoch": "deployment-specific",
  "sourceChangeId": "opaque-db-cdc-id",
  "eventId": "opaque-db-cdc-id",
  "businessKey": {"productId": "P-2001"},
  "operation": "u",
  "payload": {"price": 179.0, "currency": "USD"}
}

7. Deliberately wrong retry strategy

“Catch Exception, sleep one second, retry forever” conflates malformed schema, authorization denial, invalid cursor, target uniqueness violation, transient network loss and a poison event. Classify errors into retryable, quarantine/dead-letter, backfill-required and operator-action categories. Bound retries by an SLO; preserve the original event and trace identifiers for diagnosis.

Class Example Action
transient transport temporary target/network unavailability bounded retry with backoff/jitter; do not advance checkpoint
duplicate receipt already present success/no-op; checkpoint may advance
poison schema unsupported schemaVersion quarantine + alert; preserve ordering policy for affected aggregate
cursor invalid/pruned CDC ScanFailure/InvalidIdentifier stop incremental, run backfill/reinitialization
authorization db.cdc.query/target write denied operator/security fix; no blind retry storm

8. Production judgment

Decision surface Production questions
delivery semantics Where can replay occur—CDC polling, Kafka Connect, broker, consumer, driver retry—and what durable receipt/idempotency key makes each replay harmless?
retention/checkpoint window How long can consumers be down before source history may be pruned? How is earliest/current cursor validity monitored and tested?
identity Which business keys survive restore/import and cross systems? Are key constraints present at source and target?
ordering Which business rules require ordering only within a transaction, per aggregate, or globally? Do not infer loss from gaps in Neo4j txId.
schema evolution How are FULL/DIFF mode changes, added/removed properties, event-envelope versions and target-model migrations rolled out compatibly?
latency/SLO What are source commit → capture → broker → consumer → target durable p50/p95/p99 and lag-age distributions?
transactions/retries Is checkpoint advancement ordered after durable side effects? Are ambiguous retries and non-idempotent external actions controlled?
CPU/disk/network What transaction-log enrichment/storage overhead, broker retention, serialization size, batch size and network egress are acceptable?
security/privacy CDC can reveal all changed data in a database; who has boosted db.cdc.query access, how are topics ACLed/encrypted, and how is sensitive payload retention governed?
backup/recovery After restore/snapshot/pause-resume, how is consumer state reinitialized and how are stale downstream changes reconciled?
observability Can you correlate source tx metadata, CDC id/txId/seq, Kafka partition/offset, consumer eventId, target receipt and request trace?
testing/failure injection Have duplicate, crash-after-side-effect, invalid cursor, retention gap, schema change, broker restart, target outage and feedback-loop guards been tested?
version/edition/tier Are server/driver/connector versions and self-managed Enterprise vs AuraDB BC/VDC availability recorded? Are deprecated cdc.* names absent?
cost/migration Does near-real-time synchronization justify Enterprise/Aura/Kafka ops cost versus simpler polling/batch/read-model alternatives?

9. Cleanup/reset

Cypher 25 · optional Community target cleanup
CYPHER 25
MATCH (n) WHERE n.labTag = 'ch20' DETACH DELETE n;
DROP CONSTRAINT ch20_product_replica_id IF EXISTS;
DROP CONSTRAINT ch20_customer_replica_id IF EXISTS;
DROP CONSTRAINT ch20_order_replica_id IF EXISTS;
DROP CONSTRAINT ch20_receipt_id IF EXISTS;
Terminal · simulator cleanup
rm -f .atlasmart_ch20_state.json .atlasmart_ch20_checkpoint.json cdc_sim.py

Check your understanding

  1. What crash window forces duplicate replay?
  2. Why put SyncReceipt in the same target transaction as the business mutation?
  3. What does a CDC ScanFailure after long downtime mean operationally?
  4. Does txId provide a globally contiguous event counter?
  5. How do you roll out a breaking event schema?
Review the answers

1. Downstream commit succeeds but checkpoint/offset persistence has not yet happened.

2. So a committed business mutation cannot exist without durable evidence that lets a replay become a no-op.

3. The cursor points to transaction-log history no longer available; switch to an explicit backfill/reinitialization path.

4. No. Gaps can be normal and independent databases have separate histories.

5. Use an explicit schemaVersion/transformer, compatibility window or dual representation, tests, reconciliation and a rollback path.

Summary and next step

An at-least-once pipeline is reliable when duplicate application is boring, gaps are detected rather than guessed away, and model changes are versioned. Lesson 3 moves the same correctness contract into Kafka Connect, where connector offsets, topic partitioning, serialization and feedback-loop topology add new failure surfaces.

Authoritative references

  • Current Neo4j versions — Current database release 2026.07.1 and current 5.26 LTS patch 5.26.30.
  • CDC introduction — Current CDC availability and purpose; CDC is a change feed, not an exact database-copy mechanism.
  • CDC on self-managed Neo4j — Enterprise-only enablement, OFF/DIFF/FULL txLogEnrichment modes, security, retention, disk and unrecorded-change boundaries.
  • CDC on Aura — Current AuraDB Business Critical / Virtual Dedicated Cloud enablement, pause/resume and snapshot-reset behavior.
  • CDC procedures — db.cdc.earliest/current/query semantics, exclusive cursors, txId/seq, metadata and change identifiers.
  • CDC event schema — Current node/relationship event envelope including operation, keys, before/after state and metadata.
  • CDC selectors — Server-side entity/operation/label/type/key/property and metadata filtering.
  • CDC examples — Official cursor-management examples and empty-result/current-cursor retention guidance.
  • CDC backup/restore behavior — Why restores invalidate old change identifiers and require consumer state/backfill/reinitialization planning.
  • CDC troubleshooting — Disabled, scan-failure and invalid-identifier error mechanisms including retention and wrong-database cursors.
  • CDC known issues — ElementId and change-identifier instability across restore/import/copy/pause-resume operations.
  • CDC changelog — Current db.cdc.* namespace and deprecation history of cdc.* procedures.
  • Neo4j Connector for Kafka — Current source/sink connector architecture and CDC/query source strategies.
  • Kafka source CDC quickstart — Current CDC source configuration, txLogEnrichment prerequisite, event fields, offsets, topics and failure modes.
  • Kafka connector installation — Current connector release 5.5.3 and Kafka Connect plugin deployment model.
  • Kafka source query quickstart — Community-compatible polling source alternative and its delete/soft-delete limitation.
  • Kafka sink CDC quickstart — Neo4j-to-Neo4j CDC sink pattern, source header requirement, key constraints and loop warning.
  • Python driver 6.3 API — Official driver 6.3 and Bolt compatibility used for optional application examples.

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.