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

Neo4j CDC Concepts, Change Identifiers, Event Semantics, Capture Scope, and Consumer Checkpoints

Build the mental model for current Neo4j CDC: transaction-log enrichment, opaque change identifiers, event envelopes, selectors, security, checkpoint advancement, retention boundaries, and a free deterministic AtlasMart simulator.

Advanced230–310 minutesCDC envelope + simulator 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 wants inventory, recommendations and customer-service systems to react to graph changes within seconds. The tempting design is “read the transaction log and publish every row exactly once.” Neo4j CDC does expose transaction-log-derived graph changes, but a correct integration starts by separating database-local cursor semantics from consumer delivery semantics. The CDC procedure tells you which graph changes are available; your integration decides when a downstream effect is durable and when a checkpoint may advance.

Core mental model

Treat db.cdc.query as a resumable change reader over retained transaction-log history. Treat the returned id as an opaque exclusive cursor. Treat the downstream pipeline as at-least-once unless you can prove a stronger end-to-end protocol.

Learning outcomes

01

Explain how txLogEnrichment OFF/DIFF/FULL makes graph changes available to current Neo4j CDC and where the feature is licensed.

02

Interpret db.cdc.earliest/current/query, exclusive change identifiers, txId, seq, metadata and node/relationship event envelopes.

03

Use selectors and business-key constraints without confusing elementId, txId or change IDs with portable identity.

04

Design a durable checkpoint rule that advances only after downstream state is committed and survives duplicate replay.

05

Run the free deterministic AtlasMart event-log simulator and distinguish its teaching ordinal from real Neo4j cursor semantics.

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. CDC capture begins in the transaction log

Self-managed Neo4j leaves CDC OFF by default. Enabling DIFF records property/label/type differences; FULL records complete before/after entity state. The setting is per database. Changing DIFF↔FULL changes the event shape immediately, so consumers must be version- and mode-aware. Import/load operations that bypass the transaction layer are not CDC events.

Mode Captured state Operational consequence
OFF no CDC enrichment no usable CDC stream; disabling also breaks old cursor continuity
DIFF removals/updates/additions only smaller payload than FULL but consumer must reconstruct if it needs complete state
FULL complete before/after state simpler downstream projection; higher transaction-log volume
Cypher 25 · optional Enterprise enablement
// OPTIONAL — self-managed Enterprise 2026.07.1, run against system database.
CYPHER 25
ALTER DATABASE neo4j SET OPTION txLogEnrichment "FULL";
SHOW DATABASES YIELD name, options
WHERE name = 'neo4j'
RETURN name, options;
// Expected invariant: options contains txLogEnrichment: "FULL".

2. Three procedures define the current cursor contract

db.cdc.earliest() returns the earliest currently available cursor. db.cdc.current() returns an exclusive cursor for the latest committed transaction; with Cypher 25 from Neo4j 2026.06 it also returns txCommitTime. db.cdc.query(from, selectors) returns changes strictly after from. Every returned record has its own id, so a consumer can advance incrementally.

Cypher 25 · optional real CDC probe
// OPTIONAL — Enterprise/Aura supported CDC database.
// Transaction A: capture an exclusive starting cursor.
CYPHER 25
CALL db.cdc.current() YIELD id, txCommitTime
RETURN id, txCommitTime;

// Transaction B: make one AtlasMart change.
CYPHER 25
MERGE (p:Product {productId:'P-2001'})
SET p.name='Trail Camera Pro', p.price=189.00, p.labTag='ch20';

// Transaction C: use the id from Transaction A as $from.
CYPHER 25
CALL db.cdc.query($from, [
  {select:'n', labels:['Product'], operation:'c'},
  {select:'n', labels:['Product'], operation:'u'}
])
YIELD id, txId, seq, metadata, event
RETURN id, txId, seq, metadata.txCommitTime AS committedAt,
       event.operation AS op, event.keys AS keys,
       event.state.before AS before, event.state.after AS after;
