mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-22 04:55:39 +00:00
feat(node): add blob v2 fetch and field helpers (#4155)
this PR blob v2 field helpers and reads to the Node SDK.
`blob()` marks a field as blob v2 and lets you set the storage
thresholds. Inputs can be bytes, a URI, or a data/uri struct.
Queries return descriptors. `fetchBlobs()` reads the bytes by row ID,
and `fetchBlobFiles()` gives you lazy handles for full or range reads.
`blobColumns()` lists the blob fields, including nested ones.
Fetch uses the table’s current checkout. It preserves order, duplicates,
and nulls. Holding row IDs across compaction still requires stable row
IDs.
```javascript
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 header = await handle!.readRange(0n, 65536n);
```
### Testing
- cover input validation, thresholds, nested fields, fetch ordering,
nulls, and range reads.
This commit is contained in:
+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,
|
||||
|
||||
+75
-16
@@ -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
|
||||
@@ -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);
|
||||
});
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user