diff --git a/BENCHMARKS.md b/BENCHMARKS.md index 94a7c0a..c42be84 100644 --- a/BENCHMARKS.md +++ b/BENCHMARKS.md @@ -221,6 +221,28 @@ build: 21084.6 ms (4743 vectors/s) · exact scan: 39 QPS, p50 24739 µs | 256 | 0.9990 | 3731 | 254 | 697 | +### After: unranked keyword scores, top-k merge (rankings unchanged) + +A fusion study (`search_harness --fusion-study`) showed that capping the +keyword candidate pool is **not** a safe optimisation: against the current +full-corpus normalisation the final top-10 overlap is only 0.83-0.92 and the +first result changes for 10-35% of queries, for only a 2x saving. So the fusion +semantics were left alone and the same answer made cheaper: fusion needs every +keyword score but not their ranking, so BM25 now returns them unsorted from a +dense accumulator (it hashed every posting and then sorted every match), and +the merge selects its top k instead of sorting every candidate. Steady-state +p50: **0.24 -> 0.07 ms** (1K), **2.1 -> 0.49 ms** (10K), **23 -> 4.65 ms** +(100K) — **79x / 100x / 190x** faster than the v2.3.0 baseline, with identical +results. + +### End to end: `HDF5Memory::hybrid_search` (k = 10, weights 0.7 / 0.3) + +| N | ingest ms | cold index build ms | checkpoint ms | open ms | first query after open ms | p50 ms | p99 ms | QPS | +|---:|---:|---:|---:|---:|---:|---:|---:|---:| +| 1000 | 11 | 112 | 4.0 | 1.1 | 1.4 | 0.07 | 0.08 | 14077.9 | +| 10000 | 104 | 1487 | 33.8 | 13.7 | 13.9 | 0.49 | 0.51 | 2020.9 | +| 100000 | 1376 | 20285 | 728.9 | 353.1 | 142.2 | 4.65 | 4.78 | 214.7 | + ## Vector Search Latency Brute-force cosine similarity over 384-dimensional embeddings (OpenAI text-embedding-3-small size). diff --git a/CHANGELOG.md b/CHANGELOG.md index 5efef92..7cb31dd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -43,6 +43,14 @@ on every kernel call. - `clawhdf5-ann`: `HnswIndex::graph_to_bytes` / `from_graph_bytes` — graph-only serialization (checksummed, every neighbour id and level validated on load). +- `clawhdf5-agent`: a further 4-5x on `hybrid_search` with **identical + rankings** (p50 now 0.07 / 0.49 / 4.65 ms at 1K / 10K / 100K — 79x / 100x / + 190x faster than v2.3.0). Fusion needs every keyword score but not their + ranking: new `BM25Index::scores` returns them unsorted from a dense + accumulator (it hashed every posting, then sorted every match), and + `merge_vector_keyword` selects its top k instead of sorting every candidate. + Capping the keyword candidate pool was measured and rejected: it changes the + top-10 for most queries (`search_harness --fusion-study`). - `clawhdf5-agent`: BM25 results are deterministic (ties break by record id), top-k uses a bounded heap, and the "WAND early termination" that computed a bound and then ignored it is gone. IDF is computed per query. diff --git a/crates/clawhdf5-agent/src/bm25.rs b/crates/clawhdf5-agent/src/bm25.rs index 1aff2d7..fb6d603 100644 --- a/crates/clawhdf5-agent/src/bm25.rs +++ b/crates/clawhdf5-agent/src/bm25.rs @@ -83,35 +83,15 @@ impl BM25Index { /// Uses Block-Max WAND for early termination when remaining documents /// cannot beat the current top-k threshold. pub fn search(&self, query: &str, k: usize) -> Vec<(usize, f32)> { - if self.num_docs == 0 || k == 0 { + if k == 0 { return Vec::new(); } - - // Term-at-a-time accumulation. IDF is computed here rather than cached - // at build time: it depends on the live document count, which changes - // with every incremental add/remove, and costs one `ln` per query term. - let mut scores: HashMap = HashMap::new(); - for token in tokenize(query) { - let Some(postings) = self.inverted.get(token.as_str()) else { - continue; - }; - let df = postings.len() as f32; - let idf = ((self.num_docs as f32 - df + 0.5) / (df + 0.5) + 1.0).ln(); - for &(doc_id, freq) in postings { - let dl = self.doc_lengths[doc_id] as f32; - let freq_f = freq as f32; - let tf = (freq_f * (self.k1 + 1.0)) - / (freq_f + self.k1 * (1.0 - self.b + self.b * dl / self.avg_dl)); - *scores.entry(doc_id).or_insert(0.0) += idf * tf; - } - } - // Top-k with a bounded min-heap: O(matches * log k) instead of sorting // every match. Ties break towards the lower doc id so results are - // deterministic (the accumulator is a HashMap). + // deterministic. let mut heap: BinaryHeap)>> = - BinaryHeap::with_capacity(k + 1); - for (doc_id, score) in scores { + BinaryHeap::with_capacity(k.min(1024) + 1); + for (doc_id, score) in self.scores(query) { heap.push(Reverse((HeapScore(score), Reverse(doc_id)))); if heap.len() > k { heap.pop(); @@ -125,6 +105,47 @@ impl BM25Index { results } + /// The BM25 score of **every** matching document, in doc-id order, unsorted + /// by score. Score fusion normalises over the whole matching set, so it + /// needs all of these but not their ranking; producing a ranked list of + /// every match (`search(query, corpus_len)`) spent most of its time sorting. + pub fn scores(&self, query: &str) -> Vec<(usize, f32)> { + if self.num_docs == 0 { + return Vec::new(); + } + // Term-at-a-time accumulation into a dense array: a common term has a + // posting per document, and hashing each one dominated query time. + // IDF is computed here rather than cached at build time: it depends on + // the live document count, which changes with every incremental + // add/remove, and costs one `ln` per query term. + let mut acc = vec![0.0f32; self.doc_lengths.len()]; + let mut matched = false; + for token in tokenize(query) { + let Some(postings) = self.inverted.get(token.as_str()) else { + continue; + }; + matched = true; + let df = postings.len() as f32; + let idf = ((self.num_docs as f32 - df + 0.5) / (df + 0.5) + 1.0).ln(); + for &(doc_id, freq) in postings { + let dl = self.doc_lengths[doc_id] as f32; + let freq_f = freq as f32; + let tf = (freq_f * (self.k1 + 1.0)) + / (freq_f + self.k1 * (1.0 - self.b + self.b * dl / self.avg_dl)); + acc[doc_id] += idf * tf; + } + } + if !matched { + return Vec::new(); + } + // Every contribution is strictly positive (idf = ln(1 + x), x > 0), so + // a zero entry is a document no query term touched. + acc.into_iter() + .enumerate() + .filter(|&(_, score)| score > 0.0) + .collect() + } + /// Number of document slots (live or not) the index covers. Ids are /// positions in the document list it mirrors. pub fn len(&self) -> usize { @@ -564,6 +585,20 @@ mod tests { } } + #[test] + fn scores_is_the_unranked_form_of_a_full_search() { + let mut state = 99u64; + let docs: Vec = (0..200).map(|_| random_doc(&mut state)).collect(); + let tombstones: Vec = (0..200).map(|i| u8::from(i % 7 == 0)).collect(); + let index = BM25Index::build(&docs, &tombstones); + for query in ["alpha", "beta gamma x1", "missing", ""] { + let mut all = index.scores(query); + all.sort_by(|a, b| b.1.total_cmp(&a.1).then(a.0.cmp(&b.0))); + assert_eq!(all, index.search(query, docs.len()), "{query:?}"); + assert!(all.iter().all(|(id, _)| tombstones[*id] == 0)); + } + } + #[test] fn ties_break_towards_the_lower_doc_id() { let docs: Vec = (0..6).map(|_| "same text".to_string()).collect(); diff --git a/crates/clawhdf5-agent/src/hybrid.rs b/crates/clawhdf5-agent/src/hybrid.rs index 20a9f2d..c943ed0 100644 --- a/crates/clawhdf5-agent/src/hybrid.rs +++ b/crates/clawhdf5-agent/src/hybrid.rs @@ -58,7 +58,7 @@ pub fn hybrid_search( vector_search::cosine_similarity_batch(query_embedding, vectors, tombstones) } }; - let kw_scores = bm25_index.search(query_text, vectors.len()); + let kw_scores = bm25_index.scores(query_text); merge_vector_keyword(vec_scores, kw_scores, vector_weight, keyword_weight, k) } @@ -92,13 +92,22 @@ pub fn merge_vector_keyword( let mut results: Vec<(usize, f32)> = merged.into_iter().collect(); // Index tie-break: `merged` is a HashMap, so without it the ties that - // survive `truncate` differ from run to run. - results.sort_by(|a, b| { + // survive differ from run to run. + let by_score_then_id = |a: &(usize, f32), b: &(usize, f32)| { b.1.partial_cmp(&a.1) .unwrap_or(std::cmp::Ordering::Equal) .then(a.0.cmp(&b.0)) - }); - results.truncate(k); + }; + // Only the top k are wanted: partition them out, then order just those, + // instead of sorting every candidate (the keyword side can be the corpus). + if k == 0 { + return Vec::new(); + } + if results.len() > k { + results.select_nth_unstable_by(k - 1, by_score_then_id); + results.truncate(k); + } + results.sort_by(by_score_then_id); results } @@ -344,6 +353,25 @@ mod tests { assert_eq!(result[0].1, 1.0); } + #[test] + fn merge_top_k_matches_a_full_sort() { + // Many ties (scores repeat) so the index tie-break is exercised. + let vec_scores: Vec<(usize, f32)> = (0..300).map(|i| (i, ((i * 7) % 13) as f32)).collect(); + let kw_scores: Vec<(usize, f32)> = (100..500).map(|i| (i, ((i * 5) % 11) as f32)).collect(); + let everything = + merge_vector_keyword(vec_scores.clone(), kw_scores.clone(), 0.7, 0.3, 10_000); + assert_eq!(everything.len(), 500); + assert!( + everything + .windows(2) + .all(|w| { w[0].1 > w[1].1 || (w[0].1 == w[1].1 && w[0].0 < w[1].0) }) + ); + for k in [0, 1, 7, 50, 499, 500, 501] { + let top = merge_vector_keyword(vec_scores.clone(), kw_scores.clone(), 0.7, 0.3, k); + assert_eq!(top, everything[..k.min(500)], "k = {k}"); + } + } + #[test] fn normalize_scores_all_equal() { let matched = normalize_scores(&[(0, 0.4), (1, 0.4)]); diff --git a/crates/clawhdf5-agent/src/search.rs b/crates/clawhdf5-agent/src/search.rs index 30b679a..616b0e3 100644 --- a/crates/clawhdf5-agent/src/search.rs +++ b/crates/clawhdf5-agent/src/search.rs @@ -35,7 +35,9 @@ impl HDF5Memory { .into_iter() .map(|(id, dist)| (id, 1.0 - dist)) .collect(); - let kw_scores = bm25.search(query_text, self.cache.len()); + // Fusion normalises over every keyword match, so it needs all + // the scores — but not ranked. + let kw_scores = bm25.scores(query_text); hybrid::merge_vector_keyword( vec_scores, kw_scores, diff --git a/crates/clawhdf5-bench/src/bin/search_harness.rs b/crates/clawhdf5-bench/src/bin/search_harness.rs index 4e50c75..3c43af7 100644 --- a/crates/clawhdf5-bench/src/bin/search_harness.rs +++ b/crates/clawhdf5-bench/src/bin/search_harness.rs @@ -369,10 +369,101 @@ fn bench_end_to_end(n: usize, json: &mut Vec) { })); } +// --------------------------------------------------------------------------- +// Fusion study: does capping the keyword candidate pool change the ranking? +// --------------------------------------------------------------------------- + +/// `hybrid_search` min-max normalises each signal over the candidates it is +/// given. The vector stage supplies a pool of `max(8k, 64)`; the keyword stage +/// supplies *every* matching record, which is what now dominates query time. +/// This compares the current fusion with one whose keyword stage is capped to +/// a pool, reporting how often the final top-k agree and what each costs. +fn fusion_study(n: usize) { + use clawhdf5_agent::bm25::BM25Index; + use clawhdf5_agent::hybrid::merge_vector_keyword; + + let data = make_dataset(n, 0xE2E ^ n as u64); + let mut rng = Rng(7); + let texts: Vec = (0..n) + .map(|i| text_for(data.cluster_of[i], i, &mut rng)) + .collect(); + let query_texts: Vec = data + .query_cluster + .iter() + .enumerate() + .map(|(i, c)| text_for(*c, i, &mut rng)) + .collect(); + let bm25 = BM25Index::build(&texts, &vec![0u8; n]); + let index = HnswIndex::build_with_metric( + &data.vectors, + HNSW_M, + HNSW_EF_CONSTRUCTION, + DistanceMetric::Cosine, + ); + + let vec_pool = (K * 8).max(64); + println!("\n### Fusion study, N = {n} (k = {K}, weights 0.7 / 0.3, vector pool {vec_pool})\n"); + println!( + "| keyword pool | top-{K} overlap vs full | identical top-{K} | same #1 | keyword+merge µs |" + ); + println!("|---:|---:|---:|---:|---:|"); + + let fuse = |q: usize, kw_pool: usize| -> (Vec, Duration) { + let vec_scores: Vec<(usize, f32)> = index + .search(&data.queries[q], vec_pool, vec_pool) + .into_iter() + .map(|(id, d)| (id, 1.0 - d)) + .collect(); + let t = Instant::now(); + let kw = bm25.search(&query_texts[q], kw_pool); + let merged = merge_vector_keyword(vec_scores, kw, 0.7, 0.3, K); + let took = t.elapsed(); + (merged.into_iter().map(|(id, _)| id).collect(), took) + }; + + let full: Vec<(Vec, Duration)> = (0..N_QUERIES).map(|q| fuse(q, n)).collect(); + let full_time: Duration = full.iter().map(|f| f.1).sum(); + println!( + "| all ({n}) | 1.0000 | 100.0% | 100.0% | {:.0} |", + micros(full_time) / N_QUERIES as f64 + ); + for pool in [vec_pool, vec_pool * 4, 1000] { + if pool >= n { + continue; + } + let (mut overlap, mut identical, mut same_first) = (0usize, 0usize, 0usize); + let mut time = Duration::ZERO; + for (q, (want, _)) in full.iter().enumerate() { + let (got, took) = fuse(q, pool); + time += took; + overlap += got.iter().filter(|id| want.contains(id)).count(); + identical += usize::from(&got == want); + same_first += usize::from(got.first() == want.first()); + } + println!( + "| {pool} | {:.4} | {:.1}% | {:.1}% | {:.0} |", + overlap as f64 / (K * N_QUERIES) as f64, + 100.0 * identical as f64 / N_QUERIES as f64, + 100.0 * same_first as f64 / N_QUERIES as f64, + micros(time) / N_QUERIES as f64 + ); + } +} + fn main() { let args: Vec = std::env::args().skip(1).collect(); let full = args.iter().any(|a| a == "--full"); let ann_only = args.iter().any(|a| a == "--ann-only"); + if args.iter().any(|a| a == "--fusion-study") { + for &n in if full { + &[10_000, 100_000][..] + } else { + &[10_000][..] + } { + fusion_study(n); + } + return; + } if args.iter().any(|a| a == "--uniform") { UNIFORM.store(true, std::sync::atomic::Ordering::Relaxed); println!("(uniform random data)");