mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-11 07:42:26 +00:00
Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
af17703268 |
@@ -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()});
|
||||
```
|
||||
|
||||
|
||||
@@ -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
@@ -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
@@ -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(),
|
||||
|
||||
@@ -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())
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user