Skip to content

Commit 477d6c2

Browse files
ragnorccursoragent
andauthored
feat: expose uncommitted delete transactions (lance-format#6781)
## Summary - Add `DeleteBuilder::execute_uncommitted()` so deletes can be staged and committed later with `CommitBuilder`. - Return an `UncommittedDelete` wrapper containing the delete `Transaction`, `affected_rows`, and `num_deleted_rows` so staged commits can preserve row-level conflict rebasing. - Factor delete transaction construction into a helper shared by committed and uncommitted delete paths. ## Test plan - `cargo fmt --all` - `cargo test -p lance dataset::write::delete::tests::test_delete_execute_uncommitted_preserves_affected_rows_for_rebase --lib` - `cargo test -p lance dataset::write::delete::tests::test_delete_false_predicate_still_commits --lib` - `cargo test -p lance dataset::write::delete::tests::test_concurrent_delete_with_retries --lib` - `cargo test -p lance dataset::write::delete::tests::test_delete_concurrency --lib` - `cargo test -p lance dataset::write::delete::tests --lib` Fixes lance-format#6658. --------- Co-authored-by: Cursor <cursoragent@cursor.com>
1 parent 6621199 commit 477d6c2

3 files changed

Lines changed: 163 additions & 17 deletions

File tree

rust/lance/src/dataset.rs

Lines changed: 2 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -137,7 +137,8 @@ pub use write::update::{UpdateBuilder, UpdateJob};
137137
#[allow(deprecated)]
138138
pub use write::{
139139
AutoCleanupParams, CommitBuilder, DeleteBuilder, DeleteResult, ExternalBlobMode, InsertBuilder,
140-
WriteDestination, WriteMode, WriteParams, WriteProgressFn, WriteStats, write_fragments,
140+
UncommittedDelete, WriteDestination, WriteMode, WriteParams, WriteProgressFn, WriteStats,
141+
write_fragments,
141142
};
142143

143144
pub(crate) const INDICES_DIR: &str = "_indices";

rust/lance/src/dataset/write.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -54,7 +54,7 @@ pub mod update;
5454

5555
pub use super::progress::{WriteProgressFn, WriteStats};
5656
pub use commit::CommitBuilder;
57-
pub use delete::{DeleteBuilder, DeleteResult};
57+
pub use delete::{DeleteBuilder, DeleteResult, UncommittedDelete};
5858
pub use insert::InsertBuilder;
5959

6060
/// The destination to write data to.

rust/lance/src/dataset/write/delete.rs

Lines changed: 160 additions & 15 deletions
Original file line numberDiff line numberDiff line change
@@ -31,6 +31,21 @@ pub struct DeleteResult {
3131
pub num_deleted_rows: u64,
3232
}
3333

34+
/// Result of a staged delete operation.
35+
///
36+
/// The returned transaction can be committed later with [`CommitBuilder`].
37+
/// Pass `affected_rows` to [`CommitBuilder::with_affected_rows`] when present
38+
/// to preserve row-level conflict resolution for concurrent deletes and updates.
39+
#[derive(Debug, Clone)]
40+
pub struct UncommittedDelete {
41+
/// The transaction to commit.
42+
pub transaction: Transaction,
43+
/// The row addresses affected by the delete, if available.
44+
pub affected_rows: Option<RowAddrTreeMap>,
45+
/// The number of rows that were deleted.
46+
pub num_deleted_rows: u64,
47+
}
48+
3449
/// Apply deletions to fragments based on a RoaringTreemap of row IDs.
3550
///
3651
/// Returns the set of modified fragments and removed fragments, if any.
@@ -157,6 +172,56 @@ impl DeleteBuilder {
157172

158173
execute_with_retry(job, self.dataset, config).await
159174
}
175+
176+
/// Execute the delete operation without committing the transaction.
177+
///
178+
/// Use [`CommitBuilder`] to commit the returned transaction.
179+
///
180+
/// # Example: Delete rows from a dataset
181+
///
182+
/// ```rust
183+
/// use lance::dataset::{CommitBuilder, DeleteBuilder};
184+
///
185+
/// # use std::sync::Arc;
186+
/// # use lance::Result;
187+
/// # use lance::dataset::Dataset;
188+
/// # async fn example(dataset: Arc<Dataset>) -> Result<()> {
189+
/// let staged_delete = DeleteBuilder::new(dataset.clone(), "age > 65")
190+
/// .execute_uncommitted()
191+
/// .await?;
192+
/// let mut commit_builder = CommitBuilder::new(dataset);
193+
/// if let Some(affected_rows) = staged_delete.affected_rows {
194+
/// commit_builder = commit_builder.with_affected_rows(affected_rows);
195+
/// }
196+
/// commit_builder
197+
/// .execute(staged_delete.transaction)
198+
/// .await?;
199+
/// # Ok(())
200+
/// # }
201+
/// ```
202+
pub async fn execute_uncommitted(self) -> Result<UncommittedDelete> {
203+
let job = DeleteJob {
204+
dataset: self.dataset,
205+
filter: self.filter,
206+
};
207+
let data = job.execute_impl().await?;
208+
let DeleteData {
209+
updated_fragments,
210+
deleted_fragment_ids,
211+
affected_rows,
212+
num_deleted_rows,
213+
} = data;
214+
let transaction = job.build_transaction(
215+
job.dataset.as_ref(),
216+
updated_fragments,
217+
deleted_fragment_ids,
218+
);
219+
Ok(UncommittedDelete {
220+
transaction,
221+
affected_rows,
222+
num_deleted_rows,
223+
})
224+
}
160225
}
161226