Boundary: no cursor arithmetic

Do not decode, increment, compare lexically, or move a change identifier between databases. A restored/copied/imported database has a different history. The same warning applies after Aura pause/resume or snapshot restore.

3. Read the event envelope as evidence

Field What it tells you What it does not guarantee
id cursor associated with this change record portable/global identity across databases or restores
txId source transaction identifier contiguous sequence; gaps can be normal
seq ordering among changes in the same transaction global business ordering across independent producers
metadata commit/start time, users, connection, txMetadata, capture mode and database/server context authorization filtering of event content
event.operation c/u/d create/update/delete exactly-once delivery
event.eventType n/r node or relationship target model equivalence
event.keys logical-key values derived from applicable key constraints presence unless the source schema defines them
event.elementId source internal element identity at that database history stability through restore/import/copy/pause-resume

4. Selectors reduce traffic; they do not replace checkpoint discipline

Selectors can filter nodes, relationships, operations, labels/types, key properties, changed fields and transaction metadata. A strict selector may return no rows for a long period while unrelated transactions continue. Official examples therefore capture current before querying and may advance to that current cursor when no matching rows were returned, preventing a silent cursor from aging out of retention.

Cypher 25 · selector examples
CYPHER 25
CALL db.cdc.query($from, [
  {select:'n', labels:['Product'], operation:'u', changesTo:['price']},
  {select:'r', type:'CONTAINS', operation:'c'}
])
YIELD id, txId, seq, metadata, event
RETURN id, txId, seq, metadata.txCommitTime AS committedAt, event;

5. Security is wider than ordinary MATCH privileges

CDC can expose all captured changes in the database and is not reduced to the ordinary entities a user can read through graph privileges. db.cdc.query therefore requires admin or deliberately granted execute + boosted execute privileges plus database access. Treat the stream as a privileged data-exfiltration surface: topic ACLs, TLS, secret rotation, retention and privacy controls belong in the threat model.

Cypher 25 · optional least-privilege CDC procedure role
// OPTIONAL Enterprise RBAC example. db.cdc.query exposes all matching database changes
// rather than being reduced to the caller's ordinary entity-level graph visibility.
CYPHER 25
GRANT ACCESS ON DATABASE neo4j TO atlasmart_cdc_reader;
GRANT EXECUTE PROCEDURE db.cdc.query ON DBMS TO atlasmart_cdc_reader;
GRANT EXECUTE BOOSTED PROCEDURE db.cdc.query ON DBMS TO atlasmart_cdc_reader;

6. Mandatory free simulator: learn the semantics without faking CDC

Save the following as cdc_sim.py. It uses a synthetic ordinal only so the exercise can deterministically inject a missing event. Real Neo4j consumers do not have this ordinal; they detect retention/cursor failure through db.cdc.earliest/current/query behavior and CDC error evidence.

Python · cdc_sim.py
#!/usr/bin/env python3
"""AtlasMart Chapter 20 deterministic CDC simulator — Python standard library only."""
import argparse, json, os, sys
from pathlib import Path

