diff --git a/docs/src/js/interfaces/OptimizeOptions.md b/docs/src/js/interfaces/OptimizeOptions.md index 700632342..110eb3813 100644 --- a/docs/src/js/interfaces/OptimizeOptions.md +++ b/docs/src/js/interfaces/OptimizeOptions.md @@ -26,7 +26,8 @@ const olderThan = new Date(); olderThan.setDate(olderThan.getDate() - 1)); tbl.optimize({cleanupOlderThan: olderThan}); -// Delete all versions except the current version +// Delete versions committed before this point. Versions created by the +// optimize call itself are newer than the cutoff and will be retained. tbl.optimize({cleanupOlderThan: new Date()}); ``` diff --git a/nodejs/__test__/table.test.ts b/nodejs/__test__/table.test.ts index e8b97bf77..43a027642 100644 --- a/nodejs/__test__/table.test.ts +++ b/nodejs/__test__/table.test.ts @@ -52,6 +52,7 @@ import { Operator, instanceOfFullTextQuery, } from "../lancedb/query"; +import { LocalTable } from "../lancedb/table"; describe.each([arrow15, arrow16, arrow17, arrow18])( "Given a table", @@ -2518,7 +2519,7 @@ describe("when optimizing a dataset", () => { it("cleanups old versions", async () => { const stats = await table.optimize({ cleanupOlderThan: new Date() }); expect(stats.prune.bytesRemoved).toBeGreaterThan(0); - expect(stats.prune.oldVersionsRemoved).toBe(3); + expect(stats.prune.oldVersionsRemoved).toBe(2); }); it("delete unverified", async () => { @@ -2539,6 +2540,24 @@ describe("when optimizing a dataset", () => { }); }); +it("passes cleanupOlderThan to the native binding as an absolute timestamp", async () => { + const optimize = jest.fn().mockResolvedValue({ + compaction: { + filesAdded: 0, + filesRemoved: 0, + fragmentsAdded: 0, + fragmentsRemoved: 0, + }, + prune: { bytesRemoved: 0, oldVersionsRemoved: 0 }, + }); + const table = new LocalTable({ optimize } as never); + const cutoff = new Date("2020-01-02T03:04:05.678Z"); + + await table.optimize({ cleanupOlderThan: cutoff, deleteUnverified: true }); + + expect(optimize).toHaveBeenCalledWith(cutoff.getTime(), true); +}); + describe.each([arrow15, arrow16, arrow17, arrow18])( "when optimizing a dataset", // biome-ignore lint/suspicious/noExplicitAny: diff --git a/nodejs/lancedb/table.ts b/nodejs/lancedb/table.ts index d1fb8acd8..02d4fee83 100644 --- a/nodejs/lancedb/table.ts +++ b/nodejs/lancedb/table.ts @@ -147,7 +147,8 @@ export interface OptimizeOptions { * olderThan.setDate(olderThan.getDate() - 1)); * tbl.optimize({cleanupOlderThan: olderThan}); * - * // Delete all versions except the current version + * // Delete versions committed before this point. Versions created by the + * // optimize call itself are newer than the cutoff and will be retained. * tbl.optimize({cleanupOlderThan: new Date()}); */ cleanupOlderThan: Date; @@ -1445,16 +1446,8 @@ export class LocalTable extends Table { } async optimize(options?: Partial): Promise { - let cleanupOlderThanMs; - if ( - options?.cleanupOlderThan !== undefined && - options?.cleanupOlderThan !== null - ) { - cleanupOlderThanMs = - new Date().getTime() - options.cleanupOlderThan.getTime(); - } return await this.inner.optimize( - cleanupOlderThanMs, + options?.cleanupOlderThan?.getTime(), options?.deleteUnverified, ); } diff --git a/nodejs/src/table.rs b/nodejs/src/table.rs index db74d38fa..12ade78e0 100644 --- a/nodejs/src/table.rs +++ b/nodejs/src/table.rs @@ -7,7 +7,7 @@ use chrono::{DateTime, Utc}; use lancedb::ipc::{ipc_file_to_batches, ipc_file_to_schema}; use lancedb::table::{ - AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration, + AddDataMode, ColumnAlteration as LanceColumnAlteration, FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken, NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable, }; @@ -638,22 +638,20 @@ impl Table { #[napi(catch_unwind)] pub async fn optimize( &self, - older_than_ms: Option, + before_timestamp_ms: Option, delete_unverified: Option, ) -> napi::Result { let inner = self.inner_ref()?; - let older_than = if let Some(ms) = older_than_ms { - if ms == i64::MIN { - return Err(napi::Error::from_reason(format!( - "older_than_ms can not be {}", - i32::MIN, - ))); - } - Duration::try_milliseconds(ms) - } else { - None - }; + let before_timestamp = before_timestamp_ms + .map(|ms| { + DateTime::from_timestamp_millis(ms).ok_or_else(|| { + napi::Error::from_reason(format!( + "cleanupOlderThan timestamp is out of range: {ms}" + )) + }) + }) + .transpose()?; let compaction_stats = inner .optimize(OptimizeAction::Compact { @@ -664,16 +662,22 @@ impl Table { .default_error()? .compaction .unwrap(); - let prune_stats = inner - .optimize(OptimizeAction::Prune { - older_than, - delete_unverified, - error_if_tagged_old_versions: None, - }) - .await - .default_error()? - .prune - .unwrap(); + let prune_stats = if let Some(before_timestamp) = before_timestamp { + inner + .optimize_prune_before(before_timestamp, delete_unverified, None) + .await + } else { + inner + .optimize(OptimizeAction::Prune { + older_than: None, + delete_unverified, + error_if_tagged_old_versions: None, + }) + .await + } + .default_error()? + .prune + .unwrap(); inner .optimize(lancedb::table::OptimizeAction::Index( OptimizeOptions::default(), diff --git a/rust/lancedb/src/table.rs b/rust/lancedb/src/table.rs index 7f139c4cb..5adc3bf80 100644 --- a/rust/lancedb/src/table.rs +++ b/rust/lancedb/src/table.rs @@ -1739,6 +1739,33 @@ impl Table { self.inner.optimize(action).await } + /// Prune versions committed before an absolute timestamp. + /// + /// This is an internal entry point for language bindings whose public API + /// accepts an absolute cleanup cutoff. + #[doc(hidden)] + pub async fn optimize_prune_before( + &self, + before_timestamp: chrono::DateTime, + delete_unverified: Option, + error_if_tagged_old_versions: Option, + ) -> Result { + let native = self.as_native().ok_or_else(|| Error::NotSupported { + message: "optimize is not supported on LanceDB cloud.".into(), + })?; + let prune = optimize::cleanup_old_versions_before( + native, + before_timestamp, + delete_unverified, + error_if_tagged_old_versions, + ) + .await?; + Ok(OptimizeStats { + compaction: None, + prune: Some(prune), + }) + } + /// Add new columns to the table, providing values to fill in. pub fn add_columns(&self) -> AddColumnsBuilder { AddColumnsBuilder::new(self.inner.clone()) diff --git a/rust/lancedb/src/table/optimize.rs b/rust/lancedb/src/table/optimize.rs index 4ad58cffb..e3c171f27 100644 --- a/rust/lancedb/src/table/optimize.rs +++ b/rust/lancedb/src/table/optimize.rs @@ -8,7 +8,8 @@ use std::sync::Arc; -use lance::dataset::cleanup::RemovalStats; +use chrono::{DateTime, Utc}; +use lance::dataset::cleanup::{CleanupPolicyBuilder, RemovalStats}; use lance::dataset::optimize::{CompactionMetrics, IndexRemapperOptions, compact_files}; use lance::index::DatasetIndexExt; use lance_index::optimize::OptimizeOptions; @@ -139,6 +140,25 @@ pub(crate) async fn cleanup_old_versions( .await?) } +/// Remove dataset versions committed before an absolute timestamp. +pub(crate) async fn cleanup_old_versions_before( + table: &NativeTable, + before_timestamp: DateTime, + delete_unverified: Option, + error_if_tagged_old_versions: Option, +) -> Result { + table.dataset.ensure_mutable()?; + let dataset = table.dataset.get().await?; + let mut policy = CleanupPolicyBuilder::default().before_timestamp(before_timestamp); + if let Some(delete_unverified) = delete_unverified { + policy = policy.delete_unverified(delete_unverified); + } + if let Some(error_if_tagged_old_versions) = error_if_tagged_old_versions { + policy = policy.error_if_tagged_old_versions(error_if_tagged_old_versions); + } + Ok(dataset.cleanup_with_policy(policy.build()).await?) +} + /// Compact files in the dataset. /// /// This can be run after making several small appends to optimize the table