mirror of
https://github.com/qdrant/qdrant.git
synced 2026-07-23 11:11:00 -05:00
refactor: move deferred-point ownership into ID tracker (#9062)
* refactor: move deferred-point ownership into the ID tracker Re-implements the idea from #8512 against current `dev`. Deferred-point state (`deferred_internal_id` + `deferred_deleted_count`) moves out of `Segment.deferred_point_status` and the cached `SparseVectorIndex.deferred_internal_id` field into `PointMappings`, exposed through `IdTrackerRead`. The threshold is set once at `MutableIdTracker::open` time; `PointMappings::drop` now maintains the deleted counter inline (with double-delete protection), removing the manual increment in `delete_point_internal` and the `calculate_deleted_deferred_point_count` rescan. Read paths consume the threshold through the id tracker: - The segment read view drops the `deferred_point_status` field and `with_view` no longer threads it in; `read_view/{deferred,info}.rs` call `self.id_tracker.deferred_*()` directly. - `SparseVectorIndex` no longer stores its own copy and its `update_vector` / search debug-assert read from `self.id_tracker.borrow().deferred_internal_id()`. - `VectorQueryContext.deferred_internal_id` and the `SegmentQueryContext::get_vector_context` parameter are gone; the three downstream readers (`plain_vector_index`, sparse search, sparse `update_vector`) consult their own id tracker. `PointMappingsRefEnum` centralises the dispatch: - `iter_internal_with_behavior(DeferredBehavior)` replaces ad-hoc branches in `iter_filtered_points` impls. - `external_iter_cutoff(DeferredBehavior)` covers iterators sourced outside the mapping (field-index outputs in `struct_payload_index::iter_filtered_points`). - The internal `deferred_internal_id()` accessor is private; the raw threshold no longer leaks to consumers. - `iter_from_visible` / `iter_random_visible` read the mapping's own threshold; callers that previously passed `DeferredBehavior::apply(...)` now branch on `deferred_behavior.include_all_points()` (scroll / order_by) or simply drop the argument (sampling / facet). `PayloadIndexRead::query_points` drops the now-redundant `deferred_internal_id` parameter; `iter_filtered_points` takes `DeferredBehavior` directly so HNSW build/search can request `IncludeAll` while normal reads request `Exclude`. RocksDB-related parts of the original PR are skipped — that tracker is already gone from `dev`. Tests adapted: sites that mutated `segment.deferred_point_status` directly now construct a parallel non-deferred segment via `create_deferred_segment(..., 0)` for comparison; `test_deleted_deferred_point_count` reads counters through the id tracker. See `docs/plans/deferred-points-owned-by-id-tracker.md` for the design write-up. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * fix(benches): drop stale deferred_internal_id arg from query_points calls The boolean / range / conditional bench files weren't built by `cargo test -p segment`, so they slipped through. `cargo clippy --workspace --all-targets` catches them. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * refactor: drop id_tracker / point_mappings args from iter_filtered_points Both impls already hold an id tracker on `self`: - `StructPayloadIndexReadView` carries `id_tracker: &'a I`, so `self.id_tracker.point_mappings()` borrows from `'a` and the lazy iterator chain keeps working unchanged. - `PlainPayloadIndex` carries `id_tracker: Arc<AtomicRefCell<...>>`, where the mapping borrow is local; collect into a `Vec` and return `into_iter()`. PlainPayloadIndex::iter_filtered_points has no direct callers — only `query_points` was using it — so eager collection is a non-issue. While here, take `self` by value on `iter_internal_visible`, `iter_from_visible`, `iter_random_visible`, `iter_internal_with_behavior`, and `external_iter_cutoff`. `PointMappingsRefEnum` is `Copy`; this matches the existing `iter_internal` / `iter_from` / `iter_random` shape and lets the iterator outlive a local `let point_mappings = ...;` binding. The HNSW `condition_points` helper drops its now-unused `id_tracker` parameter. All callers (sampling, scroll, order_by, facet ×2, hnsw build/search) just drop the two arguments. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * chore: ignore /docs/plans/ and untrack the previously-committed plan `docs/plans/` is a scratch directory for per-feature planning notes — not something we want under source control. Add it to `.gitignore` and drop the deferred-points plan that slipped into history; the design is captured in the PR description. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * refactor: replace external_iter_cutoff with filter_deferred iterator wrapper Instead of exposing a raw `Option<PointOffsetType>` cutoff that every caller has to apply with their own `.filter(...)`, give `PointMappingsRefEnum` an iterator wrapper: fn filter_deferred<I: Iterator<Item = PointOffsetType>>( self, iter: I, deferred_behavior: DeferredBehavior, ) -> impl Iterator<Item = PointOffsetType> It returns the iterator unchanged for `IncludeAll` (or when the mapping has no threshold) and otherwise wraps it in a cutoff `.filter`, dispatched via `itertools::Either` so the no-cutoff path stays allocation-free. The struct payload index's `iter_filtered_points` swaps its open-coded filter for a single `point_mappings.filter_deferred(...)` call. The deferred threshold no longer leaks out of `PointMappingsRefEnum`. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * refactor: move deferred wrapping out of peek_top_all, gate it as test-only `BatchFilteredSearcher::peek_top_all` baked the deferred cutoff into its iterator construction, which was the last place outside `PointMappingsRefEnum` that knew about the threshold. Split the deleted-iteration concern out into a new accessor: fn iter_not_deleted(&self) -> impl Iterator<Item = PointOffsetType> + 'a It borrows `&'a BitSlice` directly (not via `&self`), so callers can chain `filter_deferred` and then move `self` into `peek_top_iter` without lifetime conflicts. Sparse + plain vector index call sites now do: let iter = id_tracker .point_mappings() .filter_deferred(searcher.iter_not_deleted(), DeferredBehavior::Exclude); searcher.peek_top_iter(iter, &is_stopped) leaving `BatchFilteredSearcher` completely ignorant of deferred state. With deferred handling lifted out, `peek_top_all` itself is now used only by tests (3 inline `#[cfg(test)] mod tests`, 1 integration test, 1 bench) — gate it under `#[cfg(feature = "testing")]` to match `new_for_test`. Production code goes through the `iter_not_deleted` + `filter_deferred` + `peek_top_iter` composition. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * refactor: optimize Segment::retrieve and thread user data through read_vectors Two interlocking changes that together collapse the per-point lookups and intermediate allocations in `Segment::retrieve` down to one external-to-internal pass. ## `IdTrackerRead::resolve_external_ids` (new default trait method) Single-pass translation of a `&[PointIdType]` slice into two parallel vectors `(Vec<PointIdType>, Vec<PointOffsetType>)`. Folds deferred filtering (compare offset against the threshold inline — no separate `point_is_deferred` lookup) and missing-id errors (eager `PointIdError`) into resolution. Lives on the trait so the deferred threshold never leaks out of the id tracker; the parallel-vector shape lets a future batched payload / vector fetcher consume `&offsets` straight without unzipping. The `appendable_flag` guard previously in `point_is_deferred` is gone: non-appendable trackers always carry `deferred_internal_id() == None` (set only via `MutableIdTracker::open`, guarded by the segment constructor), so the check was load-bearing nowhere. ## User-data threading through `read_vectors` `VectorStorageRead::read_vectors` now takes `IntoIterator<Item = (U, PointOffsetType)>` and yields `(U, PointOffsetType, CowVector)`. The user-data tag rides alongside each offset all the way through, so callers can map results back into a parallel array without keeping a separate `offset → ...` lookup table. - Default trait impl: one-line per-key loop. - Dense impl: `unzip()` into parallel `(Vec<U>, Vec<PointOffsetType>)` in a single pass — same allocation count as before, just U riding alongside. - Enum delegations (`VectorStorageEnum`, `VectorStorageReadEnum`) forward unchanged. - `for_each_in_batch` and below stay untouched. `SegmentReadView::vectors_by_offsets<U: Copy>` becomes a lazy filter chain — no parallel `Vec<(orig_idx, offset)>` allocation. The dead `SegmentReadView::read_vectors` helper is removed. ## `Segment::retrieve` end-to-end Per N points / V vectors / payload: | Operation | Before | After | |----------------------------|---------------------|-------| | `id_tracker.internal_id` | N × (1 + V + 1) | N | | `id_tracker.external_id` | N × V | 0 | | `point_is_deferred` | N (when applicable) | 0 | | `offset_to_id` HashMap | N entries | none | | `Vec` in `vectors_by_offsets` | 1 | 0 | The vectors stage passes the external id as `read_vectors`'s user data — the callback gets `id` directly without any index lookup. The payload stage uses `payload_by_offset` against the already-resolved offsets. The shape is also batch-friendly: swapping in a future `IdTrackerRead::batch_internal_id` or `payload_index.batch_get_payload` needs no changes outside the two call sites. ## Behavioural notes - Missing-id now errors eagerly inside resolution, instead of in the vectors stage (`WithVector::Bool(true)` / `Selector`) or payload stage (`with_payload.enable`). The previous `WithVector::Bool(false)` + no-payload path silently inserted an empty record; that is now also an error. None of the existing callers (search post-processing, external retrieve API, the deferred-points test on tests/mod.rs:1179) pass non-existent ids. - Added a per-payload `check_stopped`; the vectors stage already had `stop_if` on its iterator chain. - `vector_by_offset` (the single-element helper) passes `()` as the no-op user data. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * style: apply rustfmt to optimised retrieve / read_vectors paths Pre-push hook failure on the previous commit was rustfmt. Same content, formatted. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> * do not error out on missing points in retrieve --------- Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
committed by
generall
parent
1634d430e2
commit
03e15b5ed4
1
.gitignore
vendored
1
.gitignore
vendored
@@ -21,3 +21,4 @@ venv
|
||||
.env
|
||||
.python-version
|
||||
.claude/
|
||||
/docs/plans/
|
||||
|
||||
@@ -47,7 +47,7 @@ pub fn plain_boolean_query_points(c: &mut Criterion) {
|
||||
b.iter(|| {
|
||||
let filter = random_bool_filter(&mut rng);
|
||||
result_size += plain_index
|
||||
.query_points(&filter, &hw_counter, &is_stopped, None)
|
||||
.query_points(&filter, &hw_counter, &is_stopped)
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
@@ -77,7 +77,7 @@ pub fn struct_boolean_query_points(c: &mut Criterion) {
|
||||
b.iter(|| {
|
||||
let filter = random_bool_filter(&mut rng);
|
||||
result_size += struct_index
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
@@ -131,7 +131,7 @@ pub fn keyword_index_boolean_query_points(c: &mut Criterion) {
|
||||
b.iter(|| {
|
||||
let filter = random_bool_filter(&mut rng);
|
||||
result_size += index
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
|
||||
@@ -38,7 +38,7 @@ fn conditional_plain_search_benchmark(c: &mut Criterion) {
|
||||
b.iter(|| {
|
||||
let filter = random_must_filter(&mut rng, 2);
|
||||
result_size += plain_index
|
||||
.query_points(&filter, &hw_counter, &is_stopped, None)
|
||||
.query_points(&filter, &hw_counter, &is_stopped)
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
@@ -56,7 +56,7 @@ fn conditional_plain_search_benchmark(c: &mut Criterion) {
|
||||
b.iter(|| {
|
||||
let filter = random_must_filter(&mut rng, 1);
|
||||
result_size += plain_index
|
||||
.query_points(&filter, &hw_counter, &is_stopped, None)
|
||||
.query_points(&filter, &hw_counter, &is_stopped)
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
@@ -156,7 +156,7 @@ fn conditional_struct_search_benchmark(c: &mut Criterion) {
|
||||
b.iter(|| {
|
||||
let filter = random_must_filter(&mut rng, 2);
|
||||
result_size += struct_index
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
|
||||
@@ -111,7 +111,7 @@ fn range_filtering(c: &mut Criterion) {
|
||||
|| random_range_filter(&mut rng, FLT_KEY),
|
||||
|filter| {
|
||||
result_size += index
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
@@ -125,7 +125,7 @@ fn range_filtering(c: &mut Criterion) {
|
||||
|| random_range_filter(&mut rng, INT_KEY),
|
||||
|filter| {
|
||||
result_size += index
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
@@ -154,7 +154,7 @@ fn range_filtering(c: &mut Criterion) {
|
||||
|| random_range_filter(&mut rng, FLT_KEY),
|
||||
|filter| {
|
||||
result_size += index
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
@@ -168,7 +168,7 @@ fn range_filtering(c: &mut Criterion) {
|
||||
|| random_range_filter(&mut rng, INT_KEY),
|
||||
|filter| {
|
||||
result_size += index
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.len();
|
||||
query_count += 1;
|
||||
|
||||
@@ -92,7 +92,6 @@ fn sparse_vector_index_build_benchmark(c: &mut Criterion) {
|
||||
path: index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
assert_eq!(sparse_vector_index.indexed_vector_count(), NUM_VECTORS);
|
||||
@@ -109,7 +108,6 @@ fn sparse_vector_index_build_benchmark(c: &mut Criterion) {
|
||||
path: index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
|
||||
@@ -113,7 +113,6 @@ fn sparse_vector_index_search_benchmark_impl(
|
||||
path: mmap_index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || pb.inc(1),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
pb.finish_and_clear();
|
||||
|
||||
@@ -77,7 +77,7 @@ fn benchmark<const IO_URING: bool, const VECTORS: usize, const BATCH: usize>(c:
|
||||
id_tracker.deleted_point_bitslice(),
|
||||
10,
|
||||
)
|
||||
.peek_top_all(&DEFAULT_STOPPED, None)
|
||||
.peek_top_all(&DEFAULT_STOPPED)
|
||||
.expect("points scored")
|
||||
},
|
||||
BatchSize::SmallInput,
|
||||
|
||||
@@ -6,7 +6,7 @@ use common::bitvec::BitSlice;
|
||||
use common::counter::hardware_accumulator::HwMeasurementAcc;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::cow::SimpleCow;
|
||||
use common::types::{PointOffsetType, ScoreType};
|
||||
use common::types::ScoreType;
|
||||
use sparse::common::types::{DimId, DimWeight};
|
||||
|
||||
use crate::data_types::tiny_map;
|
||||
@@ -142,11 +142,7 @@ impl<'a> SegmentQueryContext<'a> {
|
||||
self.query_context.available_point_count()
|
||||
}
|
||||
|
||||
pub fn get_vector_context(
|
||||
&self,
|
||||
vector_name: &VectorName,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> VectorQueryContext<'_> {
|
||||
pub fn get_vector_context(&self, vector_name: &VectorName) -> VectorQueryContext<'_> {
|
||||
VectorQueryContext {
|
||||
search_optimized_threshold_kb: self.query_context.search_optimized_threshold_kb,
|
||||
is_stopped: Some(&self.query_context.is_stopped),
|
||||
@@ -159,7 +155,6 @@ impl<'a> SegmentQueryContext<'a> {
|
||||
.copied(),
|
||||
deleted_points: self.deleted_points,
|
||||
hardware_counter: self.hardware_counter.fork(),
|
||||
deferred_internal_id,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -197,8 +192,6 @@ pub struct VectorQueryContext<'a> {
|
||||
deleted_points: Option<&'a BitSlice>,
|
||||
|
||||
hardware_counter: HardwareCounterCell,
|
||||
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
}
|
||||
|
||||
impl VectorQueryContext<'_> {
|
||||
@@ -249,10 +242,6 @@ impl VectorQueryContext<'_> {
|
||||
pub fn is_require_idf(&self) -> bool {
|
||||
self.idf.is_some() && self.indexed_vectors.is_some()
|
||||
}
|
||||
|
||||
pub fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
self.deferred_internal_id
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "testing")]
|
||||
@@ -265,7 +254,6 @@ impl Default for VectorQueryContext<'_> {
|
||||
indexed_vectors: None,
|
||||
deleted_points: None,
|
||||
hardware_counter: HardwareCounterCell::new(),
|
||||
deferred_internal_id: None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -83,7 +83,6 @@ pub fn fixture_sparse_index_from_iter<I: InvertedIndex>(
|
||||
path: index_dir,
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})?;
|
||||
|
||||
assert_eq!(
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
use atomic_refcell::AtomicRef;
|
||||
use common::bitvec::{BitSlice, BitSliceExt as _};
|
||||
use common::types::PointOffsetType;
|
||||
use common::types::{DeferredBehavior, PointOffsetType};
|
||||
use itertools::Either;
|
||||
use self_cell::self_cell;
|
||||
|
||||
use super::tracker_enum::IdTrackerEnum;
|
||||
@@ -77,12 +78,10 @@ impl<'a> PointMappingsRefEnum<'a> {
|
||||
)
|
||||
}
|
||||
|
||||
/// Iterate over all internal IDs. Optionally filter all deferred points.
|
||||
pub fn iter_internal_visible(
|
||||
&self,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> Box<dyn Iterator<Item = PointOffsetType> + '_> {
|
||||
match deferred_internal_id {
|
||||
/// Iterate over all internal IDs, filtering deferred points using the
|
||||
/// mapping's own threshold.
|
||||
pub fn iter_internal_visible(self) -> Box<dyn Iterator<Item = PointOffsetType> + 'a> {
|
||||
match self.deferred_internal_id() {
|
||||
None => self.iter_internal(),
|
||||
Some(deferred_internal_id) => Box::new(
|
||||
self.iter_internal()
|
||||
@@ -91,13 +90,50 @@ impl<'a> PointMappingsRefEnum<'a> {
|
||||
}
|
||||
}
|
||||
|
||||
/// Iterate starting from a given ID. Optionally filter all deferred points.
|
||||
/// Iterate over all internal IDs, with deferred filtering selected by
|
||||
/// `deferred_behavior`:
|
||||
/// - [`DeferredBehavior::Exclude`] applies the mapping's own threshold;
|
||||
/// - [`DeferredBehavior::IncludeAll`] yields every point regardless of the
|
||||
/// threshold.
|
||||
pub fn iter_internal_with_behavior(
|
||||
self,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> Box<dyn Iterator<Item = PointOffsetType> + 'a> {
|
||||
if deferred_behavior.include_all_points() {
|
||||
self.iter_internal()
|
||||
} else {
|
||||
self.iter_internal_visible()
|
||||
}
|
||||
}
|
||||
|
||||
/// Wrap an iterator of internal IDs so that points at or above the
|
||||
/// mapping's deferred threshold are excluded.
|
||||
///
|
||||
/// For [`DeferredBehavior::IncludeAll`] — or when the mapping has no
|
||||
/// deferred threshold — the iterator is returned unchanged. Intended for
|
||||
/// iterators sourced outside the mapping (e.g., field-index outputs) where
|
||||
/// the threshold isn't applied implicitly.
|
||||
pub fn filter_deferred<I>(
|
||||
self,
|
||||
iter: I,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> impl Iterator<Item = PointOffsetType>
|
||||
where
|
||||
I: Iterator<Item = PointOffsetType>,
|
||||
{
|
||||
match deferred_behavior.apply(self.deferred_internal_id()) {
|
||||
None => Either::Left(iter),
|
||||
Some(cutoff) => Either::Right(iter.filter(move |&id| id < cutoff)),
|
||||
}
|
||||
}
|
||||
|
||||
/// Iterate starting from a given ID, filtering deferred points using the
|
||||
/// mapping's own threshold.
|
||||
pub fn iter_from_visible(
|
||||
&self,
|
||||
self,
|
||||
external_id: Option<PointIdType>,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> Box<dyn Iterator<Item = (PointIdType, PointOffsetType)> + '_> {
|
||||
match deferred_internal_id {
|
||||
) -> Box<dyn Iterator<Item = (PointIdType, PointOffsetType)> + 'a> {
|
||||
match self.deferred_internal_id() {
|
||||
None => self.iter_from(external_id),
|
||||
Some(deferred_internal_id) => Box::new(
|
||||
self.iter_from(external_id)
|
||||
@@ -106,11 +142,12 @@ impl<'a> PointMappingsRefEnum<'a> {
|
||||
}
|
||||
}
|
||||
|
||||
/// Iterate over internal IDs in random order, filtering deferred points
|
||||
/// using the mapping's own threshold.
|
||||
pub fn iter_random_visible(
|
||||
&self,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> Box<dyn Iterator<Item = (PointIdType, PointOffsetType)> + '_> {
|
||||
match deferred_internal_id {
|
||||
self,
|
||||
) -> Box<dyn Iterator<Item = (PointIdType, PointOffsetType)> + 'a> {
|
||||
match self.deferred_internal_id() {
|
||||
None => self.iter_random(),
|
||||
Some(deferred_internal_id) => Box::new(
|
||||
self.iter_random()
|
||||
@@ -120,6 +157,21 @@ impl<'a> PointMappingsRefEnum<'a> {
|
||||
),
|
||||
}
|
||||
}
|
||||
|
||||
/// Deferred threshold attached to this mapping, if any.
|
||||
///
|
||||
/// Compressed mappings (used by immutable / read-only-immutable id trackers)
|
||||
/// never carry a deferred threshold.
|
||||
///
|
||||
/// Kept private so callers go through the dispatch helpers
|
||||
/// ([`Self::iter_internal_with_behavior`], [`Self::external_iter_cutoff`])
|
||||
/// instead of leaking the raw threshold.
|
||||
fn deferred_internal_id(self) -> Option<PointOffsetType> {
|
||||
match self {
|
||||
PointMappingsRefEnum::Plain(m) => m.deferred_internal_id(),
|
||||
PointMappingsRefEnum::Compressed(_) => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
self_cell! {
|
||||
|
||||
@@ -100,4 +100,18 @@ impl<S: UniversalRead> IdTrackerRead for ReadOnlyIdTrackerEnum<S> {
|
||||
ReadOnlyIdTrackerEnum::Immutable(id_tracker) => id_tracker.iter_internal_versions(),
|
||||
}
|
||||
}
|
||||
|
||||
fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
match self {
|
||||
ReadOnlyIdTrackerEnum::Appendable(id_tracker) => id_tracker.deferred_internal_id(),
|
||||
ReadOnlyIdTrackerEnum::Immutable(id_tracker) => id_tracker.deferred_internal_id(),
|
||||
}
|
||||
}
|
||||
|
||||
fn deferred_deleted_count(&self) -> usize {
|
||||
match self {
|
||||
ReadOnlyIdTrackerEnum::Appendable(id_tracker) => id_tracker.deferred_deleted_count(),
|
||||
ReadOnlyIdTrackerEnum::Immutable(id_tracker) => id_tracker.deferred_deleted_count(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -111,6 +111,22 @@ impl IdTrackerRead for IdTrackerEnum {
|
||||
IdTrackerEnum::InMemoryIdTracker(id_tracker) => id_tracker.iter_internal_versions(),
|
||||
}
|
||||
}
|
||||
|
||||
fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
match self {
|
||||
IdTrackerEnum::MutableIdTracker(id_tracker) => id_tracker.deferred_internal_id(),
|
||||
IdTrackerEnum::ImmutableIdTracker(id_tracker) => id_tracker.deferred_internal_id(),
|
||||
IdTrackerEnum::InMemoryIdTracker(id_tracker) => id_tracker.deferred_internal_id(),
|
||||
}
|
||||
}
|
||||
|
||||
fn deferred_deleted_count(&self) -> usize {
|
||||
match self {
|
||||
IdTrackerEnum::MutableIdTracker(id_tracker) => id_tracker.deferred_deleted_count(),
|
||||
IdTrackerEnum::ImmutableIdTracker(id_tracker) => id_tracker.deferred_deleted_count(),
|
||||
IdTrackerEnum::InMemoryIdTracker(id_tracker) => id_tracker.deferred_deleted_count(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl IdTracker for IdTrackerEnum {
|
||||
|
||||
@@ -2,7 +2,7 @@ use std::fmt;
|
||||
use std::path::PathBuf;
|
||||
|
||||
use common::bitvec::{BitSlice, BitSliceExt as _};
|
||||
use common::types::PointOffsetType;
|
||||
use common::types::{DeferredBehavior, PointOffsetType};
|
||||
use rand::rngs::StdRng;
|
||||
use rand::{RngExt, SeedableRng};
|
||||
|
||||
@@ -174,4 +174,53 @@ pub trait IdTrackerRead {
|
||||
fn iter_internal_versions(
|
||||
&self,
|
||||
) -> Box<dyn Iterator<Item = (PointOffsetType, SeqNumberType)> + '_>;
|
||||
|
||||
/// Internal-id threshold above which points are hidden from reads.
|
||||
///
|
||||
/// Only appendable trackers can carry a non-`None` value.
|
||||
fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
None
|
||||
}
|
||||
|
||||
/// Number of soft-deleted points at or above the deferred threshold.
|
||||
fn deferred_deleted_count(&self) -> usize {
|
||||
0
|
||||
}
|
||||
|
||||
/// Translate external point ids into two parallel vectors of `(ids,
|
||||
/// offsets)` in a single pass.
|
||||
///
|
||||
/// Applies deferred filtering according to `deferred_behavior` inline
|
||||
/// (no separate `point_is_deferred` lookup), missing points will be ignored
|
||||
///
|
||||
/// The parallel-vector return shape lets downstream batched fetchers
|
||||
/// consume `&offsets` directly.
|
||||
///
|
||||
/// Centralising this here keeps the deferred threshold from leaking out
|
||||
/// of the id tracker — callers go through this entry point instead of
|
||||
/// reading `deferred_internal_id()` themselves.
|
||||
fn resolve_external_ids(
|
||||
&self,
|
||||
point_ids: &[PointIdType],
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> (Vec<PointIdType>, Vec<PointOffsetType>) {
|
||||
// Non-appendable trackers never carry a deferred threshold (it's set
|
||||
// only via `MutableIdTracker::open`, guarded by `appendable_flag` in
|
||||
// the segment constructor), so we don't need an appendable check here.
|
||||
let deferred_cutoff = deferred_behavior.apply(self.deferred_internal_id());
|
||||
|
||||
let mut ids = Vec::with_capacity(point_ids.len());
|
||||
let mut offsets = Vec::with_capacity(point_ids.len());
|
||||
for &point_id in point_ids {
|
||||
let Some(internal_id) = self.internal_id(point_id) else {
|
||||
continue;
|
||||
};
|
||||
if deferred_cutoff.is_some_and(|cutoff| internal_id >= cutoff) {
|
||||
continue;
|
||||
}
|
||||
ids.push(point_id);
|
||||
offsets.push(internal_id);
|
||||
}
|
||||
(ids, offsets)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -100,6 +100,14 @@ impl IdTrackerRead for InMemoryIdTracker {
|
||||
"in memory id tracker"
|
||||
}
|
||||
|
||||
fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
self.mappings.deferred_internal_id()
|
||||
}
|
||||
|
||||
fn deferred_deleted_count(&self) -> usize {
|
||||
self.mappings.deferred_deleted_count()
|
||||
}
|
||||
|
||||
fn iter_internal_versions(
|
||||
&self,
|
||||
) -> Box<dyn Iterator<Item = (PointOffsetType, SeqNumberType)> + '_> {
|
||||
|
||||
@@ -128,12 +128,15 @@ fn write_mapping_changes<W: Write>(
|
||||
/// Returns loaded point mappings and the number of bytes read from the file.
|
||||
///
|
||||
/// If the file ends with an incomplete entry, it is truncated from the file.
|
||||
pub(super) fn load_mappings(mappings_path: &Path) -> OperationResult<(PointMappings, u64)> {
|
||||
pub(super) fn load_mappings(
|
||||
mappings_path: &Path,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<(PointMappings, u64)> {
|
||||
let file = OneshotFile::open(mappings_path)?;
|
||||
let file_len = file.metadata()?.len();
|
||||
let mut reader = BufReader::new(file);
|
||||
|
||||
let mappings = read_mappings(&mut reader)?;
|
||||
let mappings = read_mappings(&mut reader, deferred_internal_id)?;
|
||||
|
||||
let read_to = reader.stream_position()?;
|
||||
reader.into_inner().drop_cache()?;
|
||||
@@ -198,7 +201,10 @@ where
|
||||
/// Read point mappings from the given reader
|
||||
///
|
||||
/// Returns loaded point mappings.
|
||||
pub(super) fn read_mappings<R>(reader: R) -> OperationResult<PointMappings>
|
||||
pub(super) fn read_mappings<R>(
|
||||
reader: R,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<PointMappings>
|
||||
where
|
||||
R: Read + Seek,
|
||||
{
|
||||
@@ -285,6 +291,7 @@ where
|
||||
internal_to_external,
|
||||
external_to_internal_num,
|
||||
external_to_internal_uuid,
|
||||
deferred_internal_id,
|
||||
);
|
||||
|
||||
Ok(mappings)
|
||||
|
||||
@@ -69,7 +69,10 @@ pub struct MutableIdTracker {
|
||||
}
|
||||
|
||||
impl MutableIdTracker {
|
||||
pub fn open(segment_path: impl Into<PathBuf>) -> OperationResult<Self> {
|
||||
pub fn open(
|
||||
segment_path: impl Into<PathBuf>,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<Self> {
|
||||
let segment_path = segment_path.into();
|
||||
|
||||
let (mappings_path, versions_path) =
|
||||
@@ -93,11 +96,18 @@ impl MutableIdTracker {
|
||||
}
|
||||
|
||||
let (mappings, mappings_expected_len) = if has_mappings {
|
||||
load_mappings(&mappings_path).map_err(|err| {
|
||||
load_mappings(&mappings_path, deferred_internal_id).map_err(|err| {
|
||||
OperationError::service_error(format!("Failed to load ID tracker mappings: {err}"))
|
||||
})?
|
||||
} else {
|
||||
(PointMappings::default(), 0)
|
||||
let mappings = PointMappings::new(
|
||||
Default::default(),
|
||||
Default::default(),
|
||||
Default::default(),
|
||||
Default::default(),
|
||||
deferred_internal_id,
|
||||
);
|
||||
(mappings, 0)
|
||||
};
|
||||
|
||||
let internal_to_version = if has_versions {
|
||||
@@ -210,6 +220,14 @@ impl IdTrackerRead for MutableIdTracker {
|
||||
fn name(&self) -> &'static str {
|
||||
"mutable id tracker"
|
||||
}
|
||||
|
||||
fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
self.mappings.deferred_internal_id()
|
||||
}
|
||||
|
||||
fn deferred_deleted_count(&self) -> usize {
|
||||
self.mappings.deferred_deleted_count()
|
||||
}
|
||||
}
|
||||
|
||||
impl IdTracker for MutableIdTracker {
|
||||
|
||||
@@ -46,6 +46,14 @@ impl IdTrackerRead for ReadOnlyAppendableIdTracker {
|
||||
"read-only appendable id tracker"
|
||||
}
|
||||
|
||||
fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
self.mappings.deferred_internal_id()
|
||||
}
|
||||
|
||||
fn deferred_deleted_count(&self) -> usize {
|
||||
self.mappings.deferred_deleted_count()
|
||||
}
|
||||
|
||||
fn iter_internal_versions(
|
||||
&self,
|
||||
) -> Box<dyn Iterator<Item = (PointOffsetType, SeqNumberType)> + '_> {
|
||||
|
||||
@@ -27,7 +27,7 @@ const DEFAULT_VERSION: SeqNumberType = 42;
|
||||
fn test_iterator() {
|
||||
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
|
||||
|
||||
let mut id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap();
|
||||
let mut id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap();
|
||||
|
||||
id_tracker.set_link(200.into(), 0).unwrap();
|
||||
id_tracker.set_link(100.into(), 1).unwrap();
|
||||
@@ -97,7 +97,7 @@ fn test_load_store() {
|
||||
(id_tracker.mappings, id_tracker.internal_to_version)
|
||||
};
|
||||
|
||||
let mut loaded_id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap();
|
||||
let mut loaded_id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap();
|
||||
|
||||
assert_eq!(
|
||||
old_versions.len(),
|
||||
@@ -155,7 +155,7 @@ fn test_store_load_mutated() {
|
||||
(dropped_points, custom_version)
|
||||
};
|
||||
|
||||
let id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap();
|
||||
let id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap();
|
||||
for (index, point) in TEST_POINTS.iter().enumerate() {
|
||||
let internal_id = index as PointOffsetType;
|
||||
|
||||
@@ -251,7 +251,7 @@ fn test_point_deletion_persists_reload() {
|
||||
};
|
||||
|
||||
// Point should still be gone
|
||||
let id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap();
|
||||
let id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap();
|
||||
assert_eq!(id_tracker.internal_id(point_to_delete), None);
|
||||
|
||||
old_mappings
|
||||
@@ -297,28 +297,28 @@ fn test_point_mappings_de_serialization_single() {
|
||||
fn test_point_mappings_deserializing_special() {
|
||||
// Empty reader creates empty mappings
|
||||
let buf = Cursor::new(b"");
|
||||
assert_eq!(read_mappings(buf).unwrap().total_point_count(), 0);
|
||||
assert_eq!(read_mappings(buf, None).unwrap().total_point_count(), 0);
|
||||
|
||||
// Corrupt if reading invalid type byte
|
||||
let buf = Cursor::new(b"\x00");
|
||||
assert!(
|
||||
read_mappings(buf)
|
||||
read_mappings(buf, None)
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("Corrupted ID tracker mapping storage")
|
||||
);
|
||||
|
||||
let buf = Cursor::new(b"malformed!");
|
||||
assert!(read_mappings(buf).is_err());
|
||||
assert!(read_mappings(buf, None).is_err());
|
||||
|
||||
// Empty if change is not fully written
|
||||
let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00");
|
||||
assert_eq!(read_mappings(buf).unwrap().total_point_count(), 0);
|
||||
assert_eq!(read_mappings(buf, None).unwrap().total_point_count(), 0);
|
||||
|
||||
// Exactly one entry
|
||||
let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00\x00\x00\x00\x02\x00\x00\x00");
|
||||
assert_eq!(
|
||||
read_mappings(buf)
|
||||
read_mappings(buf, None)
|
||||
.unwrap()
|
||||
.internal_id(&PointIdType::NumId(1)),
|
||||
Some(2)
|
||||
@@ -327,7 +327,7 @@ fn test_point_mappings_deserializing_special() {
|
||||
// Exactly one entry and a malformed second one
|
||||
let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00\x00\x00\x00\x02\x00\x00\x00\x00\x01\x00");
|
||||
assert!(
|
||||
read_mappings(buf)
|
||||
read_mappings(buf, None)
|
||||
.unwrap_err()
|
||||
.to_string()
|
||||
.contains("Corrupted ID tracker mapping storage")
|
||||
@@ -335,7 +335,7 @@ fn test_point_mappings_deserializing_special() {
|
||||
|
||||
// Exactly one entry and an incomplete second one
|
||||
let buf = Cursor::new(b"\x01\x01\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x00\x03\x01\x00");
|
||||
let mappings = read_mappings(buf).unwrap();
|
||||
let mappings = read_mappings(buf, None).unwrap();
|
||||
assert_eq!(mappings.total_point_count(), 1);
|
||||
assert_eq!(mappings.internal_id(&PointIdType::NumId(1)), Some(0));
|
||||
}
|
||||
@@ -397,7 +397,7 @@ fn test_point_mappings_truncation() {
|
||||
.unwrap();
|
||||
assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 13);
|
||||
assert_eq!(
|
||||
load_mappings(&mappings_path)
|
||||
load_mappings(&mappings_path, None)
|
||||
.unwrap()
|
||||
.0
|
||||
.internal_id(&PointIdType::NumId(1)),
|
||||
@@ -413,7 +413,7 @@ fn test_point_mappings_truncation() {
|
||||
.unwrap();
|
||||
assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 14);
|
||||
assert_eq!(
|
||||
load_mappings(&mappings_path)
|
||||
load_mappings(&mappings_path, None)
|
||||
.unwrap()
|
||||
.0
|
||||
.internal_id(&PointIdType::NumId(1)),
|
||||
@@ -429,7 +429,7 @@ fn test_point_mappings_truncation() {
|
||||
.unwrap();
|
||||
assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 16);
|
||||
assert_eq!(
|
||||
load_mappings(&mappings_path)
|
||||
load_mappings(&mappings_path, None)
|
||||
.unwrap()
|
||||
.0
|
||||
.internal_id(&PointIdType::NumId(1)),
|
||||
@@ -444,7 +444,7 @@ fn test_point_mappings_truncation() {
|
||||
).unwrap();
|
||||
assert_eq!(fs::metadata(&mappings_path).unwrap().len(), 28);
|
||||
assert_eq!(
|
||||
load_mappings(&mappings_path)
|
||||
load_mappings(&mappings_path, None)
|
||||
.unwrap()
|
||||
.0
|
||||
.internal_id(&PointIdType::NumId(1)),
|
||||
@@ -468,7 +468,8 @@ fn make_in_memory_tracker_from_memory() -> InMemoryIdTracker {
|
||||
}
|
||||
|
||||
fn make_mutable_tracker(path: &Path) -> MutableIdTracker {
|
||||
let mut id_tracker = MutableIdTracker::open(path).expect("failed to open mutable ID tracker");
|
||||
let mut id_tracker =
|
||||
MutableIdTracker::open(path, None).expect("failed to open mutable ID tracker");
|
||||
|
||||
for value in TEST_POINTS.iter() {
|
||||
let internal_id = id_tracker.total_point_count() as PointOffsetType;
|
||||
@@ -537,7 +538,7 @@ fn simple_id_tracker_vs_mutable_tracker_congruence() {
|
||||
let segment_dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
|
||||
let db = open_db(segment_dir.path(), &[DB_VECTOR_CF]).unwrap();
|
||||
|
||||
let mut mutable_id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap();
|
||||
let mut mutable_id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap();
|
||||
let mut simple_id_tracker = SimpleIdTracker::open(db).unwrap();
|
||||
|
||||
// Insert 100 random points into id_tracker
|
||||
@@ -612,7 +613,7 @@ fn simple_id_tracker_vs_mutable_tracker_congruence() {
|
||||
mutable_id_tracker.mapping_flusher()().unwrap();
|
||||
mutable_id_tracker.versions_flusher()().unwrap();
|
||||
drop(mutable_id_tracker);
|
||||
let mutable_id_tracker = MutableIdTracker::open(segment_dir.path()).unwrap();
|
||||
let mutable_id_tracker = MutableIdTracker::open(segment_dir.path(), None).unwrap();
|
||||
|
||||
check_trackers(&simple_id_tracker, &mutable_id_tracker);
|
||||
}
|
||||
|
||||
@@ -35,6 +35,14 @@ pub struct PointMappings {
|
||||
// Having two separate maps allows us iterating only over one type at a time without having to filter.
|
||||
external_to_internal_num: BTreeMap<u64, PointOffsetType>,
|
||||
external_to_internal_uuid: BTreeMap<Uuid, PointOffsetType>,
|
||||
|
||||
/// Points with internal id >= this value are hidden from reads.
|
||||
/// Only set for appendable segments with deferred points.
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
|
||||
/// Number of deleted deferred points. Maintained incrementally so we can
|
||||
/// derive the visible deferred count without re-scanning the deleted bitslice.
|
||||
deferred_deleted_count: usize,
|
||||
}
|
||||
|
||||
impl PointMappings {
|
||||
@@ -43,12 +51,25 @@ impl PointMappings {
|
||||
internal_to_external: Vec<PointIdType>,
|
||||
external_to_internal_num: BTreeMap<u64, PointOffsetType>,
|
||||
external_to_internal_uuid: BTreeMap<Uuid, PointOffsetType>,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> Self {
|
||||
let deferred_deleted_count = deferred_internal_id
|
||||
.map(|deferred_from| {
|
||||
let total = deleted.len();
|
||||
if total <= deferred_from as usize {
|
||||
0
|
||||
} else {
|
||||
deleted[deferred_from as usize..total].count_ones()
|
||||
}
|
||||
})
|
||||
.unwrap_or(0);
|
||||
Self {
|
||||
deleted,
|
||||
internal_to_external,
|
||||
external_to_internal_num,
|
||||
external_to_internal_uuid,
|
||||
deferred_internal_id,
|
||||
deferred_deleted_count,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -109,7 +130,22 @@ impl PointMappings {
|
||||
}
|
||||
|
||||
if let Some(internal_id) = &internal_id {
|
||||
let was_already_deleted = *self
|
||||
.deleted
|
||||
.get(*internal_id as usize)
|
||||
.as_deref()
|
||||
.unwrap_or(&true);
|
||||
self.deleted.set(*internal_id as usize, true);
|
||||
|
||||
// Count newly-deleted deferred points so we can report visible deferred totals
|
||||
// without rescanning the deleted bitslice.
|
||||
if !was_already_deleted
|
||||
&& self
|
||||
.deferred_internal_id
|
||||
.is_some_and(|deferred_from| *internal_id >= deferred_from)
|
||||
{
|
||||
self.deferred_deleted_count += 1;
|
||||
}
|
||||
}
|
||||
|
||||
internal_id
|
||||
@@ -268,6 +304,14 @@ impl PointMappings {
|
||||
self.internal_to_external.len()
|
||||
}
|
||||
|
||||
pub(crate) fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
self.deferred_internal_id
|
||||
}
|
||||
|
||||
pub(crate) fn deferred_deleted_count(&self) -> usize {
|
||||
self.deferred_deleted_count
|
||||
}
|
||||
|
||||
/// Generate a random [`PointMappings`].
|
||||
#[cfg(test)]
|
||||
pub fn random(rand: &mut StdRng, total_size: u32) -> Self {
|
||||
@@ -326,6 +370,8 @@ impl PointMappings {
|
||||
internal_to_external,
|
||||
external_to_internal_num,
|
||||
external_to_internal_uuid,
|
||||
deferred_internal_id: None,
|
||||
deferred_deleted_count: 0,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -349,6 +395,8 @@ impl PointMappings {
|
||||
internal_to_external,
|
||||
external_to_internal_num,
|
||||
external_to_internal_uuid,
|
||||
deferred_internal_id: _,
|
||||
deferred_deleted_count: _,
|
||||
} = self;
|
||||
|
||||
let deleted_bytes = deleted.capacity().div_ceil(u8::BITS as usize);
|
||||
|
||||
@@ -8,7 +8,7 @@ use common::cow::BoxCow;
|
||||
#[cfg(target_os = "linux")]
|
||||
use common::cpu::linux_low_thread_priority;
|
||||
use common::progress_tracker::ProgressTracker;
|
||||
use common::types::PointOffsetType;
|
||||
use common::types::{DeferredBehavior, PointOffsetType};
|
||||
use fs_err as fs;
|
||||
use log::{debug, trace};
|
||||
use rand::Rng;
|
||||
@@ -459,7 +459,6 @@ impl HNSWIndex {
|
||||
|
||||
let points_to_index = condition_points(
|
||||
payload_block.condition,
|
||||
id_tracker_ref.deref(),
|
||||
&payload_index_ref,
|
||||
&vector_storage_ref,
|
||||
stopped,
|
||||
@@ -601,7 +600,6 @@ impl HNSWIndex {
|
||||
/// Get list of points for indexing, associated with payload block filtering condition
|
||||
fn condition_points(
|
||||
condition: FieldCondition,
|
||||
id_tracker: &IdTrackerEnum,
|
||||
payload_index: &StructPayloadIndex,
|
||||
vector_storage: &VectorStorageEnum,
|
||||
stopped: &AtomicBool,
|
||||
@@ -612,18 +610,14 @@ fn condition_points(
|
||||
|
||||
let deleted_bitslice = vector_storage.deleted_vector_bitslice();
|
||||
|
||||
let point_mappings = id_tracker.point_mappings();
|
||||
|
||||
payload_index.with_view(|v| {
|
||||
let cardinality_estimation = v.estimate_cardinality(&filter, &disposed_hw_counter)?;
|
||||
Ok(v.iter_filtered_points(
|
||||
&filter,
|
||||
id_tracker,
|
||||
&point_mappings,
|
||||
&cardinality_estimation,
|
||||
&disposed_hw_counter,
|
||||
stopped,
|
||||
None,
|
||||
DeferredBehavior::IncludeAll,
|
||||
)?
|
||||
.filter(|&point_id| !deleted_bitslice.get_bit(point_id as usize).unwrap_or(false))
|
||||
.collect())
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
use common::bitvec::BitSlice;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::cow::BoxCow;
|
||||
use common::types::{PointOffsetType, ScoredPointOffset};
|
||||
use common::types::{DeferredBehavior, PointOffsetType, ScoredPointOffset};
|
||||
|
||||
use super::HNSWIndex;
|
||||
use crate::common::operation_error::OperationResult;
|
||||
@@ -298,21 +298,17 @@ impl HNSWIndex {
|
||||
let hw_counter = &vector_query_context.hardware_counter();
|
||||
let is_stopped = &vector_query_context.is_stopped();
|
||||
|
||||
let id_tracker = self.id_tracker.borrow();
|
||||
let payload_index = self.payload_index.borrow();
|
||||
let point_mappings = id_tracker.point_mappings();
|
||||
// Assume query is already estimated to be small enough so we can iterate over all matched ids
|
||||
let filtered_points: Vec<PointOffsetType> = payload_index.with_view(|v| {
|
||||
let query_cardinality = v.estimate_cardinality(filter, hw_counter)?;
|
||||
v.iter_filtered_points(
|
||||
filter,
|
||||
&*id_tracker,
|
||||
&point_mappings,
|
||||
&query_cardinality,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
// No deferred filtering here since it's HNSW index.
|
||||
None,
|
||||
DeferredBehavior::IncludeAll,
|
||||
)
|
||||
.map(|it| it.collect())
|
||||
})?;
|
||||
|
||||
@@ -347,21 +347,34 @@ impl<'a> BatchFilteredSearcher<'a> {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn peek_top_all(
|
||||
self,
|
||||
is_stopped: &AtomicBool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<Vec<Vec<ScoredPointOffset>>> {
|
||||
let iter = self
|
||||
.filters
|
||||
/// Iterator over every internal point ID that isn't soft-deleted in this
|
||||
/// searcher's `point_deleted` bitslice.
|
||||
///
|
||||
/// Does not apply deferred-point filtering — wrap with
|
||||
/// `PointMappingsRefEnum::filter_deferred` (or compose otherwise) before
|
||||
/// passing to [`Self::peek_top_iter`] when deferred awareness is needed.
|
||||
///
|
||||
/// The returned iterator borrows the underlying bitslice (lifetime `'a`),
|
||||
/// independent of `&self`, so it can be composed and then passed into
|
||||
/// `peek_top_iter(self, ...)` which consumes the searcher.
|
||||
pub fn iter_not_deleted(&self) -> impl Iterator<Item = PointOffsetType> + 'a {
|
||||
self.filters
|
||||
.point_deleted
|
||||
.iter_zeros()
|
||||
.map(|p| p as PointOffsetType)
|
||||
.take_while(|&point_id| {
|
||||
// Early exit if we hit the max point ID (e.g. a deferred point).
|
||||
point_id < deferred_internal_id.unwrap_or(PointOffsetType::MAX)
|
||||
});
|
||||
}
|
||||
|
||||
/// Score every non-deleted point without deferred filtering.
|
||||
///
|
||||
/// Production paths compose `iter_not_deleted` with
|
||||
/// `PointMappingsRefEnum::filter_deferred` and call
|
||||
/// [`Self::peek_top_iter`] directly.
|
||||
#[cfg(feature = "testing")]
|
||||
pub fn peek_top_all(
|
||||
self,
|
||||
is_stopped: &AtomicBool,
|
||||
) -> OperationResult<Vec<Vec<ScoredPointOffset>>> {
|
||||
let iter = self.iter_not_deleted();
|
||||
self.peek_top_iter(iter, is_stopped)
|
||||
}
|
||||
|
||||
|
||||
@@ -4,7 +4,7 @@ use std::sync::atomic::AtomicBool;
|
||||
|
||||
use ahash::AHashMap;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::types::{PointOffsetType, ScoreType};
|
||||
use common::types::{DeferredBehavior, PointOffsetType, ScoreType};
|
||||
use serde_json::Value;
|
||||
|
||||
use super::field_index::numeric_index::NumericFieldIndexRead;
|
||||
@@ -13,7 +13,6 @@ use super::query_optimization::rescore_formula::FormulaScorer;
|
||||
use super::query_optimization::rescore_formula::parsed_formula::ParsedFormula;
|
||||
use crate::common::Flusher;
|
||||
use crate::common::operation_error::OperationResult;
|
||||
use crate::id_tracker::{IdTrackerRead, PointMappingsRefEnum};
|
||||
use crate::index::field_index::{CardinalityEstimation, PayloadBlockCondition};
|
||||
use crate::json_path::JsonPath;
|
||||
use crate::payload_storage::FilterContext;
|
||||
@@ -66,7 +65,6 @@ pub trait PayloadIndexRead {
|
||||
filter: &Filter,
|
||||
hw_counter: &HardwareCounterCell,
|
||||
is_stopped: &AtomicBool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<Vec<PointOffsetType>>;
|
||||
|
||||
/// Return number of points, indexed by this field
|
||||
@@ -106,19 +104,16 @@ pub trait PayloadIndexRead {
|
||||
|
||||
/// Iterate point offsets that match the filter.
|
||||
///
|
||||
/// Generic over `I: IdTrackerRead` so callers pass their concrete tracker
|
||||
/// without dynamic dispatch; the iterator return uses RPITIT so each impl
|
||||
/// keeps its own zero-cost concrete chain.
|
||||
#[allow(clippy::too_many_arguments)]
|
||||
fn iter_filtered_points<'a, I: IdTrackerRead>(
|
||||
/// The iterator return uses RPITIT so each impl keeps its own zero-cost
|
||||
/// concrete chain. The id tracker is read from `&self`, so impls reach it
|
||||
/// through their own field rather than receiving a separate parameter.
|
||||
fn iter_filtered_points<'a>(
|
||||
&'a self,
|
||||
filter: &'a Filter,
|
||||
id_tracker: &'a I,
|
||||
point_mappings: &'a PointMappingsRefEnum<'a>,
|
||||
query_cardinality: &'a CardinalityEstimation,
|
||||
hw_counter: &'a HardwareCounterCell,
|
||||
is_stopped: &'a AtomicBool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> OperationResult<impl Iterator<Item = PointOffsetType> + 'a>;
|
||||
|
||||
/// Iterate conditions for payload blocks with minimum size of `threshold`
|
||||
|
||||
@@ -7,7 +7,7 @@ use ahash::AHashMap;
|
||||
use atomic_refcell::AtomicRefCell;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::iterator_ext::IteratorExt;
|
||||
use common::types::{PointOffsetType, ScoreType};
|
||||
use common::types::{DeferredBehavior, PointOffsetType, ScoreType};
|
||||
use fs_err as fs;
|
||||
use schemars::_serde_json::Value;
|
||||
|
||||
@@ -15,7 +15,7 @@ use super::field_index::FieldIndex;
|
||||
use super::payload_config::PayloadFieldSchemaWithIndexType;
|
||||
use crate::common::Flusher;
|
||||
use crate::common::operation_error::{OperationError, OperationResult};
|
||||
use crate::id_tracker::{IdTrackerEnum, IdTrackerRead, PointMappingsRefEnum};
|
||||
use crate::id_tracker::{IdTrackerEnum, IdTrackerRead};
|
||||
use crate::index::field_index::facet_index::FacetIndexEnum;
|
||||
use crate::index::field_index::numeric_index::{NumericFieldIndex, NumericFieldIndexRead};
|
||||
use crate::index::field_index::{CardinalityEstimation, FacetIndex, PayloadBlockCondition};
|
||||
@@ -111,12 +111,11 @@ impl PayloadIndexRead for PlainPayloadIndex {
|
||||
filter: &Filter,
|
||||
hw_counter: &HardwareCounterCell,
|
||||
is_stopped: &AtomicBool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<Vec<PointOffsetType>> {
|
||||
let filter_context = self.filter_context(filter, hw_counter)?;
|
||||
let id_tracker = self.id_tracker.borrow();
|
||||
let point_mappings = id_tracker.point_mappings();
|
||||
let all_points_iter = point_mappings.iter_internal_visible(deferred_internal_id);
|
||||
let all_points_iter = point_mappings.iter_internal_visible();
|
||||
Ok(all_points_iter
|
||||
.stop_if(is_stopped)
|
||||
.filter(|id| filter_context.check(*id))
|
||||
@@ -190,21 +189,26 @@ impl PayloadIndexRead for PlainPayloadIndex {
|
||||
))
|
||||
}
|
||||
|
||||
fn iter_filtered_points<'a, I: IdTrackerRead>(
|
||||
fn iter_filtered_points<'a>(
|
||||
&'a self,
|
||||
filter: &'a Filter,
|
||||
_id_tracker: &'a I,
|
||||
point_mappings: &'a PointMappingsRefEnum<'a>,
|
||||
_query_cardinality: &'a CardinalityEstimation,
|
||||
hw_counter: &'a HardwareCounterCell,
|
||||
is_stopped: &'a AtomicBool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> OperationResult<impl Iterator<Item = PointOffsetType> + 'a> {
|
||||
let filter_context = self.filter_context(filter, hw_counter)?;
|
||||
let all_points_iter = point_mappings.iter_internal_visible(deferred_internal_id);
|
||||
Ok(all_points_iter
|
||||
// `self.id_tracker` is an `Arc<AtomicRefCell<_>>`, so the mapping borrow is
|
||||
// local; collect eagerly to detach the iterator from the borrow.
|
||||
let matched: Vec<PointOffsetType> = self
|
||||
.id_tracker
|
||||
.borrow()
|
||||
.point_mappings()
|
||||
.iter_internal_with_behavior(deferred_behavior)
|
||||
.stop_if(is_stopped)
|
||||
.filter(move |id| filter_context.check(*id)))
|
||||
.filter(|id| filter_context.check(*id))
|
||||
.collect();
|
||||
Ok(matched.into_iter())
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -4,7 +4,7 @@ use std::sync::Arc;
|
||||
|
||||
use atomic_refcell::AtomicRefCell;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::types::{PointOffsetType, ScoredPointOffset, TelemetryDetail};
|
||||
use common::types::{DeferredBehavior, PointOffsetType, ScoredPointOffset, TelemetryDetail};
|
||||
use parking_lot::Mutex;
|
||||
use sparse::common::types::DimId;
|
||||
|
||||
@@ -137,16 +137,20 @@ impl VectorIndexRead for PlainVectorIndex {
|
||||
query_context.hardware_counter(),
|
||||
)?;
|
||||
|
||||
let deferred_internal_id = query_context.deferred_internal_id();
|
||||
|
||||
let mut search_results = match filter {
|
||||
Some(filter) => {
|
||||
let filtered_ids_vec = self.payload_index.borrow().with_view(|v| {
|
||||
v.query_points(filter, &hw_counter, &is_stopped, deferred_internal_id)
|
||||
})?;
|
||||
let filtered_ids_vec = self
|
||||
.payload_index
|
||||
.borrow()
|
||||
.with_view(|v| v.query_points(filter, &hw_counter, &is_stopped))?;
|
||||
batch_searcher.peek_top_iter(filtered_ids_vec.iter().copied(), &is_stopped)?
|
||||
}
|
||||
None => batch_searcher.peek_top_all(&is_stopped, deferred_internal_id)?,
|
||||
None => {
|
||||
let iter = id_tracker
|
||||
.point_mappings()
|
||||
.filter_deferred(batch_searcher.iter_not_deleted(), DeferredBehavior::Exclude);
|
||||
batch_searcher.peek_top_iter(iter, &is_stopped)?
|
||||
}
|
||||
};
|
||||
|
||||
for (search_result, query_vector) in search_results.iter_mut().zip(query_vectors) {
|
||||
|
||||
@@ -8,7 +8,6 @@ use atomic_refcell::AtomicRefCell;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::generic_consts::Random;
|
||||
use common::storage_version::StorageVersion as _;
|
||||
use common::types::PointOffsetType;
|
||||
use fs_err as fs;
|
||||
use sparse::common::scores_memory_pool::ScoresMemoryPool;
|
||||
use sparse::common::sparse_vector::SparseVector;
|
||||
@@ -37,7 +36,6 @@ pub struct SparseVectorIndex<TInvertedIndex: InvertedIndex> {
|
||||
searches_telemetry: SparseSearchesTelemetry,
|
||||
indices_tracker: IndicesTracker,
|
||||
scores_memory_pool: ScoresMemoryPool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
}
|
||||
|
||||
/// Getters for internals, used for testing.
|
||||
@@ -72,7 +70,6 @@ pub struct SparseVectorIndexOpenArgs<'a, F: FnMut()> {
|
||||
pub path: &'a Path,
|
||||
pub stopped: &'a AtomicBool,
|
||||
pub tick_progress: F,
|
||||
pub deferred_internal_id: Option<PointOffsetType>,
|
||||
}
|
||||
|
||||
impl<TInvertedIndex: InvertedIndex> SparseVectorIndex<TInvertedIndex> {
|
||||
@@ -86,7 +83,6 @@ impl<TInvertedIndex: InvertedIndex> SparseVectorIndex<TInvertedIndex> {
|
||||
path,
|
||||
stopped,
|
||||
tick_progress,
|
||||
deferred_internal_id,
|
||||
} = args;
|
||||
|
||||
let config_path = SparseIndexConfig::get_config_path(path);
|
||||
@@ -145,7 +141,6 @@ impl<TInvertedIndex: InvertedIndex> SparseVectorIndex<TInvertedIndex> {
|
||||
searches_telemetry,
|
||||
indices_tracker,
|
||||
scores_memory_pool,
|
||||
deferred_internal_id,
|
||||
})
|
||||
}
|
||||
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::types::{PointOffsetType, ScoredPointOffset};
|
||||
use common::types::{DeferredBehavior, PointOffsetType, ScoredPointOffset};
|
||||
use itertools::Itertools;
|
||||
use sparse::common::sparse_vector::SparseVector;
|
||||
use sparse::index::inverted_index::InvertedIndex;
|
||||
@@ -71,14 +71,10 @@ impl<TInvertedIndex: InvertedIndex> SparseVectorIndex<TInvertedIndex> {
|
||||
// `prefiltered_points` always contains visible points only so we don't need additional filtering here.
|
||||
Some(filtered_points) => filtered_points.iter().copied(),
|
||||
None => {
|
||||
let filtered_points = self.payload_index.borrow().with_view(|v| {
|
||||
v.query_points(
|
||||
filter,
|
||||
&hw_counter,
|
||||
&is_stopped,
|
||||
vector_query_context.deferred_internal_id(),
|
||||
)
|
||||
})?;
|
||||
let filtered_points = self
|
||||
.payload_index
|
||||
.borrow()
|
||||
.with_view(|v| v.query_points(filter, &hw_counter, &is_stopped))?;
|
||||
*prefiltered_points = Some(filtered_points);
|
||||
prefiltered_points.as_ref().unwrap().iter().copied()
|
||||
}
|
||||
@@ -86,7 +82,10 @@ impl<TInvertedIndex: InvertedIndex> SparseVectorIndex<TInvertedIndex> {
|
||||
searcher.peek_top_iter(filtered_points, &is_stopped)?
|
||||
}
|
||||
None => {
|
||||
searcher.peek_top_all(&is_stopped, vector_query_context.deferred_internal_id())?
|
||||
let iter = id_tracker
|
||||
.point_mappings()
|
||||
.filter_deferred(searcher.iter_not_deleted(), DeferredBehavior::Exclude);
|
||||
searcher.peek_top_iter(iter, &is_stopped)?
|
||||
}
|
||||
};
|
||||
let res = results.pop().expect("single element results");
|
||||
@@ -119,14 +118,10 @@ impl<TInvertedIndex: InvertedIndex> SparseVectorIndex<TInvertedIndex> {
|
||||
// so no additional filtering is required in that case.
|
||||
Some(filtered_points) => filtered_points.iter(),
|
||||
None => {
|
||||
let filtered_points = self.payload_index.borrow().with_view(|v| {
|
||||
v.query_points(
|
||||
filter,
|
||||
&hw_counter,
|
||||
&is_stopped,
|
||||
vector_query_context.deferred_internal_id(),
|
||||
)
|
||||
})?;
|
||||
let filtered_points = self
|
||||
.payload_index
|
||||
.borrow()
|
||||
.with_view(|v| v.query_points(filter, &hw_counter, &is_stopped))?;
|
||||
*prefiltered_points = Some(filtered_points);
|
||||
prefiltered_points.as_ref().unwrap().iter()
|
||||
}
|
||||
|
||||
@@ -14,6 +14,7 @@ use crate::common::operation_error::{OperationError, OperationResult, check_proc
|
||||
use crate::data_types::named_vectors::CowVector;
|
||||
use crate::data_types::query_context::VectorQueryContext;
|
||||
use crate::data_types::vectors::{QueryVector, VectorInternal, VectorRef};
|
||||
use crate::id_tracker::IdTrackerRead;
|
||||
use crate::index::sparse_index::indices_tracker::IndicesTracker;
|
||||
use crate::index::sparse_index::sparse_index_config::{SparseIndexConfig, SparseIndexType};
|
||||
use crate::index::{VectorIndex, VectorIndexRead};
|
||||
@@ -34,12 +35,6 @@ impl<TInvertedIndex: InvertedIndex> VectorIndexRead for SparseVectorIndex<TInver
|
||||
let mut results = Vec::with_capacity(vectors.len());
|
||||
let mut prefiltered_points = None;
|
||||
|
||||
debug_assert_eq!(
|
||||
self.deferred_internal_id,
|
||||
query_context.deferred_internal_id(),
|
||||
"SparseIndex and VectorQueryContext deferred_internal_id consistency violated."
|
||||
);
|
||||
|
||||
for vector in vectors {
|
||||
check_process_stopped(&query_context.is_stopped())?;
|
||||
|
||||
@@ -172,7 +167,9 @@ impl<TInvertedIndex: InvertedIndex> VectorIndex for SparseVectorIndex<TInvertedI
|
||||
}
|
||||
|
||||
let point_is_deferred = self
|
||||
.deferred_internal_id
|
||||
.id_tracker
|
||||
.borrow()
|
||||
.deferred_internal_id()
|
||||
.is_some_and(|deferred| id >= deferred);
|
||||
|
||||
if point_is_deferred {
|
||||
|
||||
@@ -6,11 +6,11 @@ use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::counter::iterator_hw_measurement::HwMeasurementIteratorExt;
|
||||
use common::either_variant::EitherVariant;
|
||||
use common::iterator_ext::IteratorExt;
|
||||
use common::types::{PointOffsetType, ScoreType};
|
||||
use common::types::{DeferredBehavior, PointOffsetType, ScoreType};
|
||||
|
||||
use super::StructPayloadIndexReadView;
|
||||
use crate::common::operation_error::OperationResult;
|
||||
use crate::id_tracker::{IdTrackerRead, PointMappingsRefEnum};
|
||||
use crate::id_tracker::IdTrackerRead;
|
||||
use crate::index::PayloadIndexRead;
|
||||
use crate::index::field_index::numeric_index::NumericFieldIndexRead;
|
||||
use crate::index::field_index::{
|
||||
@@ -68,20 +68,16 @@ where
|
||||
filter: &Filter,
|
||||
hw_counter: &HardwareCounterCell,
|
||||
is_stopped: &AtomicBool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<Vec<PointOffsetType>> {
|
||||
// Assume query is already estimated to be small enough so we can iterate over all matched ids
|
||||
let query_cardinality = self.estimate_cardinality(filter, hw_counter)?;
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
let result = self
|
||||
.iter_filtered_points(
|
||||
filter,
|
||||
self.id_tracker,
|
||||
&point_mappings,
|
||||
&query_cardinality,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
deferred_internal_id,
|
||||
DeferredBehavior::Exclude,
|
||||
)?
|
||||
.collect();
|
||||
Ok(result)
|
||||
@@ -143,18 +139,18 @@ where
|
||||
))
|
||||
}
|
||||
|
||||
fn iter_filtered_points<'b, IT: IdTrackerRead>(
|
||||
fn iter_filtered_points<'b>(
|
||||
&'b self,
|
||||
filter: &'b Filter,
|
||||
id_tracker: &'b IT,
|
||||
point_mappings: &'b PointMappingsRefEnum<'b>,
|
||||
query_cardinality: &'b CardinalityEstimation,
|
||||
hw_counter: &'b HardwareCounterCell,
|
||||
is_stopped: &'b AtomicBool,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> OperationResult<impl Iterator<Item = PointOffsetType> + 'b> {
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
|
||||
if query_cardinality.primary_clauses.is_empty() {
|
||||
let full_scan_iterator = point_mappings.iter_internal_visible(deferred_internal_id);
|
||||
let full_scan_iterator = point_mappings.iter_internal_with_behavior(deferred_behavior);
|
||||
|
||||
let struct_filtered_context = self.struct_filtered_context(filter, hw_counter)?;
|
||||
// Worst case: query expected to return few matches, but index can't be used
|
||||
@@ -165,7 +161,7 @@ where
|
||||
Ok(EitherVariant::A(matched_points))
|
||||
} else {
|
||||
// CPU-optimized strategy here: points are made unique before applying other filters.
|
||||
let mut visited_list = self.visited_pool.get(id_tracker.total_point_count());
|
||||
let mut visited_list = self.visited_pool.get(self.id_tracker.total_point_count());
|
||||
|
||||
// If even one iterator is None, we should replace the whole thing with
|
||||
// an iterator over all ids.
|
||||
@@ -180,14 +176,12 @@ where
|
||||
.iter_conditions()
|
||||
.all(|condition| query_cardinality.is_primary(condition));
|
||||
|
||||
let joined_primary_iterator = primary_iterators
|
||||
.into_iter()
|
||||
// Filter out deferred points.
|
||||
// This iterator (and each primary iterator too) can yield items in non sorted order, depending on the type of index and primary condition.
|
||||
.flatten()
|
||||
.filter(move |&internal_id| {
|
||||
internal_id < deferred_internal_id.unwrap_or(PointOffsetType::MAX)
|
||||
})
|
||||
// Primary clause iterators come from field indexes and don't go through
|
||||
// the mapping, so deferred filtering must be applied to them explicitly.
|
||||
// Each primary iterator (and the flattened stream) can yield items in
|
||||
// non-sorted order depending on the field-index type and primary condition.
|
||||
let joined_primary_iterator = point_mappings
|
||||
.filter_deferred(primary_iterators.into_iter().flatten(), deferred_behavior)
|
||||
.stop_if(is_stopped);
|
||||
|
||||
return Ok(if all_conditions_are_primary {
|
||||
@@ -212,7 +206,7 @@ where
|
||||
// and applying full filter.
|
||||
let struct_filtered_context = self.struct_filtered_context(filter, hw_counter)?;
|
||||
|
||||
let id_tracker_iterator = point_mappings.iter_internal_visible(deferred_internal_id);
|
||||
let id_tracker_iterator = point_mappings.iter_internal_with_behavior(deferred_behavior);
|
||||
|
||||
let iter = id_tracker_iterator
|
||||
.stop_if(is_stopped)
|
||||
|
||||
@@ -55,7 +55,7 @@ fn smoke_view_over_in_memory_backends() {
|
||||
// `query_points` over an empty filter on an empty tracker returns nothing.
|
||||
let empty_filter = Filter::default();
|
||||
let result = view
|
||||
.query_points(&empty_filter, &hw_counter, &is_stopped, None)
|
||||
.query_points(&empty_filter, &hw_counter, &is_stopped)
|
||||
.expect("query_points");
|
||||
assert!(result.is_empty(), "no points in tracker");
|
||||
|
||||
|
||||
@@ -20,7 +20,6 @@ impl Segment {
|
||||
payload_storage: payload_storage.deref(),
|
||||
vector_data: &self.vector_data,
|
||||
segment_config: &self.segment_config,
|
||||
deferred_point_status: self.deferred_point_status.as_ref(),
|
||||
appendable_flag: self.appendable_flag,
|
||||
};
|
||||
|
||||
|
||||
@@ -53,7 +53,6 @@ impl Segment {
|
||||
segment_type: _,
|
||||
segment_config,
|
||||
error_status: _,
|
||||
deferred_point_status: _,
|
||||
} = self;
|
||||
|
||||
let sparse_names = &segment_config.sparse_vector_data;
|
||||
|
||||
@@ -21,7 +21,6 @@ use std::sync::Arc;
|
||||
use atomic_refcell::AtomicRefCell;
|
||||
use common::is_alive_lock::IsAliveLock;
|
||||
use common::storage_version::StorageVersion;
|
||||
use common::types::PointOffsetType;
|
||||
use parking_lot::Mutex;
|
||||
use uuid::Uuid;
|
||||
|
||||
@@ -87,18 +86,6 @@ pub struct Segment {
|
||||
/// Last unhandled error
|
||||
/// If not None, all update operations will be aborted until original operation is performed properly
|
||||
pub error_status: Option<SegmentFailedState>,
|
||||
pub(crate) deferred_point_status: Option<DeferredPointStatus>,
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct DeferredPointStatus {
|
||||
/// Points with internal id >= this value are hidden from reads.
|
||||
/// Available for appendable segments only.
|
||||
pub(crate) deferred_internal_id: PointOffsetType,
|
||||
|
||||
/// Amount of deleted deferred points. Must kept track of properly to be able
|
||||
/// to calculate the amount of available deferred and visible points.
|
||||
pub(crate) deferred_deleted_count: usize,
|
||||
}
|
||||
|
||||
pub struct VectorData {
|
||||
|
||||
@@ -15,11 +15,11 @@ where
|
||||
TVD: VectorDataRead,
|
||||
{
|
||||
pub(super) fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
self.deferred_point_status.map(|s| s.deferred_internal_id)
|
||||
self.id_tracker.deferred_internal_id()
|
||||
}
|
||||
|
||||
pub(super) fn deferred_deleted_count(&self) -> Option<usize> {
|
||||
self.deferred_point_status.map(|s| s.deferred_deleted_count)
|
||||
pub(super) fn deferred_deleted_count(&self) -> usize {
|
||||
self.id_tracker.deferred_deleted_count()
|
||||
}
|
||||
|
||||
pub fn deferred_point_count(&self) -> usize {
|
||||
@@ -28,7 +28,7 @@ where
|
||||
.id_tracker
|
||||
.total_point_count()
|
||||
.saturating_sub(internal_id as usize)
|
||||
.saturating_sub(self.deferred_deleted_count().unwrap_or_default()),
|
||||
.saturating_sub(self.deferred_deleted_count()),
|
||||
None => 0,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2,7 +2,7 @@ use std::collections::{BTreeSet, HashMap};
|
||||
use std::sync::atomic::AtomicBool;
|
||||
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::types::PointOffsetType;
|
||||
use common::types::{DeferredBehavior, PointOffsetType};
|
||||
use itertools::Itertools;
|
||||
|
||||
use crate::common::operation_error::{OperationError, OperationResult, check_process_stopped};
|
||||
@@ -67,17 +67,14 @@ where
|
||||
|
||||
if use_iterative_approach {
|
||||
// Go over the filtered points and aggregate the values (read from other indexes).
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
let points = self
|
||||
.payload_index
|
||||
.iter_filtered_points(
|
||||
filter,
|
||||
self.id_tracker,
|
||||
&point_mappings,
|
||||
&filter_cardinality,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
self.deferred_internal_id(),
|
||||
DeferredBehavior::Exclude,
|
||||
)?
|
||||
.filter(|&point_id| !self.id_tracker.is_deleted_point(point_id));
|
||||
facet_index.for_points_values(points, hw_counter, |_point_id, iter| {
|
||||
@@ -150,18 +147,15 @@ where
|
||||
let filter_cardinality = self
|
||||
.payload_index
|
||||
.estimate_cardinality(filter, hw_counter)?;
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
|
||||
let points = self
|
||||
.payload_index
|
||||
.iter_filtered_points(
|
||||
filter,
|
||||
self.id_tracker,
|
||||
&point_mappings,
|
||||
&filter_cardinality,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
self.deferred_internal_id(),
|
||||
DeferredBehavior::Exclude,
|
||||
)?
|
||||
.filter(|&point_id| !self.id_tracker.is_deleted_point(point_id));
|
||||
facet_index.for_points_values(points, hw_counter, |_point_id, iter| {
|
||||
|
||||
@@ -83,7 +83,7 @@ where
|
||||
num_indexed_vectors: num_indexed_vectors_total,
|
||||
num_points: self.id_tracker.available_point_count(),
|
||||
num_deferred_points: Some(self.deferred_point_count()),
|
||||
num_deleted_deferred_points: Some(self.deferred_deleted_count().unwrap_or_default()),
|
||||
num_deleted_deferred_points: Some(self.deferred_deleted_count()),
|
||||
num_deleted_vectors: self.id_tracker.deleted_point_count(),
|
||||
vectors_size_bytes, // Considers vector storage, but not indices.
|
||||
payloads_size_bytes, // Considers payload storage, but not indices.
|
||||
|
||||
@@ -18,8 +18,8 @@ use crate::index::field_index::FieldIndex;
|
||||
use crate::index::struct_payload_index::StructPayloadIndexReadView;
|
||||
use crate::payload_storage::PayloadStorageRead;
|
||||
use crate::payload_storage::payload_storage_enum::PayloadStorageEnum;
|
||||
use crate::segment::VectorData;
|
||||
use crate::segment::vector_data_read::VectorDataRead;
|
||||
use crate::segment::{DeferredPointStatus, VectorData};
|
||||
use crate::types::{PointIdType, SegmentConfig, SeqNumberType, VectorNameBuf};
|
||||
use crate::vector_storage::VectorStorageEnum;
|
||||
|
||||
@@ -41,7 +41,6 @@ where
|
||||
pub(crate) payload_storage: &'s TPayloadStorage,
|
||||
pub(crate) vector_data: &'s HashMap<VectorNameBuf, TVectorData>,
|
||||
pub(crate) segment_config: &'s SegmentConfig,
|
||||
pub(crate) deferred_point_status: Option<&'s DeferredPointStatus>,
|
||||
pub(crate) appendable_flag: bool,
|
||||
}
|
||||
|
||||
|
||||
@@ -46,19 +46,14 @@ where
|
||||
|
||||
let start_from = order_by.start_from();
|
||||
|
||||
let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id());
|
||||
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
let values_ids_iterator = self
|
||||
.payload_index
|
||||
.iter_filtered_points(
|
||||
condition,
|
||||
self.id_tracker,
|
||||
&point_mappings,
|
||||
&cardinality_estimation,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
effective_deferred_id,
|
||||
deferred_behavior,
|
||||
)?
|
||||
.flat_map(|internal_id| {
|
||||
// Repeat a point for as many values as it has.
|
||||
|
||||
@@ -2,6 +2,7 @@ use std::sync::atomic::AtomicBool;
|
||||
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::iterator_ext::IteratorExt;
|
||||
use common::types::DeferredBehavior;
|
||||
use rand::seq::{IteratorRandom, SliceRandom};
|
||||
|
||||
use crate::common::operation_error::OperationResult;
|
||||
@@ -22,7 +23,7 @@ where
|
||||
pub fn read_by_random_id(&self, limit: usize) -> Vec<PointIdType> {
|
||||
self.id_tracker
|
||||
.point_mappings()
|
||||
.iter_random_visible(self.deferred_internal_id())
|
||||
.iter_random_visible()
|
||||
.map(|x| x.0)
|
||||
.take(limit)
|
||||
.collect()
|
||||
@@ -35,7 +36,6 @@ where
|
||||
is_stopped: &AtomicBool,
|
||||
hw_counter: &HardwareCounterCell,
|
||||
) -> OperationResult<Vec<PointIdType>> {
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
let cardinality_estimation = self
|
||||
.payload_index
|
||||
.estimate_cardinality(condition, hw_counter)?;
|
||||
@@ -43,12 +43,10 @@ where
|
||||
.payload_index
|
||||
.iter_filtered_points(
|
||||
condition,
|
||||
self.id_tracker,
|
||||
&point_mappings,
|
||||
&cardinality_estimation,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
self.deferred_internal_id(),
|
||||
DeferredBehavior::Exclude,
|
||||
)?
|
||||
.filter_map(|internal_id| self.id_tracker.external_id(internal_id));
|
||||
|
||||
@@ -69,7 +67,7 @@ where
|
||||
Ok(self
|
||||
.id_tracker
|
||||
.point_mappings()
|
||||
.iter_random_visible(self.deferred_internal_id())
|
||||
.iter_random_visible()
|
||||
.stop_if(is_stopped)
|
||||
.filter(move |(_, internal_id)| filter_context.check(*internal_id))
|
||||
.map(|(external_id, _)| external_id)
|
||||
|
||||
@@ -59,12 +59,14 @@ where
|
||||
limit: Option<usize>,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> Vec<PointIdType> {
|
||||
let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id());
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
let iter = if deferred_behavior.include_all_points() {
|
||||
point_mappings.iter_from(offset)
|
||||
} else {
|
||||
point_mappings.iter_from_visible(offset)
|
||||
};
|
||||
|
||||
self.id_tracker
|
||||
.point_mappings()
|
||||
.iter_from_visible(offset, effective_deferred_id)
|
||||
.map(|x| x.0)
|
||||
iter.map(|x| x.0)
|
||||
.take(limit.unwrap_or(usize::MAX))
|
||||
.collect()
|
||||
}
|
||||
@@ -78,23 +80,18 @@ where
|
||||
hw_counter: &HardwareCounterCell,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> OperationResult<Vec<PointIdType>> {
|
||||
let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id());
|
||||
|
||||
let cardinality_estimation = self
|
||||
.payload_index
|
||||
.estimate_cardinality(condition, hw_counter)?;
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
|
||||
let ids_iterator = self
|
||||
.payload_index
|
||||
.iter_filtered_points(
|
||||
condition,
|
||||
self.id_tracker,
|
||||
&point_mappings,
|
||||
&cardinality_estimation,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
effective_deferred_id,
|
||||
deferred_behavior,
|
||||
)?
|
||||
.filter_map(|internal_id| {
|
||||
let external_id = self.id_tracker.external_id(internal_id)?;
|
||||
@@ -121,13 +118,15 @@ where
|
||||
hw_counter: &HardwareCounterCell,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> OperationResult<Vec<PointIdType>> {
|
||||
let effective_deferred_id = deferred_behavior.apply(self.deferred_internal_id());
|
||||
|
||||
let filter_context = self.payload_index.filter_context(condition, hw_counter)?;
|
||||
Ok(self
|
||||
.id_tracker
|
||||
.point_mappings()
|
||||
.iter_from_visible(offset, effective_deferred_id)
|
||||
let point_mappings = self.id_tracker.point_mappings();
|
||||
let iter = if deferred_behavior.include_all_points() {
|
||||
point_mappings.iter_from(offset)
|
||||
} else {
|
||||
point_mappings.iter_from_visible(offset)
|
||||
};
|
||||
|
||||
Ok(iter
|
||||
.stop_if(is_stopped)
|
||||
.filter(move |(_, internal_id)| filter_context.check(*internal_id))
|
||||
.map(|(external_id, _)| external_id)
|
||||
|
||||
@@ -2,13 +2,14 @@ use std::sync::atomic::AtomicBool;
|
||||
|
||||
use ahash::AHashMap;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::iterator_ext::IteratorExt;
|
||||
use common::types::{DeferredBehavior, ScoredPointOffset};
|
||||
|
||||
use crate::common::operation_error::{OperationError, OperationResult};
|
||||
use crate::common::{check_query_vectors, check_stopped};
|
||||
use crate::data_types::query_context::{QueryContext, QueryIdfStats, SegmentQueryContext};
|
||||
use crate::data_types::segment_record::{NamedVectorsOwned, SegmentRecord};
|
||||
use crate::data_types::vectors::{QueryVector, VectorInternal, VectorStructInternal};
|
||||
use crate::data_types::vectors::{QueryVector, VectorStructInternal};
|
||||
use crate::id_tracker::IdTrackerRead;
|
||||
use crate::index::{PayloadIndexRead, VectorIndexRead};
|
||||
use crate::payload_storage::PayloadStorageRead;
|
||||
@@ -37,89 +38,87 @@ where
|
||||
is_stopped: &AtomicBool,
|
||||
deferred_behavior: DeferredBehavior,
|
||||
) -> OperationResult<AHashMap<ExtendedPointId, SegmentRecord>> {
|
||||
let mut records = AHashMap::with_capacity(point_ids.len());
|
||||
// Stage 1: resolve external → internal once, into two parallel vectors.
|
||||
// The id tracker owns this: deferred filtering happens inline (no
|
||||
// `point_is_deferred` lookup). The parallel-vector shape lets
|
||||
// a future batched payload / vector fetcher consume `&offsets`
|
||||
// straight without unzipping first.
|
||||
let (resolved_ids, resolved_offsets) = self
|
||||
.id_tracker
|
||||
.resolve_external_ids(point_ids, deferred_behavior);
|
||||
debug_assert_eq!(resolved_ids.len(), resolved_offsets.len());
|
||||
|
||||
// Filter out deferred points. This is done in two stages to prevent cloning `point_ids`
|
||||
// and iterating more than needed but still satisfy Rust's ownership constraints.
|
||||
let behavior_allows_filtering = !deferred_behavior.include_all_points();
|
||||
let filter_deferred = self.has_deferred_points() && behavior_allows_filtering;
|
||||
let filtered_point_ids = filter_deferred.then(|| {
|
||||
point_ids
|
||||
.iter()
|
||||
.filter(|&&point_id| !self.point_is_deferred(point_id))
|
||||
.copied()
|
||||
.collect::<Vec<_>>()
|
||||
});
|
||||
// Stage 2: pre-allocate one record per resolved point. The `vectors`
|
||||
// slot is initialised here according to `with_vector`, so the
|
||||
// `WithVector::Bool(false)` path needs no separate clearing pass.
|
||||
let needs_vectors = match with_vector {
|
||||
WithVector::Bool(true) | WithVector::Selector(_) => true,
|
||||
WithVector::Bool(false) => false,
|
||||
};
|
||||
let mut records: AHashMap<ExtendedPointId, SegmentRecord> = resolved_ids
|
||||
.iter()
|
||||
.map(|&id| {
|
||||
let record = SegmentRecord {
|
||||
id,
|
||||
vectors: needs_vectors.then(NamedVectorsOwned::default),
|
||||
payload: None,
|
||||
};
|
||||
(id, record)
|
||||
})
|
||||
.collect();
|
||||
|
||||
// Stage two: select the correct slice and shadow `point_ids`.
|
||||
let point_ids = filtered_point_ids.as_deref().unwrap_or(point_ids);
|
||||
|
||||
let mut update_record_vector =
|
||||
|vector_name: &VectorNameBuf,
|
||||
point_id: PointIdType,
|
||||
vector_internal: VectorInternal| {
|
||||
let point_record = records
|
||||
.entry(point_id)
|
||||
.or_insert_with(|| SegmentRecord::empty(point_id));
|
||||
|
||||
point_record
|
||||
.vectors
|
||||
.get_or_insert_with(NamedVectorsOwned::default)
|
||||
.push((vector_name.clone(), vector_internal));
|
||||
// Stage 3: vectors. The external id rides along as the read's user
|
||||
// data — it comes back unchanged in the callback, so no extra
|
||||
// `offset → id` lookup is needed.
|
||||
if needs_vectors {
|
||||
let mut process_vectors = |vector_name: &VectorNameBuf| -> OperationResult<()> {
|
||||
let keys = resolved_ids
|
||||
.iter()
|
||||
.zip(&resolved_offsets)
|
||||
.map(|(&id, &offset)| (id, offset))
|
||||
.stop_if(is_stopped);
|
||||
self.vectors_by_offsets(vector_name, keys, hw_counter, |id, _offset, vec| {
|
||||
if let Some(record) = records.get_mut(&id) {
|
||||
record
|
||||
.vectors
|
||||
.as_mut()
|
||||
.expect("needs_vectors path keeps vectors as Some")
|
||||
.push((vector_name.clone(), vec));
|
||||
}
|
||||
})
|
||||
};
|
||||
|
||||
match with_vector {
|
||||
WithVector::Bool(true) => {
|
||||
for vector_name in self.vector_data.keys() {
|
||||
self.read_vectors(
|
||||
vector_name,
|
||||
point_ids,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
|point_id, vec| {
|
||||
update_record_vector(vector_name, point_id, vec);
|
||||
},
|
||||
)?;
|
||||
match with_vector {
|
||||
WithVector::Bool(true) => {
|
||||
for vector_name in self.vector_data.keys() {
|
||||
process_vectors(vector_name)?;
|
||||
}
|
||||
}
|
||||
}
|
||||
WithVector::Bool(false) => {
|
||||
// Do not display empty `vectors: {}` if disabled.
|
||||
for &point_id in point_ids {
|
||||
let point_record = records
|
||||
.entry(point_id)
|
||||
.or_insert_with(|| SegmentRecord::empty(point_id));
|
||||
point_record.vectors = None;
|
||||
}
|
||||
}
|
||||
WithVector::Selector(selector) => {
|
||||
for vector_name in selector {
|
||||
self.read_vectors(
|
||||
vector_name,
|
||||
point_ids,
|
||||
hw_counter,
|
||||
is_stopped,
|
||||
|point_id, vec| {
|
||||
update_record_vector(vector_name, point_id, vec);
|
||||
},
|
||||
)?;
|
||||
WithVector::Selector(names) => {
|
||||
for vector_name in names {
|
||||
process_vectors(vector_name)?;
|
||||
}
|
||||
}
|
||||
WithVector::Bool(false) => unreachable!("guarded by needs_vectors"),
|
||||
}
|
||||
}
|
||||
|
||||
for &point_id in point_ids {
|
||||
let payload = if with_payload.enable {
|
||||
if let Some(selector) = &with_payload.payload_selector {
|
||||
Some(selector.process(self.payload(point_id, hw_counter)?))
|
||||
} else {
|
||||
Some(self.payload(point_id, hw_counter)?)
|
||||
// Stage 4: payload. Use the already-resolved offsets to skip another
|
||||
// external→internal lookup per point. The per-iteration shape here
|
||||
// mirrors what a future batched payload fetcher would consume —
|
||||
// `&resolved_offsets` becomes its input directly.
|
||||
if with_payload.enable {
|
||||
for (&id, &offset) in resolved_ids.iter().zip(&resolved_offsets) {
|
||||
check_stopped(is_stopped)?;
|
||||
let payload = self.payload_by_offset(offset, hw_counter)?;
|
||||
let payload = match &with_payload.payload_selector {
|
||||
Some(selector) => selector.process(payload),
|
||||
None => payload,
|
||||
};
|
||||
if let Some(record) = records.get_mut(&id) {
|
||||
record.payload = Some(payload);
|
||||
}
|
||||
} else {
|
||||
None
|
||||
};
|
||||
let point_record = records
|
||||
.entry(point_id)
|
||||
.or_insert_with(|| SegmentRecord::empty(point_id));
|
||||
point_record.payload = payload;
|
||||
}
|
||||
}
|
||||
|
||||
Ok(records)
|
||||
@@ -226,8 +225,7 @@ where
|
||||
.vector_data
|
||||
.get(vector_name)
|
||||
.ok_or_else(|| OperationError::vector_name_not_exists(vector_name))?;
|
||||
let vector_query_context =
|
||||
query_context.get_vector_context(vector_name, self.deferred_internal_id());
|
||||
let vector_query_context = query_context.get_vector_context(vector_name);
|
||||
let internal_results = vector_data.vector_index().search(
|
||||
query_vectors,
|
||||
filter,
|
||||
|
||||
@@ -1,8 +1,5 @@
|
||||
use std::sync::atomic::AtomicBool;
|
||||
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::generic_consts::Random;
|
||||
use common::iterator_ext::IteratorExt;
|
||||
use common::types::PointOffsetType;
|
||||
|
||||
use crate::common::check_vector_name;
|
||||
@@ -35,9 +32,9 @@ where
|
||||
let mut result = None;
|
||||
self.vectors_by_offsets(
|
||||
vector_name,
|
||||
std::iter::once(point_offset),
|
||||
std::iter::once(((), point_offset)),
|
||||
hw_counter,
|
||||
|_, vector_internal| {
|
||||
|(), _, vector_internal| {
|
||||
result = Some(vector_internal);
|
||||
},
|
||||
)?;
|
||||
@@ -45,12 +42,19 @@ where
|
||||
}
|
||||
|
||||
/// Retrieve multiple vectors by internal ID.
|
||||
pub fn vectors_by_offsets(
|
||||
///
|
||||
/// Each input is tagged with caller-supplied user data `U` (any `Copy`
|
||||
/// type — the external id, the input position, …). The data is threaded
|
||||
/// back to the callback unchanged, so callers can map results into a
|
||||
/// parallel input array without keeping a separate `offset → ...` lookup
|
||||
/// table. Deleted points are filtered out lazily — entries with deleted
|
||||
/// vectors or deleted points are simply not delivered to the callback.
|
||||
pub fn vectors_by_offsets<U: Copy>(
|
||||
&self,
|
||||
vector_name: &VectorName,
|
||||
point_offsets: impl IntoIterator<Item = PointOffsetType>,
|
||||
keys: impl IntoIterator<Item = (U, PointOffsetType)>,
|
||||
hw_counter: &HardwareCounterCell,
|
||||
mut callback: impl FnMut(PointOffsetType, VectorInternal),
|
||||
mut callback: impl FnMut(U, PointOffsetType, VectorInternal),
|
||||
) -> OperationResult<()> {
|
||||
check_vector_name(vector_name, self.segment_config)?;
|
||||
let vector_data = self
|
||||
@@ -61,7 +65,9 @@ where
|
||||
let total_vectors = vector_storage.total_vector_count();
|
||||
|
||||
let id_tracker = self.id_tracker;
|
||||
let non_deleted_offsets = point_offsets.into_iter().filter(|&point_offset| {
|
||||
// The user-data tag rides alongside each offset, so the deletion
|
||||
// filter stays lazy — no parallel `Vec<(orig_idx, offset)>` needed.
|
||||
let live_keys = keys.into_iter().filter(|&(_, point_offset)| {
|
||||
if total_vectors <= point_offset as usize {
|
||||
debug_assert!(
|
||||
false,
|
||||
@@ -70,59 +76,23 @@ where
|
||||
);
|
||||
return false;
|
||||
}
|
||||
|
||||
let is_vector_deleted = vector_storage.is_deleted_vector(point_offset);
|
||||
let is_point_deleted = id_tracker.is_deleted_point(point_offset);
|
||||
!is_vector_deleted && !is_point_deleted
|
||||
});
|
||||
|
||||
vector_storage.read_vectors::<Random>(non_deleted_offsets, |point_offset, cow_vector| {
|
||||
if vector_storage.is_on_disk() {
|
||||
hw_counter
|
||||
.vector_io_read()
|
||||
.incr_delta(cow_vector.estimate_size_in_bytes());
|
||||
}
|
||||
callback(point_offset, cow_vector.to_owned());
|
||||
});
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Retrieve named vectors for a list of external point IDs, invoking the
|
||||
/// callback for each (point_id, vector) pair.
|
||||
pub fn read_vectors(
|
||||
&self,
|
||||
vector_name: &VectorName,
|
||||
point_ids: &[PointIdType],
|
||||
hw_counter: &HardwareCounterCell,
|
||||
is_stopped: &AtomicBool,
|
||||
mut callback: impl FnMut(PointIdType, VectorInternal),
|
||||
) -> OperationResult<()> {
|
||||
let mut error = None;
|
||||
let internal_ids = point_ids
|
||||
.iter()
|
||||
.copied()
|
||||
.stop_if(is_stopped)
|
||||
.filter_map(|point_id| match self.lookup_internal_id(point_id) {
|
||||
Ok(point_offset) => Some(point_offset),
|
||||
Err(err) => {
|
||||
error = Some(err);
|
||||
None
|
||||
}
|
||||
});
|
||||
self.vectors_by_offsets(
|
||||
vector_name,
|
||||
internal_ids,
|
||||
hw_counter,
|
||||
|point_offset, vector_internal| {
|
||||
if let Some(point_id) = self.id_tracker.external_id(point_offset) {
|
||||
callback(point_id, vector_internal);
|
||||
vector_storage.read_vectors::<Random, U>(
|
||||
live_keys,
|
||||
|user_data, point_offset, cow_vector| {
|
||||
if vector_storage.is_on_disk() {
|
||||
hw_counter
|
||||
.vector_io_read()
|
||||
.incr_delta(cow_vector.estimate_size_in_bytes());
|
||||
}
|
||||
callback(user_data, point_offset, cow_vector.to_owned());
|
||||
},
|
||||
)?;
|
||||
if let Some(err) = error {
|
||||
return Err(err);
|
||||
}
|
||||
);
|
||||
|
||||
Ok(())
|
||||
}
|
||||
|
||||
|
||||
@@ -315,21 +315,11 @@ impl Segment {
|
||||
|
||||
let mut id_tracker = self.id_tracker.borrow_mut();
|
||||
|
||||
let is_point_already_deleted = id_tracker.is_deleted_point(internal_id);
|
||||
|
||||
// `drop_internal` updates the id tracker's deferred-deleted counter when
|
||||
// the point sits at or above the deferred threshold, with double-delete
|
||||
// protection inside `PointMappings::drop`.
|
||||
id_tracker.drop_internal(internal_id)?;
|
||||
|
||||
let deferred_point_status = self.deferred_point_status.as_mut();
|
||||
|
||||
// Increase counter for deleted points.
|
||||
if let Some(deferred_point_status) = deferred_point_status
|
||||
&& internal_id >= deferred_point_status.deferred_internal_id
|
||||
// Don't count the deletion of the same point twice
|
||||
&& !is_point_already_deleted
|
||||
{
|
||||
deferred_point_status.deferred_deleted_count += 1;
|
||||
}
|
||||
|
||||
// Before, we propagated point deletions to also delete its vectors. This turns
|
||||
// out to be problematic because this sometimes makes us lose vector data
|
||||
// because we cannot control the order of segment flushes.
|
||||
@@ -583,29 +573,6 @@ impl Segment {
|
||||
pub fn fix_id_tracker_inconsistencies(&mut self) -> OperationResult<Vec<PointOffsetType>> {
|
||||
self.id_tracker.borrow_mut().fix_inconsistencies()
|
||||
}
|
||||
|
||||
/// Calculates the amount of deleted deferred points by iterating over all points in the ID tracker. Therefore this operation
|
||||
/// can be expensive and should only be run once at segment creation.
|
||||
pub(crate) fn calculate_deleted_deferred_point_count(&self) -> usize {
|
||||
let Some(deferred_from) = self.deferred_internal_id() else {
|
||||
return 0;
|
||||
};
|
||||
|
||||
let id_tracker = self.id_tracker.borrow();
|
||||
let total_points = id_tracker.total_point_count();
|
||||
|
||||
if total_points < deferred_from as usize {
|
||||
return 0;
|
||||
}
|
||||
|
||||
id_tracker.deleted_point_bitslice()[deferred_from as usize..total_points].count_ones()
|
||||
}
|
||||
|
||||
pub(crate) fn deferred_internal_id(&self) -> Option<PointOffsetType> {
|
||||
self.deferred_point_status
|
||||
.as_ref()
|
||||
.map(|i| i.deferred_internal_id)
|
||||
}
|
||||
}
|
||||
|
||||
fn restore_snapshot_in_place(snapshot_path: &Path) -> OperationResult<()> {
|
||||
|
||||
@@ -7,7 +7,7 @@ use ahash::AHashSet;
|
||||
use common::counter::hardware_counter::HardwareCounterCell;
|
||||
use common::tar_ext;
|
||||
use common::tar_unpack::tar_unpack_file;
|
||||
use common::types::DeferredBehavior;
|
||||
use common::types::{DeferredBehavior, PointOffsetType};
|
||||
use fs_err as fs;
|
||||
use fs_err::File;
|
||||
use ordered_float::OrderedFloat;
|
||||
@@ -32,7 +32,7 @@ use crate::entry::entry_point::{
|
||||
NonAppendableSegmentEntry as _, ReadSegmentEntry as _, SegmentEntry as _,
|
||||
};
|
||||
use crate::entry::{SnapshotEntry as _, StorageSegmentEntry as _};
|
||||
use crate::id_tracker::IdTracker;
|
||||
use crate::id_tracker::{IdTracker, IdTrackerRead};
|
||||
use crate::index::sparse_index::sparse_index_config::{SparseIndexConfig, SparseIndexType};
|
||||
use crate::json_path::JsonPath;
|
||||
use crate::segment_constructor::simple_segment_constructor::{
|
||||
@@ -910,7 +910,10 @@ fn create_deferred_segment(
|
||||
// Now we should have deferred points
|
||||
assert_eq!(segment.has_deferred_points(), n_deferred > 0);
|
||||
if n_deferred > 0 {
|
||||
assert_eq!(segment.deferred_internal_id(), Some(n_vectors as u32));
|
||||
assert_eq!(
|
||||
segment.id_tracker.borrow().deferred_internal_id(),
|
||||
Some(n_vectors as u32)
|
||||
);
|
||||
}
|
||||
|
||||
// Points 1 to n_vectors should NOT be deferred
|
||||
@@ -1007,7 +1010,7 @@ fn test_dense_deferred_points() {
|
||||
"Segment should still have deferred points after reopening"
|
||||
);
|
||||
assert_eq!(
|
||||
segment.deferred_internal_id(),
|
||||
segment.id_tracker.borrow().deferred_internal_id(),
|
||||
Some(13),
|
||||
"Deferred internal ID should still be `DEFERRED_POINTS_ID` after reopening"
|
||||
);
|
||||
@@ -1071,7 +1074,7 @@ fn test_deferred_point_estimation_with_filter() {
|
||||
|
||||
// For consistency we also test that the same cardinality is estimated if no deferred points exist.
|
||||
if n_deferred == 0 {
|
||||
assert_eq!(segment.deferred_internal_id(), None);
|
||||
assert_eq!(segment.id_tracker.borrow().deferred_internal_id(), None);
|
||||
let estimation = segment
|
||||
.estimate_point_count(Some(&filter), &hw_counter)
|
||||
.unwrap();
|
||||
@@ -1301,7 +1304,7 @@ fn test_deferred_point_facets() {
|
||||
);
|
||||
|
||||
let dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
|
||||
let mut segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred);
|
||||
let segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred);
|
||||
|
||||
let request = FacetParams {
|
||||
key: key.clone(),
|
||||
@@ -1314,14 +1317,17 @@ fn test_deferred_point_facets() {
|
||||
.facet(&request, &AtomicBool::new(false), &hw_counter)
|
||||
.unwrap();
|
||||
|
||||
let old_status = segment.deferred_point_status.take();
|
||||
if n_deferred > 0 {
|
||||
assert!(old_status.is_some());
|
||||
}
|
||||
let facet_res = segment
|
||||
// Compare against the same point set without deferred mode by
|
||||
// rebuilding the segment with `deferred_internal_id = None`.
|
||||
let no_deferred_dir = Builder::new()
|
||||
.prefix("segment_dir_no_deferred")
|
||||
.tempdir()
|
||||
.unwrap();
|
||||
let no_deferred_segment =
|
||||
create_deferred_segment(&no_deferred_dir, 5, N_POINTS + n_deferred, 0);
|
||||
let facet_res = no_deferred_segment
|
||||
.facet(&request, &AtomicBool::new(false), &hw_counter)
|
||||
.unwrap();
|
||||
segment.deferred_point_status = old_status;
|
||||
|
||||
let expected_deferred = if filter.is_some() {
|
||||
n_deferred.div_ceil(3)
|
||||
@@ -1446,7 +1452,7 @@ fn assert_deferred_points_excluded<F, R, T>(
|
||||
log::debug!(" => deferred points = {n_deferred}; filter-set ID = {filter_set_id}",);
|
||||
|
||||
let dir = Builder::new().prefix("segment_dir").tempdir().unwrap();
|
||||
let mut segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred);
|
||||
let segment = create_deferred_segment(&dir, 5, N_POINTS, n_deferred);
|
||||
|
||||
// Search with deferred mode
|
||||
let search_res_deferred = operation(&segment, filter_set.filter.as_ref());
|
||||
@@ -1460,29 +1466,34 @@ fn assert_deferred_points_excluded<F, R, T>(
|
||||
}
|
||||
|
||||
// Disable deferred points and search again.
|
||||
if need_rebuilt_segment {
|
||||
// Don't run this on windows because this test is already extremely slow.
|
||||
// Recreating the segment here would double that time.
|
||||
if cfg!(target_os = "windows") {
|
||||
drop(segment);
|
||||
dir.close().unwrap();
|
||||
continue;
|
||||
}
|
||||
|
||||
let dir = Builder::new().prefix("segment_dir_2").tempdir().unwrap();
|
||||
segment = create_deferred_segment(&dir, 5, N_POINTS + n_deferred, 0);
|
||||
} else {
|
||||
segment.deferred_point_status = None;
|
||||
// On Windows segment creation is extremely IO-heavy; skip the rebuild path
|
||||
// for tests where it noticeably slows the suite down.
|
||||
if need_rebuilt_segment && cfg!(target_os = "windows") {
|
||||
drop(segment);
|
||||
dir.close().unwrap();
|
||||
continue;
|
||||
}
|
||||
|
||||
let search_res_normal = operation(&segment, filter_set.filter.as_ref());
|
||||
// Deferred state is owned by the id tracker and only set at segment
|
||||
// construction time, so we rebuild a fresh non-deferred segment with the
|
||||
// same total point count to compare against.
|
||||
let no_deferred_dir = Builder::new()
|
||||
.prefix("segment_dir_no_deferred")
|
||||
.tempdir()
|
||||
.unwrap();
|
||||
let no_deferred_segment =
|
||||
create_deferred_segment(&no_deferred_dir, 5, N_POINTS + n_deferred, 0);
|
||||
|
||||
let search_res_normal = operation(&no_deferred_segment, filter_set.filter.as_ref());
|
||||
assert_eq!(
|
||||
search_res_normal.len(),
|
||||
filter_set.expected_visible + (filter_set.expected_deferred)(n_deferred)
|
||||
);
|
||||
|
||||
drop(segment);
|
||||
drop(no_deferred_segment);
|
||||
dir.close().unwrap();
|
||||
no_deferred_dir.close().unwrap();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1510,7 +1521,7 @@ fn test_deleted_deferred_point_count() {
|
||||
assert_eq!(segment.available_point_count_without_deferred(), N_POINTS);
|
||||
|
||||
for d in 0..n_deferred {
|
||||
let delete_id = segment.deferred_internal_id().unwrap() + d as u32;
|
||||
let delete_id = segment.id_tracker.borrow().deferred_internal_id().unwrap() + d as u32;
|
||||
segment
|
||||
.delete_point_internal(delete_id, &hw_counter)
|
||||
.unwrap();
|
||||
@@ -1521,7 +1532,7 @@ fn test_deleted_deferred_point_count() {
|
||||
n_deferred.checked_sub(deleted_count).unwrap()
|
||||
);
|
||||
assert_eq!(
|
||||
segment.calculate_deleted_deferred_point_count(),
|
||||
segment.id_tracker.borrow().deferred_deleted_count(),
|
||||
deleted_count,
|
||||
);
|
||||
|
||||
@@ -1536,7 +1547,7 @@ fn test_deleted_deferred_point_count() {
|
||||
);
|
||||
|
||||
assert_eq!(
|
||||
segment.calculate_deleted_deferred_point_count(),
|
||||
segment.id_tracker.borrow().deferred_deleted_count(),
|
||||
deleted_count
|
||||
);
|
||||
|
||||
@@ -1545,7 +1556,10 @@ fn test_deleted_deferred_point_count() {
|
||||
|
||||
// We delete all deferred points in the segment.
|
||||
assert_eq!(segment.deferred_point_count(), 0);
|
||||
assert_eq!(segment.calculate_deleted_deferred_point_count(), n_deferred);
|
||||
assert_eq!(
|
||||
segment.id_tracker.borrow().deferred_deleted_count(),
|
||||
n_deferred
|
||||
);
|
||||
assert_eq!(segment.available_point_count_without_deferred(), N_POINTS);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -169,7 +169,6 @@ impl Segment {
|
||||
path: &vector_index_path,
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})?;
|
||||
|
||||
// Register the new storage with the payload index so `has_vector`
|
||||
|
||||
@@ -86,7 +86,9 @@ impl SegmentBuilder {
|
||||
let temp_dir = create_temp_dir(temp_dir)?;
|
||||
|
||||
let id_tracker = if segment_config.is_appendable() {
|
||||
IdTrackerEnum::MutableIdTracker(create_mutable_id_tracker(temp_dir.path())?)
|
||||
// Deferred state is applied when the freshly built segment is reloaded
|
||||
// via `load_segment`. The transient builder tracker doesn't need it.
|
||||
IdTrackerEnum::MutableIdTracker(create_mutable_id_tracker(temp_dir.path(), None)?)
|
||||
} else {
|
||||
IdTrackerEnum::InMemoryIdTracker(InMemoryIdTracker::new())
|
||||
};
|
||||
@@ -658,8 +660,6 @@ impl SegmentBuilder {
|
||||
path: &vector_index_path,
|
||||
stopped,
|
||||
tick_progress: || (),
|
||||
// We don't use the `index` returned here so we always set deferred to `None`. It's been loaded properly later.
|
||||
deferred_internal_id: None,
|
||||
})?;
|
||||
|
||||
if sparse_vector_config.storage_type.is_on_disk() {
|
||||
|
||||
@@ -39,9 +39,7 @@ use crate::index::sparse_index::sparse_vector_index::{
|
||||
use crate::index::struct_payload_index::StructPayloadIndex;
|
||||
use crate::payload_storage::mmap_payload_storage::MmapPayloadStorage;
|
||||
use crate::payload_storage::payload_storage_enum::PayloadStorageEnum;
|
||||
use crate::segment::{
|
||||
DeferredPointStatus, SEGMENT_STATE_FILE, Segment, SegmentVersion, VectorData,
|
||||
};
|
||||
use crate::segment::{SEGMENT_STATE_FILE, Segment, SegmentVersion, VectorData};
|
||||
use crate::types::{
|
||||
Distance, HnswGlobalConfig, Indexes, PayloadStorageType, SegmentConfig, SegmentState,
|
||||
SegmentType, SeqNumberType, SparseVectorStorageType, VectorDataConfig, VectorName,
|
||||
@@ -223,8 +221,11 @@ pub(crate) fn create_payload_storage(
|
||||
Ok(payload_storage)
|
||||
}
|
||||
|
||||
pub(crate) fn create_mutable_id_tracker(segment_path: &Path) -> OperationResult<MutableIdTracker> {
|
||||
MutableIdTracker::open(segment_path)
|
||||
pub(crate) fn create_mutable_id_tracker(
|
||||
segment_path: &Path,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<MutableIdTracker> {
|
||||
MutableIdTracker::open(segment_path, deferred_internal_id)
|
||||
}
|
||||
|
||||
pub(crate) fn get_payload_index_path(segment_path: &Path) -> PathBuf {
|
||||
@@ -411,7 +412,8 @@ fn create_segment(
|
||||
let use_mutable_id_tracker =
|
||||
appendable_flag || !immutable_id_tracker::mappings_path(segment_path).is_file();
|
||||
let started = Instant::now();
|
||||
let id_tracker = create_segment_id_tracker(use_mutable_id_tracker, segment_path)?;
|
||||
let id_tracker =
|
||||
create_segment_id_tracker(use_mutable_id_tracker, segment_path, deferred_internal_id)?;
|
||||
log_load_timing(segment_path, "id_tracker", started);
|
||||
|
||||
let mut vector_storages = HashMap::new();
|
||||
@@ -556,7 +558,6 @@ fn create_segment(
|
||||
path: &vector_index_path,
|
||||
stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id,
|
||||
})?);
|
||||
log_load_timing(
|
||||
segment_path,
|
||||
@@ -582,7 +583,7 @@ fn create_segment(
|
||||
SegmentType::Plain
|
||||
};
|
||||
|
||||
let mut segment = Segment {
|
||||
Ok(Segment {
|
||||
uuid,
|
||||
initial_version,
|
||||
version,
|
||||
@@ -598,22 +599,13 @@ fn create_segment(
|
||||
payload_storage,
|
||||
segment_config: config.clone(),
|
||||
error_status: None,
|
||||
deferred_point_status: None,
|
||||
};
|
||||
|
||||
if let Some(deferred_internal_id) = deferred_internal_id {
|
||||
segment.deferred_point_status = Some(DeferredPointStatus {
|
||||
deferred_internal_id,
|
||||
deferred_deleted_count: segment.calculate_deleted_deferred_point_count(),
|
||||
});
|
||||
}
|
||||
|
||||
Ok(segment)
|
||||
})
|
||||
}
|
||||
|
||||
fn create_segment_id_tracker(
|
||||
mutable_id_tracker: bool,
|
||||
segment_path: &Path,
|
||||
deferred_internal_id: Option<PointOffsetType>,
|
||||
) -> OperationResult<Arc<AtomicRefCell<IdTrackerEnum>>> {
|
||||
if !mutable_id_tracker {
|
||||
return Ok(sp(IdTrackerEnum::ImmutableIdTracker(
|
||||
@@ -622,7 +614,7 @@ fn create_segment_id_tracker(
|
||||
}
|
||||
|
||||
Ok(sp(IdTrackerEnum::MutableIdTracker(
|
||||
create_mutable_id_tracker(segment_path)?,
|
||||
create_mutable_id_tracker(segment_path, deferred_internal_id)?,
|
||||
)))
|
||||
}
|
||||
|
||||
|
||||
@@ -242,21 +242,22 @@ where
|
||||
.expect("Vector not found")
|
||||
}
|
||||
|
||||
fn read_vectors<P: AccessPattern>(
|
||||
fn read_vectors<P: AccessPattern, U: Copy>(
|
||||
&self,
|
||||
keys: impl IntoIterator<Item = PointOffsetType>,
|
||||
mut callback: impl FnMut(PointOffsetType, CowVector<'_>),
|
||||
keys: impl IntoIterator<Item = (U, PointOffsetType)>,
|
||||
mut callback: impl FnMut(U, PointOffsetType, CowVector<'_>),
|
||||
) {
|
||||
let point_offsets: Vec<_> = keys.into_iter().collect();
|
||||
// Split into parallel arrays in one pass: `for_each_in_batch` needs an
|
||||
// offsets slice (it chunks it for batched reads), but we still want
|
||||
// `user_data[idx]` available inside the callback.
|
||||
let (user_data, point_offsets): (Vec<U>, Vec<PointOffsetType>) = keys.into_iter().unzip();
|
||||
|
||||
self.vectors
|
||||
.as_ref()
|
||||
.unwrap()
|
||||
.for_each_in_batch(&point_offsets, |idx, vector| {
|
||||
let point_offset = point_offsets[idx];
|
||||
let vector = CowVector::from(T::slice_to_float_cow(Cow::Borrowed(vector)));
|
||||
|
||||
callback(point_offset, vector);
|
||||
callback(user_data[idx], point_offsets[idx], vector);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -495,7 +496,7 @@ mod tests {
|
||||
2,
|
||||
);
|
||||
let res = searcher
|
||||
.peek_top_all(&DEFAULT_STOPPED, None)
|
||||
.peek_top_all(&DEFAULT_STOPPED)
|
||||
.unwrap()
|
||||
.into_iter()
|
||||
.exactly_one()
|
||||
@@ -640,7 +641,7 @@ mod tests {
|
||||
5,
|
||||
);
|
||||
let closest = searcher
|
||||
.peek_top_all(&DEFAULT_STOPPED, None)
|
||||
.peek_top_all(&DEFAULT_STOPPED)
|
||||
.unwrap()
|
||||
.into_iter()
|
||||
.exactly_one()
|
||||
|
||||
@@ -105,22 +105,26 @@ impl<S: UniversalRead> VectorStorageRead for VectorStorageReadEnum<S> {
|
||||
}
|
||||
}
|
||||
|
||||
fn read_vectors<P: AccessPattern>(
|
||||
fn read_vectors<P: AccessPattern, U: Copy>(
|
||||
&self,
|
||||
keys: impl IntoIterator<Item = PointOffsetType>,
|
||||
callback: impl FnMut(PointOffsetType, CowVector<'_>),
|
||||
keys: impl IntoIterator<Item = (U, PointOffsetType)>,
|
||||
callback: impl FnMut(U, PointOffsetType, CowVector<'_>),
|
||||
) {
|
||||
match self {
|
||||
VectorStorageReadEnum::Dense(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::DenseByte(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::DenseHalf(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::DenseChunked(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::DenseChunkedByte(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::DenseChunkedHalf(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::MultiDenseChunked(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::MultiDenseChunkedByte(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::MultiDenseChunkedHalf(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::Sparse(s) => s.read_vectors::<P>(keys, callback),
|
||||
VectorStorageReadEnum::Dense(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageReadEnum::DenseByte(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageReadEnum::DenseHalf(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageReadEnum::DenseChunked(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageReadEnum::DenseChunkedByte(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageReadEnum::DenseChunkedHalf(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageReadEnum::MultiDenseChunked(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageReadEnum::MultiDenseChunkedByte(s) => {
|
||||
s.read_vectors::<P, U>(keys, callback)
|
||||
}
|
||||
VectorStorageReadEnum::MultiDenseChunkedHalf(s) => {
|
||||
s.read_vectors::<P, U>(keys, callback)
|
||||
}
|
||||
VectorStorageReadEnum::Sparse(s) => s.read_vectors::<P, U>(keys, callback),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -119,7 +119,7 @@ fn do_test_delete_points(storage: &mut VectorStorageEnum) {
|
||||
5,
|
||||
);
|
||||
let closest = searcher
|
||||
.peek_top_all(&DEFAULT_STOPPED, None)
|
||||
.peek_top_all(&DEFAULT_STOPPED)
|
||||
.unwrap()
|
||||
.into_iter()
|
||||
.exactly_one()
|
||||
|
||||
@@ -182,7 +182,7 @@ fn do_test_delete_points(vector_dim: usize, vec_count: usize, storage: &mut Vect
|
||||
5,
|
||||
);
|
||||
let closest = searcher
|
||||
.peek_top_all(&DEFAULT_STOPPED, None)
|
||||
.peek_top_all(&DEFAULT_STOPPED)
|
||||
.unwrap()
|
||||
.pop()
|
||||
.unwrap();
|
||||
|
||||
@@ -93,15 +93,21 @@ pub trait VectorStorageRead {
|
||||
/// Get the vector by the given key with potential optimizations for sequential reads.
|
||||
fn get_vector<P: AccessPattern>(&self, key: PointOffsetType) -> CowVector<'_>;
|
||||
|
||||
/// Get multiple vectors by the given keys
|
||||
/// Get multiple vectors by the given keys.
|
||||
/// Potentially optimized for internal parallel reads.
|
||||
fn read_vectors<P: AccessPattern>(
|
||||
///
|
||||
/// Each input is tagged with caller-supplied user data `U` (e.g., an
|
||||
/// external point id, or the input position). The data is threaded
|
||||
/// straight back to the callback alongside its offset and vector, so
|
||||
/// callers can map results into a parallel input array without keeping a
|
||||
/// separate `offset → ...` lookup table.
|
||||
fn read_vectors<P: AccessPattern, U: Copy>(
|
||||
&self,
|
||||
keys: impl IntoIterator<Item = PointOffsetType>,
|
||||
mut callback: impl FnMut(PointOffsetType, CowVector<'_>),
|
||||
keys: impl IntoIterator<Item = (U, PointOffsetType)>,
|
||||
mut callback: impl FnMut(U, PointOffsetType, CowVector<'_>),
|
||||
) {
|
||||
for key in keys {
|
||||
callback(key, self.get_vector::<P>(key));
|
||||
for (user_data, key) in keys {
|
||||
callback(user_data, key, self.get_vector::<P>(key));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -778,47 +784,53 @@ impl VectorStorageRead for VectorStorageEnum {
|
||||
}
|
||||
}
|
||||
|
||||
fn read_vectors<P: AccessPattern>(
|
||||
fn read_vectors<P: AccessPattern, U: Copy>(
|
||||
&self,
|
||||
keys: impl IntoIterator<Item = PointOffsetType>,
|
||||
callback: impl FnMut(PointOffsetType, CowVector<'_>),
|
||||
keys: impl IntoIterator<Item = (U, PointOffsetType)>,
|
||||
callback: impl FnMut(U, PointOffsetType, CowVector<'_>),
|
||||
) {
|
||||
match self {
|
||||
VectorStorageEnum::DenseVolatile(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseVolatile(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
#[cfg(test)]
|
||||
VectorStorageEnum::DenseVolatileByte(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseVolatileByte(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
#[cfg(test)]
|
||||
VectorStorageEnum::DenseVolatileHalf(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseMemmap(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseMemmapByte(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseMemmapHalf(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseVolatileHalf(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::DenseMemmap(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::DenseMemmapByte(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::DenseMemmapHalf(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
|
||||
#[cfg(target_os = "linux")]
|
||||
VectorStorageEnum::DenseUring(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseUring(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
#[cfg(target_os = "linux")]
|
||||
VectorStorageEnum::DenseUringByte(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseUringByte(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
#[cfg(target_os = "linux")]
|
||||
VectorStorageEnum::DenseUringHalf(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseUringHalf(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
|
||||
VectorStorageEnum::DenseAppendableMemmap(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseAppendableMemmapByte(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseAppendableMemmapHalf(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::SparseVolatile(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::SparseMmap(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::MultiDenseVolatile(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::DenseAppendableMemmap(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::DenseAppendableMemmapByte(v) => {
|
||||
v.read_vectors::<P, U>(keys, callback)
|
||||
}
|
||||
VectorStorageEnum::DenseAppendableMemmapHalf(v) => {
|
||||
v.read_vectors::<P, U>(keys, callback)
|
||||
}
|
||||
VectorStorageEnum::SparseVolatile(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::SparseMmap(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::MultiDenseVolatile(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
#[cfg(test)]
|
||||
VectorStorageEnum::MultiDenseVolatileByte(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::MultiDenseVolatileByte(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
#[cfg(test)]
|
||||
VectorStorageEnum::MultiDenseVolatileHalf(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::MultiDenseAppendableMemmap(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::MultiDenseVolatileHalf(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::MultiDenseAppendableMemmap(v) => {
|
||||
v.read_vectors::<P, U>(keys, callback)
|
||||
}
|
||||
VectorStorageEnum::MultiDenseAppendableMemmapByte(v) => {
|
||||
v.read_vectors::<P>(keys, callback)
|
||||
v.read_vectors::<P, U>(keys, callback)
|
||||
}
|
||||
VectorStorageEnum::MultiDenseAppendableMemmapHalf(v) => {
|
||||
v.read_vectors::<P>(keys, callback)
|
||||
v.read_vectors::<P, U>(keys, callback)
|
||||
}
|
||||
VectorStorageEnum::EmptyDense(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::EmptySparse(v) => v.read_vectors::<P>(keys, callback),
|
||||
VectorStorageEnum::EmptyDense(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
VectorStorageEnum::EmptySparse(v) => v.read_vectors::<P, U>(keys, callback),
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -113,7 +113,7 @@ fn exact_search_test() {
|
||||
let px = payload_index_ptr.borrow();
|
||||
let filter = Filter::new_must(Condition::Field(block.condition.clone()));
|
||||
let points = px
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap();
|
||||
for point in points {
|
||||
coverage.insert(point, coverage.get(&point).unwrap_or(&0) + 1);
|
||||
|
||||
@@ -136,7 +136,7 @@ fn _test_filterable_hnsw(
|
||||
for block in &blocks {
|
||||
let filter = Filter::new_must(Condition::Field(block.condition.clone()));
|
||||
let points = px
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &stopped))
|
||||
.unwrap();
|
||||
for point in points {
|
||||
coverage.insert(point, coverage.get(&point).unwrap_or(&0) + 1);
|
||||
|
||||
@@ -146,7 +146,7 @@ fn test_filtering_context_consistency() {
|
||||
let nested_filter_0 = Filter::new_must(nested_condition_0);
|
||||
let (res0, check_res0) = index.with_view(|v| {
|
||||
let res0 = v
|
||||
.query_points(&nested_filter_0, &hw_counter, &is_stopped, None)
|
||||
.query_points(&nested_filter_0, &hw_counter, &is_stopped)
|
||||
.unwrap();
|
||||
let filter_context = v.filter_context(&nested_filter_0, &hw_counter).unwrap();
|
||||
let check_res0: Vec<_> = (0..NUM_POINTS as PointOffsetType)
|
||||
@@ -187,7 +187,7 @@ fn test_filtering_context_consistency() {
|
||||
|
||||
let (res1, check_res1) = index.with_view(|v| {
|
||||
let res1 = v
|
||||
.query_points(&nested_filter_1, &hw_counter, &is_stopped, None)
|
||||
.query_points(&nested_filter_1, &hw_counter, &is_stopped)
|
||||
.unwrap();
|
||||
let filter_context = v.filter_context(&nested_filter_1, &hw_counter).unwrap();
|
||||
let check_res1: Vec<_> = (0..NUM_POINTS as PointOffsetType)
|
||||
@@ -225,7 +225,7 @@ fn test_filtering_context_consistency() {
|
||||
|
||||
let (res2, check_res2) = index.with_view(|v| {
|
||||
let res2 = v
|
||||
.query_points(&nested_filter_2, &hw_counter, &is_stopped, None)
|
||||
.query_points(&nested_filter_2, &hw_counter, &is_stopped)
|
||||
.unwrap();
|
||||
let filter_context = v.filter_context(&nested_filter_2, &hw_counter).unwrap();
|
||||
let check_res2: Vec<_> = (0..NUM_POINTS as PointOffsetType)
|
||||
@@ -273,7 +273,7 @@ fn test_filtering_context_consistency() {
|
||||
|
||||
let (res3, check_res3) = index.with_view(|v| {
|
||||
let res3 = v
|
||||
.query_points(&nested_filter_3, &hw_counter, &is_stopped, None)
|
||||
.query_points(&nested_filter_3, &hw_counter, &is_stopped)
|
||||
.unwrap();
|
||||
let filter_context = v.filter_context(&nested_filter_3, &hw_counter).unwrap();
|
||||
let check_res3: Vec<_> = (0..NUM_POINTS as PointOffsetType)
|
||||
|
||||
@@ -636,7 +636,7 @@ fn test_is_empty_conditions(test_segments: &TestSegments) -> Result<()> {
|
||||
.plain_segment
|
||||
.payload_index
|
||||
.borrow()
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap();
|
||||
|
||||
let real_number = plain_result.len();
|
||||
@@ -646,7 +646,7 @@ fn test_is_empty_conditions(test_segments: &TestSegments) -> Result<()> {
|
||||
.struct_segment
|
||||
.payload_index
|
||||
.borrow()
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped, None))
|
||||
.with_view(|v| v.query_points(&filter, &hw_counter, &is_stopped))
|
||||
.unwrap()
|
||||
.into_iter()
|
||||
// null index does not track deleted points, so we need to filter them out here. In callsites,
|
||||
|
||||
@@ -184,7 +184,6 @@ fn sparse_index_discover_test() {
|
||||
path: index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
@@ -219,7 +218,7 @@ fn sparse_index_discover_test() {
|
||||
|
||||
let query_context = QueryContext::default();
|
||||
let segment_query_context = query_context.get_segment_query_context();
|
||||
let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME, None);
|
||||
let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME);
|
||||
|
||||
let sparse_search_result = sparse_index
|
||||
.search(&[&sparse_query], None, top, None, &vector_context)
|
||||
@@ -302,7 +301,6 @@ fn sparse_index_hardware_measurement_test() {
|
||||
path: index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
@@ -312,7 +310,7 @@ fn sparse_index_hardware_measurement_test() {
|
||||
|
||||
let query_context = QueryContext::default();
|
||||
let segment_query_context = query_context.get_segment_query_context();
|
||||
let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME, None);
|
||||
let vector_context = segment_query_context.get_vector_context(SPARSE_VECTOR_NAME);
|
||||
|
||||
let cpu_usage = query_context.hardware_usage_accumulator().get_cpu();
|
||||
assert_eq!(cpu_usage, 0);
|
||||
|
||||
@@ -230,7 +230,6 @@ fn sparse_vector_index_consistent_with_storage() {
|
||||
path: mmap_index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
@@ -257,7 +256,6 @@ fn sparse_vector_index_consistent_with_storage() {
|
||||
path: mmap_index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap();
|
||||
|
||||
@@ -694,7 +692,6 @@ fn check_persistence<TInvertedIndex: InvertedIndex>(
|
||||
path: inverted_index_dir.path(),
|
||||
stopped: &stopped,
|
||||
tick_progress: || (),
|
||||
deferred_internal_id: None,
|
||||
})
|
||||
.unwrap()
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user