EVENTS = [
 {"epoch":"atlasmart-source-v1","ordinal":0,"id":"sim-2000","txId":500,"seq":0,"entity":"Product","key":"P-2001","op":"c","after":{"name":"Trail Camera Pro","price":189.0,"schemaVersion":1}},
 {"epoch":"atlasmart-source-v1","ordinal":1,"id":"sim-2001","txId":500,"seq":1,"entity":"Product","key":"P-2002","op":"c","after":{"name":"Field Battery","price":49.0,"schemaVersion":1}},
 {"epoch":"atlasmart-source-v1","ordinal":2,"id":"sim-2002","txId":502,"seq":0,"entity":"Customer","key":"C-2001","op":"c","after":{"name":"Mina Rahimi","tier":"GOLD","schemaVersion":1}},
 {"epoch":"atlasmart-source-v1","ordinal":3,"id":"sim-2003","txId":503,"seq":0,"entity":"Order","key":"O-2001","op":"c","after":{"customerId":"C-2001","status":"PAID","schemaVersion":1}},
 {"epoch":"atlasmart-source-v1","ordinal":4,"id":"sim-2004","txId":505,"seq":0,"entity":"Product","key":"P-2001","op":"u","after":{"name":"Trail Camera Pro","price":179.0,"schemaVersion":2,"currency":"USD"}},
 {"epoch":"atlasmart-source-v1","ordinal":5,"id":"sim-2005","txId":506,"seq":0,"entity":"Product","key":"P-2002","op":"d","after":None},
]
SNAPSHOT = {
 "Product":{"P-2001":{"name":"Trail Camera Pro","price":179.0,"schemaVersion":2,"currency":"USD"}},
 "Customer":{"C-2001":{"name":"Mina Rahimi","tier":"GOLD","schemaVersion":1}},
 "Order":{"O-2001":{"customerId":"C-2001","status":"PAID","schemaVersion":1}},
}
STATE = Path('.atlasmart_ch20_state.json')
CHECKPOINT = Path('.atlasmart_ch20_checkpoint.json')

def load(path, default):
    return json.loads(path.read_text()) if path.exists() else default

def save(path, obj):
    tmp = path.with_suffix(path.suffix + '.tmp')
    tmp.write_text(json.dumps(obj, indent=2, sort_keys=True))
    os.replace(tmp, path)

def reset():
    for p in (STATE, CHECKPOINT):
        if p.exists(): p.unlink()
    print('RESET state=empty checkpoint=none')

def durable_apply(state, ev):
    # Receipt and business-state mutation are persisted together in one state document.
    if ev['id'] in state['receipts']:
        return 'duplicate'
    bucket = state['entities'].setdefault(ev['entity'], {})
    if ev['op'] == 'd': bucket.pop(ev['key'], None)
    else: bucket[ev['key']] = ev['after']
    state['receipts'].append(ev['id'])
    save(STATE, state)
    return 'applied'

def consume(drop=None, duplicate=None, crash_after=None):
    state = load(STATE, {'receipts':[], 'entities':{}})
    cp = load(CHECKPOINT, {'epoch':'atlasmart-source-v1','ordinal':-1,'id':None})
    if cp['epoch'] != 'atlasmart-source-v1':
        print('GAP epoch-mismatch -> BACKFILL_REQUIRED'); return 3
    stream = [e for e in EVENTS if e['ordinal'] > cp['ordinal'] and e['id'] != drop]
    if duplicate:
        match = next((e for e in stream if e['id'] == duplicate), None)
        if match: stream.insert(stream.index(match)+1, dict(match))
    expected = cp['ordinal'] + 1
    for ev in stream:
        # An at-least-once source may replay an event that the target already applied.
        # A duplicate older than the next expected ordinal is harmless; a replay at the
        # expected ordinal advances the checkpoint after we verify the receipt exists.
        if ev['id'] in state['receipts']:
            if ev['ordinal'] > expected:
                print(f'GAP expectedOrdinal={expected} got={ev["ordinal"]} -> BACKFILL_REQUIRED')
                return 4
            print(f'DUPLICATE id={ev["id"]} txId={ev["txId"]} seq={ev["seq"]} key={ev["key"]}')
            if ev['ordinal'] == expected:
                cp = {'epoch':ev['epoch'], 'ordinal':ev['ordinal'], 'id':ev['id']}
                save(CHECKPOINT, cp)
                expected = ev['ordinal'] + 1
            continue
        if ev['ordinal'] != expected:
            print(f'GAP expectedOrdinal={expected} got={ev["ordinal"]} -> BACKFILL_REQUIRED')
            return 4
        result = durable_apply(state, ev)
        print(f'{result.upper()} id={ev["id"]} txId={ev["txId"]} seq={ev["seq"]} key={ev["key"]}')
        if crash_after == ev['id']:
            print('CRASH after durable target apply, before checkpoint advance')
            return 9
        # Offset/checkpoint is advanced only after downstream state is durable.
        cp = {'epoch':ev['epoch'], 'ordinal':ev['ordinal'], 'id':ev['id']}
        save(CHECKPOINT, cp)
        expected = ev['ordinal'] + 1
    print(f'CHECKPOINT ordinal={cp["ordinal"]} id={cp["id"]}')
    return 0

