Compare commits

...
9 changed files with 163 additions and 22 deletions
+3
View File
@@ -587,6 +587,9 @@ Modeled after ``VACUUM`` in PostgreSQL.
you have added or modified 100,000 or more records or run more than 20 data you have added or modified 100,000 or more records or run more than 20 data
modification operations. modification operations.
Cleanup retention must exceed the longest expected write. Retention shorter
than 10 minutes requires `deleteUnverified: true` and exclusive write access.
#### Parameters #### Parameters
* **options?**: `Partial`<[`OptimizeOptions`](../interfaces/OptimizeOptions.md)> * **options?**: `Partial`<[`OptimizeOptions`](../interfaces/OptimizeOptions.md)>
+6 -3
View File
@@ -16,7 +16,10 @@ cleanupOlderThan: Date;
If set then all versions older than the given date If set then all versions older than the given date
be removed. The current version will never be removed. be removed. The current version will never be removed.
The default is 7 days The default is 7 days. The resulting retention period must be longer than
the longest expected write. Values shorter than 10 minutes require
`deleteUnverified: true` and are only safe when no other process can write
to the dataset.
#### Example #### Example
@@ -26,8 +29,8 @@ const olderThan = new Date();
olderThan.setDate(olderThan.getDate() - 1)); olderThan.setDate(olderThan.getDate() - 1));
tbl.optimize({cleanupOlderThan: olderThan}); tbl.optimize({cleanupOlderThan: olderThan});
// Delete all versions except the current version // With exclusive access, delete all versions except the current version
tbl.optimize({cleanupOlderThan: new Date()}); tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true});
``` ```
*** ***
+9 -1
View File
@@ -2196,7 +2196,15 @@ describe("when optimizing a dataset", () => {
}); });
it("cleanups old versions", async () => { it("cleanups old versions", async () => {
const stats = await table.optimize({ cleanupOlderThan: new Date() }); await expect(
table.optimize({ cleanupOlderThan: new Date() }),
).rejects.toThrow("at least 10 minutes");
expect(await table.version()).toBe(2);
const stats = await table.optimize({
cleanupOlderThan: new Date(),
deleteUnverified: true,
});
expect(stats.prune.bytesRemoved).toBeGreaterThan(0); expect(stats.prune.bytesRemoved).toBeGreaterThan(0);
expect(stats.prune.oldVersionsRemoved).toBe(3); expect(stats.prune.oldVersionsRemoved).toBe(3);
}); });
+9 -3
View File
@@ -129,15 +129,18 @@ export interface OptimizeOptions {
/** /**
* If set then all versions older than the given date * If set then all versions older than the given date
* be removed. The current version will never be removed. * be removed. The current version will never be removed.
* The default is 7 days * The default is 7 days. The resulting retention period must be longer than
* the longest expected write. Values shorter than 10 minutes require
* `deleteUnverified: true` and are only safe when no other process can write
* to the dataset.
* @example * @example
* // Delete all versions older than 1 day * // Delete all versions older than 1 day
* const olderThan = new Date(); * const olderThan = new Date();
* olderThan.setDate(olderThan.getDate() - 1)); * olderThan.setDate(olderThan.getDate() - 1));
* tbl.optimize({cleanupOlderThan: olderThan}); * tbl.optimize({cleanupOlderThan: olderThan});
* *
* // Delete all versions except the current version * // With exclusive access, delete all versions except the current version
* tbl.optimize({cleanupOlderThan: new Date()}); * tbl.optimize({cleanupOlderThan: new Date(), deleteUnverified: true});
*/ */
cleanupOlderThan: Date; cleanupOlderThan: Date;
/** /**
@@ -744,6 +747,9 @@ export abstract class Table {
* optimize should be run frequently. A good rule of thumb is to run optimize if * optimize should be run frequently. A good rule of thumb is to run optimize if
* you have added or modified 100,000 or more records or run more than 20 data * you have added or modified 100,000 or more records or run more than 20 data
* modification operations. * modification operations.
*
* Cleanup retention must exceed the longest expected write. Retention shorter
* than 10 minutes requires `deleteUnverified: true` and exclusive write access.
*/ */
abstract optimize(options?: Partial<OptimizeOptions>): Promise<OptimizeStats>; abstract optimize(options?: Partial<OptimizeOptions>): Promise<OptimizeStats>;
/** List all indices that have been created with {@link Table.createIndex} */ /** List all indices that have been created with {@link Table.createIndex} */
+2
View File
@@ -6,6 +6,7 @@ use std::collections::HashMap;
use chrono::{DateTime, Utc}; use chrono::{DateTime, Utc};
use lancedb::ipc::{ipc_file_to_batches, ipc_file_to_schema}; use lancedb::ipc::{ipc_file_to_batches, ipc_file_to_schema};
use lancedb::table::optimize::validate_cleanup_options;
use lancedb::table::{ use lancedb::table::{
AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration, AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration,
FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken, FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken,
@@ -559,6 +560,7 @@ impl Table {
} else { } else {
None None
}; };
validate_cleanup_options(older_than, delete_unverified).default_error()?;
let compaction_stats = inner let compaction_stats = inner
.optimize(OptimizeAction::Compact { .optimize(OptimizeAction::Compact {
+39 -11
View File
@@ -117,6 +117,23 @@ _MODEL_BACKED_TOKENIZER_ERRORS = (
"Failed to initialize default tokenizer", "Failed to initialize default tokenizer",
) )
_MIN_SAFE_CLEANUP_AGE = timedelta(minutes=10)
def _validate_cleanup_options(
older_than: Optional[timedelta], delete_unverified: bool
) -> None:
if (
older_than is not None
and older_than < _MIN_SAFE_CLEANUP_AGE
and not delete_unverified
):
raise ValueError(
"cleanup age must be at least 10 minutes unless delete_unverified is "
"true; short cleanup windows can remove a manifest still needed by an "
"in-progress write"
)
def _add_unique_note(exception: BaseException, note: str) -> None: def _add_unique_note(exception: BaseException, note: str) -> None:
existing_notes = getattr(exception, "__notes__", ()) or () existing_notes = getattr(exception, "__notes__", ()) or ()
@@ -1767,7 +1784,9 @@ class Table(ABC):
---------- ----------
older_than: timedelta, default None older_than: timedelta, default None
The minimum age of the version to delete. If None, then this defaults The minimum age of the version to delete. If None, then this defaults
to two weeks. to two weeks. This must be longer than the longest expected write.
Values shorter than 10 minutes require `delete_unverified=True` and
are only safe when no other process can write to the dataset.
delete_unverified: bool, default False delete_unverified: bool, default False
Because they may be part of an in-progress transaction, files newer Because they may be part of an in-progress transaction, files newer
than 7 days old are not deleted by default. If you are sure that than 7 days old are not deleted by default. If you are sure that
@@ -1835,9 +1854,11 @@ class Table(ABC):
Parameters Parameters
---------- ----------
cleanup_older_than: timedelta, optional default 7 days cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. Set All files belonging to versions older than this will be removed. The
to 0 days to remove all versions except the latest. The latest version latest version is never removed. This must be longer than the longest
is never removed. expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
delete_unverified: bool, default False delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not in-progress operation (e.g. appending new data) and these files will not
@@ -3786,7 +3807,9 @@ class LanceTable(Table):
---------- ----------
older_than: timedelta, default None older_than: timedelta, default None
The minimum age of the version to delete. If None, then this defaults The minimum age of the version to delete. If None, then this defaults
to two weeks. to two weeks. This must be longer than the longest expected write.
Values shorter than 10 minutes require `delete_unverified=True` and
are only safe when no other process can write to the dataset.
delete_unverified: bool, default False delete_unverified: bool, default False
Because they may be part of an in-progress transaction, files newer Because they may be part of an in-progress transaction, files newer
than 7 days old are not deleted by default. If you are sure that than 7 days old are not deleted by default. If you are sure that
@@ -3799,6 +3822,7 @@ class LanceTable(Table):
The stats of the cleanup operation, including how many bytes were The stats of the cleanup operation, including how many bytes were
freed. freed.
""" """
_validate_cleanup_options(older_than, delete_unverified)
return self.to_lance().cleanup_old_versions( return self.to_lance().cleanup_old_versions(
older_than, delete_unverified=delete_unverified older_than, delete_unverified=delete_unverified
) )
@@ -3844,9 +3868,11 @@ class LanceTable(Table):
Parameters Parameters
---------- ----------
cleanup_older_than: timedelta, optional default 7 days cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. Set All files belonging to versions older than this will be removed. The
to 0 days to remove all versions except the latest. The latest version latest version is never removed. This must be longer than the longest
is never removed. expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
delete_unverified: bool, default False delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not in-progress operation (e.g. appending new data) and these files will not
@@ -6038,9 +6064,11 @@ class AsyncTable:
Parameters Parameters
---------- ----------
cleanup_older_than: timedelta, optional default 7 days cleanup_older_than: timedelta, optional default 7 days
All files belonging to versions older than this will be removed. Set All files belonging to versions older than this will be removed. The
to 0 days to remove all versions except the latest. The latest version latest version is never removed. This must be longer than the longest
is never removed. expected write. Values shorter than 10 minutes require
`delete_unverified=True` and are only safe when no other process can
write to the dataset.
delete_unverified: bool, default False delete_unverified: bool, default False
Files leftover from a failed transaction may appear to be part of an Files leftover from a failed transaction may appear to be part of an
in-progress operation (e.g. appending new data) and these files will not in-progress operation (e.g. appending new data) and these files will not
+11 -1
View File
@@ -2980,6 +2980,9 @@ def test_compact_cleanup(tmp_db: DBConnection):
stats = table.cleanup_old_versions() stats = table.cleanup_old_versions()
assert stats.bytes_removed == 0 assert stats.bytes_removed == 0
with pytest.raises(ValueError, match="at least 10 minutes"):
table.cleanup_old_versions(older_than=timedelta(0))
stats = table.cleanup_old_versions(older_than=timedelta(0), delete_unverified=True) stats = table.cleanup_old_versions(older_than=timedelta(0), delete_unverified=True)
assert stats.bytes_removed > 0 assert stats.bytes_removed > 0
assert table.version == 4 assert table.version == 4
@@ -3396,7 +3399,14 @@ async def test_optimize(mem_db_async: AsyncConnection):
assert stats.prune.bytes_removed == 0 assert stats.prune.bytes_removed == 0
assert stats.prune.old_versions_removed == 0 assert stats.prune.old_versions_removed == 0
stats = await table.optimize(cleanup_older_than=timedelta(seconds=0)) version_before_rejected_cleanup = await table.version()
with pytest.raises(ValueError, match="at least 10 minutes"):
await table.optimize(cleanup_older_than=timedelta(seconds=0))
assert await table.version() == version_before_rejected_cleanup
stats = await table.optimize(
cleanup_older_than=timedelta(seconds=0), delete_unverified=True
)
assert stats.prune.bytes_removed > 0 assert stats.prune.bytes_removed > 0
assert stats.prune.old_versions_removed == 3 assert stats.prune.old_versions_removed == 3
+2
View File
@@ -19,6 +19,7 @@ use arrow::{
}; };
use lancedb::blob::{BlobFile, BlobRangeRequest}; use lancedb::blob::{BlobFile, BlobRangeRequest};
use lancedb::index::scalar::FtsIndexBuilder; use lancedb::index::scalar::FtsIndexBuilder;
use lancedb::table::optimize::validate_cleanup_options;
use lancedb::table::{ use lancedb::table::{
AddDataMode, ColumnAlteration, Duration, FieldMetadataUpdate, FtsToken as LanceDbFtsToken, AddDataMode, ColumnAlteration, Duration, FieldMetadataUpdate, FtsToken as LanceDbFtsToken,
NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable, NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable,
@@ -1217,6 +1218,7 @@ impl Table {
} else { } else {
None None
}; };
validate_cleanup_options(older_than, delete_unverified).infer_error()?;
future_into_py(self_.py(), async move { future_into_py(self_.py(), async move {
let compaction_stats = inner let compaction_stats = inner
.optimize(OptimizeAction::Compact { .optimize(OptimizeAction::Compact {
+82 -3
View File
@@ -18,7 +18,33 @@ pub use chrono::Duration;
pub use lance::dataset::optimize::CompactionOptions; pub use lance::dataset::optimize::CompactionOptions;
use super::NativeTable; use super::NativeTable;
use crate::error::Result; use crate::error::{Error, Result};
const MIN_SAFE_CLEANUP_AGE_MINUTES: i64 = 10;
/// Validate version-cleanup options before starting an optimization operation.
///
/// This is public for the language bindings, which run compaction and cleanup as
/// separate operations and must reject unsafe cleanup options before compaction
/// changes the table.
#[doc(hidden)]
pub fn validate_cleanup_options(
older_than: Option<Duration>,
delete_unverified: Option<bool>,
) -> Result<()> {
let minimum_age =
Duration::try_minutes(MIN_SAFE_CLEANUP_AGE_MINUTES).expect("minimum cleanup age is valid");
if older_than.is_some_and(|age| age < minimum_age) && delete_unverified != Some(true) {
return Err(Error::InvalidInput {
message: format!(
"cleanup age must be at least {MIN_SAFE_CLEANUP_AGE_MINUTES} minutes unless \
delete_unverified is true; short cleanup windows can remove a manifest still \
needed by an in-progress write"
),
});
}
Ok(())
}
/// Optimize the dataset. /// Optimize the dataset.
/// ///
@@ -60,7 +86,9 @@ pub enum OptimizeAction {
/// ///
/// Once a version is pruned it can no longer be checked out. /// Once a version is pruned it can no longer be checked out.
Prune { Prune {
/// The duration of time to keep versions of the dataset. /// The duration of time to keep versions of the dataset. This should be longer than the
/// longest expected write. Values shorter than 10 minutes require `delete_unverified` to
/// be true and are only safe when no other process can write to the dataset.
older_than: Option<Duration>, older_than: Option<Duration>,
/// Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default. /// Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default.
/// If you are sure that there are no in-progress transactions, then you can set this to True to delete all files older than `older_than`. /// If you are sure that there are no in-progress transactions, then you can set this to True to delete all files older than `older_than`.
@@ -164,6 +192,15 @@ pub(crate) async fn execute_optimize(
table: &NativeTable, table: &NativeTable,
action: OptimizeAction, action: OptimizeAction,
) -> Result<OptimizeStats> { ) -> Result<OptimizeStats> {
if let OptimizeAction::Prune {
older_than,
delete_unverified,
..
} = &action
{
validate_cleanup_options(*older_than, *delete_unverified)?;
}
let mut stats = OptimizeStats { let mut stats = OptimizeStats {
compaction: None, compaction: None,
prune: None, prune: None,
@@ -222,7 +259,7 @@ mod tests {
use crate::connect; use crate::connect;
use crate::index::{Index, scalar::BTreeIndexBuilder}; use crate::index::{Index, scalar::BTreeIndexBuilder};
use crate::query::ExecutableQuery; use crate::query::ExecutableQuery;
use crate::table::{CompactionOptions, OptimizeAction, OptimizeStats}; use crate::table::{CompactionOptions, Duration, OptimizeAction, OptimizeStats};
use futures::TryStreamExt; use futures::TryStreamExt;
#[tokio::test] #[tokio::test]
@@ -380,6 +417,48 @@ mod tests {
assert_eq!(all_values, expected); assert_eq!(all_values, expected);
} }
#[tokio::test]
async fn test_optimize_rejects_unsafe_cleanup_age() {
let conn = connect("memory://").execute().await.unwrap();
let batch = RecordBatch::try_new(
Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, false)])),
vec![Arc::new(Int32Array::from_iter_values(0..10))],
)
.unwrap();
let table = conn
.create_table("test_unsafe_prune", batch.clone())
.execute()
.await
.unwrap();
table.add(batch).execute().await.unwrap();
let versions_before = table
.list_versions()
.await
.unwrap()
.into_iter()
.map(|version| version.version)
.collect::<Vec<_>>();
let err = table
.optimize(OptimizeAction::Prune {
older_than: Some(Duration::zero()),
delete_unverified: Some(false),
error_if_tagged_old_versions: None,
})
.await
.unwrap_err();
assert!(err.to_string().contains("at least 10 minutes"));
let versions_after = table
.list_versions()
.await
.unwrap()
.into_iter()
.map(|version| version.version)
.collect::<Vec<_>>();
assert_eq!(versions_after, versions_before);
}
#[tokio::test] #[tokio::test]
async fn test_optimize_index() { async fn test_optimize_index() {
let conn = connect("memory://").execute().await.unwrap(); let conn = connect("memory://").execute().await.unwrap();