mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-11 15:52:17 +00:00
Compare commits
6
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
74c003af59 | ||
|
|
695e7ae956 | ||
|
|
67532dc4ce | ||
|
|
04836120d4 | ||
|
|
160b72bbd9 | ||
|
|
53c6495a32 |
@@ -6,19 +6,162 @@
|
|||||||
//! This module contains the implementation of optimization operations that help
|
//! This module contains the implementation of optimization operations that help
|
||||||
//! maintain good performance for LanceDB tables.
|
//! maintain good performance for LanceDB tables.
|
||||||
|
|
||||||
use std::sync::Arc;
|
use std::{collections::HashSet, sync::Arc};
|
||||||
|
|
||||||
|
use arrow_schema::DataType;
|
||||||
use lance::dataset::cleanup::RemovalStats;
|
use lance::dataset::cleanup::RemovalStats;
|
||||||
use lance::dataset::optimize::{CompactionMetrics, IndexRemapperOptions, compact_files};
|
use lance::dataset::optimize::{
|
||||||
use lance::index::DatasetIndexExt;
|
CompactionMetrics, CompactionPlan, CompactionPlanner, IndexRemapperOptions,
|
||||||
|
compact_files_with_planner, plan_compaction,
|
||||||
|
};
|
||||||
|
use lance::index::{DatasetIndexExt, DatasetIndexInternalExt};
|
||||||
|
use lance_index::IndexType;
|
||||||
|
use lance_index::metrics::NoOpMetricsCollector;
|
||||||
use lance_index::optimize::OptimizeOptions;
|
use lance_index::optimize::OptimizeOptions;
|
||||||
use log::info;
|
use log::{debug, info};
|
||||||
|
|
||||||
|
use super::NativeTable;
|
||||||
|
use crate::error::{Error, Result};
|
||||||
pub use chrono::Duration;
|
pub use chrono::Duration;
|
||||||
pub use lance::dataset::optimize::CompactionOptions;
|
pub use lance::dataset::optimize::CompactionOptions;
|
||||||
|
|
||||||
use super::NativeTable;
|
const MAX_ARROW_FIXED_SIZE_LIST_CHILD_INDEX: u64 = u32::MAX as u64;
|
||||||
use crate::error::Result;
|
|
||||||
|
struct PrecomputedCompactionPlanner {
|
||||||
|
plan: CompactionPlan,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[async_trait::async_trait]
|
||||||
|
impl CompactionPlanner for PrecomputedCompactionPlanner {
|
||||||
|
async fn plan(&self, dataset: &lance::Dataset) -> lance::Result<CompactionPlan> {
|
||||||
|
debug_assert_eq!(dataset.manifest().version, self.plan.read_version());
|
||||||
|
Ok(self.plan.clone())
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn sq_vector_dimension(data_type: &DataType) -> Option<i32> {
|
||||||
|
match data_type {
|
||||||
|
DataType::FixedSizeList(_, dimension) => Some(*dimension),
|
||||||
|
DataType::List(field) => match field.data_type() {
|
||||||
|
DataType::FixedSizeList(_, dimension) => Some(*dimension),
|
||||||
|
_ => None,
|
||||||
|
},
|
||||||
|
_ => None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn segment_touches_fragments(
|
||||||
|
fragment_ids: Option<impl IntoIterator<Item = u32>>,
|
||||||
|
affected_fragments: &HashSet<u64>,
|
||||||
|
) -> bool {
|
||||||
|
fragment_ids.is_none_or(|fragment_ids| {
|
||||||
|
fragment_ids
|
||||||
|
.into_iter()
|
||||||
|
.any(|fragment| affected_fragments.contains(&(fragment as u64)))
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
fn oversized_sq_partition(
|
||||||
|
partition_sizes: impl IntoIterator<Item = u64>,
|
||||||
|
dimension: u64,
|
||||||
|
) -> Option<u64> {
|
||||||
|
partition_sizes.into_iter().find(|partition_size| {
|
||||||
|
partition_size
|
||||||
|
.checked_mul(dimension)
|
||||||
|
.is_none_or(|child_len| child_len > MAX_ARROW_FIXED_SIZE_LIST_CHILD_INDEX)
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
|
async fn validate_sq_index_remapping(
|
||||||
|
dataset: &lance::Dataset,
|
||||||
|
options: &CompactionOptions,
|
||||||
|
has_custom_remapper: bool,
|
||||||
|
) -> Result<CompactionPlan> {
|
||||||
|
let plan = plan_compaction(dataset, options).await?;
|
||||||
|
if options.defer_index_remap || has_custom_remapper || dataset.manifest().uses_stable_row_ids()
|
||||||
|
{
|
||||||
|
return Ok(plan);
|
||||||
|
}
|
||||||
|
|
||||||
|
if plan.tasks.is_empty() {
|
||||||
|
return Ok(plan);
|
||||||
|
}
|
||||||
|
let affected_fragments: HashSet<u64> = plan
|
||||||
|
.tasks
|
||||||
|
.iter()
|
||||||
|
.flat_map(|task| task.fragments.iter().map(|fragment| fragment.id))
|
||||||
|
.collect();
|
||||||
|
|
||||||
|
for segment in dataset.load_indices().await?.iter() {
|
||||||
|
if !segment_touches_fragments(
|
||||||
|
segment
|
||||||
|
.fragment_bitmap
|
||||||
|
.as_ref()
|
||||||
|
.map(|fragment_bitmap| fragment_bitmap.iter()),
|
||||||
|
&affected_fragments,
|
||||||
|
) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
let Some(field_id) = segment.fields.first() else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
let Some(field) = dataset.schema().field_by_id(*field_id) else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
let data_type = field.data_type();
|
||||||
|
let Some(dimension) = sq_vector_dimension(&data_type) else {
|
||||||
|
continue;
|
||||||
|
};
|
||||||
|
let dimension = u64::try_from(dimension).map_err(|_| Error::InvalidInput {
|
||||||
|
message: format!(
|
||||||
|
"SQ index '{}' has an invalid vector dimension of {}",
|
||||||
|
segment.name, dimension
|
||||||
|
),
|
||||||
|
})?;
|
||||||
|
|
||||||
|
let field_path = dataset.schema().field_path(*field_id)?;
|
||||||
|
let vector_index = match dataset
|
||||||
|
.open_vector_index(&field_path, &segment.uuid, &NoOpMetricsCollector)
|
||||||
|
.await
|
||||||
|
{
|
||||||
|
Ok(vector_index) => vector_index,
|
||||||
|
Err(error) => {
|
||||||
|
// The default remapper also treats an index it cannot open as
|
||||||
|
// unusable and drops it, so it cannot reach the SQ take kernel.
|
||||||
|
debug!(
|
||||||
|
"Skipping remap preflight for index segment {} because it could not be opened: {}",
|
||||||
|
segment.uuid, error
|
||||||
|
);
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
if !matches!(
|
||||||
|
vector_index.index_type(),
|
||||||
|
IndexType::IvfSq | IndexType::IvfHnswSq
|
||||||
|
) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
|
||||||
|
let partition_sizes = (0..vector_index.total_partitions())
|
||||||
|
.map(|partition_id| vector_index.partition_size(partition_id) as u64);
|
||||||
|
if let Some(partition_size) = oversized_sq_partition(partition_sizes, dimension) {
|
||||||
|
let max_partition_size = MAX_ARROW_FIXED_SIZE_LIST_CHILD_INDEX / dimension;
|
||||||
|
return Err(Error::InvalidInput {
|
||||||
|
message: format!(
|
||||||
|
"Cannot compact table because SQ index '{}' segment {} has a partition with {} vectors of dimension {}. Arrow's fixed-size-list take kernel cannot remap partitions whose child array exceeds {} values. Recreate the index with default IVF partitioning, or enough partitions to keep each partition at or below {} vectors, before compacting.",
|
||||||
|
segment.name,
|
||||||
|
segment.uuid,
|
||||||
|
partition_size,
|
||||||
|
dimension,
|
||||||
|
MAX_ARROW_FIXED_SIZE_LIST_CHILD_INDEX,
|
||||||
|
max_partition_size,
|
||||||
|
),
|
||||||
|
});
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(plan)
|
||||||
|
}
|
||||||
|
|
||||||
/// Optimize the dataset.
|
/// Optimize the dataset.
|
||||||
///
|
///
|
||||||
@@ -144,7 +287,7 @@ pub(crate) async fn cleanup_old_versions(
|
|||||||
/// This can be run after making several small appends to optimize the table
|
/// This can be run after making several small appends to optimize the table
|
||||||
/// for faster reads.
|
/// for faster reads.
|
||||||
///
|
///
|
||||||
/// This calls into [lance::dataset::optimize::compact_files].
|
/// This calls into [lance::dataset::optimize::compact_files_with_planner].
|
||||||
pub(crate) async fn compact_files_impl(
|
pub(crate) async fn compact_files_impl(
|
||||||
table: &NativeTable,
|
table: &NativeTable,
|
||||||
options: CompactionOptions,
|
options: CompactionOptions,
|
||||||
@@ -152,7 +295,9 @@ pub(crate) async fn compact_files_impl(
|
|||||||
) -> Result<CompactionMetrics> {
|
) -> Result<CompactionMetrics> {
|
||||||
table.dataset.ensure_mutable()?;
|
table.dataset.ensure_mutable()?;
|
||||||
let mut dataset = (*table.dataset.get().await?).clone();
|
let mut dataset = (*table.dataset.get().await?).clone();
|
||||||
let metrics = compact_files(&mut dataset, options, remap_options).await?;
|
let plan = validate_sq_index_remapping(&dataset, &options, remap_options.is_some()).await?;
|
||||||
|
let planner = PrecomputedCompactionPlanner { plan };
|
||||||
|
let metrics = compact_files_with_planner(&mut dataset, remap_options, &planner).await?;
|
||||||
table.dataset.update(dataset);
|
table.dataset.update(dataset);
|
||||||
Ok(metrics)
|
Ok(metrics)
|
||||||
}
|
}
|
||||||
@@ -214,13 +359,14 @@ pub(crate) async fn execute_optimize(
|
|||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
|
use super::{oversized_sq_partition, segment_touches_fragments, sq_vector_dimension};
|
||||||
use arrow_array::{
|
use arrow_array::{
|
||||||
Array, FixedSizeListArray, Float32Array, Int32Array, RecordBatch, StringArray,
|
Array, FixedSizeListArray, Float32Array, Int32Array, RecordBatch, StringArray,
|
||||||
};
|
};
|
||||||
use arrow_schema::{DataType, Field, Schema};
|
use arrow_schema::{DataType, Field, Schema};
|
||||||
use lance_arrow::FixedSizeListArrayExt;
|
use lance_arrow::FixedSizeListArrayExt;
|
||||||
use rstest::rstest;
|
use rstest::rstest;
|
||||||
use std::sync::Arc;
|
use std::{collections::HashSet, sync::Arc};
|
||||||
|
|
||||||
use crate::connect;
|
use crate::connect;
|
||||||
use crate::database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS;
|
use crate::database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS;
|
||||||
@@ -230,6 +376,172 @@ mod tests {
|
|||||||
use crate::table::{CompactionOptions, OptimizeAction, OptimizeStats};
|
use crate::table::{CompactionOptions, OptimizeAction, OptimizeStats};
|
||||||
use futures::TryStreamExt;
|
use futures::TryStreamExt;
|
||||||
|
|
||||||
|
/// Regression test for https://github.com/lancedb/lancedb/issues/2866.
|
||||||
|
#[test]
|
||||||
|
fn test_detect_oversized_sq_partition() {
|
||||||
|
const DIMENSION: u64 = 4095;
|
||||||
|
let safe_size = u32::MAX as u64 / DIMENSION;
|
||||||
|
|
||||||
|
assert_eq!(
|
||||||
|
oversized_sq_partition([safe_size, safe_size + 1], DIMENSION),
|
||||||
|
Some(safe_size + 1)
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
oversized_sq_partition([safe_size / 2, safe_size / 2 + 1], DIMENSION),
|
||||||
|
None
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_segment_touches_fragments_includes_unknown_coverage() {
|
||||||
|
let affected_fragments = HashSet::from([7]);
|
||||||
|
|
||||||
|
assert!(segment_touches_fragments(
|
||||||
|
Some([7_u32]),
|
||||||
|
&affected_fragments
|
||||||
|
));
|
||||||
|
assert!(!segment_touches_fragments(
|
||||||
|
Some([8_u32]),
|
||||||
|
&affected_fragments
|
||||||
|
));
|
||||||
|
assert!(segment_touches_fragments(
|
||||||
|
None::<[u32; 0]>,
|
||||||
|
&affected_fragments
|
||||||
|
));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_sq_vector_dimension_includes_multivectors() {
|
||||||
|
const DIMENSION: i32 = 4095;
|
||||||
|
let vector = DataType::FixedSizeList(
|
||||||
|
Arc::new(Field::new("item", DataType::Float32, false)),
|
||||||
|
DIMENSION,
|
||||||
|
);
|
||||||
|
let multivector = DataType::List(Arc::new(Field::new("item", vector.clone(), false)));
|
||||||
|
|
||||||
|
assert_eq!(sq_vector_dimension(&vector), Some(DIMENSION));
|
||||||
|
assert_eq!(sq_vector_dimension(&multivector), Some(DIMENSION));
|
||||||
|
assert_eq!(sq_vector_dimension(&DataType::Float32), None);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[tokio::test]
|
||||||
|
async fn test_compact_legacy_vector_index_with_unknown_fragment_coverage() {
|
||||||
|
use lance::Dataset;
|
||||||
|
use lance::index::vector::{IndexFileVersion, VectorIndexParams};
|
||||||
|
use lance::index::{DatasetIndexExt, DatasetIndexInternalExt};
|
||||||
|
use lance_index::IndexType;
|
||||||
|
use lance_linalg::distance::MetricType;
|
||||||
|
use lance_table::io::commit::write_manifest_file_to_path;
|
||||||
|
use object_store::ObjectStoreExt;
|
||||||
|
|
||||||
|
const INDEX_NAME: &str = "legacy_ivf_pq";
|
||||||
|
const ROWS: i32 = 128;
|
||||||
|
const DIMENSION: i32 = 8;
|
||||||
|
|
||||||
|
let tmpdir = tempfile::tempdir().unwrap();
|
||||||
|
let conn = connect(tmpdir.path().to_str().unwrap())
|
||||||
|
.execute()
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let vectors = FixedSizeListArray::try_new_from_values(
|
||||||
|
Float32Array::from_iter_values((0..ROWS).flat_map(|row| {
|
||||||
|
(0..DIMENSION).map(move |offset| (row * DIMENSION + offset) as f32)
|
||||||
|
})),
|
||||||
|
DIMENSION,
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
let schema = Arc::new(Schema::new(vec![Field::new(
|
||||||
|
"vector",
|
||||||
|
vectors.data_type().clone(),
|
||||||
|
false,
|
||||||
|
)]));
|
||||||
|
let batch = RecordBatch::try_new(schema, vec![Arc::new(vectors)]).unwrap();
|
||||||
|
let table = conn
|
||||||
|
.create_table("test_legacy_vector_compact", batch.clone())
|
||||||
|
.execute()
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
table.add(batch.clone()).execute().await.unwrap();
|
||||||
|
table.add(batch).execute().await.unwrap();
|
||||||
|
let mut dataset = (*table.dataset().unwrap().get().await.unwrap()).clone();
|
||||||
|
let mut params = VectorIndexParams::ivf_pq(1, 8, 1, MetricType::L2, 10);
|
||||||
|
params.version(IndexFileVersion::Legacy);
|
||||||
|
dataset
|
||||||
|
.create_index(
|
||||||
|
&["vector"],
|
||||||
|
IndexType::Vector,
|
||||||
|
Some(INDEX_NAME.to_string()),
|
||||||
|
¶ms,
|
||||||
|
true,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
table.dataset().unwrap().update(dataset);
|
||||||
|
|
||||||
|
let dataset = table.dataset().unwrap().get().await.unwrap();
|
||||||
|
let mut segments = dataset.load_indices_by_name(INDEX_NAME).await.unwrap();
|
||||||
|
assert_eq!(segments.len(), 1);
|
||||||
|
let original = segments.pop().unwrap();
|
||||||
|
let field_path = dataset.schema().field_path(original.fields[0]).unwrap();
|
||||||
|
let vector_index = dataset
|
||||||
|
.open_vector_index(
|
||||||
|
&field_path,
|
||||||
|
&original.uuid,
|
||||||
|
&lance_index::metrics::NoOpMetricsCollector,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert_eq!(vector_index.index_type(), IndexType::IvfPq);
|
||||||
|
assert_eq!(
|
||||||
|
vector_index.statistics().unwrap()["index_file_version"],
|
||||||
|
"Legacy"
|
||||||
|
);
|
||||||
|
let mut legacy = original.clone();
|
||||||
|
legacy.fragment_bitmap = None;
|
||||||
|
let object_store = dataset.object_store(None).await.unwrap();
|
||||||
|
let mut manifest = dataset.manifest().clone();
|
||||||
|
manifest.index_section = None;
|
||||||
|
manifest.transaction_section = None;
|
||||||
|
let manifest_path = dataset.manifest_location().path.clone();
|
||||||
|
object_store.inner.delete(&manifest_path).await.unwrap();
|
||||||
|
write_manifest_file_to_path(
|
||||||
|
object_store.as_ref(),
|
||||||
|
&mut manifest,
|
||||||
|
Some(vec![legacy]),
|
||||||
|
&manifest_path,
|
||||||
|
None,
|
||||||
|
)
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
let legacy_dataset = Dataset::open(dataset.uri()).await.unwrap();
|
||||||
|
table.dataset().unwrap().update(legacy_dataset);
|
||||||
|
|
||||||
|
let dataset = table.dataset().unwrap().get().await.unwrap();
|
||||||
|
let description_error = dataset
|
||||||
|
.describe_indices(None)
|
||||||
|
.await
|
||||||
|
.err()
|
||||||
|
.expect("legacy coverage should be unknown to index descriptions");
|
||||||
|
assert!(
|
||||||
|
description_error
|
||||||
|
.to_string()
|
||||||
|
.contains("Fragment bitmap is required")
|
||||||
|
);
|
||||||
|
|
||||||
|
let stats = table
|
||||||
|
.optimize(OptimizeAction::Compact {
|
||||||
|
options: CompactionOptions {
|
||||||
|
target_rows_per_fragment: 1_000,
|
||||||
|
..Default::default()
|
||||||
|
},
|
||||||
|
remap_options: None,
|
||||||
|
})
|
||||||
|
.await
|
||||||
|
.unwrap();
|
||||||
|
assert!(stats.compaction.unwrap().fragments_removed > 0);
|
||||||
|
assert_eq!(table.count_rows(None).await.unwrap(), ROWS as usize * 3);
|
||||||
|
}
|
||||||
|
|
||||||
#[tokio::test]
|
#[tokio::test]
|
||||||
async fn test_optimize_compact_simple() {
|
async fn test_optimize_compact_simple() {
|
||||||
let conn = connect("memory://").execute().await.unwrap();
|
let conn = connect("memory://").execute().await.unwrap();
|
||||||
|
|||||||
Reference in New Issue
Block a user