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

Kafka/Connector Patterns for Ingesting to or Emitting from Neo4j and Avoiding Feedback Loops

Connect Neo4j with Kafka deliberately: compare CDC and query source strategies, understand Kafka Connect offsets and serialization, design sink upserts, monitor lag, and prevent source/sink feedback loops.

Advanced230–320 minutesKafka/Connect topology 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 now wants two Kafka paths: graph changes should feed analytics/search, while upstream commerce events should update Neo4j. Kafka Connect can run both directions, but the direction and source strategy matter. A Neo4j source connector using CDC has different edition, delete and offset semantics from the Community-compatible query polling strategy. A sink writing back into the same CDC-enabled source can create an infinite loop.

Topology rule

Draw the data-flow arrows before writing connector JSON. Every arrow needs an owner, identity key, offset/checkpoint, retry policy, security boundary and a rule that prevents the arrow from feeding its own source.

Learning outcomes

01

Distinguish Neo4j Kafka source CDC strategy from query polling and choose based on deletes, edition and latency needs.

02

Explain Kafka Connect source offsets, start-from behavior, single-task CDC source behavior and schema-aware serialization.

03

Design idempotent sink mapping with business-key constraints rather than blind CREATE statements.

04

Prevent Neo4j→Kafka→Neo4j feedback loops through topology separation/origin governance.

05

Measure connector lag, task errors, broker offsets, target receipts and end-to-end age instead of only connector RUNNING status.

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. Source strategy is a correctness choice

Source strategy Availability Strength Important limitation
CDC Enterprise or supported Aura tier near-real-time create/update/delete with before/after change semantics requires CDC enablement, retained change history and privileged access
Query polling Community/Enterprise/Aura where query connector supported custom projection; easy when a monotonic tracking field exists cannot discover a hard-deleted entity after it is gone unless soft-delete/outbox design preserves it
Delete boundary

If downstream correctness requires physical deletes, polling a graph after deletion cannot reconstruct what vanished. Use CDC or represent deletion as an explicit durable event/soft-delete.

2. Pin the connector and understand its process boundary

Current Neo4j Connector for Kafka is 5.5.3. It runs as a Kafka Connect plugin, not inside the Neo4j server. The CDC source connector is configured with neo4j.source-strategy=CDC and always runs one source task; setting tasks.max above one does not parallelize it. Kafka Connect persists source offsets, which become the connector-managed checkpoint.

JSON · illustrative current CDC source connector
{
  "name": "AtlasMartNeo4jSourceCdc",
  "config": {
    "connector.class": "org.neo4j.connectors.kafka.source.Neo4jConnector",
    "key.converter": "org.apache.kafka.connect.json.JsonConverter",
    "key.converter.schemas.enable": true,
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": true,
    "neo4j.uri": "neo4j://neo4j:7687",
    "neo4j.authentication.type": "BASIC",
    "neo4j.authentication.basic.username": "neo4j",
    "neo4j.authentication.basic.password": "${file:/run/secrets/neo4j.properties:password}",
    "neo4j.database": "neo4j",
    "neo4j.source-strategy": "CDC",
    "neo4j.start-from": "NOW",
    "neo4j.cdc.poll-interval": "1s",
    "neo4j.cdc.poll-duration": "5s",
    "neo4j.cdc.topic.atlasmart-product.patterns.0.pattern": "(:Product)",
    "neo4j.cdc.topic.atlasmart-product.patterns.0.operation": "CREATE",
    "neo4j.cdc.topic.atlasmart-product.patterns.1.pattern": "(:Product)",
    "neo4j.cdc.topic.atlasmart-product.patterns.1.operation": "UPDATE",
    "neo4j.cdc.topic.atlasmart-product.patterns.2.pattern": "(:Product)",
    "neo4j.cdc.topic.atlasmart-product.patterns.2.operation": "DELETE"
  }
}
Secret note

