Compare commits

..
Author SHA1 Message Date
Gatefixer af17703268 fix(node): preserve optimize cleanup timestamp 2026-09-10 18:34:22 +00:00
21 changed files with 374 additions and 1564 deletions
Generated
+241 -252
View File
File diff suppressed because it is too large Load Diff
+15 -15
View File
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
lance = { "version" = "=12.0.0-beta.17", default-features = false, "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=12.0.0-beta.17", default-features = false, "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=12.0.0-beta.17", default-features = false, "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=12.0.0-beta.17", "tag" = "v12.0.0-beta.17", "git" = "https://github.com/lance-format/lance.git" }
lance = { "version" = "=12.0.0-beta.16", default-features = false, "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=12.0.0-beta.16", default-features = false, "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=12.0.0-beta.16", default-features = false, "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
lancedb = { path = "rust/lancedb", default-features = false }
ahash = "0.8"
# Note that this one does not include pyarrow
@@ -60,7 +60,7 @@ log = "0.4"
metrics = "0.24"
metrics-util = "0.19"
moka = { version = "0.12", features = ["future"] }
object_store = "0.14.1"
object_store = "0.13.2"
pin-project = "1.0.7"
rand = "0.9"
snafu = "0.8"
-62
View File
@@ -1,62 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / BlobFile
# Class: BlobFile
A lazy handle to blob bytes. Create one with [Table.fetchBlobFiles](Table.md#fetchblobfiles).
## Methods
### read()
```ts
read(): Promise<Buffer>
```
Reads from the cursor to the end and advances the cursor.
A second call returns an empty buffer. [BlobFile.readRange](BlobFile.md#readrange) does
not move the cursor.
#### Returns
`Promise`&lt;`Buffer`&gt;
***
### readRange()
```ts
readRange(start, end): Promise<Buffer>
```
Reads the half-open byte range `[start, end)`.
Fails when `end` is past the blob size. Does not move the cursor.
#### Parameters
* **start**: `bigint`
* **end**: `bigint`
#### Returns
`Promise`&lt;`Buffer`&gt;
***
### size()
```ts
size(): bigint
```
Returns the blob size in bytes.
#### Returns
`bigint`
-62
View File
@@ -137,20 +137,6 @@ containing the new version number of the table after altering the columns.
***
### blobColumns()
```ts
abstract blobColumns(): Promise<string[]>
```
Blob v2 columns, including nested dotted paths.
#### Returns
`Promise`&lt;`string`[]&gt;
***
### branches()
```ts
@@ -513,54 +499,6 @@ Drop an index from the table.
***
### fetchBlobFiles()
```ts
abstract fetchBlobFiles(column, rowIds): Promise<(null | BlobFile)[]>
```
Opens lazy blob handles for `column` at the given row IDs using the
table's current checkout.
Preserves input order, duplicates, and nulls. Use this for large payloads.
See [Table.fetchBlobs](Table.md#fetchblobs) for row-ID validity across versions.
#### Parameters
* **column**: `string`
* **rowIds**: readonly (`number` \| `bigint`)[]
#### Returns
`Promise`&lt;(`null` \| [`BlobFile`](BlobFile.md))[]&gt;
***
### fetchBlobs()
```ts
abstract fetchBlobs(column, rowIds): Promise<(null | Buffer)[]>
```
Bytes for `column` at row IDs from [Query.withRowId](Query.md#withrowid).
Reads the table's current checkout. IDs from another version can fail after
compaction unless stable row ids are enabled. Results keep input order and
duplicates. Null blobs are `null`. Empty blobs are empty buffers.
#### Parameters
* **column**: `string`
* **rowIds**: readonly (`number` \| `bigint`)[]
#### Returns
`Promise`&lt;(`null` \| `Buffer`)[]&gt;
***
### flushLsm()
```ts
-55
View File
@@ -1,55 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / blob
# Function: blob()
```ts
function blob(name, options): Field
```
Declares a `lance.blob.v2` column.
Query results are descriptors, not payload bytes. Use [Table.fetchBlobs](../classes/Table.md#fetchblobs)
or [Table.fetchBlobFiles](../classes/Table.md#fetchblobfiles) to read bytes.
## Parameters
* **name**: `string`
* **options**: [`BlobOptions`](../type-aliases/BlobOptions.md) = `{}`
## Returns
`Field`
## Example
```ts
import { readFile } from "node:fs/promises";
import { Field, Int64, Schema } from "apache-arrow";
import { blob, connect } from "@lancedb/lancedb";
const db = await connect("./data");
const video = await readFile("clip.mp4");
const table = await db.createTable(
"videos",
[{ id: 1n, video }],
{
schema: new Schema([
new Field("id", new Int64()),
blob("video"),
]),
},
);
const rows = await table.query().select(["id"]).withRowId().toArray();
const rowIds = rows.map((row) => row._rowid as bigint);
const bytes = await table.fetchBlobs("video", rowIds);
const [handle] = await table.fetchBlobFiles("video", rowIds);
const size = handle!.size();
const header = await handle!.readRange(0n, size < 65536n ? size : 65536n);
```
-22
View File
@@ -1,22 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / isBlobField
# Function: isBlobField()
```ts
function isBlobField(field): boolean
```
Checks for the `lance.blob.v2` extension marker. Does not validate the
field's storage type.
## Parameters
* **field**: `Field`&lt;`any`&gt;
## Returns
`boolean`
-4
View File
@@ -19,7 +19,6 @@
## Classes
- [AutoQuery](classes/AutoQuery.md)
- [BlobFile](classes/BlobFile.md)
- [BooleanQuery](classes/BooleanQuery.md)
- [BoostQuery](classes/BoostQuery.md)
- [BranchContents](classes/BranchContents.md)
@@ -144,7 +143,6 @@
- [AnalyzePlanDistributedMetrics](type-aliases/AnalyzePlanDistributedMetrics.md)
- [BaseTokenizer](type-aliases/BaseTokenizer.md)
- [BlobOptions](type-aliases/BlobOptions.md)
- [Data](type-aliases/Data.md)
- [DataLike](type-aliases/DataLike.md)
- [FieldLike](type-aliases/FieldLike.md)
@@ -160,11 +158,9 @@
## Functions
- [RecordBatchIterator](functions/RecordBatchIterator.md)
- [blob](functions/blob.md)
- [connect](functions/connect.md)
- [connectNamespace](functions/connectNamespace.md)
- [instrumentLanceDbMetrics](functions/instrumentLanceDbMetrics.md)
- [isBlobField](functions/isBlobField.md)
- [makeArrowTable](functions/makeArrowTable.md)
- [packBits](functions/packBits.md)
- [permutationBuilder](functions/permutationBuilder.md)
+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()});
```
-48
View File
@@ -1,48 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / BlobOptions
# Type Alias: BlobOptions
```ts
type BlobOptions: object;
```
## Type declaration
### dedicatedSizeThreshold?
```ts
optional dedicatedSizeThreshold: number;
```
Max payload bytes stored in a packed sidecar before a dedicated file. Must
be a positive safe integer.
### inlineSizeThreshold?
```ts
optional inlineSizeThreshold: number;
```
Max payload bytes kept inline in the data file. Zero is allowed. Must be a
safe integer.
### nullable?
```ts
optional nullable: boolean;
```
Defaults to true.
### packFileSizeThreshold?
```ts
optional packFileSizeThreshold: number;
```
Max bytes in one packed sidecar before starting another. Must be a positive
safe integer.
+1 -1
View File
@@ -28,7 +28,7 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<arrow.version>15.0.0</arrow.version>
<lance-core.version>12.0.0-beta.17</lance-core.version>
<lance-core.version>12.0.0-beta.16</lance-core.version>
<spotless.skip>false</spotless.skip>
<spotless.version>2.30.0</spotless.version>
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
-185
View File
@@ -1,185 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
import { Field, Int64, List, Schema, Struct, Utf8 } from "apache-arrow";
import { makeArrowTable } from "../lancedb/arrow";
import { BlobFile, blob, coerceBlobValue, isBlobField } from "../lancedb/blob";
describe("blob()", () => {
it("marks the field as lance.blob.v2", () => {
const field = blob("image", { nullable: false });
expect(field.nullable).toBe(false);
expect(isBlobField(field)).toBe(true);
expect(field.metadata.get("ARROW:extension:name")).toBe("lance.blob.v2");
});
it("writes encoding thresholds as field metadata", () => {
const field = blob("video", {
inlineSizeThreshold: 1024,
dedicatedSizeThreshold: 2 * 1024 * 1024,
packFileSizeThreshold: 64 * 1024 * 1024,
});
expect(
field.metadata.get("lance-encoding:blob-inline-size-threshold"),
).toBe("1024");
expect(
field.metadata.get("lance-encoding:blob-dedicated-size-threshold"),
).toBe(String(2 * 1024 * 1024));
expect(
field.metadata.get("lance-encoding:blob-pack-file-size-threshold"),
).toBe(String(64 * 1024 * 1024));
});
it("rejects invalid thresholds", () => {
expect(() => blob("image", { inlineSizeThreshold: -1 })).toThrow(
/inlineSizeThreshold must be non-negative/,
);
expect(() => blob("image", { dedicatedSizeThreshold: 0 })).toThrow(
/dedicatedSizeThreshold must be positive/,
);
expect(() => blob("image", { packFileSizeThreshold: 1.5 })).toThrow(
/packFileSizeThreshold must be a safe integer/,
);
expect(() =>
blob("image", { dedicatedSizeThreshold: Number.MAX_SAFE_INTEGER + 1 }),
).toThrow(/dedicatedSizeThreshold must be a safe integer/);
});
});
describe("coerceBlobValue", () => {
it.each([
["Buffer", Buffer.from("x"), { data: Buffer.from("x"), uri: null }],
[
"Uint8Array",
new Uint8Array([120]),
{ data: new Uint8Array([120]), uri: null },
],
["URI string", "s3://bucket/key", { data: null, uri: "s3://bucket/key" }],
[
"data struct",
{ data: Buffer.from("y") },
{ data: Buffer.from("y"), uri: null },
],
[
"uri struct",
{ uri: "s3://bucket/key" },
{ data: null, uri: "s3://bucket/key" },
],
["null", null, null],
])("accepts %s", (_name, input, expected) => {
expect(coerceBlobValue(input)).toEqual(expected);
});
it.each([
["empty URI", "", /uri cannot be empty/],
["object without data or uri", { position: 0 }, /data' or 'uri/],
[
"Int16Array",
new Int16Array([1]),
/Blob data must be Buffer or Uint8Array/,
],
[
"both data and uri",
{ data: Buffer.from("y"), uri: "s3://bucket/key" },
/exactly one of 'data' or 'uri'/,
],
[
"neither data nor uri",
{ data: null, uri: null },
/exactly one of 'data' or 'uri'/,
],
])("rejects %s", (_name, input, message) => {
expect(() => coerceBlobValue(input)).toThrow(message);
});
});
describe("BlobFile", () => {
it("rejects constructing BlobFile without a native handle", () => {
expect(() => new (BlobFile as unknown as { new (): BlobFile })()).toThrow(
/fetchBlobFiles/,
);
});
});
describe("makeArrowTable blob columns", () => {
it("coerces Buffer input onto a blob field", () => {
const schema = new Schema([
new Field("id", new Int64(), true),
blob("image"),
]);
const table = makeArrowTable([{ id: 1n, image: Buffer.from("hello") }], {
schema,
});
expect(isBlobField(table.schema.fields[1])).toBe(true);
const image = table.getChild("image")!;
expect(image.nullCount).toBe(0);
expect(image.getChild("uri")!.get(0)).toBeNull();
expect(image.getChild("data")!.nullCount).toBe(0);
expect(Buffer.from(image.getChild("data")!.get(0)!).toString()).toBe(
"hello",
);
});
it("coerces Buffer elements inside a list and keeps null slots", () => {
const schema = new Schema([
new Field("id", new Int64(), true),
new Field("images", new List(blob("image")), true),
]);
const table = makeArrowTable(
[
{ id: 1n, images: [Buffer.from("a"), Buffer.from("bb")] },
{ id: 2n, images: null },
{ id: 3n, images: [Buffer.from("c"), null] },
{ id: 4n, images: [] },
],
{ schema },
);
const images = table.getChild("images")!;
expect(images.nullCount).toBe(1);
const rows = images.toArray();
expect(rows[1]).toBeNull();
expect(Array.from(rows[3] as Iterable<unknown>)).toHaveLength(0);
const first = Array.from(rows[0] as Iterable<{ data: Uint8Array | null }>);
expect(Buffer.from(first[0].data!).toString()).toBe("a");
expect(Buffer.from(first[1].data!).toString()).toBe("bb");
const third = Array.from(
rows[2] as Iterable<{ data: Uint8Array | null } | null>,
);
expect(Buffer.from(third[0]!.data!).toString()).toBe("c");
expect(third[1]).toBeNull();
});
it("coerces Buffer fields inside list structs", () => {
const schema = new Schema([
new Field("id", new Int64(), true),
new Field(
"items",
new List(
new Field(
"item",
new Struct([new Field("name", new Utf8(), true), blob("image")]),
true,
),
),
true,
),
]);
const table = makeArrowTable(
[
{
id: 1n,
items: [{ name: "one", image: Buffer.from("alpha") }],
},
],
{ schema },
);
const items = Array.from(
table.getChild("items")!.toArray()[0] as Iterable<{
name: string;
image: { data: Uint8Array | null };
}>,
);
expect(items[0].name).toBe("one");
expect(Buffer.from(items[0].image.data!).toString()).toBe("alpha");
});
});
+20 -272
View File
@@ -18,7 +18,6 @@ import {
Query,
Table,
VectorQuery,
blob,
connect,
tokenize,
} from "../lancedb";
@@ -53,6 +52,7 @@ import {
Operator,
instanceOfFullTextQuery,
} from "../lancedb/query";
import { LocalTable } from "../lancedb/table";
describe.each([arrow15, arrow16, arrow17, arrow18])(
"Given a table",
@@ -2402,276 +2402,6 @@ describe("when dealing with versioning", () => {
});
});
describe("when dealing with blob columns", () => {
let tmpDir: tmp.DirResult;
beforeEach(() => {
tmpDir = tmp.dirSync({ unsafeCleanup: true });
});
afterEach(() => {
tmpDir.removeCallback();
});
it("discovers blob columns", async () => {
const { table } = await openBlobTable();
expect(await table.blobColumns()).toEqual(["image"]);
});
it("preserves order, duplicates, and nulls", async () => {
const { table, rowIds } = await openBlobTable();
const [alphaId, betaId, nullId] = rowIds;
const bytes = await table.fetchBlobs("image", [
betaId,
alphaId,
betaId,
nullId,
]);
expect(bytes.map((b) => (b == null ? null : b.toString()))).toEqual([
"beta",
"alpha",
"beta",
null,
]);
const files = await table.fetchBlobFiles("image", [
betaId,
nullId,
alphaId,
]);
expect(files.map((f) => f == null)).toEqual([false, true, false]);
});
it("reads full blob contents", async () => {
const { table, rowIds, alpha, beta } = await openBlobTable();
const bytes = await table.fetchBlobs("image", rowIds);
expect(bytes[0]!.equals(alpha)).toBe(true);
expect(bytes[1]!.equals(beta)).toBe(true);
const files = await table.fetchBlobFiles("image", rowIds);
expect(files[0]!.size()).toBe(BigInt(alpha.length));
expect(Buffer.from(await files[0]!.read()).toString()).toBe("alpha");
expect(Buffer.from(await files[1]!.read()).toString()).toBe("beta");
});
it("reads a half-open range", async () => {
const { table, rowIds } = await openBlobTable();
const files = await table.fetchBlobFiles("image", rowIds);
expect(Buffer.from(await files[0]!.readRange(0n, 2n)).toString()).toBe(
"al",
);
});
it("readRange does not move the cursor", async () => {
const { table, rowIds, alpha } = await openBlobTable();
const [handle] = await table.fetchBlobFiles("image", rowIds);
expect((await handle!.readRange(1n, 3n)).toString()).toBe("lp");
expect(await handle!.read()).toEqual(alpha);
expect(await handle!.read()).toEqual(Buffer.alloc(0));
});
it("fails when readRange end is past the blob size", async () => {
const { table, rowIds, alpha } = await openBlobTable();
const files = await table.fetchBlobFiles("image", rowIds);
await expect(
files[0]!.readRange(0n, BigInt(alpha.length + 1)),
).rejects.toThrow(/exceeds blob size/);
});
it("rejects fetchBlobs on a non-blob column", async () => {
const { table, rowIds } = await openBlobTable();
await expect(table.fetchBlobs("id", rowIds)).rejects.toThrow(/blob/i);
});
it("discovers and fetches nested blob columns", async () => {
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), true),
new Field("info", new Struct([blob("image")]), true),
]);
const payload = Buffer.from("nested");
const table = await db.createTable(
"nested_blobs",
[{ id: 1n, info: { image: payload } }],
{ schema },
);
expect(await table.blobColumns()).toEqual(["info.image"]);
const rows = await table.query().withRowId().toArray();
const bytes = await table.fetchBlobs("info.image", [
rows[0]._rowid as bigint,
]);
expect(bytes[0]!.equals(payload)).toBe(true);
});
it("creates and adds list blob columns", async () => {
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), true),
new Field("images", new List(blob("image")), true),
]);
const alpha = Buffer.from("alpha");
const beta = Buffer.from("beta");
const gamma = Buffer.from("gamma");
const table = await db.createTable(
"list_blobs",
[{ id: 1n, images: [alpha, beta] }],
{ schema },
);
await table.add([
{ id: 2n, images: null },
{ id: 3n, images: [gamma, null] },
{ id: 4n, images: [] },
]);
expect(await table.blobColumns()).toEqual(["images.image"]);
const rows = await table.query().toArray();
const byId = new Map(rows.map((row) => [Number(row.id), row]));
expect(descriptorSizes(byId.get(1)!.images)).toEqual([
alpha.length,
beta.length,
]);
expect(byId.get(2)!.images).toBeNull();
expect(descriptorSizes(byId.get(3)!.images)).toEqual([gamma.length, null]);
expect(Array.from(byId.get(4)!.images as Iterable<unknown>)).toHaveLength(
0,
);
});
it("creates and adds list struct blob columns", async () => {
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), true),
new Field(
"items",
new List(
new Field(
"item",
new Struct([new Field("name", new Utf8(), true), blob("image")]),
true,
),
),
true,
),
]);
const alpha = Buffer.from("nested-alpha");
const beta = Buffer.from("nested-beta");
const table = await db.createTable(
"list_struct_blobs",
[{ id: 1n, items: [{ name: "one", image: alpha }] }],
{ schema },
);
await table.add([
{
id: 2n,
items: [
{ name: "two", image: beta },
{ name: "three", image: null },
],
},
]);
const rows = await table.query().toArray();
const byId = new Map(rows.map((row) => [Number(row.id), row]));
expect(
descriptorSizes(
Array.from(byId.get(1)!.items as Iterable<{ image: unknown }>).map(
(item) => item.image,
),
),
).toEqual([alpha.length]);
expect(
descriptorSizes(
Array.from(byId.get(2)!.items as Iterable<{ image: unknown }>).map(
(item) => item.image,
),
),
).toEqual([beta.length, null]);
});
it("rejects blob fields inside a fixed-size list", async () => {
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), true),
new Field("frames", new FixedSizeList(2, blob("frame")), true),
]);
await expect(
db.createTable(
"fsl_blobs",
[{ id: 1n, frames: [Buffer.from("a"), Buffer.from("b")] }],
{ schema },
),
).rejects.toThrow(
"Blob fields inside FixedSizeList are not supported. Use List instead.",
);
});
it("rejects blob fields inside a nested fixed-size list", async () => {
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), true),
new Field(
"clip",
new Struct([
new Field("frames", new FixedSizeList(2, blob("frame")), true),
]),
true,
),
]);
await expect(
db.createTable(
"nested_fsl_blobs",
[
{
id: 1n,
clip: { frames: [Buffer.from("a"), Buffer.from("b")] },
},
],
{ schema },
),
).rejects.toThrow(
"Blob fields inside FixedSizeList are not supported. Use List instead.",
);
});
it("rejects an Arrow table with blob fields inside a fixed-size list", async () => {
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), true),
new Field("frames", new FixedSizeList(2, blob("frame")), true),
]);
await expect(
db.createTable("fsl_blobs_ipc", new ArrowTable(schema)),
).rejects.toThrow(
"Blob fields inside FixedSizeList are not supported. Use List instead.",
);
});
function descriptorSizes(values: unknown): (number | null)[] {
return Array.from(
values as Iterable<{ size?: bigint | number } | null>,
).map((value) => (value == null ? null : Number(value.size)));
}
async function openBlobTable() {
const db = await connect(tmpDir.name);
const schema = new Schema([
new Field("id", new Int64(), true),
blob("image"),
]);
const alpha = Buffer.from("alpha");
const beta = Buffer.from("beta");
const table = await db.createTable(
"blobs",
[
{ id: 1n, image: alpha },
{ id: 2n, image: beta },
{ id: 3n, image: null },
],
{ schema },
);
const rows = await table.query().withRowId().toArray();
const rowIdById = new Map(
rows.map((r) => [Number(r.id), r._rowid as bigint]),
);
const rowIds = [1, 2, 3].map((id) => rowIdById.get(id)!);
return { table, rowIds, alpha, beta };
}
});
describe("when dealing with tags", () => {
let tmpDir: tmp.DirResult;
beforeEach(() => {
@@ -2789,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 () => {
@@ -2810,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>
+1 -102
View File
@@ -40,7 +40,6 @@ import {
} from "apache-arrow";
import { Buffers } from "apache-arrow/data";
import { typedArrayToArrowType } from "./arrow_type";
import { coerceBlobValue, isBlobField } from "./blob";
import { type EmbeddingFunction } from "./embedding/embedding_function";
import {
EmbeddingFunctionConfig,
@@ -431,14 +430,12 @@ export function makeArrowTable(
throw new Error("A schema must be provided if data is empty");
} else {
schema = new Schema(schema.fields, schemaMetadata);
validateBlobSchema(schema);
return new ArrowTable(schema);
}
}
let inferredSchema = inferSchema(data, schema, opt);
inferredSchema = new Schema(inferredSchema.fields, schemaMetadata);
validateBlobSchema(inferredSchema);
const finalColumns: Record<string, Vector> = {};
for (const field of inferredSchema.fields) {
@@ -448,35 +445,6 @@ export function makeArrowTable(
return new ArrowTable(inferredSchema, finalColumns);
}
function validateBlobSchema(schema: Schema): void {
for (const field of schema.fields) {
validateBlobField(field);
}
}
function validateBlobField(field: Field): void {
if (
isFixedSizeList(field.type) &&
containsBlobField(field.type.children[0])
) {
throw new Error(
"Blob fields inside FixedSizeList are not supported. Use List instead.",
);
}
for (const child of field.type.children ?? []) {
validateBlobField(child);
}
}
function containsBlobField(field: Field): boolean {
if (isBlobField(field)) {
return true;
}
return (field.type.children ?? []).some((child: Field) =>
containsBlobField(child),
);
}
function isObject(value: unknown): value is Record<string, unknown> {
return (
typeof value === "object" &&
@@ -512,32 +480,6 @@ function transposeData(
path: string[] = [],
): Vector {
const valuesPath = [...path, field.name];
if (isBlobField(field) && field.type instanceof Struct) {
const blobRows = data.map((datum) =>
coerceBlobValue(valueAtPath(datum, valuesPath)),
);
const childVectors = field.type.children.map((child) => {
const values = blobRows.map((row) =>
row == null ? null : (row[child.name as "data" | "uri"] ?? null),
);
return makeVector(values, child.type, undefined, child.nullable);
});
const nullCount = blobRows.filter((row) => row === null).length;
const structData = makeData({
type: field.type,
length: blobRows.length,
nullCount,
nullBitmap:
nullCount > 0
? arrowUtil.packBools(blobRows.map((row) => row !== null))
: undefined,
children: childVectors.map((v) => v.data[0]),
});
return arrowMakeVector(structData);
}
if (isList(field.type) && containsBlobField(field.type.children[0])) {
return transposeListData(data, field, valuesPath);
}
const values = data.map((datum) => valueAtPath(datum, valuesPath));
if (field.type instanceof Struct) {
const childFields = field.type.children;
@@ -553,7 +495,7 @@ function transposeData(
nullCount > 0
? arrowUtil.packBools(values.map((value) => value !== null))
: undefined,
children: childVectors.map((v) => v.data[0]),
children: childVectors as unknown as ArrowData<DataType>[],
});
return arrowMakeVector(structData);
} else {
@@ -561,48 +503,6 @@ function transposeData(
}
}
function transposeListData(
data: Record<string, unknown>[],
field: Field,
valuesPath: string[],
): Vector {
const listType = field.type as List;
const childField = listType.children[0];
const lists = data.map((datum) => valueAtPath(datum, valuesPath));
const flattened: Record<string, unknown>[] = [];
const validity: boolean[] = [];
const offsets: number[] = [0];
for (const list of lists) {
if (list == null) {
validity.push(false);
offsets.push(flattened.length);
continue;
}
if (!Array.isArray(list)) {
throw new Error(`expected an array for list field '${field.name}'`);
}
validity.push(true);
for (const element of list) {
flattened.push({ [childField.name]: element });
}
offsets.push(flattened.length);
}
const childVector = transposeData(flattened, childField, []);
const nullCount = validity.filter((valid) => !valid).length;
return arrowMakeVector(
makeData({
type: listType,
length: lists.length,
nullCount,
nullBitmap: nullCount > 0 ? arrowUtil.packBools(validity) : undefined,
valueOffsets: Int32Array.from(offsets),
child: childVector.data[0],
}),
);
}
/**
* Create an empty Arrow table with the provided schema
*/
@@ -1052,7 +952,6 @@ export async function fromTableToBuffer(
schema = sanitizeSchema(schema);
}
const tableWithEmbeddings = await applyEmbeddings(table, embeddings, schema);
validateBlobSchema(tableWithEmbeddings.schema);
const writer = RecordBatchFileWriter.writeAll(tableWithEmbeddings);
return Buffer.from(await writer.toUint8Array());
}
-236
View File
@@ -1,236 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
import { Field, LargeBinary, Struct, Utf8 } from "apache-arrow";
import { BlobFile as NativeBlobFile } from "./native";
const BLOB_V2_EXTENSION_NAME = "lance.blob.v2";
const INLINE_SIZE_THRESHOLD_KEY = "lance-encoding:blob-inline-size-threshold";
const DEDICATED_SIZE_THRESHOLD_KEY =
"lance-encoding:blob-dedicated-size-threshold";
const PACK_FILE_SIZE_THRESHOLD_KEY =
"lance-encoding:blob-pack-file-size-threshold";
export type BlobInput = {
data: Buffer | Uint8Array | null;
uri: string | null;
};
export type BlobOptions = {
/** Defaults to true. */
nullable?: boolean;
/**
* Max payload bytes kept inline in the data file. Zero is allowed. Must be a
* safe integer.
*/
inlineSizeThreshold?: number;
/**
* Max payload bytes stored in a packed sidecar before a dedicated file. Must
* be a positive safe integer.
*/
dedicatedSizeThreshold?: number;
/**
* Max bytes in one packed sidecar before starting another. Must be a positive
* safe integer.
*/
packFileSizeThreshold?: number;
};
/**
* Declares a `lance.blob.v2` column.
*
* Query results are descriptors, not payload bytes. Use {@link Table.fetchBlobs}
* or {@link Table.fetchBlobFiles} to read bytes.
*
* @example
* ```ts
* import { readFile } from "node:fs/promises";
* import { Field, Int64, Schema } from "apache-arrow";
* import { blob, connect } from "@lancedb/lancedb";
*
* const db = await connect("./data");
* const video = await readFile("clip.mp4");
* const table = await db.createTable(
* "videos",
* [{ id: 1n, video }],
* {
* schema: new Schema([
* new Field("id", new Int64()),
* blob("video"),
* ]),
* },
* );
*
* const rows = await table.query().select(["id"]).withRowId().toArray();
* const rowIds = rows.map((row) => row._rowid as bigint);
* const bytes = await table.fetchBlobs("video", rowIds);
*
* const [handle] = await table.fetchBlobFiles("video", rowIds);
* const size = handle!.size();
* const header = await handle!.readRange(0n, size < 65536n ? size : 65536n);
* ```
*/
export function blob(name: string, options: BlobOptions = {}): Field {
const metadata = new Map<string, string>([
["ARROW:extension:name", BLOB_V2_EXTENSION_NAME],
]);
setThreshold(
metadata,
INLINE_SIZE_THRESHOLD_KEY,
"inlineSizeThreshold",
options.inlineSizeThreshold,
0,
);
setThreshold(
metadata,
DEDICATED_SIZE_THRESHOLD_KEY,
"dedicatedSizeThreshold",
options.dedicatedSizeThreshold,
1,
);
setThreshold(
metadata,
PACK_FILE_SIZE_THRESHOLD_KEY,
"packFileSizeThreshold",
options.packFileSizeThreshold,
1,
);
return new Field(
name,
new Struct([
new Field("data", new LargeBinary(), true),
new Field("uri", new Utf8(), true),
]),
options.nullable ?? true,
metadata,
);
}
/**
* Checks for the `lance.blob.v2` extension marker. Does not validate the
* field's storage type.
*/
export function isBlobField(field: Field): boolean {
return field.metadata?.get("ARROW:extension:name") === BLOB_V2_EXTENSION_NAME;
}
/**
* A lazy handle to blob bytes. Create one with {@link Table.fetchBlobFiles}.
*
* @hideconstructor
*/
export class BlobFile {
private readonly inner: NativeBlobFile;
private constructor(inner: NativeBlobFile) {
if (!(inner instanceof NativeBlobFile)) {
throw new Error("BlobFile handles come from Table.fetchBlobFiles");
}
this.inner = inner;
}
/** @ignore */
static fromNative(inner: NativeBlobFile): BlobFile {
return new BlobFile(inner);
}
/** Returns the blob size in bytes. */
size(): bigint {
return this.inner.size();
}
/**
* Reads from the cursor to the end and advances the cursor.
*
* A second call returns an empty buffer. {@link BlobFile.readRange} does
* not move the cursor.
*/
read(): Promise<Buffer> {
return this.inner.read();
}
/**
* Reads the half-open byte range `[start, end)`.
*
* Fails when `end` is past the blob size. Does not move the cursor.
*/
readRange(start: bigint, end: bigint): Promise<Buffer> {
return this.inner.readRange(start, end);
}
}
export function coerceBlobValue(value: unknown): BlobInput | null {
if (value == null) {
return null;
}
if (isBlobBytes(value)) {
return { data: value, uri: null };
}
if (ArrayBuffer.isView(value)) {
throw new Error("Blob data must be Buffer or Uint8Array");
}
if (typeof value === "string") {
if (value === "") {
throw new Error("Blob uri cannot be empty");
}
return { data: null, uri: value };
}
if (typeof value === "object") {
const record = value as Record<string, unknown>;
if (!("data" in record) && !("uri" in record)) {
throw new Error(
"Blob struct values must include a 'data' or 'uri' field",
);
}
const uri = record.uri;
if (uri === "") {
throw new Error("Blob uri cannot be empty");
}
if (uri != null && typeof uri !== "string") {
throw new Error(`Blob uri must be a string or null, got ${typeof uri}`);
}
const data = record.data;
if (data != null && !isBlobBytes(data)) {
throw new Error("Blob data must be Buffer, Uint8Array, or null");
}
const bytes = (data as Buffer | Uint8Array | null | undefined) ?? null;
const uriValue = uri ?? null;
if ((bytes == null) === (uriValue == null)) {
throw new Error(
"Blob struct values must set exactly one of 'data' or 'uri'",
);
}
return { data: bytes, uri: uriValue };
}
throw new Error(
"Blob column values must be Buffer, Uint8Array, a URI string, null, or { data?, uri? }",
);
}
function isBlobBytes(value: unknown): value is Buffer | Uint8Array {
return Buffer.isBuffer(value) || value instanceof Uint8Array;
}
function setThreshold(
metadata: Map<string, string>,
key: string,
optionName: string,
value: number | undefined,
minimum: number,
): void {
if (value === undefined) {
return;
}
if (!Number.isSafeInteger(value)) {
throw new Error(`${optionName} must be a safe integer`);
}
if (value < minimum) {
throw new Error(
minimum <= 0
? `${optionName} must be non-negative`
: `${optionName} must be positive`,
);
}
metadata.set(key, String(value));
}
-3
View File
@@ -77,9 +77,6 @@ export {
VectorColumnOptions,
} from "./arrow";
export { blob, isBlobField, BlobFile } from "./blob";
export type { BlobOptions } from "./blob";
export {
Connection,
CreateTableOptions,
+19 -85
View File
@@ -17,7 +17,6 @@ import {
tableFromIPC,
} from "./arrow";
import { BlobFile } from "./blob";
import { EmbeddingFunctionConfig, getRegistry } from "./embedding/registry";
import { IndexOptions } from "./indices";
import { Job } from "./job";
@@ -148,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;
@@ -511,35 +511,6 @@ export abstract class Table {
*/
abstract takeRowIds(rowIds: readonly (bigint | number)[]): TakeQuery;
/**
* Blob v2 columns, including nested dotted paths.
*/
abstract blobColumns(): Promise<string[]>;
/**
* Bytes for `column` at row IDs from {@link Query.withRowId}.
*
* Reads the table's current checkout. IDs from another version can fail after
* compaction unless stable row ids are enabled. Results keep input order and
* duplicates. Null blobs are `null`. Empty blobs are empty buffers.
*/
abstract fetchBlobs(
column: string,
rowIds: readonly (bigint | number)[],
): Promise<(Buffer | null)[]>;
/**
* Opens lazy blob handles for `column` at the given row IDs using the
* table's current checkout.
*
* Preserves input order, duplicates, and nulls. Use this for large payloads.
* See {@link Table.fetchBlobs} for row-ID validity across versions.
*/
abstract fetchBlobFiles(
column: string,
rowIds: readonly (bigint | number)[],
): Promise<(BlobFile | null)[]>;
/**
* Create a search query to find the nearest neighbors
* of the given query
@@ -1190,34 +1161,23 @@ export class LocalTable extends Table {
}
takeRowIds(rowIds: readonly (bigint | number)[]): TakeQuery {
return new TakeQuery(this.inner.takeRowIds(rowIdsToBigInts(rowIds)));
}
const ids = rowIds.map((id) => {
if (typeof id === "bigint") {
return id;
}
if (!Number.isInteger(id)) {
throw new Error("Row id must be an integer (or bigint)");
}
if (id < 0) {
throw new Error("Row id cannot be negative");
}
if (!Number.isSafeInteger(id)) {
throw new Error("Row id is too large for number; use bigint instead");
}
return BigInt(id);
});
blobColumns(): Promise<string[]> {
return this.inner.blobColumns();
}
async fetchBlobs(
column: string,
rowIds: readonly (bigint | number)[],
): Promise<(Buffer | null)[]> {
const values = await this.inner.fetchBlobs(column, rowIdsToBigInts(rowIds));
// N-API Option maps missing values to undefined. Collapse those to null.
return values.map((value) => value ?? null);
}
async fetchBlobFiles(
column: string,
rowIds: readonly (bigint | number)[],
): Promise<(BlobFile | null)[]> {
const files = await this.inner.fetchBlobFiles(
column,
rowIdsToBigInts(rowIds),
);
// N-API Option maps missing values to undefined. Collapse those to null.
return files.map((file) =>
file == null ? null : BlobFile.fromNative(file),
);
return new TakeQuery(this.inner.takeRowIds(ids));
}
query(): Query {
@@ -1486,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,
);
}
@@ -1774,21 +1726,3 @@ export class Branches {
)) as unknown as CherryPickResult;
}
}
function rowIdsToBigInts(rowIds: readonly (bigint | number)[]): bigint[] {
return rowIds.map((id) => {
if (typeof id === "bigint") {
return id;
}
if (!Number.isInteger(id)) {
throw new Error("Row id must be an integer (or bigint)");
}
if (id < 0) {
throw new Error("Row id cannot be negative");
}
if (!Number.isSafeInteger(id)) {
throw new Error("Row id is too large for number; use bigint instead");
}
return BigInt(id);
});
}
-95
View File
@@ -1,95 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
use std::ops::Range;
use std::sync::Arc;
use arrow_array::{Array, LargeBinaryArray};
use lancedb::blob::BlobFile as LanceBlobFile;
use napi::bindgen_prelude::*;
use napi_derive::napi;
use crate::error::convert_error;
#[napi]
pub struct BlobFile {
inner: Arc<LanceBlobFile>,
}
impl BlobFile {
pub(crate) fn new(inner: LanceBlobFile) -> Self {
Self {
inner: Arc::new(inner),
}
}
}
#[napi]
impl BlobFile {
#[napi]
pub fn size(&self) -> BigInt {
BigInt::from(self.inner.size())
}
#[napi]
pub async fn read(&self) -> napi::Result<Buffer> {
let bytes = self.inner.read().await.map_err(|err| convert_error(&err))?;
Ok(Buffer::from(bytes.as_ref()))
}
#[napi]
pub async fn read_range(&self, start: BigInt, end: BigInt) -> napi::Result<Buffer> {
let range = bigint_range(start, end)?;
let bytes = self
.inner
.read_range(range)
.await
.map_err(|err| convert_error(&err))?;
Ok(Buffer::from(bytes.as_ref()))
}
}
fn bigint_range(start: BigInt, end: BigInt) -> napi::Result<Range<u64>> {
let start = parse_u64(start, "start")?;
let end = parse_u64(end, "end")?;
if start > end {
return Err(napi::Error::from_reason(format!(
"invalid blob range: start ({start}) > end ({end})"
)));
}
Ok(start..end)
}
pub(crate) fn parse_u64(value: BigInt, name: &str) -> napi::Result<u64> {
let (negative, value, lossless) = value.get_u64();
if negative {
return Err(napi::Error::from_reason(format!(
"{name} cannot be negative"
)));
}
if !lossless {
return Err(napi::Error::from_reason(format!(
"{name} is too large to fit in u64"
)));
}
Ok(value)
}
pub(crate) fn parse_row_ids(row_ids: Vec<BigInt>) -> napi::Result<Vec<u64>> {
row_ids
.into_iter()
.map(|id| parse_u64(id, "row id"))
.collect()
}
pub(crate) fn copy_blob_buffers(array: LargeBinaryArray) -> Vec<Option<Buffer>> {
(0..array.len())
.map(|i| {
if array.is_null(i) {
None
} else {
Some(Buffer::from(array.value(i).to_vec()))
}
})
.collect()
}
-1
View File
@@ -10,7 +10,6 @@ use std::collections::HashMap;
use env_logger::Env;
use napi_derive::*;
mod blob;
mod connection;
mod error;
mod header;
+27 -62
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,
};
@@ -15,7 +15,6 @@ use napi::bindgen_prelude::*;
use napi::threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode};
use napi_derive::napi;
use crate::blob::{BlobFile, copy_blob_buffers, parse_row_ids};
use crate::error::NapiErrorExt;
use crate::index::Index;
use crate::merge::NativeMergeInsertBuilder;
@@ -330,44 +329,6 @@ impl Table {
))
}
#[napi(catch_unwind)]
pub async fn blob_columns(&self) -> napi::Result<Vec<String>> {
self.inner_ref()?.blob_columns().await.default_error()
}
#[napi(catch_unwind)]
pub async fn fetch_blobs(
&self,
column: String,
row_ids: Vec<BigInt>,
) -> napi::Result<Vec<Option<Buffer>>> {
let row_ids = parse_row_ids(row_ids)?;
let array = self
.inner_ref()?
.fetch_blobs(column.as_str(), &row_ids)
.await
.default_error()?;
Ok(copy_blob_buffers(array))
}
#[napi(catch_unwind)]
pub async fn fetch_blob_files(
&self,
column: String,
row_ids: Vec<BigInt>,
) -> napi::Result<Vec<Option<BlobFile>>> {
let row_ids = parse_row_ids(row_ids)?;
let files = self
.inner_ref()?
.fetch_blob_files(column.as_str(), &row_ids)
.await
.default_error()?;
Ok(files
.into_iter()
.map(|file| file.map(BlobFile::new))
.collect())
}
#[napi(catch_unwind)]
pub fn vector_search(&self, vector: Float32Array) -> napi::Result<VectorQuery> {
self.query()?.nearest_to(vector)
@@ -677,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 {
@@ -703,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