mirror of
https://github.com/lancedb/lancedb.git
synced 2026-08-18 03:58:26 +00:00
feat: invalidate generated columns on native delete
This commit is contained in:
@@ -89,6 +89,8 @@ pub mod write_progress;
|
||||
#[cfg(test)]
|
||||
mod append_generated_column_invalidation_contract;
|
||||
#[cfg(test)]
|
||||
mod delete_generated_column_invalidation_contract;
|
||||
#[cfg(test)]
|
||||
mod schema_metadata_updates_dependency_contract;
|
||||
#[cfg(test)]
|
||||
mod update_generated_column_invalidation_contract;
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
use std::sync::Arc;
|
||||
|
||||
use futures::FutureExt;
|
||||
use lance::dataset::DeleteBuilder;
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use lance::dataset::DeleteBuilder;
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use super::{NativeTable, Predicate};
|
||||
@@ -29,34 +29,40 @@ pub(crate) async fn execute_delete(
|
||||
predicate: Predicate<'_>,
|
||||
) -> Result<DeleteResult> {
|
||||
table.dataset.ensure_mutable()?;
|
||||
match predicate {
|
||||
Predicate::String(s) => {
|
||||
let mut dataset = (*table.dataset.get().await?).clone();
|
||||
let delete_result = dataset.delete(s).boxed().await?;
|
||||
let num_deleted_rows = delete_result.num_deleted_rows;
|
||||
let version = dataset.version().version;
|
||||
table.dataset.update(dataset);
|
||||
Ok(DeleteResult {
|
||||
num_deleted_rows,
|
||||
version,
|
||||
})
|
||||
}
|
||||
Predicate::Expr(expr) => {
|
||||
let dataset = table.dataset.get().await?;
|
||||
let delete_result = DeleteBuilder::from_expr(Arc::clone(&dataset), expr.clone())
|
||||
.execute()
|
||||
.await?;
|
||||
let num_deleted_rows = delete_result.num_deleted_rows;
|
||||
let version = delete_result.new_dataset.version().version;
|
||||
table.dataset.update(
|
||||
Arc::try_unwrap(delete_result.new_dataset).unwrap_or_else(|arc| (*arc).clone()),
|
||||
);
|
||||
Ok(DeleteResult {
|
||||
num_deleted_rows,
|
||||
version,
|
||||
})
|
||||
}
|
||||
|
||||
// One exact dataset supplies binding-snapshot planning, the DeleteBuilder,
|
||||
// and its transaction basis. Do not call table schema()/version() or another
|
||||
// get(). Conflicts are not caught/replanned here.
|
||||
let dataset = table.dataset.get().await?;
|
||||
|
||||
// String preserves the legacy Dataset::delete zero-retry baseline; Expr
|
||||
// retains DeleteBuilder defaults until a generated patch is attached.
|
||||
let mut builder = match predicate {
|
||||
Predicate::String(s) => DeleteBuilder::new(Arc::clone(&dataset), s).conflict_retries(0),
|
||||
Predicate::Expr(expr) => DeleteBuilder::from_expr(Arc::clone(&dataset), expr.clone()),
|
||||
};
|
||||
|
||||
if let Some(schema_metadata_updates) =
|
||||
super::generated_column_invalidation::plan_native_delete_generated_column_invalidation(
|
||||
dataset.as_ref(),
|
||||
)?
|
||||
{
|
||||
// Exact-basis fence: never retry an old generated patch on latest.
|
||||
builder = builder
|
||||
.with_schema_metadata_updates(schema_metadata_updates)?
|
||||
.conflict_retries(0);
|
||||
}
|
||||
|
||||
let delete_result = builder.execute().await?;
|
||||
let num_deleted_rows = delete_result.num_deleted_rows;
|
||||
let version = delete_result.new_dataset.version().version;
|
||||
table
|
||||
.dataset
|
||||
.update(Arc::try_unwrap(delete_result.new_dataset).unwrap_or_else(|arc| (*arc).clone()));
|
||||
Ok(DeleteResult {
|
||||
num_deleted_rows,
|
||||
version,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -1,12 +1,12 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! Crate-private Native wiring for generated-column invalidation (B4b / B4c).
|
||||
//! Crate-private Native wiring for generated-column invalidation (B4b / B4c / B4d).
|
||||
//!
|
||||
//! Converts the B4a pure planner into one Lance field-metadata patch for Native
|
||||
//! append and update commits. Planning is strict-decode/validate; overwrite of a
|
||||
//! table with any generated-column definition, and direct writes of generated
|
||||
//! outputs via Update, fail closed as [`Error::NotSupported`].
|
||||
//! append, update, and delete commits. Planning is strict-decode/validate;
|
||||
//! overwrite of a table with any generated-column definition, and direct writes
|
||||
//! of generated outputs via Update, fail closed as [`Error::NotSupported`].
|
||||
|
||||
use std::collections::{BTreeSet, HashMap};
|
||||
|
||||
@@ -95,6 +95,27 @@ pub(super) fn plan_native_update_generated_column_invalidation(
|
||||
Ok(Some(planned_invalidation_to_schema_metadata_updates(plan)))
|
||||
}
|
||||
|
||||
/// Plan Native delete invalidation against one exact dataset snapshot.
|
||||
///
|
||||
/// Strict-decodes and validates every present generated-column metadata value
|
||||
/// through the B4a `RowSetChanged` planner before any Delete scanner/file IO.
|
||||
/// Returns `Some(patch)` when at least one generated column would be invalidated,
|
||||
/// or `None` when the table has no generated columns. Actual zero-row Delete
|
||||
/// suppression is owned by Lance A4d, not this planner.
|
||||
pub(super) fn plan_native_delete_generated_column_invalidation(
|
||||
dataset: &Dataset,
|
||||
) -> Result<Option<SchemaMetadataUpdates>> {
|
||||
let snapshot = generated_column_binding_snapshot_from_dataset(dataset)?;
|
||||
let plan = plan_generated_column_invalidation(
|
||||
&snapshot,
|
||||
&GeneratedColumnMutationImpact::RowSetChanged,
|
||||
)?;
|
||||
if plan.is_empty() {
|
||||
return Ok(None);
|
||||
}
|
||||
Ok(Some(planned_invalidation_to_schema_metadata_updates(plan)))
|
||||
}
|
||||
|
||||
/// Convert planner replacements into one non-empty Lance field-metadata patch.
|
||||
///
|
||||
/// Each entry is keyed by stable output field ID, uses `replace: false`, and
|
||||
|
||||
Reference in New Issue
Block a user