The secret reference demonstrates indirection. Do not commit Neo4j or Kafka credentials into connector JSON. Use the secret/config provider supported by your Connect deployment and protect its storage.

3. start-from is initial positioning, not a permanent override

neo4j.start-from First-run behavior After Kafka Connect has stored an offset
NOW capture changes after connector starts stored offset normally wins
EARLIEST start at earliest available change history stored offset normally wins
USER_PROVIDED use explicit CDC cursor stored offset normally wins unless intentionally ignored/reset

Losing/resetting connector offsets can cause a replay. That is not a reason to disable replay—it is a reason to make the sink idempotent and document offset backup/reset procedures.

4. Serialization is part of compatibility

The current source quickstarts use schema-carrying Kafka Connect values. The source connector always produces messages with schemas; a schemaless converter configuration will fail. Extended/compact payload modes and Avro/Protobuf/JSON choices affect type fidelity and compatibility. Version and test the envelope separately from the target graph model.

Evidence Why capture it
topic + partition + offset broker-side ordering/progress
neo4j.source.cdc.id header proves a message came from the supported CDC source format used by the CDC sink strategy
event txId/seq/id source change correlation
business key idempotent target MERGE/update
schema id/version decoder/contract compatibility

5. Sink design: MERGE by owned identity, not by coincidence

The sink target must have constraints matching the identity fields carried by events. A source key constraint makes logical keys available in CDC events; a target constraint makes MERGE efficient and enforces uniqueness. If the target has no matching key, retries can duplicate entities. If the key is too weak, distinct source entities can collapse into one target.

Cypher 25 · target key/receipt 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;

6. Feedback-loop failure injection

A source database publishes CDC to topic A. A sink consumes topic A and writes back to the same CDC-enabled database. Those writes are new transactions, which CDC captures and republishes. The pipeline feeds itself. The official source/sink quickstart deliberately uses separate source and target databases.

Guard Mechanism
topology separation source and target are distinct databases/instances; target CDC is disabled when it is only a sink
origin metadata events carry source system/flow ID; consumers reject their own outbound origin when bidirectional sync is unavoidable
topic separation direction-specific topics and ACLs prevent accidental cross-wiring
business ownership one system owns each mutable field/aggregate; avoid “last writer wins” ping-pong
rate/loop alert detect same business key/event lineage cycling abnormally
Wrong fix

Filtering only on a label is not a loop proof if the sink writes that same label. Ownership/origin and topology must make the cycle impossible or explicitly detectable.

7. Observability: RUNNING is not healthy enough

Signal Healthy question
source offset / CDC cursor age Is the connector keeping up with source commits inside the retention window?
consumer-group lag Are partitions accumulating backlog?
task status/error log Are retries hiding poison messages or auth/schema failures?
target receipt rate Do applied + duplicate counts explain consumed messages?
end-to-end event age How old is data when target commit completes? p95/p99?
DLQ/quarantine Are failures bounded and diagnosable rather than silently skipped?

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. When is query polling preferable to CDC?
  2. What persists CDC source progress in Kafka Connect?
  3. Why can tasks.max>1 fail to improve a CDC source connector?
  4. What is the safest loop prevention?
  5. Why is connector RUNNING insufficient?
Review the answers

1. When Community compatibility/custom projection is more important and deletes can be represented safely through soft-delete/outbox/tracking semantics.

2. The connector source offset; neo4j.start-from mainly determines first positioning when no stored offset applies.

3. The current Neo4j CDC source connector runs a single source task.

4. Separate source/target ownership/topology so sink writes are not captured back into the same source flow.

5. The task can be technically running while lag, retries, poison messages, target failures or end-to-end age violate the SLO.

Summary and next step

Kafka Connect adds durable offsets, broker ordering and scalable integration, but not magical exactly-once graph semantics. Lesson 4 applies the same discipline when the source of truth is relational: the transactional outbox removes the classic dual-write race before events ever reach Neo4j.

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.