Repository navigation
perf(exec): parallelize exact cosine KNN by source (#342) - #494
Conversation
WalkthroughThe change adds instance-owned Rayon compute pools, propagates compute-thread limits through similarity execution, and parallelizes exact and filtered cosine KNN by canonical source chunks while preserving deterministic results and execution limits. ChangesDeterministic cosine parallelization
Estimated code review effort: 4 (Complex) | ~45 minutes Sequence Diagram(s)sequenceDiagram
participant GraphForge
participant similar_algorithm_with_compute
participant AlgorithmControl
participant ComputePool
participant CosineKNN
GraphForge->>similar_algorithm_with_compute: submit similarity request with limits and pool
similar_algorithm_with_compute->>AlgorithmControl: attach shared compute pool
AlgorithmControl->>CosineKNN: provide compute-thread budget
CosineKNN->>ComputePool: execute canonical source chunks
ComputePool-->>CosineKNN: return per-source scores
CosineKNN-->>GraphForge: return deterministic shaped results
Possibly related issues
Suggested labels: 🚥 Pre-merge checks | ✅ 3 | ❌ 2❌ Failed checks (1 warning, 1 inconclusive)
✅ Passed checks (3 passed)
✨ Finishing Touches 💡 1📝 Generate docstrings 💡
🧪 Generate unit tests (beta)
Comment |
94eb2cc to
ae473c3
Compare
Partition independent source rows across the instance-owned private compute pool (#337 budget) while preserving serial per-source dot-product order, deterministic merge, fingerprints, and structured limit/cancel outcomes. Keep a serial path below the documented crossover or when the policy provides one compute thread. Co-authored-by: David Spencer <DecisionNerd@users.noreply.github.com>
Allow the compute-pool similar entry arity and iterate filtered candidate chunks without a needless range index. Co-authored-by: David Spencer <DecisionNerd@users.noreply.github.com>
This comment has been minimized.
This comment has been minimized.
Bazel Bootstrap failed once on checkpoint_pin_and_open_lease_control_recovery_cleanup (WriterBusy) — outside this PR's exec/KNN diff; local cargo+bazel storage suites pass. Co-authored-by: David Spencer <DecisionNerd@users.noreply.github.com>
There was a problem hiding this comment.
Actionable comments posted: 4
🧹 Nitpick comments (4)
crates/graphforge-exec/src/algorithm_dispatch.rs (2)
173-207: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueConsider marking
AlgorithmLimitsas#[non_exhaustive].This change adds a public field to a public struct. Every external struct-literal construction breaks. The three in-repo test sites that this PR had to update confirm the literal pattern is in use. If you mark the struct
#[non_exhaustive], later field additions stay source-compatible for downstream crates, and callers move toAlgorithmLimits::default()plus builder methods.The
.max(1)clamp and the doc comment match the existingwith_batch_sizecontract, so the builder itself is consistent.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/graphforge-exec/src/algorithm_dispatch.rs` around lines 173 - 207, Mark the public AlgorithmLimits struct as #[non_exhaustive] so adding future fields does not break external struct-literal construction. Keep its existing Default implementation and with_batch_size/with_compute_threads builder methods unchanged.
247-253: 📐 Maintainability & Code Quality | 🔵 Trivial | 💤 Low valueDocument that the pool overrides the declared thread budget.
with_compute_pooloverwriteslimits.compute_threadswithpool.num_threads(). A caller that chainswith_compute_threads(n)and thenwith_compute_pool(pool)losesnwithout any signal. The override direction is the safe one, becauseselect_cosine_pathandsource_chunksmust never request more chunks than the pool has workers. State that invariant in the doc comment so a later refactor does not reverse the assignment order.📝 Proposed doc clarification
/// Attach the instance-owned private CPU pool (`#337` / `#342`). + /// + /// The pool is authoritative: this replaces any previously configured + /// `compute_threads` with `pool.num_threads()`. Chunk planning depends on + /// `compute_threads` never exceeding the pool's real worker count. #[must_use] pub(crate) fn with_compute_pool(mut self, pool: crate::SharedComputePool) -> Self {🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/graphforge-exec/src/algorithm_dispatch.rs` around lines 247 - 253, Update the doc comment for with_compute_pool to explicitly state that attaching the pool overrides the declared compute thread budget with the pool’s worker count, preserving the invariant that select_cosine_path and source_chunks never request more chunks than available pool workers.crates/graphforge-exec/src/algorithm_similar_knn.rs (2)
441-459: 🚀 Performance & Scalability | 🔵 Trivial | ⚡ Quick winConsider more chunks than workers so Rayon can balance the load.
source_chunkscreates exactly one contiguous range per worker. After the split no work stealing is possible, so the slowest chunk sets total latency.For
exact_filtered_cosine_knn_parallelthe candidate-set size varies per source. Equal source counts do not mean equal work. One source with a large candidate set serializes the tail while the other workers idle.If you emit a small multiple of
threadschunks, Rayon balances the ranges through work stealing. Determinism is unaffected, becausemerge_chunk_pairsconcatenates by chunk index and the ranges stay in ascending source order.Note that
source_chunks_cover_canonical_rangesat Lines 773-779 and thechunks: 4assertion at Lines 762-769 both pin the current one-chunk-per-worker shape, so those tests need updating with any change here.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/graphforge-exec/src/algorithm_similar_knn.rs` around lines 441 - 459, Update source_chunks to produce a small multiple of threads contiguous, ascending source ranges rather than exactly one range per worker, while preserving empty-input handling and deterministic ordering. Adjust source_chunks_cover_canonical_ranges and the chunks: 4 assertion to validate the new chunking shape and coverage, and ensure exact_filtered_cosine_knn_parallel continues merging chunks by index.
272-273: 🚀 Performance & Scalability | 🔵 Trivial | 🏗️ Heavy liftAvoid the per-source
seenallocation in the filtered path.
score_source_filteredallocates and zeroes avec![false; vectors.len()]for every source. The cost is proportional to the total vector count, not to the candidate-set size. Fornsources this isO(n²)allocation and zeroing work even when each source has only a handful of candidates.Two options preserve exact iteration order, and therefore preserve bit-for-bit results:
- Reuse one generation-stamped
Vec<u32>per chunk. Compare each entry against the current source ordinal instead of clearing the buffer.- Deduplicate
source_candidatesdirectly when the candidate set is much smaller thanvectors.len().The bounds check must stay, because it produces the
filtered KNN candidate is outside vector selectionerror that the tests assert.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/graphforge-exec/src/algorithm_similar_knn.rs` around lines 272 - 273, Update score_source_filtered to avoid allocating and clearing a vec![false; vectors.len()] for each source by reusing a generation-stamped buffer per chunk or deduplicating source_candidates when appropriate. Preserve the existing candidate iteration order and bit-for-bit results, retain the bounds check that reports “filtered KNN candidate is outside vector selection,” and keep deduplication semantics unchanged.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Inline comments:
In `@crates/graphforge-exec/src/algorithm_similar_knn.rs`:
- Around line 18-23: Validate COSINE_PARALLEL_CROSSOVER_OPS with benchmarks
comparing serial and source-parallel execution, including Rayon scheduling,
per-chunk allocation, and merge overhead on the target hosts. Replace the
coincidental 16_384 value with the smallest measured workload where parallel
execution wins, and document the supporting benchmark results alongside the
constant.
- Around line 824-857: Update
parallel_output_limits_and_cancellation_remain_atomic to use workload dimensions
that exceed COSINE_PARALLEL_CROSSOVER_OPS, such as adversarial_vectors(32, 32),
while preserving the expected errors. Add an explicit select_cosine_path
assertion confirming the parallel path is selected, following the guard used by
thread_matrix_preserves_exact_knn_cosine_and_filtered_fingerprints.
- Around line 119-138: Update the parallel chunk collection in the current KNN
implementation and exact_filtered_cosine_knn_parallel to collect all chunk
results without short-circuiting, preserving their range order. After
collection, return the first error by chunk index, and otherwise unwrap the
successful chunk results for downstream aggregation.
In `@crates/graphforge-exec/src/compute_pool.rs`:
- Around line 85-95: Update multi_thread_pool_runs_on_private_workers to inspect
each parallel iterator task’s current thread name and assert it uses the
configured private-worker name prefix, while retaining the sum assertion. Ensure
the check occurs inside the install closure so execution on Rayon’s global pool
cannot satisfy the test.
---
Nitpick comments:
In `@crates/graphforge-exec/src/algorithm_dispatch.rs`:
- Around line 173-207: Mark the public AlgorithmLimits struct as
#[non_exhaustive] so adding future fields does not break external struct-literal
construction. Keep its existing Default implementation and
with_batch_size/with_compute_threads builder methods unchanged.
- Around line 247-253: Update the doc comment for with_compute_pool to
explicitly state that attaching the pool overrides the declared compute thread
budget with the pool’s worker count, preserving the invariant that
select_cosine_path and source_chunks never request more chunks than available
pool workers.
In `@crates/graphforge-exec/src/algorithm_similar_knn.rs`:
- Around line 441-459: Update source_chunks to produce a small multiple of
threads contiguous, ascending source ranges rather than exactly one range per
worker, while preserving empty-input handling and deterministic ordering. Adjust
source_chunks_cover_canonical_ranges and the chunks: 4 assertion to validate the
new chunking shape and coverage, and ensure exact_filtered_cosine_knn_parallel
continues merging chunks by index.
- Around line 272-273: Update score_source_filtered to avoid allocating and
clearing a vec![false; vectors.len()] for each source by reusing a
generation-stamped buffer per chunk or deduplicating source_candidates when
appropriate. Preserve the existing candidate iteration order and bit-for-bit
results, retain the bounds check that reports “filtered KNN candidate is outside
vector selection,” and keep deduplication semantics unchanged.
🪄 Autofix
Fix all unresolved CodeRabbit comments on this PR:
- Push a commit to this branch (recommended)
- Create a new PR with the fixes
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: CHILL
Plan: Pro
Run ID: f58cd092-57c5-4348-aabe-543c1561059b
⛔ Files ignored due to path filters (3)
Cargo.lockis excluded by!**/*.lock,!**/*.lockdocs/development/execution-resource-policy.mdis excluded by!**/*.md,!**/docs/**docs/development/m4-entry-baseline.mdis excluded by!**/*.md,!**/docs/**
📒 Files selected for processing (13)
cargo-bazel-lock.jsoncrates/graphforge-api/src/embedding_refresh.rscrates/graphforge-api/src/lib.rscrates/graphforge-api/src/resource_policy.rscrates/graphforge-exec/Cargo.tomlcrates/graphforge-exec/src/algorithm_analyze_conductance.rscrates/graphforge-exec/src/algorithm_analyze_modularity.rscrates/graphforge-exec/src/algorithm_dispatch.rscrates/graphforge-exec/src/algorithm_similar.rscrates/graphforge-exec/src/algorithm_similar_knn.rscrates/graphforge-exec/src/compute_pool.rscrates/graphforge-exec/src/lib.rstools/bazel/drift/cargo_feature_fingerprint.json
There was a problem hiding this comment.
🧹 Nitpick comments (2)
crates/graphforge-exec/src/algorithm_similar_knn.rs (2)
272-273: 🚀 Performance & Scalability | 🔵 Trivial | ⚡ Quick winThe per-source
seenbuffer scales with the vector count, not the candidate count.
score_source_filteredallocates and zeroesvec![false; vectors.len()]for every source. Total cost is O(sources × vectors), even when each candidate list holds only a few entries. For filtered KNN the candidate lists are normally short, so this allocation dominates the actual scoring work on large selections.Size the deduplication structure by
source_candidates.len()instead.♻️ Proposed change to size deduplication by candidate count
- let mut seen = vec![false; vectors.len()]; let mut candidates = Vec::with_capacity(source_candidates.len()); + let mut seen = std::collections::HashSet::with_capacity(source_candidates.len()); for &target_index in source_candidates { checkpoint(control, work)?; - let Some(target_seen) = seen.get_mut(target_index) else { + if target_index >= vectors.len() { return Err(execution( "filtered KNN candidate is outside vector selection", )); - }; - if source_index == target_index || std::mem::replace(target_seen, true) { + } + if source_index == target_index || !seen.insert(target_index) { continue; }Confirm that the traversal order and the resulting candidate order stay identical, because
#342requires bit-for-bit parity with the serial path. The order is unchanged here, sinceseenonly gates insertion.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/graphforge-exec/src/algorithm_similar_knn.rs` around lines 272 - 273, Update the per-source `seen` buffer in `score_source_filtered` to use `source_candidates.len()` rather than `vectors.len()`, reducing allocation to the candidate count. Preserve the existing traversal and candidate insertion order so filtered KNN results remain bit-for-bit identical with the serial path.
321-325: 🚀 Performance & Scalability | 🔵 Trivial | 💤 Low valueFull sort per source costs more than a partial selection.
top_k_pairssorts every candidate, then keepsk. The exact path producessources - 1candidates per source, so the cost is O(n log n) per source while onlykresults survive.select_nth_unstable_byfollowed by a sort of the firstkelements reduces this to O(n + k log k).This change is optional. If you apply it, verify the tie order stays identical, because
#342requires bit-for-bit parity with the current output.🤖 Prompt for AI Agents
Verify each finding against current code. Fix only still-valid issues, skip the rest with a brief reason, keep changes minimal, and validate. In `@crates/graphforge-exec/src/algorithm_similar_knn.rs` around lines 321 - 325, Optimize top_k_pairs by replacing the full candidates sort with select_nth_unstable_by to partition around the top-k boundary, then sort only the retained first k candidates using the existing score-descending and index-ascending comparator. Preserve the current tie ordering exactly to maintain bit-for-bit parity.
🤖 Prompt for all review comments with AI agents
Verify each finding against current code. Fix only still-valid issues, skip the
rest with a brief reason, keep changes minimal, and validate.
Nitpick comments:
In `@crates/graphforge-exec/src/algorithm_similar_knn.rs`:
- Around line 272-273: Update the per-source `seen` buffer in
`score_source_filtered` to use `source_candidates.len()` rather than
`vectors.len()`, reducing allocation to the candidate count. Preserve the
existing traversal and candidate insertion order so filtered KNN results remain
bit-for-bit identical with the serial path.
- Around line 321-325: Optimize top_k_pairs by replacing the full candidates
sort with select_nth_unstable_by to partition around the top-k boundary, then
sort only the retained first k candidates using the existing score-descending
and index-ascending comparator. Preserve the current tie ordering exactly to
maintain bit-for-bit parity.
ℹ️ Review info
⚙️ Run configuration
Configuration used: Path: .coderabbit.yaml
Review profile: CHILL
Plan: Pro
Run ID: 98285860-7738-4677-807c-c9a801bf55f9
⛔ Files ignored due to path filters (3)
Cargo.lockis excluded by!**/*.lock,!**/*.lockdocs/development/execution-resource-policy.mdis excluded by!**/*.md,!**/docs/**docs/development/m4-entry-baseline.mdis excluded by!**/*.md,!**/docs/**
📒 Files selected for processing (13)
cargo-bazel-lock.jsoncrates/graphforge-api/src/embedding_refresh.rscrates/graphforge-api/src/lib.rscrates/graphforge-api/src/resource_policy.rscrates/graphforge-exec/Cargo.tomlcrates/graphforge-exec/src/algorithm_analyze_conductance.rscrates/graphforge-exec/src/algorithm_analyze_modularity.rscrates/graphforge-exec/src/algorithm_dispatch.rscrates/graphforge-exec/src/algorithm_similar.rscrates/graphforge-exec/src/algorithm_similar_knn.rscrates/graphforge-exec/src/compute_pool.rscrates/graphforge-exec/src/lib.rstools/bazel/drift/cargo_feature_fingerprint.json
🚧 Files skipped from review as they are similar to previous changes (12)
- crates/graphforge-exec/Cargo.toml
- crates/graphforge-exec/src/algorithm_analyze_conductance.rs
- cargo-bazel-lock.json
- crates/graphforge-api/src/embedding_refresh.rs
- crates/graphforge-exec/src/algorithm_dispatch.rs
- tools/bazel/drift/cargo_feature_fingerprint.json
- crates/graphforge-exec/src/algorithm_similar.rs
- crates/graphforge-exec/src/algorithm_analyze_modularity.rs
- crates/graphforge-exec/src/lib.rs
- crates/graphforge-api/src/resource_policy.rs
- crates/graphforge-exec/src/compute_pool.rs
- crates/graphforge-api/src/lib.rs
Raise COSINE_PARALLEL_CROSSOVER_OPS to the measured 32_768 win boundary, use worker-local checkpoint counters (no shared atomic per MAC), prefer lowest-index chunk errors, and strengthen pool/path tests per review. Co-authored-by: David Spencer <DecisionNerd@users.noreply.github.com>
Summary
Parallelize exact / filtered cosine KNN and all-score cosine similarity by canonical source through the instance-owned Rayon
ComputePoolfrom #337, preserving bit-for-bit serial numeric order and merge semantics.Closes #342.
Changes
graphforge-exec::compute_pool(instance-owned Rayon pool sized bycompute_threads)COSINE_PARALLEL_CROSSOVER_OPS = 32_768)Test plan
cargo test -p graphforge-exec --lib algorithm_similar_knncargo test -p graphforge-exec --lib compute_poolcargo clippy -p graphforge-exec --lib -- -D warningsmeasure_cosine_parallel_crossover --ignored)Notes
Land before #343 / #344 so those PRs can reuse
ComputePoolwithout duplicating the module.Need help on this PR? Tag
@codesmith-botwith what you need. Autofix is disabled.Note
Parallelize exact cosine KNN by source using a private Rayon thread pool
ComputePoolwrapper around a private RayonThreadPoolin compute_pool.rs, sized byresource_policy.compute_threadsand owned perGraphForgeinstance, avoiding the global Rayon pool.COSINE_PARALLEL_CROSSOVER_OPSthreshold in algorithm_similar_knn.rs; when estimated work exceeds it and a pool is attached, sources are chunked and processed in parallel.similar_algorithm_with_computein algorithm_similar.rs and threads the pool throughAlgorithmControlso the KNN kernel can use it; existingsimilar_algorithmandsimilar_algorithm_with_limitscallers are unaffected.AtomicUsize.compute_threads > 1and work exceeds the crossover, changing runtime behavior for existing callers who setcompute_threads.Macroscope summarized 903ac87.
Summary by CodeRabbit
New Features
Bug Fixes