fix(server): converge schema caches across HStore servers - #3237
bitflicker64 wants to merge 9 commits into
Conversation
Codecov Report❌ Patch coverage is Additional details and impacted files@@ Coverage Diff @@
## master #3237 +/- ##
============================================
+ Coverage 41.35% 41.49% +0.14%
- Complexity 7299 7376 +77
============================================
Files 802 804 +2
Lines 69688 69991 +303
Branches 9291 9337 +46
============================================
+ Hits 28816 29044 +228
- Misses 37576 37630 +54
- Partials 3296 3317 +21 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
e270112 to
ec6f9ed
Compare
A schema change made through one Server could stay invisible on the
others indefinitely: updates and status flips publish no PD event, and
the PD watch does not replay events lost while it reconnects.
Every schema change now writes a per-graph opaque version to PD at
HUGEGRAPH/{cluster}/SCHEMA_VERSION/{graphspace}/{graph}. Each Server
checks it once per schema.sync.reconcile_interval (10s by default) and
clears that graph's schema cache when it changed, so a change becomes
visible everywhere within about one interval without relying on events.
A burst of changes costs each Server at most one clear per interval.
Existing events are unchanged.
A clear caused by a remote change and a reader's cache update are now
exclusive on the graph's SchemaCaches, and a reader skips its update if
such a clear ran after it read, so data read before the clear can't
be cached after it.
New options: schema.sync.enabled (default true) and
schema.sync.reconcile_interval (0 to 3600 seconds, 0 disables checks).
Fixes apache#3235
ec6f9ed to
a4f7ae1
Compare
imbajin
left a comment
There was a problem hiding this comment.
Blocking: no newly verified runtime regression. Score: 8.5/10 after five independent review lanes and synthesis; please clarify the rollout preconditions below before release. This assessment uses code/test inspection and CI evidence, not an independently rerun multi-Server experiment; current CI is not fully green.
The convergence bound needs every Server of the graph on a release with these options, sync enabled, polling on and PD reachable. Say so in the conf template, with what to do after a rolling upgrade: older Servers write no version for schema updates.
Per review on apache/hugegraph#3237: the convergence bound needs every server on a release with these options, sync enabled, polling on and PD reachable, and a schema change or restart after a rolling upgrade.
There was a problem hiding this comment.
Changes requested: revisit the propagation design
The earlier 9/10 recommendation is withdrawn pending design agreement. The burst benchmark is useful, but does not establish that routine, low-frequency DDL needs continuous polling.
| Current PR | Alternative to evaluate | |
|---|---|---|
| Normal update/status change | Write a version; peers detect it on their next poll | Notify peers; ACK after successful cache invalidation |
| No Schema changes | Every Server checks each open graph periodically | No periodic Schema-version reads required by this mechanism |
| Missed delivery | Next successful version poll detects the change | PD retries unacknowledged updates, coalesced to the latest version per graph |
The current PR still preserves existing add notifications. The alternative is a design proposal, not an implemented guarantee.
Please confirm before implementation
- Commit and failover. establish a reliable link between the Schema commit and version update. Recover the latest version and confirmation state after PD leadership changes; duplicate delivery is acceptable, lost updates are not.
- Bounded retry. retry asynchronously without blocking DDL. Expire delivery attempts for long-offline instances; distinguish restarted instances from old sessions.
- Rejoin before serving. an expired instance must reconcile with PD and invalidate/reload stale Schema before accepting graph requests. Do not assume reconnect already guarantees this.
Please compare the smallest notification/ACK/recovery design against polling under realistic DDL frequency and graph/Server counts. Restore normal update/status notification as part of that evaluation; batch redundant invalidations within one operation where practical.
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
An unresolved critical stale-cache race and moderate lifecycle and retry-state issues must be fixed.
Get a fresh assessment by requesting another Copilot review.
Review effort: Lite
Findings: 1
What changed in this PR
This PR adds bounded schema-cache convergence across HStore servers using PD-backed per-graph versions and periodic reconciliation.
Changes:
- Adds schema version writes, reconciliation, retries, and configuration options.
- Protects cache updates during remote invalidation.
- Adds lifecycle cleanup, documentation, and comprehensive tests.
| File | Summary |
|---|---|
hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/UnitTestSuite.java |
Registers the new tests. |
hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/core/GraphManagerDropGraphTest.java |
Tests version-key cleanup during graph drops. |
hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/SchemaVersionReconcilerTest.java |
Tests reconciliation, retries, failures, and scheduling. |
hugegraph-server/hugegraph-test/src/main/java/org/apache/hugegraph/unit/cache/CachedSchemaTransactionTest.java |
Tests cache guards, options, metadata, and lifecycle behavior. |
hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties |
Documents synchronization configuration. |
hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/StandardHugeGraph.java |
Manages reconciler shutdown and version-key cleanup. |
hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/MetaManager.java |
Exposes schema-version metadata operations. |
hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/meta/managers/GraphMetaManager.java |
Implements schema-version key handling. |
hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/config/CoreOptions.java |
Defines synchronization options. |
hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/SchemaVersionReconciler.java |
Implements periodic PD reconciliation and retries. |
hugegraph-server/hugegraph-core/src/main/java/org/apache/hugegraph/backend/cache/CachedSchemaTransactionV2.java |
Integrates version bumps and guarded cache updates. |
hugegraph-server/hugegraph-api/src/main/java/org/apache/hugegraph/core/GraphManager.java |
Removes version metadata during graph drops. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
addSchema() and updateSchema() cached the written element without the generation check the read paths use. A remote clear landing between the storage write and the cache update could then be undone by an element another server had already replaced. Capture the generation before the write and update the caches under the arrayCaches lock only if it is unchanged.
070a961 skipped the cache update of a local write if the generation changed during it. The reconciler also clears on this server's own version writes, so a reader could cache the old element after that clear and the skipped update left it there: an index label stayed CREATING in the cache after its rebuild set CREATED, failing two hstore core tests. Invalidate the element and reset the cached-all flag of its type instead, so the next read loads it from storage.
imbajin
left a comment
There was a problem hiding this comment.
Blocking: yes. Summary: A schema mutation can be committed without publishing the version that tells peers to invalidate their caches, leaving them on stale schema after a writer crash. Evidence: CachedSchemaTransactionV2 writes the schema before a separate version update; SchemaVersionReconciler.run() skips when the stored version still matches applied.
| this.updateCache(schema); | ||
| this.updateCache(schema, generation); | ||
|
|
||
| this.bumpSchemaVersion(); |
There was a problem hiding this comment.
super.addSchema() before this separate PD version write; updateSchema() and removeSchema() use the same ordering. If the writer exits between those calls, peers keep the previously applied version and SchemaVersionReconciler skips them, so they can continue serving stale schema until another mutation or a peer restart. Please make the mutation notification recoverable across this crash window, for example by atomically storing a revision or publishing a repair revision during graph startup, and add a recovery test.
bitflicker64
left a comment
There was a problem hiding this comment.
Blocking: no. Summary: One new Minor finding: the version key deleted on graph drop can be written again before the graph's reconciler stops, so the key the drop path is meant to remove can survive in PD. The design question and the writer-crash window are already open in earlier reviews and are not repeated here. Evidence: static reading of GraphManager.dropGraph(), StandardHugeGraph.drop()/close() and SchemaVersionReconciler.run()/stop() at 7f5672a; latest-head CI checks all pass.
| g.clearBackend(); | ||
| try { | ||
| // The schema version kept in PD for the schema cache sync | ||
| this.metaManager.deleteSchemaVersion(graphSpace, name); |
There was a problem hiding this comment.
Minor: The version key is deleted while the graph's reconciler is still scheduled. SchemaVersionReconciler is stopped only later, in g.close() via stopReconciler(). If any bump() during the drop failed (for example the one from CachedSchemaTransactionV2.clear() inside clearBackend()), pendingWrite is true, and a tick that runs between this delete and g.close() passes the graphClosed check and writes a fresh version at SchemaVersionReconciler.run() line 134. stop() uses future.cancel(false), so a tick already running also completes its write. The key then outlives the graph. StandardHugeGraph.drop() (line 1172) has the same order: delete first, close() last. Please delete the key after the reconciler is stopped, for example by moving the delete after g.close() here and after this.close() in drop(), and extend GraphManagerDropGraphTest to cover a pending write at drop time.
Design update: push from the commit plus a Server-held lease, replacing the 10 s poll@imbajin Following your review (O(1) notification, not polling) and your design doc of 2026-09-26, I redid the design from scratch instead of tuning the poll. This comment is the proposed design and how it was reached. Nothing new is pushed to this branch yet; I would like your agreement on the design and on the questions at the end before implementing it here in stages. How the design was reached
The design
Staleness bound. If a DDL is acknowledged at time T, then from T + 12 s every Server either serves it or answers 503 on that graph. The bound depends only on the Server's own clock rate. It holds under a Server stall, half-open TCP, a PD leader change including a partitioned old leader, and crashes. Requests already running keep the schema they started with. MeasuredPrototype of the PD TXN with the record, the ring and sender,
One claim failed. A SIGSTOP of the first PD leader after cluster start leaves no new leader for about 11 s, and all 3 Servers fenced for 0.6 to 1.8 s. The cause is jraft's pre-vote, which blocks on the connect timeout while holding the node lock. PD sets that timeout from Idle cost at 20 Servers is 72,000 small echo pairs per hour whatever the number of graphs, against 720,000 PD reads per hour for this PR's current poll with 100 open graphs. A DDL costs 1 raft entry instead of 4 to 10 puts today. Master bugs found on the way (fixed by this design)
Where it differs from your doc
Not in this PR: barrier DDLA wait-for-all-Servers step before index rebuild and label or index removal is designed, but it failed every review round. The last gap is a Store write that the client gave up on but Store still applies after the barrier; only Store-side fencing closes that. Cross-Server rebuild and removal are also unsafe on master today, independent of the cache. So it is a separate proposal, as your TODO 2 frames it. In this PR those DDLs still commit atomically and reach every Server within the 12 s bound, and a Server that has seen DELETING rejects new writes to that label. Stages on this branch
No stage ships sessions without the lease and the resync that recover them. Questions
The full specification is about 20,000 words; I will put the parts you want into the PR body and the doc PR (apache/hugegraph-doc#498) once the design is agreed. The blocking point in your inline comment (a writer crash between the schema write and the version write) goes away with the atomic TXN; I will answer it there. |
|
@imbajin A second, independent review of the design above, checked against the code, led to these changes. Nothing else in the design moves.
Next I start stage A on this branch. |
HstoreSession.truncate() called resetIdByKey() for the store's graph
name. For the schema store that is {graphspace}/{graph}/m, so PD dropped
every schema id counter of the graph by prefix, while the schema itself
stays in PD meta and survives the truncate. The truncating Server kept
handing out ids from its cached range, but any other Server, or the same
one after a restart, got ids from 1 again and reused the ids of existing
schema elements.
Stop resetting the counters on truncate.
clearGraphSpace() deletes HUGEGRAPH/{cluster}/GRAPHSPACE/{gs}/AUTH by
prefix before it drops the graphs. Graph keys sit beside AUTH under the
graphspace, so the prefix also matched every key of a graph whose name
starts with AUTH, such as AUTHx, and wiped its schema and config before
the graph's own drop ran.
End the auth prefix with the path delimiter.
IdMetaStore.getId() reads the counter from the local RocksDB and then writes the next value through raft. The read is outside raft, so a leader elected a moment ago could read the counter before it had applied its predecessor's last write and hand out the same id range twice. getId() now waits for a jraft ReadIndex first: the read index is confirmed by a quorum in the current term, and the wait ends once the local applied index reaches it. EAGAIN (no entry of the new term committed yet) and EBUSY (leadership transfer) are retried 5 times, 20 ms apart; the wait is bounded by raft.rpc-timeout. Stores that are not replicated skip the wait.
jraft calls onSnapshotSave() on the state machine thread with the snapshot index set to the last applied index, and applies the next entry as soon as it returns. PD handed the whole save to a job thread, so the RocksDB checkpoint could include entries past the snapshot index, and a node installing that snapshot applied them a second time. Take the checkpoint before returning; only the compression stays on the job thread. The checkpoint is hard links under the store's write lock. A failed checkpoint now completes the snapshot once with an error; it used to report the error and then continue and report success as well.
PD set jraft's connect, default and install-snapshot rpc timeouts all from raft.rpc-timeout (10 s). A candidate opens a connection to each peer while it holds the node lock, and a peer that accepts the TCP connection but never answers the ping (a stopped process, a lost host) stalls the election for the whole connect timeout. Losing the host of the first leader after cluster start left the PD cluster without a leader for about 11 s. Add raft.rpc-connect-timeout (default 1000 ms, jraft's own default). Split out raft.rpc-install-snapshot-timeout as well, so the connect timeout can be lowered without shortening snapshot installs; its default 0 keeps the current value, raft.rpc-timeout. raft.rpc-timeout keeps its meaning for other raft requests.
|
Stage A of the design above is pushed:
Tests on a local 1 PD, 1 Store run, Java 11, 2026-09-27: PD common 83 passed, PD core 118 passed (2 skipped as before), unit 776 passed, hstore core 817 passed. Each commit adds tests that fail with its fix reverted (5 PD core failures, and the truncate and auth tests fail). What is left on this PR:
Separate from this PR: barrier DDL with Store-side fencing (Q4), and the Store schema cache fix (Q3). |
imbajin
left a comment
There was a problem hiding this comment.
Blocking: yes. Summary: The previously reported schema-write/version-publication crash window remains unresolved; this review adds a distinct ReadIndex retry window that can fail ID allocation during leader recovery. Evidence: exact-head static review; latest-head CI/status checks passed; no local tests were run.
| */ | ||
| private static final long ALIVE_PEERS_REFRESH_MS = 1000L; | ||
|
|
||
| private static final int READ_INDEX_RETRIES = 5; |
There was a problem hiding this comment.
waitReadIndex() waits for rpcTimeout (10 seconds by default), but this code schedules only five retries 20 ms apart; after roughly 100 ms, readIndex() completes with UNKNOWN_VALUE. A newly elected JRaft leader returns EAGAIN until an entry in its current term commits, so a quorum that needs longer than 100 ms makes IdMetaStore.getId() fail even while the configured timeout remains. Please retry EAGAIN/EBUSY until the timeoutMs deadline and add coverage for recovery after the current retry window.


Purpose of the PR
With several Servers on HStore, a schema change made through one Server can stay invisible on the others indefinitely. Updates and status flips publish no PD event, on purpose, because one event per
updateSchemaStatus()call is a broadcast storm. Events that are published can be lost, because the PD watch starts from "now" after a reconnect (#3157). This PR makes every Server converge within a bounded time without adding any event.The design and its open questions are in discussion #3205. The questions had no maintainer answer yet, so this PR uses conservative defaults. Each one is listed in #3235 with how to override it.
Main Changes
CachedSchemaTransactionV2writes a new opaque version of the graph to PD after every committed schema change: add, update, status flip, remove, clear, store clear or truncate. The key isHUGEGRAPH/{cluster}/SCHEMA_VERSION/{graphspace}/{graph}. It sits outsideGRAPHSPACE/{gs}/{graph}/SCHEMAbecauseclearAllSchema()deletes that prefix.SchemaVersionReconciler: one task per open graph on one daemon thread per JVM. Everyschema.sync.reconcile_intervalseconds it reads the version. If the version changed, it clears that graph's schema cache and adopts the value it read before the clear. Versions are only compared for equality.SchemaCaches. A reader skips its update if such a clear ran after its storage read, so data read before the clear is never cached after it. Cache hit paths are unchanged, except that an id-cache hit takes the lock once when it promotes the element to the array cache.StandardHugeGraph.close()stops the graph's reconciler. Dropping a graph deletes its version key, inStandardHugeGraph.drop()and, for the PD mode drop that callsclearBackend()andclose()instead, inGraphManager.dropGraph().schema.sync.enabled(defaulttrue) andschema.sync.reconcile_interval(default10, range 0 to 3600 seconds, 0 turns the checks off). Withschema.sync.enabled=falsethere is no reconciler and no PD traffic, except one delete of the version key when a graph is dropped in PD mode. The cache guard above stays on either way.Bounds, from the measurements below:
ceil(D / T) + 1clears per Server, so cluster-wide clears no longer grow with M. Notifying on every update costsM x (N - 1).Verifying these changes
SchemaVersionReconcilerTest(16 tests) covers the reconciler with an in-memory store:CachedSchemaTransactionTesthas 8 new tests:getSchema(type, id)andgetAllSchemapaths don't cache data across a remote clear injected during the read;mvn test -pl hugegraph-server/hugegraph-test -am -P unit-teston Java 11.0.32 ata4f7ae1c: 774 run, 0 failures, 1 skipped. Master83ef9f3fahas 748 run, 0 failures, 1 skipped.mvn editorconfig:checkis clean on the changed modules.notifySchemaCacheClear()inupdateSchema(), and an INFO line per remote clear so the clears can be counted.83ef9f3fae2701127)user_datato it (no event)The live runs used
e2701127. The current heada4f7ae1cadds two fixes from review on hugegraph#236, both in the graph drop path, which none of these measurements touch:drop()readsschema.sync.enabledfrom the configuration, and the PD mode drop also deletes the version key (GraphManagerDropGraphTest). A second commit,be932423, documents the rollout conditions from review here; it changes only the conf template comment.The 14.40 s is one skipped tick. Server-2's own PD client had not reconnected yet when its tick ran, about 3 s after the write landed through server-0. The next tick picked the change up. From that Server's log:
Known limits, also in #3235:
schema.sync.enabled=true,schema.sync.reconcile_intervalabove 0 and PD reachable. During a rolling upgrade an older Server writes no version for schema updates, and upgrading it later does not publish the missed change. After the last Server is upgraded, make one schema change or restart the Servers, so every Server clears its schema cache and reloads it from PD. The conf template and doc(config): add schema.sync options hugegraph-doc#498 say this.schema.sync.enabled=falsewrites no version, so its updates are not seen by the others.Does this PR potentially affect the following parts?
Documentation Status
Doc - TODO: required documentation is pending; complete it before merging.Doc - Done: documentation is included here or linked below.Doc - No Need: no user-visible documentation is affected.Documentation files in this PR or paired hugegraph-doc PR:
hugegraph-server/hugegraph-dist/src/assembly/static/conf/graphs/hugegraph.properties(commented examples).schema.sync.enabledandschema.sync.reconcile_intervalto the EN and CN config option pages.