def backfill():
    state = {'receipts':['BACKFILL@sim-current'], 'entities':SNAPSHOT}
    save(STATE, state)
    latest = EVENTS[-1]
    save(CHECKPOINT, {'epoch':latest['epoch'],'ordinal':latest['ordinal'],'id':latest['id']})
    print('BACKFILL entities=3 products=1 customers=1 orders=1')
    print(f'CHECKPOINT ordinal={latest["ordinal"]} id={latest["id"]}')

def show():
    print(json.dumps({'checkpoint':load(CHECKPOINT,None),'state':load(STATE,None)},indent=2,sort_keys=True))

p=argparse.ArgumentParser()
sub=p.add_subparsers(dest='cmd',required=True)
sub.add_parser('reset'); sub.add_parser('backfill'); sub.add_parser('show')
c=sub.add_parser('consume'); c.add_argument('--drop'); c.add_argument('--duplicate'); c.add_argument('--crash-after')
a=p.parse_args()
if a.cmd=='reset': reset()
elif a.cmd=='backfill': backfill()
elif a.cmd=='show': show()
elif a.cmd=='consume': sys.exit(consume(a.drop,a.duplicate,a.crash_after))
Terminal · normal expected output
# 1) Normal at-least-once consumer run
$ python cdc_sim.py reset
RESET state=empty checkpoint=none
$ python cdc_sim.py consume
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
APPLIED id=sim-2003 txId=503 seq=0 key=O-2001
APPLIED 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

# Notice txId jumps 500 -> 502 and 503 -> 505. That is NOT our gap detector.
# The simulator's ordinal is synthetic teaching metadata; real CDC uses opaque cursor validity/retention evidence.

7. Wrong approach → diagnosis → repair

Wrong approach Concrete failure Repair and verification
“CDC is exactly once.” crash after target commit but before cursor persistence replays the same event idempotent business key + durable receipt; run the crash/restart scenario in Lesson 2
Use elementId as cross-system key restore/import changes elementIds and downstream identity splits use stable productId/customerId/orderId and constraints
Detect loss from txId gaps system/schema transactions can create legitimate txId gaps treat cursor validity/retention as the source of truth
Persist cursor before target commit crash loses an event permanently downstream commit target effect first, then checkpoint; replay must be safe
Use cdc.query()/cdc.current() forever deprecated API may disappear from future language/runtime support use db.cdc.* and include deprecation tests in upgrade gate

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?

Check your understanding

  1. Why is db.cdc.current() called an exclusive cursor?
  2. Why can txId jump without a lost CDC event?
  3. When may a consumer advance its checkpoint?
  4. Why is elementId unsafe for cross-system identity?
  5. What is the free learning substitute for Enterprise/Aura CDC?
Review the answers

1. A later db.cdc.query(from) returns changes after the transaction represented by that cursor, not the changes in that transaction.

2. Some transaction kinds are not recorded as change events, so transaction IDs are not guaranteed contiguous.

3. Only after the intended downstream side effect is durably committed; otherwise a crash can create permanent loss.

4. Database-changing operations such as restore/import/copy/pause-resume can change it; use logical business keys.

5. The deterministic event-log simulator plus optional Community target projection; it teaches recovery semantics without claiming it is Neo4j CDC.

Summary and next step

CDC correctness begins with an opaque database-local cursor, a privileged event envelope, business-key identity and a checkpoint that advances after durable effects. Lesson 2 turns that contract into an at-least-once consumer that proves duplicate, restart, replay and schema-evolution behavior.

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.