162227
/// Job that executes the delete operation
@@ -174,6 +239,29 @@ struct DeleteData {
174239
num_deleted_rows: u64,
175240
}
176241

242+
impl DeleteJob {
243+
fn build_transaction(
244+
&self,
245+
dataset: &Dataset,
246+
updated_fragments: Vec<Fragment>,
247+
deleted_fragment_ids: Vec<u64>,
248+
) -> Transaction {
249+
let predicate = match &self.filter {
250+
ExprFilter::Sql(s) => s.clone(),
251+
ExprFilter::Datafusion(expr) => expr.to_string(),
252+
ExprFilter::Substrait(_) => {
253+
unreachable!("Substrait filters are not supported in DeleteBuilder")
254+
}
255+
};
256+
let operation = Operation::Delete {
257+
updated_fragments,
258+
deleted_fragment_ids,
259+
predicate,
260+
};
261+
Transaction::new(dataset.manifest.version, operation, None)
262+
}
263+
}
264+
177265
impl RetryExecutor for DeleteJob {
178266
type Data = DeleteData;
179267
type Result = DeleteResult;
@@ -265,24 +353,18 @@ impl RetryExecutor for DeleteJob {
265353
}
266354

267355
async fn commit(&self, dataset: Arc<Dataset>, data: Self::Data) -> Result<Self::Result> {
268-
let num_deleted_rows = data.num_deleted_rows;
269-
let predicate = match &self.filter {
270-
ExprFilter::Sql(s) => s.clone(),
271-
ExprFilter::Datafusion(expr) => expr.to_string(),
272-
ExprFilter::Substrait(_) => {
273-
unreachable!("Substrait filters are not supported in DeleteBuilder")
274-
}
275-
};
276-
let operation = Operation::Delete {
277-
updated_fragments: data.updated_fragments,
278-
deleted_fragment_ids: data.deleted_fragment_ids,
279-
predicate,
280-
};
281-
let transaction = Transaction::new(dataset.manifest.version, operation, None);
356+
let DeleteData {
357+
updated_fragments,
358+
deleted_fragment_ids,
359+
affected_rows,
360+
num_deleted_rows,
361+
} = data;
362+
let transaction =
363+
self.build_transaction(dataset.as_ref(), updated_fragments, deleted_fragment_ids);
282364

283365
let mut builder = CommitBuilder::new(dataset);
284366

285-
if let Some(affected_rows) = data.affected_rows {
367+
if let Some(affected_rows) = affected_rows {
286368
builder = builder.with_affected_rows(affected_rows);
287369
}
288370

@@ -615,6 +697,69 @@ mod tests {
615697
assert!(fragments[0].metadata.deletion_file.is_none());
616698
}
617699

700+
#[tokio::test]
701+
async fn test_delete_execute_uncommitted_preserves_affected_rows_for_rebase() {
702+
fn sequence_data(range: Range<u32>) -> RecordBatch {
703+
let schema = Arc::new(ArrowSchema::new(vec![ArrowField::new(
704+
"i",
705+
DataType::UInt32,
706+
false,
707+
)]));
708+
RecordBatch::try_new(schema, vec![Arc::new(UInt32Array::from_iter_values(range))])
709+
.unwrap()
710+
}
711+
712+
let tmp_dir = TempStrDir::default();
713+
let tmp_path = tmp_dir.as_str().to_string();
714+
715+
let dataset = InsertBuilder::new(&tmp_path)
716+
.execute(vec![sequence_data(0..100)])
717+
.await
718+
.unwrap();
719+
let initial_version = dataset.version().version;
720+
721+
let staged_delete = DeleteBuilder::new(Arc::new(dataset.clone()), "i < 10")
722+
.execute_uncommitted()
723+
.await
724+
.unwrap();
725+
726+
let dataset_before_commit = Dataset::open(&tmp_path).await.unwrap();
727+
assert_eq!(dataset_before_commit.version().version, initial_version);
728+
assert_eq!(dataset_before_commit.count_rows(None).await.unwrap(), 100);
729+
730+
assert_eq!(staged_delete.num_deleted_rows, 10);
731+
assert!(staged_delete.affected_rows.is_some());
732+
assert_eq!(staged_delete.transaction.read_version, initial_version);
733+
match &staged_delete.transaction.operation {
734+
Operation::Delete {
735+
updated_fragments,
736+
deleted_fragment_ids,
737+
predicate,
738+
} => {
739+
assert_eq!(predicate, "i < 10");
740+
assert_eq!(updated_fragments.len(), 1);
741+
assert!(deleted_fragment_ids.is_empty());
742+
}
743+
other => panic!("expected delete transaction, got {other:?}"),
744+
}
745+
746+
DeleteBuilder::new(Arc::new(dataset.clone()), "i >= 10 AND i < 20")
747+
.execute()
748+
.await
749+
.unwrap();
750+
751+
let mut commit_builder = CommitBuilder::new(&tmp_path);
752+
if let Some(affected_rows) = staged_delete.affected_rows {
753+
commit_builder = commit_builder.with_affected_rows(affected_rows);
754+
}
755+
let committed = commit_builder
756+
.execute(staged_delete.transaction)
757+
.await
758+
.unwrap();
759+
assert_eq!(committed.version().version, initial_version + 2);
760+
assert_eq!(committed.count_rows(None).await.unwrap(), 80);
761+
}
762+
618763
#[tokio::test]
619764
async fn test_concurrent_delete_with_retries() {
620765
use futures::future::try_join_all;

0 commit comments

Comments
 (0)