Skip to content

Commit cbd6fc2

Browse files
porunovclaude
andcommitted
feat: CDC-based mixed index synchronization (#4873)
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 #4873 Replaces #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>
1 parent ac0eb23 commit cbd6fc2

39 files changed

Lines changed: 5781 additions & 9 deletions

.github/workflows/ci-cdc-dummy.yml

Lines changed: 47 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,47 @@
1+
# Copyright 2026 JanusGraph Authors
2+
#
3+
# Licensed under the Apache License, Version 2.0 (the "License");
4+
# you may not use this file except in compliance with the License.
5+
# You may obtain a copy of the License at
6+
#
7+
# http://www.apache.org/licenses/LICENSE-2.0
8+
#
9+
# Unless required by applicable law or agreed to in writing, software
10+
# distributed under the License is distributed on an "AS IS" BASIS,
11+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
# See the License for the specific language governing permissions and
13+
# limitations under the License.
14+
15+
# Companion to ci-cdc.yml: when a change touches only the paths that ci-cdc.yml ignores (docs, etc.),
16+
# the real workflow does not run. This dummy reports the same status check names as success so required
17+
# checks are satisfied. Keep the matrix names in sync with ci-cdc.yml.
18+
19+
name: CI CDC
20+
21+
on:
22+
pull_request:
23+
paths:
24+
- 'docs/**'
25+
- '.github/workflows/ci-docs.yml'
26+
- '.github/ISSUE_TEMPLATE/**'
27+
- 'requirements.txt'
28+
- 'mkdocs.yml'
29+
- 'docs.Dockerfile'
30+
- '*.md'
31+
32+
jobs:
33+
tests:
34+
name: ${{ matrix.name }}
35+
runs-on: ubuntu-22.04
36+
strategy:
37+
matrix:
38+
include:
39+
- name: cdc-java8
40+
java: 8
41+
- name: cdc-java11
42+
install-args: "-Pjava-11"
43+
java: 11
44+
- name: cdc-e2e-java17
45+
java: 17
46+
steps:
47+
- run: 'echo "No build required"'

.github/workflows/ci-cdc.yml

Lines changed: 91 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,91 @@
1+
# Copyright 2026 JanusGraph Authors
2+
#
3+
# Licensed under the Apache License, Version 2.0 (the "License");
4+
# you may not use this file except in compliance with the License.
5+
# You may obtain a copy of the License at
6+
#
7+
# http://www.apache.org/licenses/LICENSE-2.0
8+
#
9+
# Unless required by applicable law or agreed to in writing, software
10+
# distributed under the License is distributed on an "AS IS" BASIS,
11+
# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
# See the License for the specific language governing permissions and
13+
# limitations under the License.
14+
15+
name: CI CDC
16+
17+
on:
18+
pull_request:
19+
paths-ignore:
20+
- 'docs/**'
21+
- '.github/workflows/ci-docs.yml'
22+
- '.github/ISSUE_TEMPLATE/**'
23+
- 'requirements.txt'
24+
- 'mkdocs.yml'
25+
- 'docs.Dockerfile'
26+
- '*.md'
27+
push:
28+
paths-ignore:
29+
- 'docs/**'
30+
- '.github/workflows/ci-docs.yml'
31+
- '.github/ISSUE_TEMPLATE/**'
32+
- 'requirements.txt'
33+
- 'mkdocs.yml'
34+
- 'docs.Dockerfile'
35+
- '*.md'
36+
branches-ignore:
37+
- 'dependabot/**'
38+
39+
env:
40+
# (The ElasticSearch test container heap is capped by JanusGraphElasticsearchContainer itself, which hard-codes
41+
# ES_JAVA_OPTS for the container -- a workflow-level ES_JAVA_OPTS env var would not reach it.)
42+
BUILD_MAVEN_OPTS: "-DskipTests=true --batch-mode --also-make"
43+
VERIFY_MAVEN_OPTS: "-Pcoverage"
44+
45+
jobs:
46+
tests:
47+
name: ${{ matrix.name }}
48+
runs-on: ubuntu-22.04
49+
strategy:
50+
fail-fast: false
51+
matrix:
52+
include:
53+
# Java 8 / 11: the module compiles and runs its unit tests plus the real Kafka + ElasticSearch
54+
# worker test (Testcontainers; Docker is available on the runners, as for the other backend suites).
55+
# Only the Debezium-dependent pipeline test sources are excluded from compilation here (Debezium 3.x
56+
# is Java 17 bytecode -- see janusgraph-cdc/pom.xml); they run only on the Java 17 job below.
57+
- name: cdc-java8
58+
java: 8
59+
- name: cdc-java11
60+
install-args: "-Pjava-11"
61+
java: 11
62+
# Java 17: the cassandra-cdc-e2e Maven profile auto-activates (Debezium 3.x requires Java 17+;
63+
# bounded below JDK 24 where cassandra-all 4.1.7 does not run), re-including and running the full
64+
# Cassandra-CDC -> Debezium -> Kafka -> ElasticSearch pipeline against Docker.
65+
- name: cdc-e2e-java17
66+
java: 17
67+
steps:
68+
- uses: actions/checkout@v4
69+
with:
70+
fetch-depth: 2
71+
- uses: actions/cache@v4
72+
with:
73+
path: ~/.m2/repository
74+
key: ${{ runner.os }}-maven-${{ hashFiles('**/pom.xml') }}
75+
restore-keys: |
76+
${{ runner.os }}-maven-
77+
- uses: actions/setup-java@v4
78+
with:
79+
java-version: ${{ matrix.java }}
80+
distribution: zulu
81+
- run: mvn clean install --projects janusgraph-cdc ${{ env.BUILD_MAVEN_OPTS }} ${{ matrix.install-args }}
82+
- run: mvn verify --projects janusgraph-cdc ${{ env.VERIFY_MAVEN_OPTS }} ${{ matrix.install-args }}
83+
- uses: actions/upload-artifact@v4
84+
with:
85+
name: jacoco-reports-cdc-${{ matrix.name }}
86+
path: target/jacoco-combined.exec
87+
- uses: codecov/codecov-action@v5
88+
env:
89+
CODECOV_TOKEN: ${{ secrets.CODECOV_TOKEN }}
90+
with:
91+
name: codecov-cdc-${{ matrix.name }}

NOTICE.txt

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,13 +22,15 @@ It also includes software from other open source projects including, but not lim
2222
* Apache Groovy [http://groovy-lang.org/]
2323
* Apache HBase [https://hbase.apache.org/]
2424
* Apache Hadoop [https://hadoop.apache.org/]
25+
* Apache Kafka [https://kafka.apache.org/]
2526
* Apache Kerby [https://github.com/apache/directory-kerby]
2627
* Apache Log4j [https://logging.apache.org/log4j]
2728
* Apache Lucene [https://lucene.apache.org/]
2829
* Apache Solr [https://lucene.apache.org/solr/]
2930
* Apache TinkerPop [https://tinkerpop.apache.org/]
3031
* Astyanax [https://github.com/Netflix/astyanax]
3132
* DataStax Driver for Apache Cassandra [https://github.com/datastax/java-driver]
33+
* Debezium [https://debezium.io/]
3234
* EasyMock [http://easymock.org/]
3335
* Elasticsearch [https://www.elastic.co/]
3436
* Google Cloud Bigtable [https://github.com/googlecloudplatform/cloud-bigtable-client]
@@ -46,6 +48,7 @@ It also includes software from other open source projects including, but not lim
4648
* Reflections8 [https://github.com/aschoerk/reflections8]
4749
* SLF4J [https://www.slf4j.org/]
4850
* Spatial4j [https://github.com/locationtech/spatial4j]
51+
* Testcontainers [https://testcontainers.com/]
4952
* Vavr [https://www.vavr.io/]
5053

5154
=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=

0 commit comments

Comments
 (0)