Skip to content
Open
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
816 changes: 703 additions & 113 deletions rust/lance-table/src/rowids/index.rs

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion rust/lance/src/dataset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3001,7 +3001,7 @@ impl Dataset {
let mut live_ids = Vec::with_capacity(ids.len());
let mut addresses = Vec::with_capacity(ids.len());
for id in ids {
if let Some(address) = row_id_index.get(*id) {
if let Some(address) = row_id_index.get(*id)? {
live_ids.push(*id);
addresses.push(u64::from(address));
}
Expand Down
2 changes: 1 addition & 1 deletion rust/lance/src/dataset/optimize.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2345,7 +2345,7 @@ async fn rewrite_files(
let captured_ids = row_ids_rx
.try_recv()
.map_err(|err| Error::internal(format!("Failed to receive row ids: {}", err)))?;
let row_addrs = captured_ids.row_addrs(None).into_owned();
let row_addrs = captured_ids.row_addrs(None)?.into_owned();
let mut serialized = Vec::with_capacity(row_addrs.serialized_size());
row_addrs.serialize_into(&mut serialized)?;
Ok(Some(serialized))
Expand Down
25 changes: 14 additions & 11 deletions rust/lance/src/dataset/rowids.rs
Original file line number Diff line number Diff line change
Expand Up @@ -330,7 +330,7 @@ mod test {
assert!(dataset.manifest.uses_stable_row_ids());

let index = get_row_id_index(&dataset).await.unwrap().unwrap();
assert!(index.get(0).is_none());
assert!(index.get(0).unwrap().is_none());

assert_eq!(dataset.manifest().next_row_id, 0);
}
Expand Down Expand Up @@ -384,7 +384,7 @@ mod test {
let index = get_row_id_index(&dataset).await.unwrap().unwrap();

let found_addresses = (0..num_rows)
.map(|i| index.get(i).unwrap())
.map(|i| index.get(i).unwrap().unwrap())
.collect::<Vec<_>>();
let expected_addresses = (0..num_rows)
.map(|i| {
Expand Down Expand Up @@ -446,8 +446,8 @@ mod test {

failing_store.clear_fail_when("get_opts", "_deletions");
let index = get_row_id_index(&dataset).await.unwrap().unwrap();
assert!(index.get(2).is_some());
assert!(index.get(3).is_none());
assert!(index.get(2).unwrap().is_some());
assert!(index.get(3).unwrap().is_none());
}

#[tokio::test]
Expand Down Expand Up @@ -530,8 +530,8 @@ mod test {
assert_eq!(dataset.manifest.fragments[0].id, 1);

let index = get_row_id_index(&dataset).await.unwrap().unwrap();
assert!(index.get(0).is_none());
assert!(index.get(num_rows).is_some());
assert!(index.get(0).unwrap().is_none());
assert!(index.get(num_rows).unwrap().is_some());
}

/// Fragment ids are a high water mark within one dataset, but a dataset
Expand Down Expand Up @@ -683,8 +683,8 @@ mod test {
assert_eq!(dataset.manifest().next_row_id, 60);

let index = get_row_id_index(&dataset).await.unwrap().unwrap();
assert!(index.get(0).is_some());
assert!(index.get(60).is_none());
assert!(index.get(0).unwrap().is_some());
assert!(index.get(60).unwrap().is_none());
}

#[tokio::test]
Expand Down Expand Up @@ -830,11 +830,14 @@ mod test {

let dataset = update_result.new_dataset;
let index = get_row_id_index(&dataset).await.unwrap().unwrap();
assert!(index.get(0).is_some());
assert!(index.get(0).unwrap().is_some());
// the updated row ids mapping to new address
assert_eq!(index.get(3), Some(RowAddress::new_from_parts(1, 0)));
assert_eq!(
index.get(3).unwrap(),
Some(RowAddress::new_from_parts(1, 0))
);
// there is no new row id
assert_eq!(index.get(5), None);
assert_eq!(index.get(5).unwrap(), None);
}

/// 100 sequential rows across 4 fragments with every third row deleted.
Expand Down
2 changes: 1 addition & 1 deletion rust/lance/src/dataset/take.rs
Original file line number Diff line number Diff line change
Expand Up @@ -555,7 +555,7 @@ impl TakeBuilder {
.as_ref()
.expect("row_ids must be set if row_addrs is not");
let addrs = if let Some(row_id_index) = get_row_id_index(&self.dataset).await? {
let resolved = row_id_index.get_many(row_ids);
let resolved = row_id_index.get_many(row_ids)?;
if self.missing_row_policy == MissingRowPolicy::Error
&& let Some(first_missing_index) =
resolved.iter().position(|address| address.is_none())
Expand Down
13 changes: 9 additions & 4 deletions rust/lance/src/dataset/utils.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,18 +117,23 @@ impl CapturedRowIds {
}
}

pub fn row_addrs(&self, index: Option<&RowIdIndex>) -> Cow<'_, RoaringTreemap> {
pub fn row_addrs(&self, index: Option<&RowIdIndex>) -> Result<Cow<'_, RoaringTreemap>> {
match self {
Self::AddressStyle(addrs) => Cow::Borrowed(addrs),
Self::AddressStyle(addrs) => Ok(Cow::Borrowed(addrs)),
Self::SequenceStyle(sequence) => {
let mut treemap = RoaringTreemap::new();
let Some(index) = index else {
panic!("RowIdIndex required for sequence style row ids")
};
for row_id in sequence.iter() {
treemap.insert(index.get(row_id).expect("row id missing from index").into());
treemap.insert(
index
.get(row_id)?
.expect("row id missing from index")
.into(),
);
}
Cow::Owned(treemap)
Ok(Cow::Owned(treemap))
}
}
}
Expand Down
2 changes: 1 addition & 1 deletion rust/lance/src/dataset/write/delete.rs
Original file line number Diff line number Diff line change
Expand Up @@ -326,7 +326,7 @@ impl RetryExecutor for DeleteJob {
Error::internal(format!("Failed to receive row ids: {}", err))
})?;
let row_id_index = get_row_id_index(&self.dataset).await?;
let removed_row_addrs = removed_row_ids.row_addrs(row_id_index.as_deref());
let removed_row_addrs = removed_row_ids.row_addrs(row_id_index.as_deref())?;

let (fragments, deleted_ids) =
apply_deletions(&self.dataset, &removed_row_addrs).await?;
Expand Down
21 changes: 13 additions & 8 deletions rust/lance/src/dataset/write/merge_insert.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2404,10 +2404,13 @@ impl MergeInsertJob {
let removed_row_ids = Arc::into_inner(deleted_rows).unwrap().into_inner().unwrap();
let removed_row_addr_vec =
if let Some(row_id_index) = get_row_id_index(&self.dataset).await? {
removed_row_ids
.iter()
.filter_map(|id| row_id_index.get(*id).map(|address| address.into()))
.collect::<Vec<_>>()
let mut addresses = Vec::with_capacity(removed_row_ids.len());
for id in &removed_row_ids {
if let Some(address) = row_id_index.get(*id)? {
addresses.push(address.into());
}
}
addresses
} else {
removed_row_ids
};
Expand Down Expand Up @@ -2517,10 +2520,12 @@ impl MergeInsertJob {

let removed_row_addr_vec =
if let Some(row_id_index) = get_row_id_index(&self.dataset).await? {
let addresses: Vec<u64> = removed_row_ids
.iter()
.filter_map(|id| row_id_index.get(*id).map(|address| address.into()))
.collect::<Vec<_>>();
let mut addresses: Vec<u64> = Vec::with_capacity(removed_row_ids.len());
for id in &removed_row_ids {
if let Some(address) = row_id_index.get(*id)? {
addresses.push(address.into());
}
}
addresses
} else {
removed_row_ids
Expand Down
2 changes: 1 addition & 1 deletion rust/lance/src/dataset/write/update.rs
Original file line number Diff line number Diff line change
Expand Up @@ -469,7 +469,7 @@ impl UpdateJob {

// Apply deletions
let row_id_index = get_row_id_index(&self.dataset).await?;
let row_addrs = removed_row_ids.row_addrs(row_id_index.as_deref());
let row_addrs = removed_row_ids.row_addrs(row_id_index.as_deref())?;
let deletions_result = self.apply_deletions(&row_addrs).await;
let (old_fragments, removed_fragment_ids) = match deletions_result {
Ok(v) => v,
Expand Down
8 changes: 6 additions & 2 deletions rust/lance/src/io/exec/rowids.rs
Original file line number Diff line number Diff line change
Expand Up @@ -136,7 +136,9 @@ impl AddRowAddrExec {
let mut builder = arrow::array::UInt64Builder::with_capacity(row_id_values.len());
for rowid in row_id_values.iter() {
if let Some(rowid) = rowid {
if let Some(row_addr) = row_id_index.get(rowid) {
if let Some(row_addr) =
row_id_index.get(rowid).map_err(DataFusionError::from)?
{
builder.append_value(row_addr.into());
} else {
return Err(DataFusionError::Internal(format!(
Expand All @@ -153,7 +155,9 @@ impl AddRowAddrExec {
// Fast path - no branching for null values
let mut rowaddrs: Vec<u64> = Vec::with_capacity(row_id_values.len());
for rowid in row_id_values.values() {
if let Some(row_addr) = row_id_index.get(*rowid) {
if let Some(row_addr) =
row_id_index.get(*rowid).map_err(DataFusionError::from)?
{
rowaddrs.push(row_addr.into());
} else {
return Err(DataFusionError::Internal(format!(
Expand Down
2 changes: 1 addition & 1 deletion rust/lance/src/io/exec/take.rs
Original file line number Diff line number Diff line change
Expand Up @@ -193,7 +193,7 @@ impl TakeStream {
let mut valid = Vec::with_capacity(row_id_array.len());

for id in row_id_array.values().iter() {
if let Some(address) = row_id_index.get(*id) {
if let Some(address) = row_id_index.get(*id)? {
addresses.push(u64::from(address));
valid.push(true);
} else {
Expand Down
Loading