Skip to content

Commit 6038455

Browse files
fulmicoton-ddPSeitz
authored andcommitted
First stab at tantivy's codec
Convert SegmentReader, InvertedIndexReader and postinglists to traits. Add special functions to pushdown certain performance methods to keep them strictly typed. We rely on a ObjectSafeCodec contraption to avoid the proliferation of generics. That object's point is to make sure we can build TermScorer with a concrete codec specific type before reboxing it. (same thing for PhraseScorer). fix performance regression: fix incorrect scorer cast for buffered union bock wand
1 parent 57fe659 commit 6038455

95 files changed

Lines changed: 2737 additions & 1662 deletions

File tree

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

benches/str_search_and_get.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -45,7 +45,7 @@ fn build_shared_indices(num_docs: usize, distribution: &str) -> BenchIndex {
4545
match distribution {
4646
"dense_random" => {
4747
for _doc_id in 0..num_docs {
48-
let suffix = rng.gen_range(0u64..1000u64);
48+
let suffix = rng.random_range(0u64..1000u64);
4949
let str_val = format!("str_{:03}", suffix);
5050

5151
writer
@@ -71,7 +71,7 @@ fn build_shared_indices(num_docs: usize, distribution: &str) -> BenchIndex {
7171
}
7272
"sparse_random" => {
7373
for _doc_id in 0..num_docs {
74-
let suffix = rng.gen_range(0u64..1000000u64);
74+
let suffix = rng.random_range(0u64..1000000u64);
7575
let str_val = format!("str_{:07}", suffix);
7676

7777
writer

common/src/bitset.rs

Lines changed: 10 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -178,13 +178,11 @@ impl TinySet {
178178
#[derive(Clone)]
179179
pub struct BitSet {
180180
tinysets: Box<[TinySet]>,
181-
len: u64,
182181
max_value: u32,
183182
}
184183
impl std::fmt::Debug for BitSet {
185184
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
186185
f.debug_struct("BitSet")
187-
.field("len", &self.len)
188186
.field("max_value", &self.max_value)
189187
.finish()
190188
}
@@ -212,7 +210,6 @@ impl BitSet {
212210
let tinybitsets = vec![TinySet::empty(); num_buckets as usize].into_boxed_slice();
213211
BitSet {
214212
tinysets: tinybitsets,
215-
len: 0,
216213
max_value,
217214
}
218215
}
@@ -230,7 +227,6 @@ impl BitSet {
230227
}
231228
BitSet {
232229
tinysets: tinybitsets,
233-
len: max_value as u64,
234230
max_value,
235231
}
236232
}
@@ -249,17 +245,19 @@ impl BitSet {
249245

250246
/// Intersect with tinysets
251247
fn intersect_update_with_iter(&mut self, other: impl Iterator<Item = TinySet>) {
252-
self.len = 0;
253248
for (left, right) in self.tinysets.iter_mut().zip(other) {
254249
*left = left.intersect(right);
255-
self.len += left.len() as u64;
256250
}
257251
}
258252

259253
/// Returns the number of elements in the `BitSet`.
260254
#[inline]
261255
pub fn len(&self) -> usize {
262-
self.len as usize
256+
self.tinysets
257+
.iter()
258+
.copied()
259+
.map(|tinyset| tinyset.len())
260+
.sum::<u32>() as usize
263261
}
264262

265263
/// Inserts an element in the `BitSet`
@@ -268,7 +266,7 @@ impl BitSet {
268266
// we do not check saturated els.
269267
let higher = el / 64u32;
270268
let lower = el % 64u32;
271-
self.len += u64::from(self.tinysets[higher as usize].insert_mut(lower));
269+
self.tinysets[higher as usize].insert_mut(lower);
272270
}
273271

274272
/// Inserts an element in the `BitSet`
@@ -277,7 +275,7 @@ impl BitSet {
277275
// we do not check saturated els.
278276
let higher = el / 64u32;
279277
let lower = el % 64u32;
280-
self.len -= u64::from(self.tinysets[higher as usize].remove_mut(lower));
278+
self.tinysets[higher as usize].remove_mut(lower);
281279
}
282280

283281
/// Returns true iff the elements is in the `BitSet`.
@@ -299,6 +297,9 @@ impl BitSet {
299297
.map(|delta_bucket| bucket + delta_bucket as u32)
300298
}
301299

300+
/// Returns the maximum number of elements in the bitset.
301+
///
302+
/// Warning: The largest element the bitset can contain is `max_value - 1`.
302303
#[inline]
303304
pub fn max_value(&self) -> u32 {
304305
self.max_value

examples/custom_collector.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -70,7 +70,7 @@ impl Collector for StatsCollector {
7070
fn for_segment(
7171
&self,
7272
_segment_local_id: u32,
73-
segment_reader: &SegmentReader,
73+
segment_reader: &dyn SegmentReader,
7474
) -> tantivy::Result<StatsSegmentCollector> {
7575
let fast_field_reader = segment_reader.fast_fields().u64(&self.field)?;
7676
Ok(StatsSegmentCollector {

examples/faceted_search_with_tweaked_score.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,7 @@ fn main() -> tantivy::Result<()> {
6565
);
6666
let top_docs_by_custom_score =
6767
// Call TopDocs with a custom tweak score
68-
TopDocs::with_limit(2).tweak_score(move |segment_reader: &SegmentReader| {
68+
TopDocs::with_limit(2).tweak_score(move |segment_reader: &dyn SegmentReader| {
6969
let ingredient_reader = segment_reader.facet_reader("ingredient").unwrap();
7070
let facet_dict = ingredient_reader.facet_dict();
7171

examples/iterating_docs_and_positions.rs

Lines changed: 1 addition & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -91,46 +91,10 @@ fn main() -> tantivy::Result<()> {
9191
}
9292
}
9393

94-
// A `Term` is a text token associated with a field.
95-
// Let's go through all docs containing the term `title:the` and access their position
96-
let term_the = Term::from_field_text(title, "the");
97-
98-
// Some other powerful operations (especially `.skip_to`) may be useful to consume these
94+
// Some other powerful operations (especially `.seek`) may be useful to consume these
9995
// posting lists rapidly.
10096
// You can check for them in the [`DocSet`](https://docs.rs/tantivy/~0/tantivy/trait.DocSet.html) trait
10197
// and the [`Postings`](https://docs.rs/tantivy/~0/tantivy/trait.Postings.html) trait
10298

103-
// Also, for some VERY specific high performance use case like an OLAP analysis of logs,
104-
// you can get better performance by accessing directly the blocks of doc ids.
105-
for segment_reader in searcher.segment_readers() {
106-
// A segment contains different data structure.
107-
// Inverted index stands for the combination of
108-
// - the term dictionary
109-
// - the inverted lists associated with each terms and their positions
110-
let inverted_index = segment_reader.inverted_index(title)?;
111-
112-
// This segment posting object is like a cursor over the documents matching the term.
113-
// The `IndexRecordOption` arguments tells tantivy we will be interested in both term
114-
// frequencies and positions.
115-
//
116-
// If you don't need all this information, you may get better performance by decompressing
117-
// less information.
118-
if let Some(mut block_segment_postings) =
119-
inverted_index.read_block_postings(&term_the, IndexRecordOption::Basic)?
120-
{
121-
loop {
122-
let docs = block_segment_postings.docs();
123-
if docs.is_empty() {
124-
break;
125-
}
126-
// Once again these docs MAY contains deleted documents as well.
127-
let docs = block_segment_postings.docs();
128-
// Prints `Docs [0, 2].`
129-
println!("Docs {docs:?}");
130-
block_segment_postings.advance();
131-
}
132-
}
133-
}
134-
13599
Ok(())
136100
}

examples/warmer.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -43,7 +43,7 @@ impl DynamicPriceColumn {
4343
}
4444
}
4545

46-
pub fn price_for_segment(&self, segment_reader: &SegmentReader) -> Option<Arc<Vec<Price>>> {
46+
pub fn price_for_segment(&self, segment_reader: &dyn SegmentReader) -> Option<Arc<Vec<Price>>> {
4747
let segment_key = (segment_reader.segment_id(), segment_reader.delete_opstamp());
4848
self.price_cache.read().unwrap().get(&segment_key).cloned()
4949
}
@@ -157,7 +157,7 @@ fn main() -> tantivy::Result<()> {
157157
let query = query_parser.parse_query("cooking")?;
158158

159159
let searcher = reader.searcher();
160-
let score_by_price = move |segment_reader: &SegmentReader| {
160+
let score_by_price = move |segment_reader: &dyn SegmentReader| {
161161
let price = price_dynamic_column
162162
.price_for_segment(segment_reader)
163163
.unwrap();

src/aggregation/accessor_helpers.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -57,7 +57,7 @@ pub(crate) fn get_numeric_or_date_column_types() -> &'static [ColumnType] {
5757

5858
/// Get fast field reader or empty as default.
5959
pub(crate) fn get_ff_reader(
60-
reader: &SegmentReader,
60+
reader: &dyn SegmentReader,
6161
field_name: &str,
6262
allowed_column_types: Option<&[ColumnType]>,
6363
) -> crate::Result<(columnar::Column<u64>, ColumnType)> {
@@ -74,7 +74,7 @@ pub(crate) fn get_ff_reader(
7474
}
7575

7676
pub(crate) fn get_dynamic_columns(
77-
reader: &SegmentReader,
77+
reader: &dyn SegmentReader,
7878
field_name: &str,
7979
) -> crate::Result<Vec<columnar::DynamicColumn>> {
8080
let ff_fields = reader.fast_fields().dynamic_column_handles(field_name)?;
@@ -90,7 +90,7 @@ pub(crate) fn get_dynamic_columns(
9090
///
9191
/// Is guaranteed to return at least one column.
9292
pub(crate) fn get_all_ff_reader_or_empty(
93-
reader: &SegmentReader,
93+
reader: &dyn SegmentReader,
9494
field_name: &str,
9595
allowed_column_types: Option<&[ColumnType]>,
9696
fallback_type: ColumnType,

src/aggregation/agg_data.rs

Lines changed: 6 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -469,7 +469,7 @@ impl AggKind {
469469
/// Build AggregationsData by walking the request tree.
470470
pub(crate) fn build_aggregations_data_from_req(
471471
aggs: &Aggregations,
472-
reader: &SegmentReader,
472+
reader: &dyn SegmentReader,
473473
segment_ordinal: SegmentOrdinal,
474474
context: AggContextParams,
475475
) -> crate::Result<AggregationsSegmentCtx> {
@@ -489,7 +489,7 @@ pub(crate) fn build_aggregations_data_from_req(
489489
fn build_nodes(
490490
agg_name: &str,
491491
req: &Aggregation,
492-
reader: &SegmentReader,
492+
reader: &dyn SegmentReader,
493493
segment_ordinal: SegmentOrdinal,
494494
data: &mut AggregationsSegmentCtx,
495495
is_top_level: bool,
@@ -728,7 +728,7 @@ fn build_nodes(
728728
let idx_in_req_data = data.push_filter_req_data(FilterAggReqData {
729729
name: agg_name.to_string(),
730730
req: filter_req.clone(),
731-
segment_reader: reader.clone(),
731+
segment_reader: reader.clone_arc(),
732732
evaluator,
733733
matching_docs_buffer,
734734
is_top_level,
@@ -745,7 +745,7 @@ fn build_nodes(
745745

746746
fn build_children(
747747
aggs: &Aggregations,
748-
reader: &SegmentReader,
748+
reader: &dyn SegmentReader,
749749
segment_ordinal: SegmentOrdinal,
750750
data: &mut AggregationsSegmentCtx,
751751
) -> crate::Result<Vec<AggRefNode>> {
@@ -764,7 +764,7 @@ fn build_children(
764764
}
765765

766766
fn get_term_agg_accessors(
767-
reader: &SegmentReader,
767+
reader: &dyn SegmentReader,
768768
field_name: &str,
769769
missing: &Option<Key>,
770770
) -> crate::Result<Vec<(Column<u64>, ColumnType)>> {
@@ -817,7 +817,7 @@ fn build_terms_or_cardinality_nodes(
817817
agg_name: &str,
818818
field_name: &str,
819819
missing: &Option<Key>,
820-
reader: &SegmentReader,
820+
reader: &dyn SegmentReader,
821821
segment_ordinal: SegmentOrdinal,
822822
data: &mut AggregationsSegmentCtx,
823823
sub_aggs: &Aggregations,

src/aggregation/bucket/filter.rs

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,5 @@
11
use std::fmt::Debug;
2+
use std::sync::Arc;
23

34
use common::BitSet;
45
use serde::{Deserialize, Deserializer, Serialize, Serializer};
@@ -402,7 +403,7 @@ pub struct FilterAggReqData {
402403
/// The filter aggregation
403404
pub req: FilterAggregation,
404405
/// The segment reader
405-
pub segment_reader: SegmentReader,
406+
pub segment_reader: Arc<dyn SegmentReader>,
406407
/// Document evaluator for the filter query (precomputed BitSet)
407408
/// This is built once when the request data is created
408409
pub evaluator: DocumentQueryEvaluator,
@@ -416,7 +417,7 @@ impl FilterAggReqData {
416417
pub(crate) fn get_memory_consumption(&self) -> usize {
417418
// Estimate: name + segment reader reference + bitset + buffer capacity
418419
self.name.len()
419-
+ std::mem::size_of::<SegmentReader>()
420+
+ std::mem::size_of::<Arc<dyn SegmentReader>>()
420421
+ self.evaluator.bitset.len() / 8 // BitSet memory (bits to bytes)
421422
+ self.matching_docs_buffer.capacity() * std::mem::size_of::<DocId>()
422423
+ std::mem::size_of::<bool>()
@@ -438,7 +439,7 @@ impl DocumentQueryEvaluator {
438439
pub(crate) fn new(
439440
query: Box<dyn Query>,
440441
schema: Schema,
441-
segment_reader: &SegmentReader,
442+
segment_reader: &dyn SegmentReader,
442443
) -> crate::Result<Self> {
443444
let max_doc = segment_reader.max_doc();
444445

src/aggregation/collector.rs

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -66,7 +66,7 @@ impl Collector for DistributedAggregationCollector {
6666
fn for_segment(
6767
&self,
6868
segment_local_id: crate::SegmentOrdinal,
69-
reader: &crate::SegmentReader,
69+
reader: &dyn SegmentReader,
7070
) -> crate::Result<Self::Child> {
7171
AggregationSegmentCollector::from_agg_req_and_reader(
7272
&self.agg,
@@ -96,7 +96,7 @@ impl Collector for AggregationCollector {
9696
fn for_segment(
9797
&self,
9898
segment_local_id: crate::SegmentOrdinal,
99-
reader: &crate::SegmentReader,
99+
reader: &dyn SegmentReader,
100100
) -> crate::Result<Self::Child> {
101101
AggregationSegmentCollector::from_agg_req_and_reader(
102102
&self.agg,
@@ -145,7 +145,7 @@ impl AggregationSegmentCollector {
145145
/// reader. Also includes validation, e.g. checking field types and existence.
146146
pub fn from_agg_req_and_reader(
147147
agg: &Aggregations,
148-
reader: &SegmentReader,
148+
reader: &dyn SegmentReader,
149149
segment_ordinal: SegmentOrdinal,
150150
context: &AggContextParams,
151151
) -> crate::Result<Self> {

0 commit comments

Comments
 (0)