Chapter 20 · Change Data Capture, Event Integration, Kafka/Connect Patterns, and Near-Real-Time Graph Synchronization
Build a Recoverable CDC Pipeline and Prove Restart, Duplicate, Gap, and Backfill Behavior
Assemble and test a recoverable synchronization runbook that proves checkpoint persistence, duplicate suppression, crash recovery, gap detection, backfill, reconciliation, security, observability, and rollback behavior.
AtlasMart is ready to call the synchronization pipeline production-capable only if it can prove recovery, not merely process a happy-path event. The acceptance drill deliberately creates the failures operators will eventually see: duplicate delivery, crash after durable side effect, missing retained history, incompatible source epoch after restore, schema evolution, target outage and a potential feedback loop. Every failure must end in a known state with an observable checkpoint and a tested remediation.
A recoverable CDC pipeline has a deterministic answer to: “What was the last source position we safely applied, what target evidence proves it, and what exact procedure do we follow if that source position no longer exists?”
Learning outcomes
Run a complete failure-injection matrix for duplicate, restart, gap and backfill behavior.
Reconcile source/snapshot state with target graph state using stable business keys and receipts.
Define cursor/offset storage, security, retention and restore procedures as operator runbook steps.
Gate connector/server/schema upgrades with replay compatibility and deprecated-procedure checks.
Choose a safe synchronization architecture from CDC, Kafka query polling, outbox and batch alternatives based on evidence.
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.
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. Build the mandatory free recovery harness
Use the exact cdc_sim.py from Lessons 1–2.
Optionally create the Community target constraints and replace
the simulator’s JSON durable_apply with the Python driver
transaction pattern. The simulator remains the source because
Community does not expose Neo4j CDC.
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 cdc_sim.py reset
rm -f .atlasmart_ch20_state.json .atlasmart_ch20_checkpoint.json
# Then run the scenario-specific command below.
2. Acceptance matrix: prove each failure mode
| Scenario | Injection | Expected evidence | Pass condition |
|---|---|---|---|
| happy path | consume full stream | APPLIED sim-2000..2005; checkpoint=sim-2005 | snapshot matches expected final state |
| duplicate delivery | --duplicate sim-2004 | one APPLIED then DUPLICATE for same eventId | final Product P-2001 price remains 179; one receipt |
| crash after target commit | --crash-after sim-2004 then restart | replay reports DUPLICATE sim-2004 | no double mutation; checkpoint later reaches sim-2005 |
| retention/missing event teaching case | --drop sim-2003 | GAP expectedOrdinal=3 got=4 | incremental processing stops; no silent skip |
| backfill | backfill after gap | BACKFILL + checkpoint sim-2005 | source snapshot and target reconcile |
| source history replaced | edit checkpoint epoch to atlasmart-source-old | GAP epoch-mismatch | consumer refuses foreign/restored history; backfill required |
| schema v2 | sim-2004 adds currency/schemaVersion=2 | transformer handles additive field | no decoder failure or loss |
3. Run the duplicate scenario
python cdc_sim.py reset
python cdc_sim.py consume --duplicate sim-2004
# Expected around sim-2004:
# APPLIED id=sim-2004 ...
# DUPLICATE id=sim-2004 ...
# CHECKPOINT ordinal=5 id=sim-2005
The duplicate is not an error if it maps to the same immutable event ID/business mutation. A rising duplicate rate can still signal offset resets, connector instability or retry pressure and should be observable.
4. Run crash → replay → recovery
$ 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.
In a real Kafka path, the analog is target commit succeeded but
the consumer/Kafka Connect offset did not commit. In a direct
db.cdc.query consumer, the analog is target commit
succeeded but the application checkpoint did not persist. Both
require safe replay.
5. Run gap → stop → backfill
$ 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
Real Neo4j CDC does not expose the simulator ordinal. The equivalent operational trigger is an invalid/pruned cursor or a restore/pause-resume history reset. On self-managed systems, transaction-log retention determines how far back changes remain queryable. On Aura, fixed log-space rotation means the available time window varies with workload. Never promise a time window from one quiet-period observation.
6. Real CDC operator runbook (optional licensed path)
| Step | Action | Evidence |
|---|---|---|
| 1 | verify database CDC mode and supported server/tier | SHOW DATABASES options or Aura CDC mode |
| 2 | record db.cdc.current before deployment/maintenance | opaque cursor + txCommitTime |
| 3 | persist consumer checkpoint only after durable side effect | checkpoint store + target receipt |
| 4 | monitor oldest/earliest available and consumer lag age | retention headroom |
| 5 | on InvalidIdentifier/ScanFailure, stop incremental writes | incident state; do not skip ahead silently |
| 6 | take/reconstruct current snapshot and reconcile target | business-key counts/hashes/invariants |
| 7 | anchor new cursor under documented snapshot/watermark procedure | new source history identity |
| 8 | resume and verify lag + duplicates + business correctness | receipts, errors, SLOs |
CYPHER 25
CALL db.cdc.earliest() YIELD id RETURN id AS earliest;
CALL db.cdc.current() YIELD id, txCommitTime RETURN id AS current, txCommitTime;
SHOW PROCEDURES YIELD name
WHERE name STARTS WITH 'db.cdc.' OR name STARTS WITH 'cdc.'
RETURN name ORDER BY name;
// Upgrade gate: application uses db.cdc.*; legacy cdc.* must not be a dependency.
7. Reconciliation is the final truth test
| Invariant | AtlasMart example | Failure interpretation |
|---|---|---|
| entity presence | P-2001 exists; P-2002 deleted in final snapshot | missing/extra target projection |
| business state | P-2001 price=179 USD | stale/out-of-order schema transform |
| referential identity | O-2001.customerId=C-2001 | mapping/backfill defect |
| receipt uniqueness | one SyncReceipt per immutable eventId | idempotency invariant broken |
| checkpoint/target relation | checkpoint never ahead of durable receipts | possible permanent event loss |
| lag age | current source position minus target applied time within SLO | overload/outage/retention risk |
8. Upgrade and migration gate
| Change | Preflight |
|---|---|
| Neo4j server upgrade | verify CDC procedures, Cypher 25 behavior, txLogEnrichment, restore semantics, connector compatibility |
| Kafka connector upgrade | read changelog, test offsets/start-from, serialization/payload mode, retries, source header and sink behavior |
| driver upgrade | run consumer integration tests and transient/error classification |
| capture mode DIFF/FULL | contract tests for both event shapes before switching |
| graph model/key migration | dual-key/version mapping, backfill, receipt compatibility and rollback window |
| restore/snapshot/pause-resume | assume old cursors/elementIds may be invalid; reinitialize consumer deliberately |
9. Choose the simplest mechanism that meets the requirement
| Requirement | Prefer | Why |
|---|---|---|
| near-real-time create/update/delete from Neo4j; licensed tier | Neo4j CDC or Kafka CDC source | native deletes + transaction-log changes |
| Community source; no hard-delete requirement | Kafka query polling or application event/outbox | works without Enterprise CDC |
| relational source of truth | transactional outbox / relational CDC → Kafka → graph | closes source dual-write race |
| nightly analytics/read model | batch snapshot/backfill | simpler operations if latency requirement allows |
| exact physical copy / DR | backup/restore/cluster mechanisms, not CDC | CDC does not replicate all database metadata/internal identity |
10. 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? |
11. Chapter 20 completion checklist
| Evidence | Pass condition |
|---|---|
| API | db.cdc.* only; deprecated cdc.* not required |
| edition | Community simulator vs Enterprise/Aura CDC clearly separated |
| identity | business keys constrained; elementId not cross-system identity |
| checkpoint | advances after durable side effect; source history identity recorded |
| duplicates | replay produces no duplicate business effect |
| restart | crash-after-commit scenario recovers |
| gap | invalid/missing history stops incremental path; backfill invoked |
| backfill | source/target counts and business invariants reconcile |
| Kafka | offsets, serialization, lag, secret and loop controls documented |
| schema | schemaVersion/mode evolution tested |
| security | CDC privileged data exposure/topic ACL/secret/TLS model documented |
| recovery | restore/pause-resume cursor invalidation included in runbook |
12. 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;
rm -f .atlasmart_ch20_state.json .atlasmart_ch20_checkpoint.json cdc_sim.py
Check your understanding
- What proves a pipeline is recoverable?
- What should happen when a CDC cursor is invalid after restore?
- Why is CDC not a backup?
- What prevents a Kafka feedback loop most reliably?
- When should AtlasMart avoid CDC entirely?
Review the answers
1. A tested relation between source checkpoint, durable target receipt/state, replay safety, gap detection, backfill and reconciliation—not merely “messages are flowing.”
2. Stop incremental assumptions, reinitialize/backfill and establish a new cursor for the new database history.
3. It captures graph changes, not every database metadata/internal detail, and its retained history can be pruned; use backup/cluster mechanisms for recovery copies.
4. Clear source/target ownership and topology separation, optionally reinforced with origin metadata and topic ACLs.
5. When batch/query/outbox mechanisms meet the actual latency/delete/recovery requirements with lower licensing and operational cost.
Summary and next step
Chapter 20 treats synchronization as a recovery protocol: opaque source cursors, stable business identity, at-least-once replay, atomic target receipts, versioned schemas, measurable lag, explicit gap/backfill transitions and no feedback loops. Chapter 21 can now add full-text retrieval knowing how near-real-time source changes become durable, testable graph/search state.
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.