mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-12 08:12:28 +00:00
Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
6702e3fec1 | ||
|
|
e0bd4b5fa1 |
Generated
+250
-240
File diff suppressed because it is too large
Load Diff
+15
-15
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
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" }
|
||||
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" }
|
||||
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.13.2"
|
||||
object_store = "0.14.1"
|
||||
pin-project = "1.0.7"
|
||||
rand = "0.9"
|
||||
snafu = "0.8"
|
||||
|
||||
@@ -0,0 +1,62 @@
|
||||
[**@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`<`Buffer`>
|
||||
|
||||
***
|
||||
|
||||
### 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`<`Buffer`>
|
||||
|
||||
***
|
||||
|
||||
### size()
|
||||
|
||||
```ts
|
||||
size(): bigint
|
||||
```
|
||||
|
||||
Returns the blob size in bytes.
|
||||
|
||||
#### Returns
|
||||
|
||||
`bigint`
|
||||
@@ -74,10 +74,10 @@ now: the column is committed with no values, and rows get them from
|
||||
[Table#refreshColumn](Table.md#refreshcolumn). Declaring one therefore costs the same on a
|
||||
large table as on an empty one.
|
||||
|
||||
A refresh also recomputes the rows whose inputs changed since they were
|
||||
computed, so a mutated input is reflected by the next refresh. While a
|
||||
declaration reads a column, that column cannot be renamed, retyped or
|
||||
dropped.
|
||||
A refresh does not revisit rows it has already filled, so mutating an
|
||||
input leaves the value computed at fill time; recomputing means dropping
|
||||
the column and declaring it again. While a declaration reads a column,
|
||||
that column cannot be renamed, retyped or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
server, and the refresh runs as a server job -- see
|
||||
@@ -137,6 +137,20 @@ 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`<`string`[]>
|
||||
|
||||
***
|
||||
|
||||
### branches()
|
||||
|
||||
```ts
|
||||
@@ -499,6 +513,54 @@ 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`<(`null` \| [`BlobFile`](BlobFile.md))[]>
|
||||
|
||||
***
|
||||
|
||||
### 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`<(`null` \| `Buffer`)[]>
|
||||
|
||||
***
|
||||
|
||||
### flushLsm()
|
||||
|
||||
```ts
|
||||
@@ -854,10 +916,10 @@ abstract refreshColumn(column): Promise<RefreshColumnResult>
|
||||
|
||||
Fill the rows of a computed column that hold no value yet.
|
||||
|
||||
Rows appended since the last refresh are filled by the next one, and
|
||||
rows whose inputs changed since they were computed are recomputed;
|
||||
everything else is left as it is. Local tables only: a remote refresh
|
||||
runs as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
|
||||
Rows appended since the last refresh are filled by the next one; rows
|
||||
already filled are left as they are, so the call is idempotent and does
|
||||
not observe a mutated input. Local tables only: a remote refresh runs
|
||||
as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
|
||||
|
||||
#### Parameters
|
||||
|
||||
|
||||
@@ -0,0 +1,55 @@
|
||||
[**@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);
|
||||
```
|
||||
@@ -0,0 +1,22 @@
|
||||
[**@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`<`any`>
|
||||
|
||||
## Returns
|
||||
|
||||
`boolean`
|
||||
@@ -19,6 +19,7 @@
|
||||
## Classes
|
||||
|
||||
- [AutoQuery](classes/AutoQuery.md)
|
||||
- [BlobFile](classes/BlobFile.md)
|
||||
- [BooleanQuery](classes/BooleanQuery.md)
|
||||
- [BoostQuery](classes/BoostQuery.md)
|
||||
- [BranchContents](classes/BranchContents.md)
|
||||
@@ -143,6 +144,7 @@
|
||||
|
||||
- [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)
|
||||
@@ -158,9 +160,11 @@
|
||||
## 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)
|
||||
|
||||
@@ -0,0 +1,48 @@
|
||||
[**@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
@@ -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.16</lance-core.version>
|
||||
<lance-core.version>12.0.0-beta.17</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>
|
||||
|
||||
@@ -0,0 +1,185 @@
|
||||
// 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");
|
||||
});
|
||||
});
|
||||
@@ -18,6 +18,7 @@ import {
|
||||
Query,
|
||||
Table,
|
||||
VectorQuery,
|
||||
blob,
|
||||
connect,
|
||||
tokenize,
|
||||
} from "../lancedb";
|
||||
@@ -2401,6 +2402,276 @@ 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(() => {
|
||||
|
||||
+102
-1
@@ -40,6 +40,7 @@ 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,
|
||||
@@ -430,12 +431,14 @@ 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) {
|
||||
@@ -445,6 +448,35 @@ 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" &&
|
||||
@@ -480,6 +512,32 @@ 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;
|
||||
@@ -495,7 +553,7 @@ function transposeData(
|
||||
nullCount > 0
|
||||
? arrowUtil.packBools(values.map((value) => value !== null))
|
||||
: undefined,
|
||||
children: childVectors as unknown as ArrowData<DataType>[],
|
||||
children: childVectors.map((v) => v.data[0]),
|
||||
});
|
||||
return arrowMakeVector(structData);
|
||||
} else {
|
||||
@@ -503,6 +561,48 @@ 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
|
||||
*/
|
||||
@@ -952,6 +1052,7 @@ 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());
|
||||
}
|
||||
|
||||
@@ -0,0 +1,236 @@
|
||||
// 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));
|
||||
}
|
||||
@@ -77,6 +77,9 @@ export {
|
||||
VectorColumnOptions,
|
||||
} from "./arrow";
|
||||
|
||||
export { blob, isBlobField, BlobFile } from "./blob";
|
||||
export type { BlobOptions } from "./blob";
|
||||
|
||||
export {
|
||||
Connection,
|
||||
CreateTableOptions,
|
||||
|
||||
+83
-24
@@ -17,6 +17,7 @@ import {
|
||||
tableFromIPC,
|
||||
} from "./arrow";
|
||||
|
||||
import { BlobFile } from "./blob";
|
||||
import { EmbeddingFunctionConfig, getRegistry } from "./embedding/registry";
|
||||
import { IndexOptions } from "./indices";
|
||||
import { Job } from "./job";
|
||||
@@ -510,6 +511,35 @@ 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
|
||||
@@ -542,10 +572,10 @@ export abstract class Table {
|
||||
* {@link Table#refreshColumn}. Declaring one therefore costs the same on a
|
||||
* large table as on an empty one.
|
||||
*
|
||||
* A refresh also recomputes the rows whose inputs changed since they were
|
||||
* computed, so a mutated input is reflected by the next refresh. While a
|
||||
* declaration reads a column, that column cannot be renamed, retyped or
|
||||
* dropped.
|
||||
* A refresh does not revisit rows it has already filled, so mutating an
|
||||
* input leaves the value computed at fill time; recomputing means dropping
|
||||
* the column and declaring it again. While a declaration reads a column,
|
||||
* that column cannot be renamed, retyped or dropped.
|
||||
*
|
||||
* On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
* server, and the refresh runs as a server job -- see
|
||||
@@ -576,10 +606,10 @@ export abstract class Table {
|
||||
/**
|
||||
* Fill the rows of a computed column that hold no value yet.
|
||||
*
|
||||
* Rows appended since the last refresh are filled by the next one, and
|
||||
* rows whose inputs changed since they were computed are recomputed;
|
||||
* everything else is left as it is. Local tables only: a remote refresh
|
||||
* runs as a server job, through {@link Table#refreshColumnAsync}.
|
||||
* Rows appended since the last refresh are filled by the next one; rows
|
||||
* already filled are left as they are, so the call is idempotent and does
|
||||
* not observe a mutated input. Local tables only: a remote refresh runs
|
||||
* as a server job, through {@link Table#refreshColumnAsync}.
|
||||
* @param {string} column The name of the computed column to fill.
|
||||
* @returns {Promise<RefreshColumnResult>} A promise that resolves to the
|
||||
* number of rows filled and the new version number of the table.
|
||||
@@ -1160,23 +1190,34 @@ export class LocalTable extends Table {
|
||||
}
|
||||
|
||||
takeRowIds(rowIds: readonly (bigint | number)[]): TakeQuery {
|
||||
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);
|
||||
});
|
||||
return new TakeQuery(this.inner.takeRowIds(rowIdsToBigInts(rowIds)));
|
||||
}
|
||||
|
||||
return new TakeQuery(this.inner.takeRowIds(ids));
|
||||
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),
|
||||
);
|
||||
}
|
||||
|
||||
query(): Query {
|
||||
@@ -1733,3 +1774,21 @@ 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);
|
||||
});
|
||||
}
|
||||
|
||||
@@ -0,0 +1,95 @@
|
||||
// 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()
|
||||
}
|
||||
@@ -10,6 +10,7 @@ use std::collections::HashMap;
|
||||
use env_logger::Env;
|
||||
use napi_derive::*;
|
||||
|
||||
mod blob;
|
||||
mod connection;
|
||||
mod error;
|
||||
mod header;
|
||||
|
||||
@@ -15,6 +15,7 @@ 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;
|
||||
@@ -329,6 +330,44 @@ 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)
|
||||
|
||||
@@ -2188,10 +2188,10 @@ class Table(ABC):
|
||||
Declaring one therefore costs the same on a large table as on an
|
||||
empty one.
|
||||
|
||||
A refresh also recomputes the rows whose inputs changed since they
|
||||
were computed, so a mutated input is reflected by the next refresh.
|
||||
While a declaration reads a column, that column cannot be renamed,
|
||||
retyped or dropped.
|
||||
A refresh does not revisit rows it has already filled, so mutating
|
||||
an input leaves the value computed at fill time; recomputing means
|
||||
dropping the column and declaring it again. While a declaration
|
||||
reads a column, that column cannot be renamed, retyped or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
server, and the refresh runs as a server job -- see
|
||||
@@ -2211,7 +2211,7 @@ class Table(ABC):
|
||||
>>> table.add_columns(computed={"doubled": "x * 2"})
|
||||
AddColumnsResult(version=2)
|
||||
>>> table.refresh_column("doubled")
|
||||
RefreshColumnResult(rows_filled=2, version=4)
|
||||
RefreshColumnResult(rows_filled=2, version=3)
|
||||
>>> table.to_arrow().sort_by("x").to_pandas()
|
||||
x doubled
|
||||
0 1 2
|
||||
@@ -2225,8 +2225,8 @@ class Table(ABC):
|
||||
|
||||
Declared with ``add_columns(computed=...)``, a column starts empty and
|
||||
gets its values here. Rows appended since the last refresh are filled
|
||||
by the next one, and rows whose inputs changed since they were computed
|
||||
are recomputed; everything else is left as it is.
|
||||
by the next one; rows already filled are left as they are, so the call
|
||||
is idempotent and does not observe a mutated input.
|
||||
|
||||
Local tables only: a remote refresh runs as a server job, through
|
||||
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
@@ -4318,14 +4318,13 @@ class LanceTable(Table):
|
||||
return LOOP.run(self._table.add_columns(transforms, computed=computed))
|
||||
|
||||
def refresh_column(self, column: str) -> "RefreshColumnResult":
|
||||
"""Fill a computed column's unfilled rows and recompute those whose
|
||||
inputs changed. See
|
||||
"""Fill a computed column's unfilled rows. See
|
||||
[`AsyncTable.refresh_column`][lancedb.AsyncTable.refresh_column]."""
|
||||
return LOOP.run(self._table.refresh_column(column))
|
||||
|
||||
def refresh_column_async(self, column: str) -> Job[RefreshColumnJobResult]:
|
||||
"""Fill a computed column's unfilled rows and recompute those whose
|
||||
inputs changed, returning a handle to the refresh job. See
|
||||
"""Fill a computed column's unfilled rows, returning a handle to the
|
||||
refresh job. See
|
||||
[`Table.refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
"""
|
||||
return Job(LOOP.run(self._table.refresh_column_async(column)))
|
||||
@@ -6313,10 +6312,10 @@ class AsyncTable:
|
||||
them from
|
||||
[`refresh_column`][lancedb.table.AsyncTable.refresh_column].
|
||||
|
||||
A refresh also recomputes the rows whose inputs changed since they
|
||||
were computed, so a mutated input is reflected by the next refresh.
|
||||
While a declaration reads a column, that column cannot be renamed,
|
||||
retyped or dropped.
|
||||
A refresh does not revisit rows it has already filled, so mutating
|
||||
an input leaves the value computed at fill time. While a
|
||||
declaration reads a column, that column cannot be renamed, retyped
|
||||
or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by
|
||||
the server. Cannot be combined with ``transforms``.
|
||||
@@ -6378,8 +6377,8 @@ class AsyncTable:
|
||||
|
||||
Declared with ``add_columns(computed=...)``, a column starts empty and
|
||||
gets its values here. Rows appended since the last refresh are filled
|
||||
by the next one, and rows whose inputs changed since they were computed
|
||||
are recomputed; everything else is left as it is.
|
||||
by the next one; rows already filled are left as they are, so the call
|
||||
is idempotent and does not observe a mutated input.
|
||||
|
||||
Local tables only: a remote refresh runs as a server job, through
|
||||
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
|
||||
@@ -4183,14 +4183,13 @@ def test_refresh_column_async_returns_job(tmp_path):
|
||||
assert result.rows_failed == 0
|
||||
assert result.rows_remaining == 0
|
||||
assert result.source_version == 2
|
||||
# The fill lands at 3; the stamp recording its inputs is published at 4.
|
||||
assert result.published_version == 4
|
||||
assert result.published_version == 3
|
||||
assert job.status() == "finished"
|
||||
assert sorted(table.to_arrow()["doubled"].to_pylist()) == [2, 4]
|
||||
|
||||
no_op = table.refresh_column_async("doubled").wait()
|
||||
assert no_op.rows_assigned == 0
|
||||
assert no_op.source_version == 4
|
||||
assert no_op.source_version == 3
|
||||
assert no_op.published_version is None
|
||||
|
||||
# Bad input raises at the call, not through the job.
|
||||
@@ -4209,6 +4208,6 @@ async def test_refresh_column_async_job_async_table(tmp_path):
|
||||
assert isinstance(result, lancedb.RefreshColumnResult)
|
||||
assert result.rows_assigned == 1
|
||||
assert result.source_version == 2
|
||||
assert result.published_version == 4
|
||||
assert result.published_version == 3
|
||||
assert await job.status() == "finished"
|
||||
assert (await table.to_arrow())["tripled"].to_pylist() == [9]
|
||||
|
||||
@@ -95,14 +95,13 @@ candle-transformers = { version = "0.9.1", optional = true }
|
||||
candle-nn = { version = "0.9.1", optional = true }
|
||||
tokenizers = { version = "0.19.1", optional = true }
|
||||
semver = { workspace = true }
|
||||
roaring = "0.11.4"
|
||||
sha2 = "0.10"
|
||||
|
||||
[dev-dependencies]
|
||||
anyhow = "1"
|
||||
lance-testing = { workspace = true }
|
||||
tempfile = { workspace = true }
|
||||
random_word = { version = "0.4.3", features = ["en"] }
|
||||
roaring = "0.11.4"
|
||||
tokio = { workspace = true, features = ["io-util", "macros", "net", "test-util"] }
|
||||
uuid = { workspace = true }
|
||||
walkdir = "2"
|
||||
|
||||
@@ -1124,9 +1124,8 @@ struct RowScope {
|
||||
|
||||
/// Whether every commit on the view after `recorded` is a fill of its
|
||||
/// computed columns: a column rewrite or data replacement touching only
|
||||
/// those fields and neither adding nor removing rows, or the freshness
|
||||
/// stamp a fill leaves on them. A version whose transaction cannot be read
|
||||
/// is not proven, so it counts as drift.
|
||||
/// those fields and neither adding nor removing rows. A version whose
|
||||
/// transaction cannot be read is not proven, so it counts as drift.
|
||||
async fn only_computed_rewrites_since(view_ds: &Dataset, recorded: u64) -> Result<bool> {
|
||||
// A fill may write any field under a computed column, so the whole
|
||||
// subtree counts, not only the root.
|
||||
@@ -1177,19 +1176,6 @@ async fn only_computed_rewrites_since(view_ds: &Dataset, recorded: u64) -> Resul
|
||||
.all(|field| computed_fields.contains(&(*field as u32)))
|
||||
})
|
||||
}
|
||||
// The stamp `refresh_column` writes after its fill (see
|
||||
// `table::freshness`): field metadata on computed columns, no data.
|
||||
Operation::UpdateConfig {
|
||||
config_updates: None,
|
||||
table_metadata_updates: None,
|
||||
schema_metadata_updates: None,
|
||||
field_metadata_updates,
|
||||
} => {
|
||||
!field_metadata_updates.is_empty()
|
||||
&& field_metadata_updates
|
||||
.keys()
|
||||
.all(|field| computed_fields.contains(&(*field as u32)))
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
if !fill {
|
||||
@@ -3373,46 +3359,6 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// Field metadata on `field` only, the commit shape of the freshness
|
||||
/// stamp `refresh_column` leaves after its fill.
|
||||
async fn commit_field_metadata(view: &MaterializedView, field: &str, key: &str) {
|
||||
let native = view.table().as_native().unwrap();
|
||||
native.dataset.reload().await.unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.update_field_metadata()
|
||||
.update(field, [(key.to_string(), "{}".to_string())])
|
||||
.unwrap()
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// The stamp is metadata on the computed column and rewrites nothing
|
||||
/// refresh certifies, so it is not drift; the same commit shape on a
|
||||
/// projected column is, like any other write to it.
|
||||
#[tokio::test]
|
||||
async fn test_a_freshness_stamp_is_not_drift() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
|
||||
commit_field_metadata(
|
||||
&view,
|
||||
"emb",
|
||||
crate::table::computed_columns::SOURCE_SIGNATURE_META_KEY,
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
|
||||
commit_field_metadata(&view, "id", "probe").await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
}
|
||||
|
||||
/// The fill job's commit rewrites only computed columns. It is the one
|
||||
/// commit on a view that is not drift: the next refresh carries on from
|
||||
/// its watermark instead of rebuilding, which would null what the fill
|
||||
@@ -3578,9 +3524,9 @@ mod tests {
|
||||
}
|
||||
|
||||
/// A SQL declaration is filled by `refresh_column` on the view, which
|
||||
/// commits a data replacement and then its freshness stamp; the next
|
||||
/// refresh continues from its watermark and keeps what the fill wrote,
|
||||
/// and only rows the view added since come back unfilled.
|
||||
/// commits a data replacement; the next refresh continues from its
|
||||
/// watermark and keeps what the fill wrote, and only rows the view added
|
||||
/// since come back unfilled.
|
||||
#[tokio::test]
|
||||
async fn test_a_sql_fill_is_not_drift() {
|
||||
use crate::materialized_view::tests::{people, sql_field};
|
||||
|
||||
@@ -75,7 +75,6 @@ mod create_index;
|
||||
pub mod datafusion;
|
||||
pub(crate) mod dataset;
|
||||
pub mod delete;
|
||||
pub mod freshness;
|
||||
pub mod lsm_stats;
|
||||
pub mod merge;
|
||||
pub mod optimize;
|
||||
@@ -779,8 +778,7 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
|
||||
message: "Function columns are supported only on LanceDB Cloud and Enterprise".into(),
|
||||
})
|
||||
}
|
||||
/// Fill a computed column's unfilled rows and recompute those whose
|
||||
/// inputs changed.
|
||||
/// Fill a computed column's unfilled rows.
|
||||
///
|
||||
/// The default returns `NotSupported`; Lance-backed tables override it.
|
||||
async fn refresh_column(&self, _column: &str) -> Result<RefreshColumnResult> {
|
||||
@@ -788,8 +786,8 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
|
||||
message: "computed columns are supported only on local tables".into(),
|
||||
})
|
||||
}
|
||||
/// Fill a computed column's unfilled rows and recompute those whose
|
||||
/// inputs changed, returning a [`Job`] tracking the operation.
|
||||
/// Fill a computed column's unfilled rows, returning a [`Job`] tracking
|
||||
/// the operation.
|
||||
async fn refresh_column_async(
|
||||
&self,
|
||||
_column: &str,
|
||||
@@ -1751,10 +1749,9 @@ impl Table {
|
||||
/// Declared with
|
||||
/// [`AddColumnsBuilder::computed`](add_columns::AddColumnsBuilder::computed),
|
||||
/// a column starts empty and gets its values here. Fragments appended
|
||||
/// since the last refresh are filled by the next one, and fragments whose
|
||||
/// inputs changed since they were computed are recomputed (see
|
||||
/// [`freshness`](crate::table::freshness)); everything else is left as
|
||||
/// it is.
|
||||
/// since the last refresh are filled by the next one; fragments already
|
||||
/// filled are left as they are, so the call is idempotent and does not
|
||||
/// observe a mutated input.
|
||||
///
|
||||
/// Local tables only: a remote refresh runs as a server job, through
|
||||
/// [`Table::refresh_column_async`].
|
||||
|
||||
@@ -60,11 +60,10 @@ impl AddColumnsBuilder {
|
||||
/// every fragment that has none -- including fragments appended since the
|
||||
/// last refresh.
|
||||
///
|
||||
/// A refresh also recomputes the rows of a fragment whose inputs changed
|
||||
/// since it was computed (see [`freshness`](super::freshness)), so a
|
||||
/// mutated input is reflected by the next refresh. An input cannot be
|
||||
/// renamed, retyped or dropped while a declaration reads it, since the
|
||||
/// expression names it.
|
||||
/// Refresh does not revisit a fragment it has filled, so mutating an input
|
||||
/// leaves the value computed at fill time; recomputing means dropping the
|
||||
/// column and declaring it again. An input cannot be renamed, retyped or
|
||||
/// dropped while a declaration reads it, since the expression names it.
|
||||
///
|
||||
/// On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
/// server, and the refresh runs as a server job -- see
|
||||
|
||||
@@ -71,22 +71,6 @@ pub const FUNCTION_BINDINGS_META_KEY: &str = "lancedb::function_bindings";
|
||||
/// Version of the schema-level Function binding envelope.
|
||||
pub const FUNCTION_BINDINGS_VERSION: u32 = 1;
|
||||
|
||||
/// Field metadata key holding `{fragment id -> input signature}` as JSON,
|
||||
/// recorded by the refresh that last computed each fragment. Outside the
|
||||
/// declaration namespace on purpose: a declaration is immutable through
|
||||
/// metadata edits, this is rewritten by every refresh. Seeded empty at
|
||||
/// declaration, so a column is tracked from birth; a column without it was
|
||||
/// declared before signatures existed.
|
||||
pub const SOURCE_SIGNATURE_META_KEY: &str = "computed_refresh.source_signature";
|
||||
|
||||
/// Field metadata key holding the definition digest a column was last
|
||||
/// computed under. A change to it makes every row stale.
|
||||
pub const DEFINITION_VERSION_META_KEY: &str = "computed_refresh.definition_version";
|
||||
|
||||
/// Field metadata key holding the table version the signature map describes:
|
||||
/// where a refresh starts following compactions to carry freshness forward.
|
||||
pub const RECORDED_AT_VERSION_META_KEY: &str = "computed_refresh.recorded_at_version";
|
||||
|
||||
/// Value of [`KIND_META_KEY`] for a column defined by a SQL expression.
|
||||
pub const SQL_KIND: &str = "sql";
|
||||
|
||||
@@ -155,7 +139,6 @@ fn computed_column_metadata(expression: &str, inputs: &[String]) -> HashMap<Stri
|
||||
INPUTS_META_KEY.to_string(),
|
||||
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
|
||||
),
|
||||
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
|
||||
])
|
||||
}
|
||||
|
||||
@@ -180,7 +163,6 @@ pub fn function_computed_column_metadata(
|
||||
INPUTS_META_KEY.to_string(),
|
||||
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
|
||||
),
|
||||
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
|
||||
])
|
||||
}
|
||||
|
||||
@@ -1313,8 +1295,7 @@ pub(crate) fn ensure_not_an_input(schema: &SchemaRef, paths: &[&str]) -> Result<
|
||||
}
|
||||
|
||||
/// Reject a write that supplies values for a computed column directly:
|
||||
/// only refresh materializes one, and only refresh decides what it
|
||||
/// recomputes.
|
||||
/// only refresh materializes one, and refresh never revisits a filled row.
|
||||
pub(crate) fn ensure_not_written<'a>(
|
||||
schema: &ArrowSchema,
|
||||
written: impl IntoIterator<Item = &'a str>,
|
||||
@@ -1489,9 +1470,7 @@ fn ensure_no_foreign_declaration(field: &ArrowField) -> Result<()> {
|
||||
/// kind, the expression, the inputs -- would bypass that validation or move
|
||||
/// a binding out from under a refresh. Drop the column and declare it again.
|
||||
pub(crate) fn is_declaration_key(key: &str) -> bool {
|
||||
key == COMPUTED_COLUMN_META_KEY
|
||||
|| key.starts_with("computed_column.")
|
||||
|| key.starts_with("computed_refresh.")
|
||||
key == COMPUTED_COLUMN_META_KEY || key.starts_with("computed_column.")
|
||||
}
|
||||
|
||||
/// Reject retyping a computed column itself.
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -134,17 +134,9 @@ pub(crate) async fn cleanup_old_versions(
|
||||
) -> Result<RemovalStats> {
|
||||
table.dataset.ensure_mutable()?;
|
||||
let dataset = table.dataset.get().await?;
|
||||
let stats = dataset
|
||||
Ok(dataset
|
||||
.cleanup_old_versions(older_than, delete_unverified, error_if_tagged_old_versions)
|
||||
.await?;
|
||||
// Computed-column signature sidecars live outside lance's directories;
|
||||
// drop the ones the surviving versions no longer reference.
|
||||
let removed =
|
||||
super::freshness::prune_sidecars(&dataset, delete_unverified.unwrap_or(false)).await?;
|
||||
if removed > 0 {
|
||||
log::debug!("removed {removed} unreferenced computed-column signature sidecars");
|
||||
}
|
||||
Ok(stats)
|
||||
.await?)
|
||||
}
|
||||
|
||||
/// Compact files in the dataset.
|
||||
|
||||
@@ -3,9 +3,9 @@
|
||||
|
||||
//! Filling computed columns.
|
||||
//!
|
||||
//! A row without a value gets one; a row that has one keeps it unless its
|
||||
//! fragment's inputs moved since it was computed, which `freshness` decides
|
||||
//! from the manifest and stamps after every fill.
|
||||
//! A row without a value gets one; a row that has one keeps it. Refresh is
|
||||
//! therefore idempotent and does not observe input mutation -- once a row is
|
||||
//! filled, changing what the expression reads leaves the stored result alone.
|
||||
//!
|
||||
//! A column's computed inputs are filled first -- the dependency graph is
|
||||
//! walked once, each reachable column filled once in dependency order, each
|
||||
@@ -31,7 +31,6 @@
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
use arrow_array::{
|
||||
Array, ArrayRef, BooleanArray, LargeBinaryArray, RecordBatch, RecordBatchOptions, StructArray,
|
||||
@@ -49,7 +48,6 @@ use lance_core::datatypes::{BlobHandling, Schema as LanceSchema};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use super::computed_columns::{BoundExpression, ComputedColumnKind, computed_column_from_field};
|
||||
use super::freshness::{self, SignatureMap, StalenessPlan};
|
||||
use super::{BaseTable, NativeTable};
|
||||
use crate::job::Job;
|
||||
use crate::{Error, Result};
|
||||
@@ -112,81 +110,28 @@ async fn execute_refresh_column_with_source(
|
||||
};
|
||||
let output_is_blob = field.is_blob_v2();
|
||||
|
||||
// Which fragments the null filter cannot speak for: their inputs moved
|
||||
// since they were computed, or the definition did. Decided once, from
|
||||
// the manifest the values are read from.
|
||||
let inputs = freshness::fields_for_paths(dataset.schema(), &bound.inputs)?;
|
||||
let definition = freshness::definition_version(&expression);
|
||||
let staleness = freshness::staleness_against(&dataset, column, &definition, &inputs).await?;
|
||||
|
||||
let mut rows_filled = 0u64;
|
||||
let mut replacements = Vec::new();
|
||||
// Fragments this refresh computed in full, signed at the version read.
|
||||
let mut computed = SignatureMap::new();
|
||||
for fragment in dataset.get_fragments() {
|
||||
let fragment_id = u32::try_from(fragment.id()).map_err(|_| Error::Runtime {
|
||||
message: format!("fragment id {} does not fit a signature map", fragment.id()),
|
||||
})?;
|
||||
// A recompute rewrites every live row, so it is staged without the
|
||||
// probe and counted as it fills; a null fill probes first, since a
|
||||
// fragment with nothing to gain is not worth a write.
|
||||
let recompute = staleness.is_dirty(fragment_id);
|
||||
let whole = recompute || {
|
||||
let (gained, unfilled) =
|
||||
count_fragment_gains(&dataset, &fragment, &bound, column).await?;
|
||||
if gained == 0 {
|
||||
continue;
|
||||
}
|
||||
rows_filled += gained;
|
||||
unfilled == u64::try_from(fragment.count_rows(None).await?).unwrap_or(u64::MAX)
|
||||
};
|
||||
if whole {
|
||||
computed.insert(
|
||||
fragment_id,
|
||||
freshness::fragment_input_signature(fragment.metadata(), &inputs)?,
|
||||
);
|
||||
let gained = count_fragment_gains(&dataset, &fragment, &bound, column).await?;
|
||||
if gained == 0 {
|
||||
continue;
|
||||
}
|
||||
let gained = Arc::new(AtomicU64::new(0));
|
||||
let values = fill_stream(
|
||||
&dataset,
|
||||
&fragment,
|
||||
bound.clone(),
|
||||
column,
|
||||
output_is_blob,
|
||||
recompute,
|
||||
gained.clone(),
|
||||
)
|
||||
.await?;
|
||||
rows_filled += gained;
|
||||
let values =
|
||||
fill_stream(&dataset, &fragment, bound.clone(), column, output_is_blob).await?;
|
||||
replacements.push(fragment.write_columns(values, &column_schema).await?);
|
||||
if recompute {
|
||||
rows_filled += gained.load(Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
let source_version = dataset.version().version;
|
||||
if replacements.is_empty() {
|
||||
// Nothing to fill; the stamp may still have something to record -- a
|
||||
// column not yet enrolled, or fragments a compaction carried.
|
||||
let mut latest = (*dataset).clone();
|
||||
let stamped = record(
|
||||
&mut latest,
|
||||
(&dataset, &staleness),
|
||||
column,
|
||||
&definition,
|
||||
&inputs,
|
||||
computed,
|
||||
)
|
||||
.await;
|
||||
if stamped.is_some() {
|
||||
table.dataset.update(latest);
|
||||
}
|
||||
return Ok(RefreshExecution {
|
||||
result: RefreshColumnResult {
|
||||
rows_filled: 0,
|
||||
version: stamped.unwrap_or(source_version),
|
||||
version: source_version,
|
||||
},
|
||||
source_version,
|
||||
published_version: stamped,
|
||||
published_version: None,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -204,17 +149,7 @@ async fn execute_refresh_column_with_source(
|
||||
)
|
||||
.await?;
|
||||
|
||||
let mut new_dataset = new_dataset;
|
||||
let version = record(
|
||||
&mut new_dataset,
|
||||
(&dataset, &staleness),
|
||||
column,
|
||||
&definition,
|
||||
&inputs,
|
||||
computed,
|
||||
)
|
||||
.await
|
||||
.unwrap_or(new_dataset.version().version);
|
||||
let version = new_dataset.version().version;
|
||||
table.dataset.update(new_dataset);
|
||||
Ok(RefreshExecution {
|
||||
result: RefreshColumnResult {
|
||||
@@ -226,31 +161,6 @@ async fn execute_refresh_column_with_source(
|
||||
})
|
||||
}
|
||||
|
||||
/// Stamp the input state the refresh computed from (see
|
||||
/// [`freshness::record_freshness`]); the version the stamp landed at, which
|
||||
/// is the last one the refresh wrote. Never fails the refresh: the values
|
||||
/// are committed, and a missing stamp only costs a recompute next time.
|
||||
async fn record(
|
||||
latest: &mut Dataset,
|
||||
pinned: (&Dataset, &StalenessPlan),
|
||||
column: &str,
|
||||
definition: &str,
|
||||
inputs: &freshness::InputFields,
|
||||
computed: SignatureMap,
|
||||
) -> Option<u64> {
|
||||
match freshness::record_freshness(latest, Some(pinned), column, definition, inputs, computed)
|
||||
.await
|
||||
{
|
||||
Ok(record) => record.version,
|
||||
Err(error) => {
|
||||
log::warn!(
|
||||
"could not record the input state computed column '{column}' was refreshed from ({error}); its fragments will recompute on the next refresh"
|
||||
);
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Refuse while a computed input still has rows a refresh of it would fill:
|
||||
/// read now, its placeholder null would be evaluated as a value and kept.
|
||||
async fn ensure_inputs_filled(
|
||||
@@ -278,9 +188,7 @@ async fn ensure_inputs_filled(
|
||||
let input_bound = super::computed_columns::bind(schema.clone(), input, expression)?;
|
||||
let mut unfilled = 0u64;
|
||||
for fragment in dataset.get_fragments() {
|
||||
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input)
|
||||
.await?
|
||||
.0;
|
||||
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input).await?;
|
||||
}
|
||||
if unfilled > 0 {
|
||||
return Err(Error::InvalidInput {
|
||||
@@ -528,35 +436,32 @@ fn blob_array_from_binary(
|
||||
/// Scans only the unfilled live rows -- deleted rows never reach the
|
||||
/// expression here, the filter having already excluded them -- and counts the
|
||||
/// non-null results. Exact, so it is both the staging decision and the
|
||||
/// fragment's contribution to `rows_filled`. Returns the gains and the rows
|
||||
/// scanned.
|
||||
/// fragment's contribution to `rows_filled`.
|
||||
async fn count_fragment_gains(
|
||||
dataset: &Dataset,
|
||||
fragment: &FileFragment,
|
||||
bound: &BoundExpression,
|
||||
column: &str,
|
||||
) -> Result<(u64, u64)> {
|
||||
) -> Result<u64> {
|
||||
let mut scanner = dataset.scan();
|
||||
scanner
|
||||
.with_fragments(vec![fragment.metadata().clone()])
|
||||
.with_row_id()
|
||||
.project(&bound.roots)?
|
||||
.filter(&format!("{} IS NULL", quote_identifier(column)))?;
|
||||
.filter(&format!("{} IS NULL", quote_identifier(column)))?
|
||||
.project(&bound.roots)?;
|
||||
configure_blob_inputs(&mut scanner, dataset.schema(), bound, None)?;
|
||||
|
||||
let mut gained = 0u64;
|
||||
let mut considered = 0u64;
|
||||
let mut batches = scanner.try_into_stream().await?;
|
||||
while let Some(batch) = batches.try_next().await? {
|
||||
let evaluated = evaluate(bound, &evaluation_batch(&batch, bound, None)?)?;
|
||||
gained += (batch.num_rows() - evaluated.null_count()) as u64;
|
||||
considered += batch.num_rows() as u64;
|
||||
}
|
||||
Ok((gained, considered))
|
||||
Ok(gained)
|
||||
}
|
||||
|
||||
/// Stream one fragment's column in physical order, filling the unfilled live
|
||||
/// rows -- every live row, for a recompute -- and keeping every other value.
|
||||
/// rows and keeping every other value.
|
||||
///
|
||||
/// Deleted rows are carried through so the values line up positionally with
|
||||
/// the fragment's data files; they are never read back, but the column file
|
||||
@@ -567,8 +472,6 @@ async fn fill_stream(
|
||||
bound: Arc<BoundExpression>,
|
||||
column: &str,
|
||||
output_is_blob: bool,
|
||||
recompute: bool,
|
||||
gained: Arc<AtomicU64>,
|
||||
) -> Result<impl Stream<Item = lance_core::Result<RecordBatch>> + Send + use<>> {
|
||||
let mut projection: Vec<String> = bound.roots.clone();
|
||||
projection.push(column.to_string());
|
||||
@@ -618,20 +521,14 @@ async fn fill_stream(
|
||||
.column_by_name(ROW_ID)
|
||||
.ok_or_else(|| missing(ROW_ID))?;
|
||||
|
||||
// Only an unfilled live row gains a value, or every live row under a
|
||||
// recompute; a deleted row has a null row id and keeps its (null) slot.
|
||||
// Only an unfilled live row gains a value; a deleted row has a null
|
||||
// row id and keeps its (null) slot.
|
||||
let unfilled = arrow::compute::is_null(existing.as_ref())?;
|
||||
let live = arrow::compute::is_not_null(row_ids.as_ref())?;
|
||||
let fill = if recompute {
|
||||
live
|
||||
} else {
|
||||
let unfilled = arrow::compute::is_null(existing.as_ref())?;
|
||||
arrow::compute::and(&unfilled, &live)?
|
||||
};
|
||||
let fill = arrow::compute::and(&unfilled, &live)?;
|
||||
let keep = arrow::compute::not(&fill)?;
|
||||
|
||||
let computed = evaluate(&bound, &evaluation_batch(&batch, &bound, Some(&keep))?)?;
|
||||
let values = arrow::compute::and(&fill, &arrow::compute::is_not_null(&computed)?)?;
|
||||
gained.fetch_add(values.true_count() as u64, Ordering::Relaxed);
|
||||
let merged = arrow_select::zip::zip(&fill, &computed, existing)?;
|
||||
let merged = if output_is_blob {
|
||||
blob_array_from_binary(&merged, projected.field(0))?
|
||||
@@ -840,8 +737,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(no_op.rows_assigned, 0);
|
||||
// The fill, then the stamp recording what it computed from.
|
||||
assert_eq!(no_op.source_version, 4);
|
||||
assert_eq!(no_op.source_version, 3);
|
||||
assert_eq!(no_op.published_version, None);
|
||||
}
|
||||
|
||||
@@ -881,8 +777,8 @@ mod tests {
|
||||
}
|
||||
|
||||
/// A row is filled only by gaining a value, so an expression yielding null
|
||||
/// settles at once instead of re-selecting the same rows forever: the
|
||||
/// second refresh finds the fragment signed and moves nothing.
|
||||
/// settles at once instead of re-selecting the same rows forever. Nothing
|
||||
/// is staged, so the version does not move either.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_converges_on_a_null_result() {
|
||||
let table = table_with("refresh_null_result", vec![1, 2, 3]).await;
|
||||
@@ -896,186 +792,28 @@ mod tests {
|
||||
|
||||
let first = table.refresh_column("maybe").await.unwrap();
|
||||
assert_eq!(first.rows_filled, 0);
|
||||
assert!(first.version > declared);
|
||||
assert_eq!(first.version, declared);
|
||||
assert_eq!(read(&table, "maybe").await, vec![None, None, None]);
|
||||
|
||||
let again = table.refresh_column("maybe").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 0);
|
||||
assert_eq!(again.version, first.version);
|
||||
assert_eq!(again.version, declared);
|
||||
}
|
||||
|
||||
/// A filled row whose input moved is recomputed: the update rewrites
|
||||
/// the row into a fragment the stamp never signed, and only that one.
|
||||
/// The contract's boundary: a filled fragment is not revisited, so
|
||||
/// mutating an input leaves the value computed at fill time.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_recomputes_a_row_whose_input_moved() {
|
||||
let table = table_with("refresh_mutation", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
|
||||
table
|
||||
.update()
|
||||
.column("x", "7")
|
||||
.only_if("x = 5")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 1);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(14)]
|
||||
);
|
||||
|
||||
let settled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(settled.rows_filled, 0);
|
||||
assert_eq!(settled.version, again.version);
|
||||
}
|
||||
|
||||
/// Each stamp is a sidecar under `_computed/`; pruning old versions
|
||||
/// removes the sidecars only they referenced, and keeps the current one.
|
||||
#[tokio::test]
|
||||
async fn test_pruning_drops_the_sidecars_of_pruned_versions() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let conn = connect(dir.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let batch = record_batch!(("x", Int32, [1, 2])).unwrap();
|
||||
let table = conn
|
||||
.create_table("sidecars", batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
async fn test_refresh_does_not_observe_input_mutation() {
|
||||
let table = table_with("refresh_mutation", vec![1]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
let sidecars = || {
|
||||
std::fs::read_dir(dir.path().join("sidecars.lance").join("_computed"))
|
||||
.unwrap()
|
||||
.count()
|
||||
};
|
||||
assert_eq!(sidecars(), 2);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
|
||||
|
||||
table
|
||||
.optimize(crate::table::OptimizeAction::Prune {
|
||||
older_than: Some(chrono::Duration::zero()),
|
||||
delete_unverified: Some(true),
|
||||
error_if_tagged_old_versions: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(sidecars(), 1);
|
||||
assert_eq!(
|
||||
table.refresh_column("doubled").await.unwrap().rows_filled,
|
||||
0
|
||||
);
|
||||
}
|
||||
|
||||
/// A deleted row is never computed and the rows that stay keep their
|
||||
/// values: a delete recomputes nothing and stamps nothing.
|
||||
#[tokio::test]
|
||||
async fn test_a_delete_recomputes_nothing() {
|
||||
let table = table_with("refresh_delete", vec![1, 2, 3]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
let filled = table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
table.delete("x = 2").await.unwrap();
|
||||
let deleted = table.version().await.unwrap();
|
||||
table.update().column("x", "3").execute().await.unwrap();
|
||||
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 0);
|
||||
assert_eq!(again.version, deleted);
|
||||
assert!(deleted > filled.version);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(6)]);
|
||||
}
|
||||
|
||||
/// Compaction copies inputs unchanged, so a fragment it builds from
|
||||
/// signed ones is fresh: the refresh recomputes nothing and only records
|
||||
/// the new fragment.
|
||||
#[tokio::test]
|
||||
async fn test_a_compaction_of_signed_fragments_recomputes_nothing() {
|
||||
let table = table_with("refresh_compact_signed", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
table
|
||||
.optimize(crate::table::OptimizeAction::Compact {
|
||||
options: crate::table::CompactionOptions::default(),
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
let compacted = table.version().await.unwrap();
|
||||
|
||||
let carried = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(carried.rows_filled, 0);
|
||||
assert_eq!(carried.version, compacted + 1);
|
||||
let settled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(settled.version, carried.version);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
}
|
||||
|
||||
/// A column declared before signatures existed has no map. Its first
|
||||
/// refresh keeps the null-fill contract and enrolls what it read from;
|
||||
/// from then on a moved input is recomputed like any other.
|
||||
#[tokio::test]
|
||||
async fn test_an_unsigned_column_is_enrolled_by_its_first_refresh() {
|
||||
let table = table_with("refresh_legacy", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
table
|
||||
.update()
|
||||
.column("x", "3")
|
||||
.only_if("x = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let native = table.as_native().unwrap();
|
||||
let mut dataset = (*native.dataset.get().await.unwrap()).clone();
|
||||
let declaration = dataset
|
||||
.schema()
|
||||
.field("doubled")
|
||||
.unwrap()
|
||||
.metadata
|
||||
.iter()
|
||||
.filter(|(key, _)| !key.starts_with("computed_refresh."))
|
||||
.map(|(key, value)| (key.clone(), value.clone()))
|
||||
.collect::<Vec<_>>();
|
||||
dataset
|
||||
.update_field_metadata()
|
||||
.replace("doubled", declaration)
|
||||
.unwrap()
|
||||
.await
|
||||
.unwrap();
|
||||
table.checkout_latest().await.unwrap();
|
||||
|
||||
// Null-fill only: the moved row keeps the value it was filled with.
|
||||
let enrolled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(enrolled.rows_filled, 0);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(4)]);
|
||||
|
||||
table
|
||||
.update()
|
||||
.column("x", "5")
|
||||
.only_if("x = 3")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 1);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(4), Some(10)]);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
|
||||
}
|
||||
|
||||
/// A row rewrite before the first refresh materializes the declared
|
||||
@@ -1093,11 +831,11 @@ mod tests {
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(6)]);
|
||||
}
|
||||
|
||||
/// A fragment compacted out of one the stamp never signed cannot vouch
|
||||
/// for any of its rows: every live row is recomputed, the moved one
|
||||
/// included.
|
||||
/// The contract holds row by row, not fragment by fragment: revisiting a
|
||||
/// fragment to fill one row must not recompute a filled row sitting beside
|
||||
/// it, even where the input behind it has since changed.
|
||||
#[tokio::test]
|
||||
async fn test_a_compaction_of_an_unsigned_fragment_recomputes_it() {
|
||||
async fn test_refresh_does_not_recompute_a_filled_row_beside_an_unfilled_one() {
|
||||
let table = table_with("refresh_mixed", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
@@ -1119,70 +857,18 @@ mod tests {
|
||||
.unwrap();
|
||||
|
||||
let result = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 3);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(4), Some(10), Some(200)]
|
||||
);
|
||||
}
|
||||
|
||||
/// The gate's reproducer: a raw lance append may carry a value for the
|
||||
/// computed column. Compaction cannot certify it, so the product is
|
||||
/// recomputed and the supplied value replaced.
|
||||
#[tokio::test]
|
||||
async fn test_raw_append_values_are_not_trusted_after_compaction() {
|
||||
use arrow_array::RecordBatchIterator;
|
||||
use lance::Dataset;
|
||||
use lance::dataset::{WriteMode, WriteParams};
|
||||
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let conn = connect(dir.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let batch = record_batch!(("x", Int32, [1, 2])).unwrap();
|
||||
let table = conn
|
||||
.create_table("raw_append", batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
let batch = record_batch!(("x", Int32, [5]), ("doubled", Int32, [Some(999_i32)])).unwrap();
|
||||
let schema = batch.schema();
|
||||
let uri = table.uri().await.unwrap();
|
||||
Dataset::write(
|
||||
RecordBatchIterator::new(vec![Ok(batch)], schema),
|
||||
&uri,
|
||||
Some(WriteParams {
|
||||
mode: WriteMode::Append,
|
||||
..Default::default()
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
table.checkout_latest().await.unwrap();
|
||||
table
|
||||
.optimize(crate::table::OptimizeAction::Compact {
|
||||
options: crate::table::CompactionOptions::default(),
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let result = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 3);
|
||||
assert_eq!(result.rows_filled, 1);
|
||||
// 2 is the mutated row keeping the value it was filled with, not 200.
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
}
|
||||
|
||||
/// An appended fragment holds no values, so compacting it into a signed
|
||||
/// one leaves the product fresh: only the appended rows are filled.
|
||||
/// Filling a fragment must not disturb the values it already holds, which
|
||||
/// is what makes a compaction-mixed fragment safe to revisit.
|
||||
#[tokio::test]
|
||||
async fn test_a_compaction_with_an_appended_fragment_fills_only_its_rows() {
|
||||
async fn test_refresh_preserves_already_filled_rows() {
|
||||
let table = table_with("refresh_preserves", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
@@ -1309,8 +995,7 @@ mod tests {
|
||||
assert_eq!(result.rows_failed, 0);
|
||||
assert_eq!(result.rows_remaining, 0);
|
||||
assert_eq!(result.source_version, 2);
|
||||
// The fill lands at 3; the stamp recording its inputs is published at 4.
|
||||
assert_eq!(result.published_version, Some(4));
|
||||
assert_eq!(result.published_version, Some(3));
|
||||
assert_eq!(job.status().await.unwrap(), "finished");
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
@@ -1374,32 +1059,31 @@ mod tests {
|
||||
assert_eq!(read(&table, "quotient").await, vec![Some(10)]);
|
||||
}
|
||||
|
||||
/// A filled row whose input moved is re-evaluated, and a row whose
|
||||
/// input did not move is not: the untouched fragment is never read, so
|
||||
/// its poison input is never reached.
|
||||
/// The gate's reproducer: an already-filled row's value must not be
|
||||
/// re-evaluated either -- its input may have mutated into one the
|
||||
/// expression chokes on.
|
||||
#[tokio::test]
|
||||
async fn test_only_a_moved_rows_value_is_re_evaluated() {
|
||||
let table = table_with("refresh_filled_poison", vec![1, 0]).await;
|
||||
async fn test_a_filled_rows_value_is_never_evaluated() {
|
||||
let table = table_with("refresh_filled_poison", vec![1, 2]).await;
|
||||
table
|
||||
.add_columns()
|
||||
.computed("quotient", "10 / coalesce(nullif(x, 0), 1)")
|
||||
.computed("quotient", "10 / x")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
table.refresh_column("quotient").await.unwrap();
|
||||
assert_eq!(read(&table, "quotient").await, vec![Some(10), Some(10)]);
|
||||
|
||||
append(&table, vec![5]).await;
|
||||
table
|
||||
.update()
|
||||
.column("x", "2")
|
||||
.column("x", "0")
|
||||
.only_if("x = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
|
||||
let result = table.refresh_column("quotient").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 2);
|
||||
assert_eq!(result.rows_filled, 1);
|
||||
assert_eq!(
|
||||
read(&table, "quotient").await,
|
||||
vec![Some(2), Some(5), Some(10)]
|
||||
|
||||
Reference in New Issue
Block a user