feat: CDC-based mixed index synchronization (#4873) - #4906
Conversation
There was a problem hiding this comment.
Pull request overview
Adds an opt-in Change-Data-Capture (CDC) pipeline to keep mixed indexes (ElasticSearch/Solr/Lucene) eventually consistent with committed graph data by asynchronously reindexing affected elements from their current state, eliminating the “dual write” divergence risk.
Changes:
- Introduces per-index CDC configuration (
index.[X].cdc.*) and commit-path logic to skip synchronous mixed-index writes in cdc-only mode. - Adds
MixedIndexUpdateApplier(reindex-from-current-state) andCdcElementChangein core, plus Cassandrastorage.cql.cdctable option support. - Adds a new
janusgraph-cdcmodule implementing a Kafka consumer worker, Debezium Cassandra JSON decoder, and extensive unit/component/E2E tests + documentation.
Reviewed changes
Copilot reviewed 33 out of 33 changed files in this pull request and generated 4 comments.
Show a summary per file
| File | Description |
|---|---|
| pom.xml | Registers new janusgraph-cdc Maven module in the build. |
| mkdocs.yml | Adds CDC operator guide page to documentation nav. |
| janusgraph-lucene/src/test/java/org/janusgraph/diskstorage/lucene/MixedIndexUpdateApplierTest.java | Validates reindex-from-current-state behavior over Lucene in cdc-only mode. |
| janusgraph-lucene/src/test/java/org/janusgraph/diskstorage/lucene/CdcSkipMutationTest.java | Verifies synchronous mixed-index write is skipped in cdc-only and retained in dual mode. |
| janusgraph-cql/src/test/java/org/janusgraph/diskstorage/cql/CQLCdcTableOptionTest.java | Unit-tests that storage.cql.cdc toggles cdc=true on edgestore DDL only. |
| janusgraph-cql/src/main/java/org/janusgraph/diskstorage/cql/CQLKeyColumnValueStore.java | Refactors CREATE TABLE building and conditionally applies Cassandra cdc=true for edgestore. |
| janusgraph-cql/src/main/java/org/janusgraph/diskstorage/cql/CQLConfigOptions.java | Adds storage.cql.cdc configuration option. |
| janusgraph-core/src/test/java/org/janusgraph/graphdb/configuration/CdcIndexConfigTest.java | Tests defaults and per-index scoping of index.[X].cdc.*. |
| janusgraph-core/src/main/java/org/janusgraph/graphdb/database/StandardJanusGraph.java | Computes cdc-only backing indexes and skips synchronous mixed-index writes for them. |
| janusgraph-core/src/main/java/org/janusgraph/graphdb/database/log/TransactionLogHeader.java | Makes TransactionLogHeader.Modification constructor public for reuse in decoding. |
| janusgraph-core/src/main/java/org/janusgraph/graphdb/database/index/MixedIndexUpdateApplier.java | Adds backend-agnostic reindex-from-current-state applier for CDC worker. |
| janusgraph-core/src/main/java/org/janusgraph/graphdb/database/index/CdcElementChange.java | Adds normalized “element changed” model consumed by the applier/worker. |
| janusgraph-core/src/main/java/org/janusgraph/graphdb/configuration/GraphDatabaseConfiguration.java | Adds index.[X].cdc.enabled and index.[X].cdc.synchronous options. |
| janusgraph-cdc/src/test/resources/cassandra-cdc.yaml | Provides Cassandra config enabling CDC for full pipeline E2E test. |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/DebeziumCassandraJsonDecoderTest.java | Tests Debezium JSON decoding against real JanusGraph-serialized bytes. |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcWorkerConvergenceTest.java | Drives worker+decoder+applier via MockConsumer to validate convergence semantics. |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcWorkerConfigurationTest.java | Tests worker configuration defaults and properties parsing. |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcKafkaElasticsearchTest.java | Testcontainers E2E for Kafka → worker → ElasticSearch convergence. |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcIndexUpdateWorkerMainTest.java | Tests runner wiring and config reflection of CDC-enabled backing indexes. |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcIndexUpdateWorkerLoopTest.java | Unit-tests polling loop semantics (dedupe/retry/commit/rewind). |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcEventDecoderTest.java | Smoke test for decoder SPI and CdcElementChange interop. |
| janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcCassandraDebeziumElasticsearchTest.java | Full Cassandra CDC → Debezium → Kafka → ElasticSearch pipeline E2E (profile-gated). |
| janusgraph-cdc/src/test/java/io/debezium/connector/cassandra/JanusGraphCdcConnectorStarter.java | Test-only bridge to start Debezium Cassandra connector embedded. |
| janusgraph-cdc/src/main/java/org/janusgraph/cdc/DebeziumCassandraJsonDecoder.java | Implements Debezium Cassandra JSON → CdcElementChange decoding. |
| janusgraph-cdc/src/main/java/org/janusgraph/cdc/CdcWorkerConfiguration.java | Defines immutable worker/Kafka configuration and consumer properties. |
| janusgraph-cdc/src/main/java/org/janusgraph/cdc/CdcIndexUpdateWorkerMain.java | Adds standalone runner that opens graph, wires decoder+applier, starts workers. |
| janusgraph-cdc/src/main/java/org/janusgraph/cdc/CdcIndexUpdateWorker.java | Implements Kafka consume/decode/dedupe/apply/retry/commit/rewind loop. |
| janusgraph-cdc/src/main/java/org/janusgraph/cdc/CdcIndexApplier.java | Functional interface to abstract index application for worker tests. |
| janusgraph-cdc/src/main/java/org/janusgraph/cdc/CdcEventDecoder.java | Decoder SPI for CDC record formats. |
| janusgraph-cdc/pom.xml | New module POM with Kafka clients + testcontainers/Debezium profile gating. |
| docs/configs/janusgraph-cfg.md | Regenerates config reference including new CDC options. |
| docs/changelog.md | Adds 1.2.0 upgrade note describing CDC mixed index synchronization. |
| docs/advanced-topics/cdc-mixed-index.md | Adds operator guide for Cassandra CDC + Debezium + Kafka + worker setup. |
💡 Add Copilot custom instructions for smarter, more guided reviews. Learn how to get started.
a7d91de to
5659caf
Compare
48c32ba to
3e6cc0e
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 35 out of 35 changed files in this pull request and generated 4 comments.
Comments suppressed due to low confidence (1)
janusgraph-lucene/src/test/java/org/janusgraph/diskstorage/lucene/MixedIndexUpdateApplierTest.java:1
- The 10-second schema/index enablement timeout is likely to be flaky under loaded CI runners (especially with filesystem-backed Lucene + BerkeleyJE). Consider increasing these timeouts (e.g., 30–60 seconds) or using a shared constant aligned with other JanusGraph index-status tests. The same concern applies to other new tests using 10-second
awaitGraphIndexStatustimeouts.
3e6cc0e to
0ee776b
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 35 out of 35 changed files in this pull request and generated 3 comments.
Comments suppressed due to low confidence (1)
janusgraph-cql/src/test/java/org/janusgraph/diskstorage/cql/CQLCdcTableOptionTest.java:1
- This assertion is fairly broad and could pass if an unrelated substring containing 'cdc' appears in the DDL. Tightening it to assert the specific table option (e.g., matching
cdc = true/WITH cdc = true) would make the test more robust and less prone to false positives.
0ee776b to
0b2ce55
Compare
0b2ce55 to
de9c758
Compare
7aaf3dc to
411925e
Compare
411925e to
daa717d
Compare
a5f4bb0 to
d81f2b7
Compare
d81f2b7 to
1e709bf
Compare
1e709bf to
285d79f
Compare
dafefda to
6df6595
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 39 out of 39 changed files in this pull request and generated 1 comment.
Suppressed comments (1)
janusgraph-cdc/src/test/java/org/janusgraph/cdc/CdcCassandraDebeziumElasticsearchTest.java:303
- This describes the opposite path from the test setup.
knowshas default MULTI multiplicity and both endpoints survive, socdcCanIdentifyDeletedRelationdeliberately skips the synchronous delete; this removal therefore validates the real tombstone/decoder/worker path, not the synchronous carve-out.
// Edge add + remove against the REAL pipeline. The ADD must flow through Cassandra-CDC -> Debezium -> Kafka
// -> worker (cdc-only skips synchronous additions). The remove is satisfied primarily by the synchronous
// relation-document deletion carve-out at commit; the decoder's tombstone-decode path (edge identity
// recovered from the column alone, pinned by the decoder unit tests) remains the idempotent safety net.
Keep mixed indexes (ElasticSearch/Solr/Lucene) eventually consistent with the graph by deriving
their updates from a Change-Data-Capture stream of the committed graph data, instead of a
synchronous second write during the transaction that can diverge on failure and leave a
permanently stale index.
Pipeline (Apache Cassandra): commit (graph data only) -> Cassandra edgestore(cdc=true) ->
Debezium -> Kafka -> CdcIndexUpdateWorker consumer group -> reindex-from-current-state -> ES bulk.
Key design points:
- Reindex-from-current-state: the worker reads each changed element's current graph state and
fully replaces its index document (reusing IndexSerializer, like transaction recovery). This is
idempotent and order-independent, so out-of-order or duplicate events still converge to the
current state -- no strict ordering required, and a stale event can never overwrite a fresh
value. Worker transactions use skipDBCacheRead(): the worker JVM's database-level cache is
never invalidated by remote writers, so reads must hit the live graph.
- Additions are never dual-written: the only synchronous write is to storage; documents are
created and refreshed downstream from the committed change stream, so they cannot diverge.
- Deletions are routed by event identifiability, decided per deleted relation at commit:
* A removed vertex document is keyed by the vertex id -- the partition key every event carries,
including a whole-row partition delete -- so it is always removed by the worker.
* A removed MULTI-multiplicity edge with a surviving endpoint leaves an ordinary column
tombstone on that endpoint's row (its mirror copy) whose column carries the full edge
identity: the worker removes the document. In particular, removing a super node with
storage.drop-whole-row-on-vertex-removal while its neighbors survive costs one partition
delete and ZERO synchronous index operations -- the per-edge work happens asynchronously in
the CDC pipeline, where eventual consistency is the contract.
* Deletions no event can identify are written synchronously by the deleting transaction, the
only place their identities exist: constrained-multiplicity edges and meta-properties keep
their relation id in the storage value region (absent from tombstones), and unidirected
edges on removed vertices, self-loops, and edges whose both endpoints die in one transaction
have no surviving mirror. These carry dual-mode-grade durability; the docs recommend
tx.log-tx (whose PRECOMMIT entry durably records the deleted identities) plus transaction
recovery to close the crash window, and note that a graph-scanning REINDEX cannot remove
documents of already-deleted elements.
* Updates (a deleted relation whose id the same transaction re-adds) never delete
synchronously: the worker rewrites the still-live document id from current state, and a
delayed synchronous whole-document delete could otherwise erase that rewrite with no later
event to restore it.
* The reverse race -- the worker reads a relation as live, a concurrent transaction deletes it
and issues its synchronous document removal, then the worker's write lands last and would
resurrect the document permanently -- is closed by post-write verification: the applier
re-reads every relation document it wrote on a fresh snapshot and removes those whose
element vanished, which decides every interleaving correctly.
- Element-keyed Kafka partitioning gives per-element ordering and horizontal scaling via a
consumer group; batches are de-duplicated and applied in one transaction spanning all backing
indexes (one ElasticSearch _bulk per index), with VERTEX changes batch-preloaded via
getVertices(...) + multiQuery().properties().
- At-least-once: offsets are committed only after a batch is durably applied; on failure the
batch is reprocessed (rewind) rather than skipped, so the index eventually catches up. The
standalone runner supervises worker liveness and exits when every worker died unexpectedly, so
a process supervisor can restart it instead of leaving a healthy-looking zombie.
Configuration (opt-in, disabled by default):
- storage.cql.cdc: emit the Cassandra cdc=true table option on the edgestore table. Opening a
store logs a warning when the option is enabled but the live table lacks cdc=true (the option
only applies at table creation; an existing table must be ALTERed) -- otherwise that
misconfiguration is a silent no-capture drift.
- index.[X].cdc.enabled / index.[X].cdc.synchronous: per-index dual mode (write synchronously AND
via CDC) or cdc-only mode (skip synchronous additions; deletions no event can identify remain
synchronous). Both are GLOBAL_OFFLINE: the commit-side filter and the worker's index discovery
are one cluster-wide contract that per-instance values could split. cdc.synchronous=false
without cdc.enabled=true logs a warning instead of being silently inert.
Components:
- janusgraph-core: per-index CDC config options, the commit-side filters in StandardJanusGraph
(cdc-only mixed-index additions are filtered out at generation time; deletions choose their
filter per relation via the event-identifiability rule above; composite indexes are never
filtered; has2iMods/WAL/lock semantics are unchanged -- with a startup warning when cdc-only
mode is configured), and MixedIndexUpdateApplier (the backend-agnostic
reindex-from-current-state engine, covering vertex, edge and property-element mixed indexes,
and issuing removals for constraint-mismatched live vertices on graphs with user-settable ids,
where id reuse could leave a stale document from a previous incarnation). The restore paths
(ElementCategory.retrieve, IndexSerializer.removeElement) now accept custom String vertex ids
in addition to Long, RelationIdentifierUtils.findRelation no longer NPEs when a relation's
adjacent vertex has been removed, and findEdgeRelations returns an empty Iterable instead of
null. CdcElementChange documents the id contract for alternative capture sources (canonical
vertex ids; raw partition-representative RelationIdentifier endpoints).
- janusgraph-cql: the storage.cql.cdc table option and the live-table drift warning (no Kafka
dependency in production code).
- janusgraph-cdc (new module; core + kafka-clients 3.9.1, which carries the fixes for
CVE-2025-27817/27818): the CdcEventDecoder SPI, DebeziumCassandraJsonDecoder (parses the
relation header directly, so IN-direction edge columns and value-less delete tombstones of
MULTI edges resolve to the correct edge identity; canonicalizes partitioned-vertex ids for
VERTEX changes; skips corrupt records -- invalid JSON/Base64, malformed keys, and string-id
keys on graphs whose id regime forbids them -- while rethrowing transient backend failures so
the batch is redelivered), the CdcWorkerConfiguration (fail-fast validation incl. rejecting a
shared group.instance.id across multiple worker threads, which would fence itself forever;
cdc.consumer.* passthrough; auto.offset.reset defaults to "earliest" and the docs explain why
"latest" risks skipping events on a rebalance before a partition's first commit; pass-through
settings colliding with managed consumer keys are warned about instead of silently ignored),
the CdcIndexUpdateWorker (two-phase shutdown, interrupt-aware retries that stop when a
shutdown arrives mid-backoff, jittered exponential backoff so multiple workers do not retry
in lockstep, per-partition batch rewind on failure that tolerates partitions revoked
mid-rewind, a paced error loop, no consumer leaks), and the standalone
CdcIndexUpdateWorkerMain runner with worker-liveness supervision.
Testing: 80 tests across the touched modules, including unit/component coverage (decoder vs real
serialized bytes incl. poison-pill skips for invalid-Base64/malformed-key/string-id-keys,
delete-envelope after=null/before fallback, payload-wrapped envelopes, IN-direction columns and
value-less edge-delete tombstones, with assertions comparing full RelationIdentifier identity
incl. endpoint ids and the real meta-property id; reindex engine over vertex/edge/property-element
indexes incl. document removal when an element loses all indexed fields, removed-endpoint edges,
custom String vertex ids, stale-document removal after id reuse under a different label, the
partitioned-vertex id contract driven through the applier to real documents, multi-backing apply
and unknown-id-in-batch; worker loop via Kafka MockConsumer incl. multi-partition rewind,
no-op-batch offset commits, run()-loop recovery after a failed batch and Error-death
observability; commit-side deletion routing: synchronous removal for both-endpoints-removed,
constrained-multiplicity and meta-property deletions, no synchronous removal for mirror-identified
deletions and same-id replacements; composite indexes under cdc-only mode; full-chain convergence
over Lucene incl. vertex/edge add/update/remove and out-of-order delivery) and two real-container
E2Es -- worker -> Kafka -> ElasticSearch (runs on Java 8, 11 and 17), and the full Cassandra-CDC
-> Debezium -> Kafka -> ElasticSearch pipeline covering the vertex AND edge lifecycle
(add/update/property-removal/delete) against real Debezium delete envelopes, including a
whole-row both-endpoints vertex drop whose edge document converges with no per-edge event. The
Debezium pipeline test is gated behind the cassandra-cdc-e2e Maven profile (auto-activated on
JDK 17-23; Debezium 3.x requires 17+ and cassandra-all 4.1.7 does not run on 24+); the default
Java 8/11 build excludes only the two Debezium-dependent sources and stays green.
CI: a dedicated workflow (.github/workflows/ci-cdc.yml) runs the cdc unit tests plus the real
Kafka+ElasticSearch worker test on Java 8 and 11 (Testcontainers' optional jna dependency is
re-added at test scope, matching janusgraph-cql/janusgraph-es), and the full real-container
suite -- including the Cassandra-CDC -> Debezium -> Kafka -> ElasticSearch pipeline -- on Java 17
with Docker, so the integration is exercised on every change and guards against regression.
Docs: advanced-topics/cdc-mixed-index.md operator guide (incl. the deletion-routing design and
its tx.log-tx recommendation, cdc_raw backpressure and the CDC disable/teardown path, systematic
RF>1 event duplication, tuning guidance for the retry budget vs Kafka's max.poll.interval.ms
during long index outages, the auto.offset.reset caveat, and document-TTL limits), a 1.2.0
changelog upgrade note, and the regenerated configuration reference.
Fixes JanusGraph#4873
Replaces JanusGraph#4874
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_011rNFck9BY9s58qW3XQ1YTK
Signed-off-by: Oleksandr Porunov <alexandr.porunov@gmail.com>
6df6595 to
eea465c
Compare
There was a problem hiding this comment.
Pull request overview
Copilot reviewed 39 out of 39 changed files in this pull request and generated no new comments.
Suppressed comments (1)
janusgraph-cdc/pom.xml:60
- This comment says the Kafka/Elasticsearch component test requires Java 17+, but the same POM later states that
CdcKafkaElasticsearchTestruns on Java 8/11/17, and the new workflow executes it on Java 8 and 11. Update the comment so it does not misstate the test's runtime requirement.
<!-- Component E2E: real ElasticSearch + Kafka (Testcontainers, require Java 17+ to run) -->
|
I believe this PR is now safe to be merged. The CDC functionality is fully gated behind these new feature flags, so existing users that don't use CDC won't have any differences. |
Keep mixed indexes (ElasticSearch/Solr/Lucene) eventually consistent with the graph by deriving their updates from a Change-Data-Capture stream of the committed graph data, instead of a synchronous second write during the transaction that can diverge on failure and leave a permanently stale index.
Pipeline (Apache Cassandra): commit (graph data only) -> Cassandra edgestore(cdc=true) ->
Debezium -> Kafka -> CdcIndexUpdateWorker consumer group -> reindex-from-current-state -> ES bulk.
Key design points:
storage.drop-whole-row-on-vertex-removalwhile its neighbors survive therefore costs one partition delete and zero synchronous index operations -- the per-edge work happens asynchronously in the CDC pipeline.tx.log-tx+ transaction recovery to close the crash window._bulkper backing index. The standalone runner supervises worker liveness and exits when all workers died so a process supervisor can restart it.Configuration (opt-in, disabled by default):
storage.cql.cdc: emit the Cassandracdc=truetable option on the edgestore table (with a startup warning when the live table lacks it, since the option only applies at table creation).index.[X].cdc.enabled/index.[X].cdc.synchronous: per-index dual mode (write synchronously AND via CDC) or cdc-only mode (skip synchronous additions; event-unidentifiable deletions remain synchronous). Both GLOBAL_OFFLINE -- the commit-side filter and the worker's index discovery are one cluster-wide contract.Components:
storage.cql.cdctable option and live-table drift warning (no Kafka dependency in production code).cdc.consumer.*passthrough), CdcIndexUpdateWorker, and the standalone CdcIndexUpdateWorkerMain runner.Testing: 80 tests across the touched modules -- decoder against real serialized bytes (incl. poison-pill skips, delete-envelope before/after images, payload-wrapped envelopes, full RelationIdentifier identity assertions), the reindex engine over vertex/edge/meta-property indexes (incl. stale-document removal after id reuse, the partitioned-vertex id contract, custom String vertex ids), the worker loop via Kafka MockConsumer (multi-partition rewind, no-op-batch offset commits, run()-loop recovery, Error-death observability), commit-side deletion routing for every class, full-chain convergence over Lucene, and two real-container E2Es: worker -> Kafka -> ElasticSearch (runs on Java 8, 11 and 17), and the full Cassandra-CDC -> Debezium -> Kafka -> ElasticSearch pipeline covering the vertex AND edge lifecycle against real Debezium delete envelopes, including a whole-row both-endpoints vertex drop. The Debezium pipeline test is gated behind the
cassandra-cdc-e2eMaven profile (auto-active on JDK 17-23); dedicated CI workflow runs Java 8/11 (unit + Kafka/ES containers) and Java 17 (full pipeline).Docs:
advanced-topics/cdc-mixed-index.mdoperator guide (deletion-routing design,tx.log-txrecommendation,cdc_rawbackpressure and CDC teardown, RF>1 duplication, retry-budget vsmax.poll.interval.mstuning,auto.offset.resetcaveat, document-TTL limits), a 1.2.0 changelog upgrade note, and the regenerated configuration reference.Fixes #4873
Replaces #4874
Thank you for contributing to JanusGraph!
In order to streamline the review of the contribution we ask you
to ensure the following steps have been taken:
For all changes:
master)?For code changes:
For documentation related changes: