Chapter 20 · Change Data Capture, Event Integration, Kafka/Connect Patterns, and Near-Real-Time Graph Synchronization
Outbox/CDC Integration with Relational Sources, Identity Mapping, and Eventual Consistency
Integrate a relational system of record with AtlasMart using the transactional outbox pattern, stable identity mapping, ordered per-aggregate events, idempotent graph projection, reconciliation, and explicit eventual-consistency SLOs.
AtlasMart Orders remains authoritative in a relational database while Neo4j powers connected customer/product/order traversals. A service updates SQL and then separately publishes “OrderPaid.” If the process crashes between those operations, SQL says PAID but the graph never receives an event. Reversing the calls merely changes which side can become wrong. The transactional outbox solves the atomicity boundary by storing the business update and an event row in the same relational transaction.
The source system owns the authoritative write. The graph is a projection that may lag. Do not pretend two independent databases share an atomic transaction unless you actually deploy a distributed transaction protocol and accept its operational cost.
Learning outcomes
Explain the dual-write race and why a transactional outbox closes it at the relational source boundary.
Design stable source event IDs and business-key mappings for Product, Customer and Order graph projections.
Separate source-event ordering, broker delivery and target-graph idempotency responsibilities.
Define eventual-consistency SLOs, reconciliation and backfill for an AtlasMart graph projection.
Plan schema evolution and deletion/tombstone semantics without creating bidirectional ownership loops.
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. The unsafe dual write
| Sequence | Crash window | Result |
|---|---|---|
| SQL COMMIT → publish Kafka | after SQL commit, before publish | source correct; graph never hears about change |
| publish Kafka → SQL COMMIT | after publish, before SQL commit | graph sees change that authoritative source rolled back |
“Retry publish in a catch block” cannot prove whether the previous publish succeeded and does not make SQL+Kafka atomic. Use a durable outbox record written with the source transaction.
2. Transactional outbox mechanism
The order row and outbox row commit atomically. A separate
publisher/CDC connector reads the outbox and emits it to Kafka.
If publication repeats, the immutable event_id lets
downstream consumers deduplicate. After successful publication,
retention/compaction of outbox rows is an operational policy—not
part of the business transaction.
BEGIN;
UPDATE orders
SET status = 'PAID', updated_at = CURRENT_TIMESTAMP
WHERE order_id = 'O-2001';
INSERT INTO outbox_events(event_id, aggregate_type, aggregate_id, event_type, schema_version, payload, created_at)
VALUES (
'evt-o-2001-paid-v1', 'Order', 'O-2001', 'OrderPaid', 1,
'{"orderId":"O-2001","customerId":"C-2001","status":"PAID"}',
CURRENT_TIMESTAMP
);
COMMIT;
-- Business row and outbox row are atomic in the relational source transaction.
| Outbox field | Purpose |
|---|---|
| event_id | globally unique immutable delivery/idempotency key |
| aggregate_type/id | business ordering and ownership key |
| event_type | semantic change, not database-table trivia |
| schema_version | decoder/migration contract |
| payload | minimum data needed or pointer/key to hydrate from source |
| created_at | lag/age evidence; not a substitute for broker offset/order |
3. Identity mapping must survive every storage engine
AtlasMart maps SQL order_id → graph
OrderReplica.orderId, not SQL row location or Neo4j
elementId. The same rule applies to customers/products. Stable
constrained keys are the bridge between a relational aggregate
and graph projection.
| Source key | Target identity | Invariant |
|---|---|---|
| customers.customer_id | CustomerReplica.customerId | unique, stable, never reused |
| products.product_id | ProductReplica.productId | unique, stable across backfill/replay |
| orders.order_id | OrderReplica.orderId | unique; status transitions versioned/ordered per order |
| outbox.event_id | SyncReceipt.eventId | unique; duplicate delivery becomes no-op |
4. Apply semantic events, not raw SQL mutations
OrderPaid describes a business fact. The graph
projection can choose its own labels/relationships while
retaining source identity and event version. This reduces
coupling to source table layouts and makes model refactors
explicit.
CYPHER 25
// Parameterized event application on Neo4j target.
MERGE (r:SyncReceipt {eventId:$eventId})
ON CREATE SET r.firstSeenAt = datetime(), r.labTag='ch20'
WITH r
MATCH (o:OrderReplica {orderId:$orderId})
SET o.status = $status,
o.customerId = $customerId,
o.schemaVersion = $schemaVersion,
o.labTag='ch20';
// In production use a transaction function and a receipt check so replay is a no-op,
// rather than mutating again unconditionally.
The abbreviated Cypher above illustrates mapping. In production, make the receipt check and business mutation one transaction (as in Lesson 2); a blind MERGE receipt followed by unconditional mutation is not sufficient idempotency for non-commutative updates.
5. Ordering and concurrency belong to the aggregate contract
If two Order events race—OrderPaid then
OrderCancelled—the target needs a source
version/sequence or a Kafka partitioning key that preserves
per-order order. Global serialization of all orders is
unnecessary. Include a monotonically increasing
aggregateVersion from the authoritative source when
transitions are non-commutative.
{
"eventId":"evt-o-2001-cancelled-v2",
"aggregateType":"Order",
"aggregateId":"O-2001",
"aggregateVersion":2,
"eventType":"OrderCancelled",
"schemaVersion":1,
"occurredAt":"2026-09-09T18:00:00Z",
"payload":{"reason":"CUSTOMER_REQUEST"}
}
6. Eventual consistency requires an SLO and reconciliation
| Evidence | Meaning |
|---|---|
| source high-water mark | latest committed source/outbox position |
| published Kafka offset | publisher progress |
| consumer offset/checkpoint | delivery progress |
| target receipt/event age | durable graph application progress |
| source-vs-target count/hash/sample | projection correctness beyond mere liveness |
| oldest unprocessed event age | user-visible staleness risk |
Define an SLO such as “99% of Order status changes are queryable in the graph within the agreed product-specific window,” then derive alerts from measured distributions. Do not invent a universal “under one second” target.
7. Backfill without racing the live stream
A safe bootstrap uses a watermark: record source boundary W, export/backfill a consistent snapshot corresponding to W, apply it idempotently, then consume events after W. Exact implementation depends on the relational CDC/outbox technology. If the source cannot give a consistent snapshot + position pair, design overlap and dedupe rather than assuming no gap.
| Phase | Acceptance condition |
|---|---|
| capture watermark | source position is durable and auditable |
| snapshot export | stable business keys + schema version recorded |
| bulk apply | target constraints valid; counts/invariants reconcile |
| incremental catch-up | start after/at W according to source semantics; overlap duplicates harmless |
| cutover | lag and reconciliation within SLO; rollback path retained |
8. Feedback loops and ownership
If graph-derived enrichment must return to the relational
system, use a different event type/topic and explicit ownership.
Never echo the same Order status field in both directions.
Include origin/flowId metadata and
reject a flow’s own emissions where bidirectional integration is
unavoidable.
9. 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
- What race does the outbox solve?
- Why is event_id different from order_id?
- What ordering usually matters for Order status?
- How do you verify eventual consistency?
- How can a backfill overlap live events safely?
Review the answers
1. The atomicity gap between committing the authoritative relational business change and durably recording an event for later publication.
2. order_id identifies the business aggregate; event_id uniquely identifies one immutable change/delivery for deduplication.
3. Per-order aggregate ordering/version, not a single global order for every AtlasMart event.
4. Measure source→target lag plus reconcile business keys/counts/invariants; connector liveness alone is insufficient.
5. Use stable identity and idempotent event receipts/upserts so replay/overlap is a no-op rather than duplication.
Summary and next step
The outbox makes event intent durable with the source transaction, while the graph consumer remains idempotent and measurable. Lesson 5 combines Neo4j CDC concepts, simulator failure injection, Kafka/outbox reasoning and backfill into one recovery acceptance test.
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.