Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
37 changes: 36 additions & 1 deletion rust/lance-table/benches/row_id_index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,7 +27,7 @@ use lance_core::utils::deletion::DeletionVector;
use lance_io::ReadBatchParams;
use lance_select::{RowAddrMask, RowAddrTreeMap};
use lance_table::format::pb;
use lance_table::rowids::FragmentRowIdIndex;
use lance_table::rowids::{FragmentRowIdIndex, segment::U64Segment};
use lance_table::{
rowids::{RowIdIndex, RowIdSequence, read_row_ids, write_row_ids},
utils::stream::{RowIdAndDeletesConfig, apply_row_id_and_deletes},
Expand Down Expand Up @@ -70,6 +70,36 @@ fn make_frag_sequences(
.collect()
}

fn make_clustered_range_frag_sequences(num_rows: u64) -> Vec<FragmentRowIdIndex> {
let rows_per_fragment = num_rows / 100;
(0_u32..100)
.map(|fragment_id| {
let start = fragment_id as u64 * rows_per_fragment;
let end = start + rows_per_fragment;
let mut ranges = Vec::new();
let mut row_id = start;
while row_id < end {
let present_end = (row_id + 250).min(end);
ranges.push(row_id..present_end);
row_id = present_end.saturating_add(250);
}
let wire_sequence = pb::RowIdSequence {
segments: ranges
.into_iter()
.map(|range| pb::U64Segment::from(U64Segment::Range(range)))
.collect(),
};
// Decode the wire ranges into the compact form used after PR #9311.
let sequence = read_row_ids(&wire_sequence.encode_to_vec()).unwrap();
FragmentRowIdIndex {
fragment_id,
row_id_sequence: Arc::new(sequence),
deletion_vector: Arc::new(DeletionVector::default()),
}
})
.collect()
}

// For range of values
// https://bheisler.github.io/criterion.rs/book/user_guide/benchmarking_with_inputs.html

Expand Down Expand Up @@ -199,6 +229,11 @@ fn bench_creation(c: &mut Criterion) {
);
}

let clustered_indices = make_clustered_range_frag_sequences(num_rows());
group.bench_function("BuildIndexClusteredRanges", |b| {
b.iter(|| std::hint::black_box(RowIdIndex::new(&clustered_indices).unwrap()));
});

group.finish();
}

Expand Down
103 changes: 103 additions & 0 deletions rust/lance-table/src/rowids/index.rs
Original file line number Diff line number Diff line change
Expand Up @@ -557,6 +557,15 @@ fn decompose_segment_no_deletions(segment: &U64Segment, start_address: u64) -> O
let coverage = range.start..=range.end - 1;
Some((coverage, (row_id_segment, address_segment)))
}
U64Segment::Ranges { .. } if !segment.is_empty() => {
// Without deletions, a live ID's position is also its physical
// offset within this segment.
let row_id_segment = segment.clone();
let address_segment =
U64Segment::Range(start_address..start_address + row_id_segment.len() as u64);
let coverage = row_id_segment.range()?;
Some((coverage, (row_id_segment, address_segment)))
}
_ if segment.is_empty() => None,
_ => {
// Non-Range segments: must iterate to build address mapping.
Expand Down Expand Up @@ -801,6 +810,100 @@ mod tests {
);
}

#[test]
fn test_merged_index_keeps_ranges_segments_compact() {
let row_id_ranges = [100..103, 108..110, 115..119];
let sequence = RowIdSequence(vec![
U64Segment::Range(50..53),
U64Segment::from_sorted_ranges(&row_id_ranges).unwrap(),
]);
let source = fragment(7, sequence);
let index = RowIdIndex::new(std::slice::from_ref(&source)).unwrap();

let merged = index.merged.as_ref().expect("small input should be merged");
assert!(matches!(
merged.get(&100).unwrap().0,
U64Segment::Ranges { .. }
));

for (row_id, row_offset) in [
(50, 0),
(51, 1),
(52, 2),
(100, 3),
(101, 4),
(102, 5),
(108, 6),
(109, 7),
(115, 8),
(116, 9),
(117, 10),
(118, 11),
] {
assert_eq!(
index.get(row_id).unwrap(),
Some(RowAddress::new_from_parts(7, row_offset)),
"row id {row_id}"
);
}
for missing in [49, 53, 99, 103, 107, 110, 114, 119, 120] {
assert_eq!(index.get(missing).unwrap(), None, "row id {missing}");
}
assert_eq!(
index.get_many(&[108, 53, 100, 108]).unwrap(),
vec![
Some(RowAddress::new_from_parts(7, 6)),
None,
Some(RowAddress::new_from_parts(7, 3)),
Some(RowAddress::new_from_parts(7, 6)),
]
);
}

#[test]
fn test_overlapping_ranges_chunks_keep_row_addresses() {
let left_ranges = [100..103, 108..110, 115..119];
let right_ranges = [103..108, 110..115];
let sources = [
fragment(
7,
RowIdSequence(vec![U64Segment::from_sorted_ranges(&left_ranges).unwrap()]),
),
fragment(
8,
RowIdSequence(vec![U64Segment::from_sorted_ranges(&right_ranges).unwrap()]),
),
];
let index = RowIdIndex::new(&sources).unwrap();

for row_id in 100..119 {
let (fragment_id, row_offset) = if (100..103).contains(&row_id)
|| (108..110).contains(&row_id)
|| (115..119).contains(&row_id)
{
let row_offset = match row_id {
100..=102 => row_id - 100,
108..=109 => row_id - 105,
115..=118 => row_id - 110,
_ => unreachable!(),
};
(7, row_offset)
} else {
let row_offset = match row_id {
103..=107 => row_id - 103,
110..=114 => row_id - 105,
_ => unreachable!(),
};
(8, row_offset)
};
assert_eq!(
index.get(row_id).unwrap(),
Some(RowAddress::new_from_parts(fragment_id, row_offset as u32)),
"row id {row_id}"
);
}
}

#[test]
fn test_deep_overlap_merges_however_many_rows_it_reads() {
// Just past the row budget in total, interleaved so every fragment
Expand Down
Loading