Compare commits

...
Author SHA1 Message Date
Gatefixer af17703268 fix(node): preserve optimize cleanup timestamp 2026-09-10 18:34:22 +00:00
6 changed files with 100 additions and 36 deletions
+2 -1
View File
@@ -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()});
```
+20 -1
View File
@@ -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: <explanation>
+3 -10
View File
@@ -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<OptimizeOptions>): Promise<OptimizeStats> {
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,
);
}
+27 -23
View File
@@ -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<i64>,
before_timestamp_ms: Option<i64>,
delete_unverified: Option<bool>,
) -> napi::Result<OptimizeStats> {
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(),
+27
View File
@@ -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<chrono::Utc>,
delete_unverified: Option<bool>,
error_if_tagged_old_versions: Option<bool>,
) -> Result<OptimizeStats> {
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())
+21 -1
View File
@@ -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<Utc>,
delete_unverified: Option<bool>,
error_if_tagged_old_versions: Option<bool>,
) -> Result<RemovalStats> {
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