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 372 additions and 1562 deletions
Generated
+239 -250
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" rust-version = "1.91.0"
[workspace.dependencies] [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 = { "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", default-features = false, "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", default-features = false, "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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.17", "tag" = "v12.0.0-beta.17", "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 } lancedb = { path = "rust/lancedb", default-features = false }
ahash = "0.8" ahash = "0.8"
# Note that this one does not include pyarrow # Note that this one does not include pyarrow
@@ -60,7 +60,7 @@ log = "0.4"
metrics = "0.24" metrics = "0.24"
metrics-util = "0.19" metrics-util = "0.19"
moka = { version = "0.12", features = ["future"] } moka = { version = "0.12", features = ["future"] }
object_store = "0.14.1" object_store = "0.13.2"
pin-project = "1.0.7" pin-project = "1.0.7"
rand = "0.9" rand = "0.9"
snafu = "0.8" 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() ### branches()
```ts ```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() ### flushLsm()
```ts ```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 ## Classes
- [AutoQuery](classes/AutoQuery.md) - [AutoQuery](classes/AutoQuery.md)
- [BlobFile](classes/BlobFile.md)
- [BooleanQuery](classes/BooleanQuery.md) - [BooleanQuery](classes/BooleanQuery.md)
- [BoostQuery](classes/BoostQuery.md) - [BoostQuery](classes/BoostQuery.md)
- [BranchContents](classes/BranchContents.md) - [BranchContents](classes/BranchContents.md)
@@ -144,7 +143,6 @@
- [AnalyzePlanDistributedMetrics](type-aliases/AnalyzePlanDistributedMetrics.md) - [AnalyzePlanDistributedMetrics](type-aliases/AnalyzePlanDistributedMetrics.md)
- [BaseTokenizer](type-aliases/BaseTokenizer.md) - [BaseTokenizer](type-aliases/BaseTokenizer.md)
- [BlobOptions](type-aliases/BlobOptions.md)
- [Data](type-aliases/Data.md) - [Data](type-aliases/Data.md)
- [DataLike](type-aliases/DataLike.md) - [DataLike](type-aliases/DataLike.md)
- [FieldLike](type-aliases/FieldLike.md) - [FieldLike](type-aliases/FieldLike.md)
@@ -160,11 +158,9 @@
## Functions ## Functions
- [RecordBatchIterator](functions/RecordBatchIterator.md) - [RecordBatchIterator](functions/RecordBatchIterator.md)
- [blob](functions/blob.md)
- [connect](functions/connect.md) - [connect](functions/connect.md)
- [connectNamespace](functions/connectNamespace.md) - [connectNamespace](functions/connectNamespace.md)
- [instrumentLanceDbMetrics](functions/instrumentLanceDbMetrics.md) - [instrumentLanceDbMetrics](functions/instrumentLanceDbMetrics.md)
- [isBlobField](functions/isBlobField.md)
- [makeArrowTable](functions/makeArrowTable.md) - [makeArrowTable](functions/makeArrowTable.md)
- [packBits](functions/packBits.md) - [packBits](functions/packBits.md)
- [permutationBuilder](functions/permutationBuilder.md) - [permutationBuilder](functions/permutationBuilder.md)
+2 -1
View File
@@ -26,7 +26,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 // 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()}); 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> <properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<arrow.version>15.0.0</arrow.version> <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.skip>false</spotless.skip>
<spotless.version>2.30.0</spotless.version> <spotless.version>2.30.0</spotless.version>
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.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, Query,
Table, Table,
VectorQuery, VectorQuery,
blob,
connect, connect,
tokenize, tokenize,
} from "../lancedb"; } from "../lancedb";
@@ -53,6 +52,7 @@ import {
Operator, Operator,
instanceOfFullTextQuery, instanceOfFullTextQuery,
} from "../lancedb/query"; } from "../lancedb/query";
import { LocalTable } from "../lancedb/table";
describe.each([arrow15, arrow16, arrow17, arrow18])( describe.each([arrow15, arrow16, arrow17, arrow18])(
"Given a table", "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", () => { describe("when dealing with tags", () => {
let tmpDir: tmp.DirResult; let tmpDir: tmp.DirResult;
beforeEach(() => { beforeEach(() => {
@@ -2789,7 +2519,7 @@ describe("when optimizing a dataset", () => {
it("cleanups old versions", async () => { it("cleanups old versions", async () => {
const stats = await table.optimize({ cleanupOlderThan: new Date() }); const stats = await table.optimize({ cleanupOlderThan: new Date() });
expect(stats.prune.bytesRemoved).toBeGreaterThan(0); expect(stats.prune.bytesRemoved).toBeGreaterThan(0);
expect(stats.prune.oldVersionsRemoved).toBe(3); expect(stats.prune.oldVersionsRemoved).toBe(2);
}); });
it("delete unverified", async () => { 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])( describe.each([arrow15, arrow16, arrow17, arrow18])(
"when optimizing a dataset", "when optimizing a dataset",
// biome-ignore lint/suspicious/noExplicitAny: <explanation> // biome-ignore lint/suspicious/noExplicitAny: <explanation>
+1 -102
View File
@@ -40,7 +40,6 @@ import {
} from "apache-arrow"; } from "apache-arrow";
import { Buffers } from "apache-arrow/data"; import { Buffers } from "apache-arrow/data";
import { typedArrayToArrowType } from "./arrow_type"; import { typedArrayToArrowType } from "./arrow_type";
import { coerceBlobValue, isBlobField } from "./blob";
import { type EmbeddingFunction } from "./embedding/embedding_function"; import { type EmbeddingFunction } from "./embedding/embedding_function";
import { import {
EmbeddingFunctionConfig, EmbeddingFunctionConfig,
@@ -431,14 +430,12 @@ export function makeArrowTable(
throw new Error("A schema must be provided if data is empty"); throw new Error("A schema must be provided if data is empty");
} else { } else {
schema = new Schema(schema.fields, schemaMetadata); schema = new Schema(schema.fields, schemaMetadata);
validateBlobSchema(schema);
return new ArrowTable(schema); return new ArrowTable(schema);
} }
} }
let inferredSchema = inferSchema(data, schema, opt); let inferredSchema = inferSchema(data, schema, opt);
inferredSchema = new Schema(inferredSchema.fields, schemaMetadata); inferredSchema = new Schema(inferredSchema.fields, schemaMetadata);
validateBlobSchema(inferredSchema);
const finalColumns: Record<string, Vector> = {}; const finalColumns: Record<string, Vector> = {};
for (const field of inferredSchema.fields) { for (const field of inferredSchema.fields) {
@@ -448,35 +445,6 @@ export function makeArrowTable(
return new ArrowTable(inferredSchema, finalColumns); 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> { function isObject(value: unknown): value is Record<string, unknown> {
return ( return (
typeof value === "object" && typeof value === "object" &&
@@ -512,32 +480,6 @@ function transposeData(
path: string[] = [], path: string[] = [],
): Vector { ): Vector {
const valuesPath = [...path, field.name]; 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)); const values = data.map((datum) => valueAtPath(datum, valuesPath));
if (field.type instanceof Struct) { if (field.type instanceof Struct) {
const childFields = field.type.children; const childFields = field.type.children;
@@ -553,7 +495,7 @@ function transposeData(
nullCount > 0 nullCount > 0
? arrowUtil.packBools(values.map((value) => value !== null)) ? arrowUtil.packBools(values.map((value) => value !== null))
: undefined, : undefined,
children: childVectors.map((v) => v.data[0]), children: childVectors as unknown as ArrowData<DataType>[],
}); });
return arrowMakeVector(structData); return arrowMakeVector(structData);
} else { } 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 * Create an empty Arrow table with the provided schema
*/ */
@@ -1052,7 +952,6 @@ export async function fromTableToBuffer(
schema = sanitizeSchema(schema); schema = sanitizeSchema(schema);
} }
const tableWithEmbeddings = await applyEmbeddings(table, embeddings, schema); const tableWithEmbeddings = await applyEmbeddings(table, embeddings, schema);
validateBlobSchema(tableWithEmbeddings.schema);
const writer = RecordBatchFileWriter.writeAll(tableWithEmbeddings); const writer = RecordBatchFileWriter.writeAll(tableWithEmbeddings);
return Buffer.from(await writer.toUint8Array()); 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, VectorColumnOptions,
} from "./arrow"; } from "./arrow";
export { blob, isBlobField, BlobFile } from "./blob";
export type { BlobOptions } from "./blob";
export { export {
Connection, Connection,
CreateTableOptions, CreateTableOptions,
+19 -85
View File
@@ -17,7 +17,6 @@ import {
tableFromIPC, tableFromIPC,
} from "./arrow"; } from "./arrow";
import { BlobFile } from "./blob";
import { EmbeddingFunctionConfig, getRegistry } from "./embedding/registry"; import { EmbeddingFunctionConfig, getRegistry } from "./embedding/registry";
import { IndexOptions } from "./indices"; import { IndexOptions } from "./indices";
import { Job } from "./job"; import { Job } from "./job";
@@ -148,7 +147,8 @@ export interface OptimizeOptions {
* 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 * // 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()}); * tbl.optimize({cleanupOlderThan: new Date()});
*/ */
cleanupOlderThan: Date; cleanupOlderThan: Date;
@@ -511,35 +511,6 @@ export abstract class Table {
*/ */
abstract takeRowIds(rowIds: readonly (bigint | number)[]): TakeQuery; 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 * Create a search query to find the nearest neighbors
* of the given query * of the given query
@@ -1190,34 +1161,23 @@ export class LocalTable extends Table {
} }
takeRowIds(rowIds: readonly (bigint | number)[]): TakeQuery { 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 new TakeQuery(this.inner.takeRowIds(ids));
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),
);
} }
query(): Query { query(): Query {
@@ -1486,16 +1446,8 @@ export class LocalTable extends Table {
} }
async optimize(options?: Partial<OptimizeOptions>): Promise<OptimizeStats> { 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( return await this.inner.optimize(
cleanupOlderThanMs, options?.cleanupOlderThan?.getTime(),
options?.deleteUnverified, options?.deleteUnverified,
); );
} }
@@ -1774,21 +1726,3 @@ export class Branches {
)) as unknown as CherryPickResult; )) 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 env_logger::Env;
use napi_derive::*; use napi_derive::*;
mod blob;
mod connection; mod connection;
mod error; mod error;
mod header; 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::ipc::{ipc_file_to_batches, ipc_file_to_schema};
use lancedb::table::{ use lancedb::table::{
AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration, AddDataMode, ColumnAlteration as LanceColumnAlteration,
FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken, FieldMetadataUpdate as LanceFieldMetadataUpdate, FtsToken as LanceDbFtsToken,
NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable, NewColumnTransform, OptimizeAction, OptimizeOptions, Ref, Table as LanceDbTable,
}; };
@@ -15,7 +15,6 @@ use napi::bindgen_prelude::*;
use napi::threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode}; use napi::threadsafe_function::{ThreadsafeFunction, ThreadsafeFunctionCallMode};
use napi_derive::napi; use napi_derive::napi;
use crate::blob::{BlobFile, copy_blob_buffers, parse_row_ids};
use crate::error::NapiErrorExt; use crate::error::NapiErrorExt;
use crate::index::Index; use crate::index::Index;
use crate::merge::NativeMergeInsertBuilder; 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)] #[napi(catch_unwind)]
pub fn vector_search(&self, vector: Float32Array) -> napi::Result<VectorQuery> { pub fn vector_search(&self, vector: Float32Array) -> napi::Result<VectorQuery> {
self.query()?.nearest_to(vector) self.query()?.nearest_to(vector)
@@ -677,22 +638,20 @@ impl Table {
#[napi(catch_unwind)] #[napi(catch_unwind)]
pub async fn optimize( pub async fn optimize(
&self, &self,
older_than_ms: Option<i64>, before_timestamp_ms: Option<i64>,
delete_unverified: Option<bool>, delete_unverified: Option<bool>,
) -> napi::Result<OptimizeStats> { ) -> napi::Result<OptimizeStats> {
let inner = self.inner_ref()?; let inner = self.inner_ref()?;
let older_than = if let Some(ms) = older_than_ms { let before_timestamp = before_timestamp_ms
if ms == i64::MIN { .map(|ms| {
return Err(napi::Error::from_reason(format!( DateTime::from_timestamp_millis(ms).ok_or_else(|| {
"older_than_ms can not be {}", napi::Error::from_reason(format!(
i32::MIN, "cleanupOlderThan timestamp is out of range: {ms}"
))); ))
} })
Duration::try_milliseconds(ms) })
} else { .transpose()?;
None
};
let compaction_stats = inner let compaction_stats = inner
.optimize(OptimizeAction::Compact { .optimize(OptimizeAction::Compact {
@@ -703,16 +662,22 @@ impl Table {
.default_error()? .default_error()?
.compaction .compaction
.unwrap(); .unwrap();
let prune_stats = inner let prune_stats = if let Some(before_timestamp) = before_timestamp {
.optimize(OptimizeAction::Prune { inner
older_than, .optimize_prune_before(before_timestamp, delete_unverified, None)
delete_unverified, .await
error_if_tagged_old_versions: None, } else {
}) inner
.await .optimize(OptimizeAction::Prune {
.default_error()? older_than: None,
.prune delete_unverified,
.unwrap(); error_if_tagged_old_versions: None,
})
.await
}
.default_error()?
.prune
.unwrap();
inner inner
.optimize(lancedb::table::OptimizeAction::Index( .optimize(lancedb::table::OptimizeAction::Index(
OptimizeOptions::default(), OptimizeOptions::default(),
+27
View File
@@ -1739,6 +1739,33 @@ impl Table {
self.inner.optimize(action).await 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. /// Add new columns to the table, providing values to fill in.
pub fn add_columns(&self) -> AddColumnsBuilder { pub fn add_columns(&self) -> AddColumnsBuilder {
AddColumnsBuilder::new(self.inner.clone()) AddColumnsBuilder::new(self.inner.clone())
+21 -1
View File
@@ -8,7 +8,8 @@
use std::sync::Arc; 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::dataset::optimize::{CompactionMetrics, IndexRemapperOptions, compact_files};
use lance::index::DatasetIndexExt; use lance::index::DatasetIndexExt;
use lance_index::optimize::OptimizeOptions; use lance_index::optimize::OptimizeOptions;
@@ -139,6 +140,25 @@ pub(crate) async fn cleanup_old_versions(
.await?) .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. /// Compact files in the dataset.
/// ///
/// This can be run after making several small appends to optimize the table /// This can be run after making several small appends to optimize the table