Skip to content

Cooperative embedding#690

Draft
jshook wants to merge 12 commits into
mainfrom
cooperative-embedding
Draft

Cooperative embedding#690
jshook wants to merge 12 commits into
mainfrom
cooperative-embedding

Conversation

@jshook

@jshook jshook commented Jul 2, 2026

Copy link
Copy Markdown
Contributor

This includes some changes we need to make the embedding surface more robust around shared resource usage and scoping. I've shared it here for visibility and will move it from draft status once our integrated testing bears fruit.

This is now rebased on top of Aaron's compaction improvement fixes, so testing of this branch is all-inclusive of both

Synopsis

Make the on-disk graph compactor embeddable (cooperative resource sharing). The shape of the interfaces types provided here were carefully selected to be useful and compatible with an embedding system like Cassandra or OpenSearch, without being overly specific to any. They were also chosen to be compatible with Java 11 onward.

Purpose

OnDiskGraphIndexCompactor currently runs on its own thread pool, is invisible while it works, and always writes to its own file. This PR adds a few small, optional extension points so a host system (e.g. a database's compaction pipeline) can drive the merge cooperatively — on the host's own threads, under the host's observation and throttling, writing straight into the host's own file. Everything is additive and @experimental; with nothing supplied, behavior and output are unchanged.

Key elements

  • Bring-your-own executor. The compactor now takes any Executor plus an explicit taskWindowSize instead of requiring a ForkJoinPool. Passing a caller-runs executor (Runnable::run) runs the whole merge on the calling thread, so a host reuses its existing compaction threads with no extra pool.
  • Progress + throttling (ProgressLimiter). One small SPI the host installs to (a) observe merge progress and (b) block/pace the bytes the merge writes against a shared budget — without jvector knowing anything about the host's limiter. Ships with two ready-made, composable implementations:
    • a leaky-bucket rate meter
    • a logging wrapper
  • No-copy output (CompactionDestination / compact(Path, startOffset)). Write the compacted graph body directly into the host's container file after a reserved header — with a commit-on-success / discard-on-failure lifecycle — instead of writing a temp file and copying it. SeekableSink is the small primitive used to address that file region.

Tests and docs/compaction.md cover the new surface.

@github-actions

github-actions Bot commented Jul 2, 2026

Copy link
Copy Markdown
Contributor

Before you submit for review:

  • Does your PR follow guidelines from CONTRIBUTIONS.md?
  • Did you summarize what this PR does clearly and concisely?
  • Did you include performance data for changes which may be performance impacting?
  • Did you include useful docs for any user-facing changes or features?
  • Did you include useful javadocs for developer oriented changes, explaining new concepts or key changes?
  • Did you rebase your branch onto the latest main for regression testing and PR submission?
  • Did you trigger regression testing via Run Bench Main and review results?
  • Did you adhere to the code formatting guidelines (TBD)
  • Did you group your changes for easy review, providing meaningful descriptions for each commit?
  • Did you ensure that all files contain the correct copyright header?
  • Did you add documentation for this feature to the release notes directory?

If you did not complete any of these, then please explain below.

…uteLayerInfoFromSources

computeLayerInfoFromSources called getNodes(0) on each source graph to count live nodes
at level 0. getNodes(0) sequentially seeks through every node record on disk to filter
out deleted entries. On a cold page cache this touches large amounts of source data before
compaction even begins, significantly delaying the start of actual graph merging.

Since every live node is present at level 0 by the HNSW invariant, the count is simply
liveNodes.get(s).cardinality() — an in-memory popcount requiring no I/O.

Also switch PQ retraining from ProductQuantization.compute() (full k-means++ init)
to basePQ.refine() (Lloyd's iterations only, warm-started from the existing codebook).
The source codebooks are already trained on the same distribution, so warm-starting
converges in far fewer passes with no recall loss.
@jshook
jshook force-pushed the cooperative-embedding branch from fa4de43 to 22d1de3 Compare July 2, 2026 22:24
dian-lun-lin and others added 11 commits July 16, 2026 15:31
Source graphs are mapped with MADV_RANDOM (correct for search-time access),
which disables kernel readahead and makes compaction's bulk phases fault one
page at a time on a cold cache (measured: retrain ~37s, pre-encode ~42s on a
disk-cold 10M-node compaction that takes ~3s/~8s warm).

Adds ReaderSupplier.prefetch(offset, length) — streams a byte range into the
page cache through a separate readahead-enabled descriptor — and
OnDiskGraphIndex.prefetchL0Records(minNode, maxNode) on top of it. Each bulk
phase warms exactly the records it is about to read, from the worker that
will read them:
- PQ retrain prefetches its (source, node)-sorted sample ranges before
  extraction (bounded by training-set size, not file size),
- code pre-encode prefetches each chunk's records at task start,
- L0 batch processing prefetches each batch's own records at task start
  (cross-source search reads are data-dependent and stay demand-faulted).

Transient cache demand is proportional to the in-flight windows, so there is
no up-front whole-file streaming pass and no memory-availability gate to
mistune: on a box that cannot hold the sources, pages are simply evicted and
reads degrade per-page to the old fault-on-demand behavior.
The L0 cross-source candidate search passed beamWidth (= 2x searchTopK) as
rerankK, doubling the approximate-phase beam over what the candidate budget
needs. A seeding-vs-beam decomposition study (7 paired disk-cold arms,
cohere-10M, median-of-3) showed the narrower beam is where the time goes:

  S=2: L0 165.2s -> 133.4s (-19%), recall 0.5718 -> 0.5660
  S=4: L0 260.3s -> 224.6s (-14%), recall 0.5733 -> 0.5659

Query latency on the merged index is unaffected (0.55ms avg both ways).
The wider beam's extra candidates were largely pruned by diversity
selection, which keeps at most degree edges per node.

The same study found warm-start seeding of these searches (from finished
neighbors' merged adjacency, reverse candidates, or upper-layer descent) is
net-negative: at matched beam width, seeded searches run 8-17% slower with
equal recall — the per-search seeding overhead exceeds the few cheap
descent hops it saves. Beam width is the whole lever.
Refinement ablations (three datasets, disk-cold, paired same-window runs)
show its recall contribution on the merged index is ~0: same-beam arms with
and without refinement land within 0.001. Skipping it saves 45-58s, ~20-25%
of total compaction time at 10M nodes. What it buys is navigability — query
latency on the merged index rises from ~0.55ms to ~0.95ms avg (p99 2.0ms to
3.9ms, cohere-10M) without it.

That is a workload tradeoff, not a correctness call: default to compaction
throughput, and let latency-sensitive pipelines opt back in with
setRefineAfterCompaction(true).
Lets an embedder run graph build/cleanup on the calling thread instead of a ForkJoinPool. Existing pool constructors are preserved as delegating overloads.
…ult pools

Facade carrying an embedder's compute/IO executors so build, PQ/NVQ, and compaction all run on it; also closes the PQRetrainer leak that hardcoded the default pools.
Widens the compress/train/encode entry points from ForkJoinPool to ParallelExecutor (keeping ForkJoinPool overloads), so encode/train can run on the calling thread. Encoding is byte-identical across executors; training is quality-equivalent (k-means seeds from ThreadLocalRandom).
…Runs()

Now that quantization accepts ParallelExecutor, the facade carries one (plus a merge Executor and IO ExecutorService) instead of a ForkJoinPool, so a memtable flush can run build + PQ/NVQ entirely on its own thread via callerRuns(). PQRetrainer and CompactionContext.computeExecutor move to ParallelExecutor too, which also keeps retrain off the all-core pool when the merge executor isn't a ForkJoinPool.
@jshook
jshook force-pushed the cooperative-embedding branch from 9398abb to c5c1b0c Compare July 17, 2026 21:15
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants