Compare commits

..

10 Commits

Author SHA1 Message Date
Wyatt Alt c463ca1503 docs: state what a computed column promises
The declaration API shipped without a runnable example, and none of the three
binding docs said when values appear, what happens to them when an input
changes, which schema operations a declaration blocks, or that the feature is
local-only. Those are the questions a caller has to answer before using it.

Adds Rust doctests on both entry points and the same semantics to the Python
and TypeScript parameter docs, plus a worked Python example. Also adds the
abstract refresh_column that both concrete Python tables already implemented,
so the surface is declared in one place and the cross-references resolve.
2026-08-12 21:01:11 -07:00
Wyatt Alt 79e9dffd06 fix: make computed columns explicitly local-only
Both operations were advertised on remote tables and neither could work. A
declaration reaches the wire as AllNulls, which RemoteTable::add_columns does
not accept, and refresh_column fell through to the BaseTable default; the
TypeScript wrapper reached the same surface. Callers got errors that named
neither the feature nor the reason.

The remote protocol has no representation for a stored expression, and what
one should look like is not settled -- the server persists a declaration under
a different vocabulary. So this states the boundary rather than guessing at a
wire format: an explicit AllNulls arm, a default that names computed columns,
and NotImplementedError raised in Python before the round trip.
2026-08-12 20:00:18 -07:00
Wyatt Alt c3efc320a6 fix(rust): refuse schema changes that invalidate a computed column
A declaration records the columns its expression reads, but nothing consulted
them: renaming an input left an expression naming a column that no longer
exists, and the failure surfaced at refresh time as a plan error rather than
at the operation that caused it. Dropping or retyping an input did the same.

alter_columns and drop_columns now reject a change to a column some
declaration reads. Nullability is not part of what an expression resolves
against, so it stays allowed. A declaration does not read itself and so
travels with its own binding, and paths compare at their root, since a change
to `metadata.age` invalidates an expression reading `metadata` just as surely.

Binding to field ids instead would leave the expression text naming the old
column, so it would need rewriting stored SQL on every rename. Refusing the
operation is what a generated column does elsewhere.
2026-08-12 20:00:18 -07:00
Wyatt Alt 8d6dea6313 fix(rust): fill computed columns row by row
Refresh selected rows with `{column} IS NULL` and rewrote them through an
UPDATE, which made the output value double as the record of whether the row
had been computed. Two consequences: an expression yielding null re-selected
the same rows on every run and reported them as filled forever, and the
target name was interpolated into SQL unquoted, so a column named
`double value` could be declared and never refreshed.

Filling is now per fragment. Each fragment that could hold an unfilled row
has the expression evaluated over its physical rows and the result written as
a standalone column file, published together in one DataReplacement. A row
that already holds a value keeps it -- the computed and current values are
merged on the is-null mask -- and a row counts as filled only when it gains a
value, so a fragment where nothing would change is never staged and a null
expression settles after one pass.

Which fragments are worth looking at comes from the manifest first: one whose
data files do not carry the field cannot hold a filled row. A fragment that
does carry it is still asked, because a row rewrite -- an update, or a
compaction folding an unfilled fragment into a filled one -- leaves nulls
behind a covering file. That case is the reason coverage alone is not the
marker; the test for it fails against a coverage-only implementation.

Names now reach the evaluator through a projection alias or a backtick-quoted
identifier, lance's dialect having no other way to spell one -- a
double-quoted name parses as a string literal.
2026-08-12 20:00:18 -07:00
Wyatt Alt 5a2d3f39e2 chore: update lance dependency to v11.0.0-beta.7
Filling a computed column needs FileFragment::write_column, which lands in
this release.
2026-08-12 20:00:18 -07:00
Wyatt Alt 219f41339d refactor(rust): tag a computed column's definition by kind
ComputedColumn carried a bare expression string, which asserts that every
computed column is a SQL expression. That holds for the only kind there is,
but it is the wrong shape for the next one: a column defined by a registered
function cannot be typed by parsing its definition, so its type and its inputs
have a different provenance than a SQL column's. A single string has nowhere
to say which it is.

The definition is now ComputedColumnKind, non-exhaustive so another kind is
additive, and the metadata carries a matching computed_column.kind tag beside
the payload. Tagging the persisted form is the point -- the Rust type stays
cheap to change and field metadata does not, and a second kind distinguished
only by which keys happen to be present would leave every reader sniffing the
shape.

An unrecognized kind reads back as Unrecognized rather than as absent. A newer
version's declaration is a computed column this one cannot evaluate, not a
plain column: reported as absent it would be redeclarable over and would fail
refresh as "not a computed column".

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-12 18:48:48 -07:00
Wyatt Alt 5982eebbb3 feat(nodejs): expose computed columns and refreshColumn
addComputedColumns declares columns defined by a SQL expression, and
refreshColumn fills a computed column's unfilled rows:

    await table.addComputedColumns([{ name: "doubled", valueSql: "x * 2" }]);
    await table.refreshColumn("doubled");
2026-08-12 18:48:48 -07:00
Wyatt Alt 5a1c839382 feat(python): expose computed columns and refresh_column
add_columns gains a computed= mapping of column name to SQL expression, and
refresh_column fills a computed column's unfilled rows. Both are available on
the sync and async tables:

    table.add_columns(computed={"doubled": "x * 2"})
    table.refresh_column("doubled")

transforms and computed are mutually exclusive, since they commit through
different paths and could half-apply.
2026-08-12 18:48:48 -07:00
Wyatt Alt e40e073a9d feat(rust): fill computed columns with refresh_column
Declaring a computed column stores its expression but computes nothing, so
until now the column stayed null with no way to fill it. refresh_column
evaluates the expression over the rows that still hold no value and commits the
results:

    table.refresh_column("doubled").await?

Rows without a value are the ones to fill, which also makes the operation
idempotent and resumable after a failure: refreshing again picks up whatever
did not land. A row whose expression evaluates to null is indistinguishable
from an unfilled one and is recomputed, which costs work but cannot change the
result.

Values written after a refresh are reachable by the next one, which is the case
that matters -- an expression column populated once and then appended to would
otherwise read null for every later row forever.
2026-08-12 18:48:48 -07:00
Wyatt Alt a47c22b26e feat(rust): declare computed columns through add_columns
A computed column is a column defined by a SQL expression rather than by
values supplied at write time, so it is added through add_columns like any
other:

    table.add_columns().computed("doubled", "x * 2").execute().await?

Declaring does not compute. The column is committed carrying its expression in
field metadata but no data, which makes declaration cost the same on a large
table as on an empty one and leaves a single code path that ever produces
values. A later refresh fills it.

The expression is the whole definition: both the result type and the input
columns are derived from it with lance-datafusion's planner, so a caller writes
neither, and the two can never disagree the way a hand-declared input list can.
Everything statically knowable is rejected at declare time rather than deferred:
an expression that does not parse, one referencing a column that does not
exist, a name already in use, and the same name declared twice in one call.

The binding lives in three field metadata keys: virtual_column marks the
column, virtual_column.expression holds it, and virtual_column.inputs holds the
parsed inputs as a JSON array. computed_columns() reads them back off a schema,
the way a SQL catalog reports a generation expression as another column of its
information schema, so introspection needs no round trip.

A transform and computed columns cannot be combined in one call, since they
commit through different transforms and could half-apply.

Only functions the query engine already knows can be named; resolving a
user-defined one needs a planner aware of the function registry, which does not
exist yet.
2026-08-12 18:48:48 -07:00
30 changed files with 2037 additions and 1021 deletions
-10
View File
@@ -69,16 +69,6 @@ jobs:
uses: actions/setup-python@v6
with:
python-version: "3.10"
- name: Add swap for Arm fat LTO
if: matrix.config.platform == 'aarch64'
shell: bash
run: |
swap_file="$RUNNER_TEMP/lancedb-swap"
sudo fallocate --length 16G "$swap_file"
sudo chmod 600 "$swap_file"
sudo mkswap "$swap_file"
sudo swapon "$swap_file"
free -h
- uses: ./.github/workflows/build_linux_wheel
with:
python-minor-version: 10
Generated
-1
View File
@@ -5495,7 +5495,6 @@ dependencies = [
"urlencoding",
"uuid",
"walkdir",
"windows-sys 0.61.2",
]
[[package]]
+53 -1
View File
@@ -69,14 +69,33 @@ abstract addColumns(newColumnTransforms): Promise<AddColumnsResult>
Add new columns with defined values.
The `{ computed }` form stores the expression rather than evaluating it
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 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.
Computed columns are local-only: LanceDB Cloud and Enterprise reject a
declaration.
#### Parameters
* **newColumnTransforms**: `Field`&lt;`any`&gt; \| `Field`&lt;`any`&gt;[] \| `Schema`&lt;`any`&gt; \| [`AddColumnsSql`](../interfaces/AddColumnsSql.md)[]
* **newColumnTransforms**:
\| `Field`&lt;`any`&gt;
\| `Field`&lt;`any`&gt;[]
\| `Schema`&lt;`any`&gt;
\| [`AddColumnsSql`](../interfaces/AddColumnsSql.md)[]
\| `object`
Either:
- An array of objects with column names and SQL expressions to calculate values
- A single Arrow Field defining one column with its data type (column will be initialized with null values)
- An array of Arrow Fields defining columns with their data types (columns will be initialized with null values)
- An Arrow Schema defining columns with their data types (columns will be initialized with null values)
- `{ computed }`, declaring columns defined by a SQL expression whose type and inputs are derived from it
#### Returns
@@ -85,6 +104,13 @@ Add new columns with defined values.
A promise that resolves to an object
containing the new version number of the table after adding the columns.
#### Example
```ts
await table.addColumns({ computed: [{ name: "doubled", valueSql: "x * 2" }] });
const { rowsFilled } = await table.refreshColumn("doubled");
```
***
### alterColumns()
@@ -718,6 +744,32 @@ for await (const batch of table.query()) {
***
### refreshColumn()
```ts
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; rows
already filled are left as they are, so the call is idempotent and does
not observe a mutated input. Local tables only.
#### Parameters
* **column**: `string`
The name of the computed column to fill.
#### Returns
`Promise`&lt;[`RefreshColumnResult`](../interfaces/RefreshColumnResult.md)&gt;
A promise that resolves to the
number of rows filled and the new version number of the table.
***
### restore()
```ts
+1
View File
@@ -105,6 +105,7 @@
- [OptimizeOptions](interfaces/OptimizeOptions.md)
- [OptimizeStats](interfaces/OptimizeStats.md)
- [QueryExecutionOptions](interfaces/QueryExecutionOptions.md)
- [RefreshColumnResult](interfaces/RefreshColumnResult.md)
- [RemovalStats](interfaces/RemovalStats.md)
- [RenameTableOptions](interfaces/RenameTableOptions.md)
- [RestNamespaceConfig](interfaces/RestNamespaceConfig.md)
@@ -0,0 +1,23 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / RefreshColumnResult
# Interface: RefreshColumnResult
## Properties
### rowsFilled
```ts
rowsFilled: number;
```
***
### version
```ts
version: number;
```
+1 -1
View File
@@ -28,7 +28,7 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<arrow.version>15.0.0</arrow.version>
<lance-core.version>11.0.0-beta.7</lance-core.version>
<lance-core.version>11.0.0-beta.6</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>
+42
View File
@@ -3340,3 +3340,45 @@ describe("LSM merge insert", () => {
await expect(table.query().useLsm(true).toArray()).rejects.toThrow();
});
});
describe("computed columns", () => {
let tmpDir: tmp.DirResult;
beforeEach(() => {
tmpDir = tmp.dirSync({ unsafeCleanup: true });
});
afterEach(() => tmpDir.removeCallback());
it("declares a column and fills it on refresh", async () => {
const db = await connect(tmpDir.name);
const table = await db.createTable("computed", [{ x: 1 }, { x: 2 }]);
await table.addColumns({
computed: [{ name: "doubled", valueSql: "x * 2" }],
});
let rows = await table.query().toArray();
expect(rows.map((r) => r.doubled)).toEqual([null, null]);
const result = await table.refreshColumn("doubled");
expect(result.rowsFilled).toBe(2);
rows = await table.query().toArray();
expect(rows.map((r) => r.doubled).sort()).toEqual([2, 4]);
});
it("fills rows added since the last refresh", async () => {
const db = await connect(tmpDir.name);
const table = await db.createTable("computed_append", [{ x: 1 }]);
await table.addColumns({
computed: [{ name: "doubled", valueSql: "x * 2" }],
});
await table.refreshColumn("doubled");
await table.add([{ x: 5 }]);
const result = await table.refreshColumn("doubled");
expect(result.rowsFilled).toBe(1);
const rows = await table.query().toArray();
expect(rows.map((r) => r.doubled).sort()).toEqual([10, 2]);
});
});
+1
View File
@@ -50,6 +50,7 @@ export {
MergeResult,
AddResult,
AddColumnsResult,
RefreshColumnResult,
AlterColumnsResult,
UpdateFieldMetadataResult,
DeleteResult,
+57 -2
View File
@@ -33,6 +33,7 @@ import {
Job,
Branches as NativeBranches,
OptimizeStats,
RefreshColumnResult,
TableStatistics,
Tags,
UpdateFieldMetadataResult,
@@ -525,18 +526,54 @@ export abstract class Table {
abstract vectorSearch(vector: IntoVector | MultiVector): VectorQuery;
/**
* Add new columns with defined values.
*
* The `{ computed }` form stores the expression rather than evaluating it
* now: the column is committed with no values, and rows get them from
* {@link Table#refreshColumn}. Declaring one therefore costs the same on a
* large table as on an empty one.
*
* 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.
*
* Computed columns are local-only: LanceDB Cloud and Enterprise reject a
* declaration.
* @param {AddColumnsSql[] | Field | Field[] | Schema} newColumnTransforms Either:
* - An array of objects with column names and SQL expressions to calculate values
* - A single Arrow Field defining one column with its data type (column will be initialized with null values)
* - An array of Arrow Fields defining columns with their data types (columns will be initialized with null values)
* - An Arrow Schema defining columns with their data types (columns will be initialized with null values)
* - `{ computed }`, declaring columns defined by a SQL expression whose type and inputs are derived from it
* @returns {Promise<AddColumnsResult>} A promise that resolves to an object
* containing the new version number of the table after adding the columns.
* @example
* ```ts
* await table.addColumns({ computed: [{ name: "doubled", valueSql: "x * 2" }] });
* const { rowsFilled } = await table.refreshColumn("doubled");
* ```
*/
abstract addColumns(
newColumnTransforms: AddColumnsSql[] | Field | Field[] | Schema,
newColumnTransforms:
| AddColumnsSql[]
| Field
| Field[]
| Schema
| { computed: AddColumnsSql[] },
): Promise<AddColumnsResult>;
/**
* Fill the rows of a computed column that hold no value yet.
*
* 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.
* @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.
*/
abstract refreshColumn(column: string): Promise<RefreshColumnResult>;
/**
* Alter the name or nullability of columns.
* @param {ColumnAlteration[]} columnAlterations One or more alterations to
@@ -1088,8 +1125,22 @@ export class LocalTable extends Table {
// TODO: Support BatchUDF
async addColumns(
newColumnTransforms: AddColumnsSql[] | Field | Field[] | Schema,
newColumnTransforms:
| AddColumnsSql[]
| Field
| Field[]
| Schema
| { computed: AddColumnsSql[] },
): Promise<AddColumnsResult> {
// Columns defined by an expression are declared, not materialized here.
if (
typeof newColumnTransforms === "object" &&
!Array.isArray(newColumnTransforms) &&
"computed" in newColumnTransforms
) {
return await this.inner.addComputedColumns(newColumnTransforms.computed);
}
// Handle single Field -> convert to array of Fields
if (newColumnTransforms instanceof Field) {
newColumnTransforms = [newColumnTransforms];
@@ -1124,6 +1175,10 @@ export class LocalTable extends Table {
throw new Error("Invalid input type for addColumns");
}
async refreshColumn(column: string): Promise<RefreshColumnResult> {
return await this.inner.refreshColumn(column);
}
async alterColumns(
columnAlterations: ColumnAlteration[],
): Promise<AlterColumnsResult> {
+39
View File
@@ -347,6 +347,30 @@ impl Table {
Ok(res.into())
}
#[napi(catch_unwind)]
pub async fn add_computed_columns(
&self,
columns: Vec<AddColumnsSql>,
) -> napi::Result<AddColumnsResult> {
let table = self.inner_ref()?;
let mut builder = table.add_columns();
for column in columns {
builder = builder.computed(column.name, column.value_sql);
}
let res = builder.execute().await.default_error()?;
Ok(res.into())
}
#[napi(catch_unwind)]
pub async fn refresh_column(&self, column: String) -> napi::Result<RefreshColumnResult> {
let res = self
.inner_ref()?
.refresh_column(column)
.await
.default_error()?;
Ok(res.into())
}
#[napi(catch_unwind)]
pub async fn add_columns_with_schema(
&self,
@@ -1196,6 +1220,21 @@ pub struct AddColumnsResult {
pub version: i64,
}
#[napi(object)]
pub struct RefreshColumnResult {
pub rows_filled: i64,
pub version: i64,
}
impl From<lancedb::table::RefreshColumnResult> for RefreshColumnResult {
fn from(value: lancedb::table::RefreshColumnResult) -> Self {
Self {
rows_filled: value.rows_filled as i64,
version: value.version as i64,
}
}
}
impl From<lancedb::table::AddColumnsResult> for AddColumnsResult {
fn from(value: lancedb::table::AddColumnsResult) -> Self {
Self {
+8
View File
@@ -335,6 +335,10 @@ class Table:
) -> list[FtsToken]: ...
async def delete(self, filter: Union[str, PyExpr]) -> DeleteResult: ...
async def add_columns(self, columns: list[tuple[str, str]]) -> AddColumnsResult: ...
async def add_computed_columns(
self, columns: list[tuple[str, str]]
) -> AddColumnsResult: ...
async def refresh_column(self, column: str) -> RefreshColumnResult: ...
async def add_columns_with_schema(self, schema: pa.Schema) -> AddColumnsResult: ...
async def alter_columns(
self, columns: list[dict[str, Any]]
@@ -680,6 +684,10 @@ class LsmWriteSpec:
class AddColumnsResult:
version: int
class RefreshColumnResult:
rows_filled: int
version: int
class AlterColumnsResult:
version: int
+13 -1
View File
@@ -958,9 +958,21 @@ class RemoteTable(Table):
def count_rows(self, filter: Optional[str] = None) -> int:
return LOOP.run(self._table.count_rows(filter))
def add_columns(self, transforms: Dict[str, str]) -> AddColumnsResult:
def add_columns(
self,
transforms: Dict[str, str] | None = None,
*,
computed: Dict[str, str] | None = None,
) -> AddColumnsResult:
if computed:
raise NotImplementedError(
"computed columns are supported only on local tables"
)
return LOOP.run(self._table.add_columns(transforms))
def refresh_column(self, column: str):
raise NotImplementedError("computed columns are supported only on local tables")
def alter_columns(
self, *alterations: Iterable[Dict[str, str]]
) -> AlterColumnsResult:
+135 -4
View File
@@ -176,6 +176,7 @@ if TYPE_CHECKING:
CompactionStats,
Tag,
AddColumnsResult,
RefreshColumnResult,
AddResult,
AlterColumnsResult,
UpdateFieldMetadataResult,
@@ -1916,7 +1917,14 @@ class Table(ABC):
@abstractmethod
def add_columns(
self, transforms: Dict[str, str] | pa.Field | List[pa.Field] | pa.Schema
self,
transforms: Dict[str, str]
| pa.Field
| List[pa.Field]
| pa.Schema
| None = None,
*,
computed: Dict[str, str] | None = None,
):
"""
Add new columns with defined values.
@@ -1930,11 +1938,68 @@ class Table(ABC):
Alternatively, a pyarrow Field or Schema can be provided to add
new columns with the specified data types. The new columns will
be initialized with null values.
computed: Dict[str, str], optional
A map of column name to a SQL expression defining the column. The
column's type and inputs are derived from the expression, so no
data type is supplied.
Unlike ``transforms``, the expression is stored rather than
evaluated now: the column is committed with no values, and rows get
them from [`refresh_column`][lancedb.table.Table.refresh_column].
Declaring one therefore costs the same on a large table as on an
empty one.
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.
Local tables only; LanceDB Cloud and Enterprise raise
``NotImplementedError``. Cannot be combined with ``transforms``.
Returns
-------
AddColumnsResult
version: the new version number of the table after adding columns.
Examples
--------
>>> import lancedb
>>> db = lancedb.connect("./.lancedb")
>>> table = db.create_table("computed_demo", [{"x": 1}, {"x": 2}])
>>> table.add_columns(computed={"doubled": "x * 2"})
AddColumnsResult(version=2)
>>> table.refresh_column("doubled")
RefreshColumnResult(rows_filled=2, version=3)
>>> table.to_arrow().sort_by("x").to_pandas()
x doubled
0 1 2
1 2 4
"""
@abstractmethod
def refresh_column(self, column: str) -> "RefreshColumnResult":
"""
Fill the rows of a computed column that hold no value yet.
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; rows already filled are left as they are, so the call
is idempotent and does not observe a mutated input.
Local tables only; LanceDB Cloud and Enterprise raise
``NotImplementedError``.
Parameters
----------
column: str
The name of the computed column to fill.
Returns
-------
RefreshColumnResult
rows_filled: the number of rows given a value.
version: the new version number of the table.
"""
@abstractmethod
@@ -3939,9 +4004,21 @@ class LanceTable(Table):
return LOOP.run(self._table.index_stats(index_name))
def add_columns(
self, transforms: Dict[str, str] | pa.field | List[pa.field] | pa.Schema
self,
transforms: Dict[str, str]
| pa.field
| List[pa.field]
| pa.Schema
| None = None,
*,
computed: Dict[str, str] | None = None,
) -> AddColumnsResult:
return LOOP.run(self._table.add_columns(transforms))
return LOOP.run(self._table.add_columns(transforms, computed=computed))
def refresh_column(self, column: str) -> "RefreshColumnResult":
"""Fill a computed column's unfilled rows. See
[`AsyncTable.refresh_column`][lancedb.AsyncTable.refresh_column]."""
return LOOP.run(self._table.refresh_column(column))
def alter_columns(
self, *alterations: Iterable[Dict[str, str]]
@@ -5856,7 +5933,14 @@ class AsyncTable:
return await self._inner.update(updates_sql, where)
async def add_columns(
self, transforms: dict[str, str] | pa.field | List[pa.field] | pa.Schema
self,
transforms: dict[str, str]
| pa.field
| List[pa.field]
| pa.Schema
| None = None,
*,
computed: dict[str, str] | None = None,
) -> AddColumnsResult:
"""
Add new columns with defined values.
@@ -5869,6 +5953,21 @@ class AsyncTable:
each row in the table, and can reference existing columns.
Alternatively, you can pass a pyarrow field or schema to add
new columns with NULLs.
computed: Dict[str, str], optional
A map of column name to a SQL expression defining the column. The
column's type and inputs are derived from the expression.
Unlike ``transforms``, the expression is stored rather than
evaluated now: the column is committed with no values, and rows get
them from
[`refresh_column`][lancedb.table.AsyncTable.refresh_column].
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.
Local tables only. Cannot be combined with ``transforms``.
Returns
-------
@@ -5882,11 +5981,43 @@ class AsyncTable:
{isinstance(f, pa.Field) for f in transforms}
):
transforms = pa.schema(transforms)
if computed:
if transforms:
raise ValueError(
"add_columns cannot take both transforms and computed columns"
)
return await self._inner.add_computed_columns(list(computed.items()))
if transforms is None:
raise ValueError("add_columns requires transforms or computed columns")
if isinstance(transforms, pa.Schema):
return await self._inner.add_columns_with_schema(transforms)
else:
return await self._inner.add_columns(list(transforms.items()))
async def refresh_column(self, column: str) -> RefreshColumnResult:
"""
Fill the rows of a computed column that hold no value yet.
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; rows already filled are left as they are, so the call
is idempotent and does not observe a mutated input.
Local tables only; LanceDB Cloud and Enterprise raise
``NotImplementedError``.
Parameters
----------
column: str
The name of the computed column to fill.
Returns
-------
RefreshColumnResult
The number of rows filled and the new version of the table.
"""
return await self._inner.refresh_column(column)
async def alter_columns(
self, *alterations: Iterable[dict[str, Any]]
) -> AlterColumnsResult:
+34
View File
@@ -3854,3 +3854,37 @@ async def test_async_search_runs_embedding_on_dedicated_executor(
assert all(name.startswith("lancedb-embedding") for name in captured_threads), (
f"embedding ran off the dedicated executor: {captured_threads}"
)
def test_computed_column_declare_and_refresh(tmp_path):
db = lancedb.connect(tmp_path)
table = db.create_table("computed", [{"x": 1}, {"x": 2}])
table.add_columns(computed={"doubled": "x * 2"})
assert table.to_arrow()["doubled"].to_pylist() == [None, None]
result = table.refresh_column("doubled")
assert result.rows_filled == 2
assert sorted(table.to_arrow()["doubled"].to_pylist()) == [2, 4]
table.add([{"x": 5}])
assert table.refresh_column("doubled").rows_filled == 1
assert sorted(table.to_arrow()["doubled"].to_pylist()) == [2, 4, 10]
def test_computed_column_rejects_transforms_and_computed_together(tmp_path):
db = lancedb.connect(tmp_path)
table = db.create_table("computed_mixed", [{"x": 1}])
with pytest.raises(ValueError):
table.add_columns({"a": "x + 1"}, computed={"b": "x * 2"})
@pytest.mark.asyncio
async def test_computed_column_async(tmp_path):
db = await lancedb.connect_async(tmp_path)
table = await db.create_table("computed_async", [{"x": 3}])
await table.add_columns(computed={"tripled": "x * 3"})
await table.refresh_column("tripled")
assert (await table.to_arrow())["tripled"].to_pylist() == [9]
+3 -1
View File
@@ -16,7 +16,8 @@ use query::{FTSQuery, HybridQuery, Query, VectorQuery};
use session::Session;
use table::{
AddColumnsResult, AddResult, AlterColumnsResult, DeleteResult, DropColumnsResult, FtsToken,
LsmWriteSpec, MergeResult, PyBlobFile, Table, UpdateFieldMetadataResult, UpdateResult,
LsmWriteSpec, MergeResult, PyBlobFile, RefreshColumnResult, Table, UpdateFieldMetadataResult,
UpdateResult,
};
pub mod arrow;
@@ -57,6 +58,7 @@ pub fn _lancedb(_py: Python, m: &Bound<'_, PyModule>) -> PyResult<()> {
m.add_class::<VectorQuery>()?;
m.add_class::<RecordBatchStream>()?;
m.add_class::<AddColumnsResult>()?;
m.add_class::<RefreshColumnResult>()?;
m.add_class::<AlterColumnsResult>()?;
m.add_class::<UpdateFieldMetadataResult>()?;
m.add_class::<AddResult>()?;
+49
View File
@@ -415,6 +415,32 @@ pub struct AddColumnsResult {
pub version: u64,
}
#[pyclass(get_all, from_py_object)]
#[derive(Clone, Debug)]
pub struct RefreshColumnResult {
pub rows_filled: u64,
pub version: u64,
}
#[pymethods]
impl RefreshColumnResult {
pub fn __repr__(&self) -> String {
format!(
"RefreshColumnResult(rows_filled={}, version={})",
self.rows_filled, self.version
)
}
}
impl From<lancedb::table::RefreshColumnResult> for RefreshColumnResult {
fn from(result: lancedb::table::RefreshColumnResult) -> Self {
Self {
rows_filled: result.rows_filled,
version: result.version,
}
}
}
#[pymethods]
impl AddColumnsResult {
pub fn __repr__(&self) -> String {
@@ -1510,6 +1536,29 @@ impl Table {
})
}
pub fn add_computed_columns(
self_: PyRef<'_, Self>,
columns: Vec<(String, String)>,
) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let mut builder = inner.add_columns();
for (name, expression) in columns {
builder = builder.computed(name, expression);
}
let result = builder.execute().await.infer_error()?;
Ok(AddColumnsResult::from(result))
})
}
pub fn refresh_column(self_: PyRef<'_, Self>, column: String) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let result = inner.refresh_column(column).await.infer_error()?;
Ok(RefreshColumnResult::from(result))
})
}
pub fn add_columns_with_schema(
self_: PyRef<'_, Self>,
schema: PyArrowType<Schema>,
-8
View File
@@ -115,11 +115,6 @@ serial_test = "3"
[target.'cfg(unix)'.dev-dependencies]
pprof = { version = "0.14", features = ["flamegraph"] }
[target.'cfg(windows)'.dependencies]
windows-sys = { version = "0.61", features = [
"Win32_Foundation",
"Win32_Storage_FileSystem",
] }
[features]
default = []
@@ -193,9 +188,6 @@ required-features = ["bedrock"]
[[example]]
name = "bench_streaming_dataloader"
[[example]]
name = "bench_open_missing_table"
[[example]]
name = "simple"
@@ -1,150 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
// Release benchmark for opening a missing table as sibling-table cardinality grows.
//
// The fixture uses real `.lance` directories and marker files. Fixture creation is
// outside the timed section. Defaults intentionally cover 1k, 10k, and 100k siblings
// with 10 warmups and 100 distinct missing-table opens per scale:
//
// ```text
// cargo run --release -p lancedb --example bench_open_missing_table
// ```
//
// `BENCH_SIBLINGS`, `BENCH_WARMUPS`, and `BENCH_TRIALS` override those defaults.
// Reduced settings are useful only as a smoke test. Performance comparisons require
// the same machine, filesystem, fixture sizes, settings, lockfile, and alternating
// baseline/candidate execution order.
use std::time::{Duration, Instant};
use anyhow::{Context, Result, bail};
use lancedb::connection::Connection;
use lancedb::{Error, connect};
use object_store::ObjectStoreExt as _;
use object_store::path::Path;
const MAX_SIBLINGS: usize = 1_000_000;
const MAX_WARMUPS: usize = 10_000;
const MAX_TRIALS: usize = 100_000;
fn env_usize(key: &str, default: usize, max: usize) -> Result<usize> {
let value = match std::env::var(key) {
Ok(value) => value
.parse()
.with_context(|| format!("invalid {key} value: {value}"))?,
Err(std::env::VarError::NotPresent) => default,
Err(error) => return Err(error).with_context(|| format!("reading {key}")),
};
if value == 0 || value > max {
bail!("{key} must be between 1 and {max}");
}
Ok(value)
}
fn sibling_counts() -> Result<Vec<usize>> {
let raw = std::env::var("BENCH_SIBLINGS").unwrap_or_else(|_| "1000,10000,100000".into());
let mut counts = raw
.split(',')
.map(|value| {
value
.trim()
.parse::<usize>()
.with_context(|| format!("invalid BENCH_SIBLINGS value: {value}"))
})
.collect::<Result<Vec<_>>>()?;
counts.sort_unstable();
counts.dedup();
if counts.is_empty() || counts[0] == 0 || counts[counts.len() - 1] > MAX_SIBLINGS {
bail!("BENCH_SIBLINGS values must be between 1 and {MAX_SIBLINGS}");
}
Ok(counts)
}
async fn add_siblings(
store: &object_store::local::LocalFileSystem,
start: usize,
end: usize,
) -> Result<()> {
for index in start..end {
let marker = Path::from(format!("sibling_{index:06}.lance/_marker"));
store
.put(&marker, bytes::Bytes::new().into())
.await
.with_context(|| format!("creating benchmark marker {marker}"))?;
}
Ok(())
}
async fn time_missing_open(db: &Connection, name: &str) -> Result<Duration> {
let started = Instant::now();
let result = db.open_table(name).execute().await;
let elapsed = started.elapsed();
match result {
Err(Error::TableNotFound { .. }) => Ok(elapsed),
Err(error) => bail!("expected TableNotFound for {name}, got {error:?}"),
Ok(_) => bail!("benchmark missing-table name unexpectedly exists: {name}"),
}
}
fn percentile(sorted: &[Duration], percentile: usize) -> Duration {
let rank = (sorted.len() * percentile).div_ceil(100).saturating_sub(1);
sorted[rank]
}
#[tokio::main]
async fn main() -> Result<()> {
let counts = sibling_counts()?;
let warmups = env_usize("BENCH_WARMUPS", 10, MAX_WARMUPS)?;
let trials = env_usize("BENCH_TRIALS", 100, MAX_TRIALS)?;
let fixture = tempfile::tempdir().context("creating benchmark fixture")?;
let database_path = fixture.path();
let fixture_store = object_store::local::LocalFileSystem::new_with_prefix(database_path)
.context("creating benchmark object store")?;
let db = connect(database_path.to_str().context("non-UTF-8 fixture path")?)
.execute()
.await?;
println!(
"config: siblings={counts:?} warmups={warmups} trials={trials} profile={} os={} arch={}",
if cfg!(debug_assertions) {
"debug"
} else {
"release"
},
std::env::consts::OS,
std::env::consts::ARCH,
);
println!("lower is better; fixture setup and teardown are excluded");
println!("| siblings | samples | p50 | p95 | max |");
println!("| ---: | ---: | ---: | ---: | ---: |");
let mut created = 0;
for sibling_count in counts {
add_siblings(&fixture_store, created, sibling_count).await?;
created = sibling_count;
for index in 0..warmups {
let name = format!("__missing_warmup_{sibling_count}_{index}");
let _ = time_missing_open(&db, &name).await?;
}
let mut samples = Vec::with_capacity(trials);
for index in 0..trials {
let name = format!("__missing_trial_{sibling_count}_{index}");
samples.push(time_missing_open(&db, &name).await?);
}
samples.sort_unstable();
println!(
"| {sibling_count} | {} | {:?} | {:?} | {:?} |",
samples.len(),
percentile(&samples, 50),
percentile(&samples, 95),
samples[samples.len() - 1],
);
}
Ok(())
}
+4 -8
View File
@@ -409,11 +409,6 @@ impl Connection {
///
/// The names will be returned in lexicographical order (ascending)
///
/// Listing databases discover physical `*.lance` entries without opening every
/// dataset. The result is a point-in-time discovery snapshot: an entry may still be
/// under creation, may contain only uncommitted storage, or may be concurrently
/// dropped before it is opened.
///
/// The parameters `page_token` and `limit` can be used to paginate the results
pub fn table_names(&self) -> TableNamesBuilder {
TableNamesBuilder::new(self.internal.clone())
@@ -461,9 +456,10 @@ impl Connection {
///
/// # Returns
/// Created [`TableRef`], or [`Error::TableNotFound`] if the table does not exist.
/// On listing databases, a committed Lance manifest is authoritative for table
/// existence. Uncommitted files or a physical `<name>.lance` directory alone do not
/// make a table openable.
/// If the table's storage is present but holds no readable dataset (for example a
/// `<name>.lance` directory left behind by an interrupted drop and re-create, which
/// [`Self::table_names`] still lists) this returns [`Error::TableCorrupted`]
/// instead.
pub fn open_table(&self, name: impl Into<String>) -> OpenTableBuilder {
OpenTableBuilder::new(
self.internal.clone(),
+5 -234
View File
@@ -25,7 +25,7 @@ use crate::database::namespace::LanceNamespaceDatabase;
use crate::error::{CreateDirSnafu, Error, Result};
use crate::io::object_store::MirroringObjectStoreWrapper;
use crate::table::NativeTable;
use crate::utils::{PatchStoreParam, validate_table_name};
use crate::utils::validate_table_name;
use lance_namespace::models::{
CreateNamespaceRequest, CreateNamespaceResponse, DescribeNamespaceRequest,
@@ -355,14 +355,6 @@ impl ListingDatabase {
url.to_string()
}
#[cfg(any(windows, test))]
fn uses_local_file_provider(object_store: &ObjectStore) -> bool {
matches!(
object_store.scheme(),
"file" | "file-object-store" | "file+uring"
)
}
async fn prepare_namespace_root(
uri: &str,
storage_options: &HashMap<String, String>,
@@ -589,17 +581,6 @@ impl ListingDatabase {
}
None => None,
};
#[cfg(windows)]
let write_store_wrapper = if Self::uses_local_file_provider(&object_store) {
// Local manifest commits need create-only rename semantics,
// including on filesystems that do not support hard links.
Some(
Arc::new(crate::io::object_store::windows::WindowsLocalFileSystemWrapper)
as Arc<dyn WrappingObjectStore>,
)
} else {
write_store_wrapper
};
let namespace_database = Self::connect_namespace_database(
&storage_base_uri,
@@ -664,22 +645,12 @@ impl ListingDatabase {
)
.await?;
#[cfg(windows)]
let write_store_wrapper = Self::uses_local_file_provider(&object_store).then(|| {
// Local manifest commits need create-only rename semantics,
// including on filesystems that do not support hard links.
Arc::new(crate::io::object_store::windows::WindowsLocalFileSystemWrapper)
as Arc<dyn WrappingObjectStore>
});
#[cfg(not(windows))]
let write_store_wrapper = None;
Ok(Self {
uri: path.to_string(),
query_string: None,
base_path,
object_store,
store_wrapper: write_store_wrapper,
store_wrapper: None,
read_consistency_interval,
storage_options: HashMap::new(),
storage_options_provider: None,
@@ -1141,12 +1112,6 @@ impl Database for ListingDatabase {
},
..Default::default()
};
let storage_params = match self.store_wrapper.clone() {
Some(wrapper) => Some(storage_params)
.patch_with_store_wrapper(wrapper)?
.expect("patching store params always returns parameters"),
None => storage_params,
};
let read_params = ReadParams {
store_options: Some(storage_params.clone()),
session: Some(self.session.clone()),
@@ -1326,37 +1291,16 @@ impl Database for ListingDatabase {
mod tests {
use super::*;
use crate::Table;
use crate::arrow::{SendableRecordBatchStream, SimpleRecordBatchStream};
use crate::connection::ConnectRequest;
use crate::data::scannable::Scannable;
use crate::database::{CreateTableMode, CreateTableRequest};
use crate::io::object_store::io_tracking::IoStatsHolder;
use crate::query::QueryRequest;
use crate::table::{AnyQuery, WriteOptions};
use arrow_array::{Int32Array, RecordBatch, StringArray};
use arrow_schema::{DataType, Field, Schema, SchemaRef};
use futures::{TryStreamExt, stream::once};
use arrow_schema::{DataType, Field, Schema};
use futures::TryStreamExt;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use tempfile::tempdir;
use tokio::sync::Barrier;
use tokio::time::timeout;
#[derive(Debug)]
struct PassthroughStoreWrapper(Arc<AtomicUsize>);
impl WrappingObjectStore for PassthroughStoreWrapper {
fn wrap(
&self,
_store_prefix: &str,
target: Arc<dyn object_store::ObjectStore>,
) -> Arc<dyn object_store::ObjectStore> {
self.0.fetch_add(1, Ordering::Relaxed);
target
}
}
async fn setup_database() -> (tempfile::TempDir, ListingDatabase) {
let tempdir = tempdir().unwrap();
@@ -1380,114 +1324,6 @@ mod tests {
(tempdir, db)
}
struct BarrierScannable {
batch: RecordBatch,
barrier: Arc<Barrier>,
}
impl Scannable for BarrierScannable {
fn schema(&self) -> SchemaRef {
self.batch.schema()
}
fn scan_as_stream(&mut self) -> SendableRecordBatchStream {
let batch = self.batch.clone();
let schema = batch.schema();
let barrier = self.barrier.clone();
Box::pin(SimpleRecordBatchStream {
schema,
stream: once(async move {
barrier.wait().await;
Ok(batch)
}),
})
}
}
fn create_request(name: &str, data: Box<dyn Scannable>) -> CreateTableRequest {
CreateTableRequest {
name: name.to_string(),
namespace_path: vec![],
data,
mode: CreateTableMode::Create,
write_options: Default::default(),
location: None,
namespace_client: None,
}
}
#[tokio::test]
async fn test_create_ignores_uncommitted_storage_without_manifest() {
let (tmp_dir, db) = setup_database().await;
let data_dir = tmp_dir.path().join("test.lance/data");
std::fs::create_dir_all(&data_dir).unwrap();
std::fs::write(data_dir.join("orphan.lance"), b"uncommitted").unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
let batch =
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1]))]).unwrap();
let table = db
.create_table(create_request("test", Box::new(batch)))
.await
.unwrap();
assert_eq!(table.count_rows(None).await.unwrap(), 1);
}
#[tokio::test]
async fn test_concurrent_create_is_arbitrated_by_manifest_commit() {
let uri = format!("memory:///concurrent-create-{}", uuid::Uuid::new_v4());
let db = crate::connect(&uri).execute().await.unwrap();
let store: Arc<dyn object_store::ObjectStore> =
Arc::new(object_store::memory::InMemory::new());
let table_url = url::Url::parse("memory:///database/test.lance").unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
let batch =
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1]))]).unwrap();
let barrier = Arc::new(Barrier::new(2));
#[allow(deprecated)]
let request = |batch, barrier| {
let mut request = create_request("test", Box::new(BarrierScannable { batch, barrier }));
request.write_options = WriteOptions {
lance_write_params: Some(lance::dataset::WriteParams {
store_params: Some(ObjectStoreParams {
object_store: Some((store.clone(), table_url.clone())),
..Default::default()
}),
commit_handler: Some(Arc::new(
lance_table::io::commit::ConditionalPutCommitHandler,
)),
..Default::default()
}),
};
request
};
let left = db
.database()
.create_table(request(batch.clone(), barrier.clone()));
let right = db.database().create_table(request(batch, barrier));
let (left, right) = timeout(Duration::from_secs(30), async { tokio::join!(left, right) })
.await
.expect("concurrent creates deadlocked");
let results = [left, right];
assert_eq!(
results.iter().filter(|result| result.is_ok()).count(),
1,
"expected one successful create, got {results:?}"
);
assert_eq!(
results
.iter()
.filter(|result| matches!(result, Err(Error::TableAlreadyExists { .. })))
.count(),
1,
"expected one manifest conflict, got {results:?}"
);
}
#[tokio::test]
async fn test_listing_database_root_ops_do_not_create_manifest() {
let tempdir = tempdir().unwrap();
@@ -1542,25 +1378,6 @@ mod tests {
assert!(!tempdir.path().join("__manifest").exists());
}
#[tokio::test]
async fn test_file_object_store_uses_local_file_provider() {
let tempdir = tempdir().unwrap();
let path = tempdir.path().to_string_lossy().replace('\\', "/");
let uri = if path.starts_with('/') {
format!("file-object-store://{path}")
} else {
format!("file-object-store:///{path}")
};
let registry = Arc::new(lance_io::object_store::ObjectStoreRegistry::default());
let (store, _) =
ObjectStore::from_uri_and_params(registry, &uri, &ObjectStoreParams::default())
.await
.unwrap();
assert_eq!(store.scheme(), "file-object-store");
assert!(ListingDatabase::uses_local_file_provider(&store));
}
/// Regression test for https://github.com/lancedb/lancedb/issues/1600.
///
/// Opening a table used to create a separate object-store client instead of
@@ -1584,13 +1401,9 @@ mod tests {
read_consistency_interval: None,
session: Some(session),
};
let mut db = ListingDatabase::connect_with_options(&request)
let db = ListingDatabase::connect_with_options(&request)
.await
.unwrap();
// A connection-level write wrapper must not prevent table opens from
// reusing the connection's registered object store.
let wrapper_calls = Arc::new(AtomicUsize::new(0));
db.store_wrapper = Some(Arc::new(PassthroughStoreWrapper(wrapper_calls.clone())));
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
db.create_table(CreateTableRequest {
@@ -1606,7 +1419,6 @@ mod tests {
.unwrap();
let before_open = registry.stats();
let wrapper_calls_before_open = wrapper_calls.load(Ordering::Relaxed);
for _ in 0..3 {
let table = db
.open_table(OpenTableRequest {
@@ -1626,7 +1438,6 @@ mod tests {
let after_open = registry.stats();
assert_eq!(after_open.misses, before_open.misses);
assert!(after_open.hits >= before_open.hits + 3);
assert!(wrapper_calls.load(Ordering::Relaxed) >= wrapper_calls_before_open + 3);
}
/// Regression test for https://github.com/lancedb/lancedb/issues/3197.
@@ -1770,46 +1581,6 @@ mod tests {
);
}
#[tokio::test]
async fn test_clone_table_uses_connection_store_wrapper() {
let (_tempdir, mut db) = setup_database().await;
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
db.create_table(CreateTableRequest {
name: "source_table".to_string(),
namespace_path: vec![],
data: Box::new(RecordBatch::new_empty(schema)) as Box<dyn Scannable>,
mode: CreateTableMode::Create,
write_options: Default::default(),
location: None,
namespace_client: None,
})
.await
.unwrap();
let source_uri = db.table_uri("source_table").unwrap();
let tracker = IoStatsHolder::default();
db.store_wrapper = Some(Arc::new(tracker.clone()));
let _ = tracker.incremental_stats();
db.clone_table(CloneTableRequest {
target_table_name: "cloned_table".to_string(),
target_namespace_path: vec![],
source_uri,
source_version: None,
source_tag: None,
is_shallow: true,
namespace_client: None,
})
.await
.unwrap();
let stats = tracker.incremental_stats();
assert!(
stats.write_iops > 0,
"clone bypassed the wrapper: {stats:?}"
);
}
#[tokio::test]
async fn test_clone_table_with_data() {
let (_tempdir, db) = setup_database().await;
+8
View File
@@ -71,6 +71,14 @@ pub enum Error {
IndexNotFound { name: String },
#[snafu(display("Embedding function '{name}' was not found. : {reason}"))]
EmbeddingFunctionNotFound { name: String, reason: String },
#[snafu(display("Column '{name}' was not found"))]
ColumnNotFound { name: String },
#[snafu(display("Column '{name}' already exists"))]
ColumnAlreadyExists { name: String },
#[snafu(display("Column '{name}' is not a computed column"))]
NotAComputedColumn { name: String },
#[snafu(display("Invalid expression for column '{column}': {message}"))]
InvalidExpression { column: String, message: String },
#[snafu(display("Table '{name}' already exists"))]
TableAlreadyExists { name: String },
-3
View File
@@ -18,9 +18,6 @@ use async_trait::async_trait;
#[cfg(test)]
pub mod io_tracking;
#[cfg(windows)]
pub(crate) mod windows;
#[derive(Debug)]
struct MirroringObjectStore {
primary: Arc<dyn ObjectStore>,
-223
View File
@@ -1,223 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Windows local filesystem compatibility for atomic manifest commits.
use std::ffi::OsStr;
use std::fmt::{Display, Formatter};
use std::os::windows::ffi::OsStrExt;
use std::path::{Path as StdPath, PathBuf};
use std::sync::Arc;
use bytes::Bytes;
use futures::stream::BoxStream;
use lance::io::WrappingObjectStore;
use object_store::{
CopyOptions, Error, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta,
ObjectStore, PutMultipartOptions, PutOptions, PutPayload, PutResult, RenameOptions,
RenameTargetMode, Result, UploadPart, path::Path,
};
use windows_sys::Win32::Foundation::{ERROR_ALREADY_EXISTS, ERROR_FILE_EXISTS};
use windows_sys::Win32::Storage::FileSystem::MoveFileExW;
const STORE_NAME: &str = "WindowsLocalFileSystem";
/// Uses the Windows move primitive for create-only renames on local stores.
///
/// `object_store` implements create-only local renames with a hard link followed
/// by a delete. Some Windows filesystems do not support hard links, but
/// `MoveFileExW` without `MOVEFILE_REPLACE_EXISTING` provides the same atomic
/// create-only rename semantics without requiring them.
#[derive(Debug, Default)]
pub struct WindowsLocalFileSystemWrapper;
impl WrappingObjectStore for WindowsLocalFileSystemWrapper {
fn wrap(&self, _store_prefix: &str, target: Arc<dyn ObjectStore>) -> Arc<dyn ObjectStore> {
Arc::new(WindowsLocalFileSystem { target })
}
}
#[derive(Debug)]
struct WindowsLocalFileSystem {
target: Arc<dyn ObjectStore>,
}
impl Display for WindowsLocalFileSystem {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "{STORE_NAME}({})", self.target)
}
}
#[async_trait::async_trait]
#[deny(clippy::missing_trait_methods)]
impl ObjectStore for WindowsLocalFileSystem {
async fn put_opts(
&self,
location: &Path,
bytes: PutPayload,
opts: PutOptions,
) -> Result<PutResult> {
self.target.put_opts(location, bytes, opts).await
}
async fn put_multipart_opts(
&self,
location: &Path,
opts: PutMultipartOptions,
) -> Result<Box<dyn MultipartUpload>> {
self.target.put_multipart_opts(location, opts).await
}
async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
self.target.get_opts(location, options).await
}
async fn get_ranges(
&self,
location: &Path,
ranges: &[std::ops::Range<u64>],
) -> Result<Vec<Bytes>> {
self.target.get_ranges(location, ranges).await
}
fn delete_stream(
&self,
locations: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
self.target.delete_stream(locations)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
self.target.list(prefix)
}
fn list_with_offset(
&self,
prefix: Option<&Path>,
offset: &Path,
) -> BoxStream<'static, Result<ObjectMeta>> {
self.target.list_with_offset(prefix, offset)
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
self.target.list_with_delimiter(prefix).await
}
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
self.target.copy_opts(from, to, options).await
}
async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
if options.target_mode != RenameTargetMode::Create {
return self.target.rename_opts(from, to, options).await;
}
let from = PathBuf::from(from.as_ref());
let to = PathBuf::from(to.as_ref());
tokio::task::spawn_blocking(move || move_file_if_not_exists(&from, &to))
.await
.map_err(|source| Error::Generic {
store: STORE_NAME,
source: Box::new(source),
})?
}
}
fn move_file_if_not_exists(from: &StdPath, to: &StdPath) -> Result<()> {
let from_wide = null_terminated_wide(from.as_os_str());
let to_wide = null_terminated_wide(to.as_os_str());
// SAFETY: both pointers reference null-terminated UTF-16 buffers that remain
// alive for the duration of this call. A zero flag value deliberately omits
// MOVEFILE_REPLACE_EXISTING, giving this operation create-only semantics.
if unsafe { MoveFileExW(from_wide.as_ptr(), to_wide.as_ptr(), 0) } != 0 {
return Ok(());
}
let source = std::io::Error::last_os_error();
let path = to.to_string_lossy().into_owned();
match source.raw_os_error().map(|code| code as u32) {
Some(ERROR_ALREADY_EXISTS | ERROR_FILE_EXISTS) => Err(Error::AlreadyExists {
path,
source: Box::new(source),
}),
_ if source.kind() == std::io::ErrorKind::NotFound => Err(Error::NotFound {
path,
source: Box::new(source),
}),
_ => Err(Error::Generic {
store: STORE_NAME,
source: Box::new(source),
}),
}
}
fn null_terminated_wide(value: &OsStr) -> Vec<u16> {
value.encode_wide().chain(Some(0)).collect()
}
#[cfg(test)]
mod tests {
use object_store::memory::InMemory;
use super::*;
#[tokio::test]
async fn create_only_rename_does_not_use_hard_links() {
let tempdir = tempfile::tempdir().unwrap();
let source_path = tempdir.path().join("staged.manifest");
let destination_path = tempdir.path().join("1.manifest");
std::fs::write(&source_path, b"manifest").unwrap();
let source = Path::from_absolute_path(&source_path).unwrap();
let destination = Path::from_absolute_path(&destination_path).unwrap();
let store = WindowsLocalFileSystem {
// The source does not exist in this inner store. Delegating the
// rename would fail, proving the wrapper uses the native move path.
target: Arc::new(InMemory::new()),
};
store
.rename_opts(
&source,
&destination,
RenameOptions::new().with_target_mode(RenameTargetMode::Create),
)
.await
.unwrap();
assert!(!source_path.exists());
assert_eq!(std::fs::read(destination_path).unwrap(), b"manifest");
}
#[tokio::test]
async fn create_only_rename_preserves_existing_destination() {
let tempdir = tempfile::tempdir().unwrap();
let source_path = tempdir.path().join("staged.manifest");
let destination_path = tempdir.path().join("1.manifest");
std::fs::write(&source_path, b"new manifest").unwrap();
std::fs::write(&destination_path, b"existing manifest").unwrap();
let source = Path::from_absolute_path(&source_path).unwrap();
let destination = Path::from_absolute_path(&destination_path).unwrap();
let store = WindowsLocalFileSystem {
target: Arc::new(InMemory::new()),
};
let error = store
.rename_opts(
&source,
&destination,
RenameOptions::new().with_target_mode(RenameTargetMode::Create),
)
.await
.unwrap_err();
assert!(matches!(error, Error::AlreadyExists { .. }));
assert_eq!(std::fs::read(source_path).unwrap(), b"new manifest");
assert_eq!(
std::fs::read(destination_path).unwrap(),
b"existing manifest"
);
}
}
+38
View File
@@ -2706,6 +2706,13 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
Ok(result)
}
// A declaration reaches here as AllNulls, which the remote protocol
// has no representation for.
NewColumnTransform::AllNulls(_) => {
return Err(Error::NotSupported {
message: "computed columns are supported only on local tables".into(),
});
}
_ => {
return Err(Error::NotSupported {
message: "Only SQL expressions are supported for adding columns".into(),
@@ -6455,6 +6462,37 @@ mod tests {
assert_eq!(result.version, if old_server { 0 } else { 43 });
}
/// Computed columns are local-only. Both halves say so here rather than
/// reaching the wire and failing somewhere less legible.
#[tokio::test]
async fn test_computed_columns_are_refused() {
let table = Table::new_with_handler("my_table", |request| -> http::Response<String> {
panic!("unexpected request: {}", request.url().path())
});
let declared = Arc::new(Schema::new(vec![Field::new(
"doubled",
DataType::Int32,
true,
)]));
let err = table
.add_columns()
.transform(NewColumnTransform::AllNulls(declared))
.execute()
.await
.unwrap_err();
assert!(
matches!(&err, Error::NotSupported { message } if message.contains("local tables")),
"{err:?}"
);
let err = table.refresh_column("doubled").await.unwrap_err();
assert!(
matches!(&err, Error::NotSupported { message } if message.contains("local tables")),
"{err:?}"
);
}
#[tokio::test]
async fn test_prewarm_index() {
let table = Table::new_with_handler("my_table", |request| {
+157 -298
View File
@@ -50,6 +50,7 @@ use crate::DistanceType;
use crate::blob::BlobRangeRequest;
use crate::data::scannable::{PeekedScannable, Scannable, estimate_write_partitions};
use crate::database::Database;
use crate::database::listing::LANCE_FILE_EXTENSION;
use crate::database::read_freshness::TableFreshness;
use crate::embeddings::{EmbeddingDefinition, EmbeddingRegistry, MemoryRegistry};
use crate::error::{Error, Result};
@@ -68,6 +69,7 @@ pub mod add_columns;
mod add_data;
pub mod branch_merge;
pub mod checkpoint;
pub mod computed_columns;
mod create_index;
pub mod datafusion;
pub(crate) mod dataset;
@@ -77,6 +79,7 @@ pub mod merge;
pub mod optimize;
mod primary_key;
pub mod query;
pub mod refresh;
pub mod schema_evolution;
pub mod update;
pub mod write_progress;
@@ -90,6 +93,9 @@ pub use branch_merge::{
MergeBranchResult, MergeBranchStatus, MergePreview, RowCountSummary,
};
pub use chrono::Duration;
pub use computed_columns::{
ComputedColumn, ComputedColumnKind, computed_column_from_field, computed_columns,
};
pub use delete::DeleteResult;
use futures::future::join_all;
pub use lance::dataset::refs::{BranchContents, Ref, TagContents, Tags as LanceTags};
@@ -97,6 +103,7 @@ pub use lance::dataset::scanner::DatasetRecordBatchStream;
pub use lance_index::optimize::OptimizeOptions;
pub use lsm_stats::{BucketStats, GenerationStats, LsmStats, MemtableStats};
pub use optimize::{CompactionOptions, OptimizeAction, OptimizeStats};
pub use refresh::RefreshColumnResult;
pub use schema_evolution::{
AddColumnsResult, AlterColumnsResult, DropColumnsResult, FieldMetadataUpdate,
UpdateFieldMetadataResult,
@@ -151,6 +158,55 @@ pub(crate) fn map_namespace_lance_error(err: lance::Error, table_name: &str) ->
}
}
/// Map a `lance::Error::DatasetNotFound` for the table at `uri` into a `lancedb::Error`.
///
/// Lance reports "there is nothing at this location" and "there is a table directory
/// here but nothing loadable inside it" with the same error. Only the first is a
/// `TableNotFound`: a `<name>.lance` directory left behind by an interrupted drop and
/// re-create is still reported by `Connection::table_names`, so callers need to be able
/// to tell "never existed" from "exists but is broken".
///
/// See <https://github.com/lancedb/lancedb/issues/3127>.
async fn map_dataset_not_found(
uri: &str,
name: &str,
params: ReadParams,
err: lance::Error,
) -> Error {
let name = name.to_string();
let source = Box::new(err);
if table_dir_exists(uri, params).await.unwrap_or(false) {
Error::TableCorrupted { name, source }
} else {
Error::TableNotFound { name, source }
}
}
/// Whether a table directory is present at `uri`, even though no dataset could be
/// loaded from it.
///
/// This looks for a `<name>.lance` entry in the parent directory, which is exactly what
/// `ListingDatabase::table_names` lists, so the two APIs agree on whether a table is
/// present. Probing `uri` itself would not work: object stores have no empty
/// directories to probe, and on a local filesystem the interesting case is precisely an
/// empty directory.
async fn table_dir_exists(uri: &str, params: ReadParams) -> Result<bool> {
let (object_store, path, _) = DatasetBuilder::from_uri(uri)
.with_read_params(params)
.build_object_store()
.await?;
// Only `*.lance` entries are ever reported as tables, so nothing else can produce
// the list-then-open mismatch this guards against.
if path.extension() != Some(LANCE_FILE_EXTENSION) {
return Ok(false);
}
let (Some(parent), Some(dir_name)) = (path.parent(), path.filename()) else {
return Ok(false);
};
let entries = object_store.read_dir(parent).await?;
Ok(entries.iter().any(|entry| entry.as_str() == dir_name))
}
/// Defines the type of column
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ColumnKind {
@@ -732,6 +788,14 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
transforms: NewColumnTransform,
read_columns: Option<Vec<String>>,
) -> Result<AddColumnsResult>;
/// 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> {
Err(Error::NotSupported {
message: "computed columns are supported only on local tables".into(),
})
}
/// Alter columns in the table.
async fn alter_columns(&self, alterations: &[ColumnAlteration]) -> Result<AlterColumnsResult>;
/// Drop columns from the table.
@@ -1624,6 +1688,29 @@ impl Table {
AddColumnsBuilder::new(self.inner.clone())
}
/// Fill the fragments of a computed column that hold no values yet.
///
/// 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; fragments already
/// filled are left as they are, so the call is idempotent and does not
/// observe a mutated input.
///
/// Local tables only.
///
/// ```
/// # use lancedb::Table;
/// # async fn refresh(table: &Table) -> Result<(), Box<dyn std::error::Error>> {
/// let result = table.refresh_column("doubled").await?;
/// println!("filled {} rows at version {}", result.rows_filled, result.version);
/// # Ok(())
/// # }
/// ```
pub async fn refresh_column(&self, column: impl AsRef<str>) -> Result<RefreshColumnResult> {
self.inner.refresh_column(column.as_ref()).await
}
/// Change a column's name or nullability.
pub async fn alter_columns(
&self,
@@ -2335,19 +2422,10 @@ impl NativeTable {
managed_versioning: Option<bool>,
) -> Result<Self> {
let params = params.unwrap_or_default();
let has_caller_store_wrapper = params
.store_options
.as_ref()
.and_then(|options| options.object_store_wrapper.as_ref())
.is_some();
// A caller wrapper must remain outside connection-level compatibility
// behavior. When there is no caller wrapper, apply the compatibility
// layer after loading so the session's registered store can be reused.
let (params, wrapper_after_load) = match write_store_wrapper {
Some(wrapper) if has_caller_store_wrapper => {
(params.patch_with_store_wrapper(wrapper)?, None)
}
wrapper => (params, wrapper),
// patch the params if we have a write store wrapper
let params = match write_store_wrapper.clone() {
Some(wrapper) => params.patch_with_store_wrapper(wrapper)?,
None => params,
};
// Build table_id from namespace + name
@@ -2379,6 +2457,8 @@ impl NativeTable {
None => false,
};
// Kept so that a `DatasetNotFound` can be re-checked against storage below.
let recovery_params = params.clone();
let mut builder = DatasetBuilder::from_uri(uri).with_read_params(params);
// Set up commit handler when managed_versioning is enabled
@@ -2397,23 +2477,10 @@ impl NativeTable {
let dataset = match builder.load().await {
Ok(dataset) => dataset,
Err(e @ lance::Error::DatasetNotFound { .. }) => {
// The manifest load is the existence check. A physical prefix may be
// from a concurrent or abandoned create, so it cannot refine this error.
return Err(Error::TableNotFound {
name: name.to_string(),
source: Box::new(e),
});
return Err(map_dataset_not_found(uri, name, recovery_params, e).await);
}
Err(e) => return Err(e.into()),
};
// Resolve the store from the session registry before applying a
// connection-level write wrapper. Wrapper identity is part of the
// registry key, so including it in ReadParams prevents reuse when the
// opened table (and its wrapped store) is short-lived.
let dataset = match wrapper_after_load {
Some(wrapper) => dataset.with_object_store_wrappers([wrapper]),
None => dataset,
};
let dataset = DatasetConsistencyWrapper::new_latest(dataset, read_consistency_interval);
let id = Self::build_id(&namespace, name);
@@ -2515,16 +2582,11 @@ impl NativeTable {
if let Some(sess) = session {
params.session(sess);
}
let has_caller_store_wrapper = params
.store_options
.as_ref()
.and_then(|options| options.object_store_wrapper.as_ref())
.is_some();
let (params, wrapper_after_load) = match write_store_wrapper {
Some(wrapper) if has_caller_store_wrapper => {
(params.patch_with_store_wrapper(wrapper)?, None)
}
wrapper => (params, wrapper),
// patch the params if we have a write store wrapper
let params = match write_store_wrapper.clone() {
Some(wrapper) => params.patch_with_store_wrapper(wrapper)?,
None => params,
};
// Build table_id from namespace + name
@@ -2548,13 +2610,6 @@ impl NativeTable {
},
e => e.into(),
})?;
// Apply the write wrapper after the session registry has resolved the
// shared store. The cloned dataset retains the wrapper for subsequent
// reads, manifest commits, and any additional base stores.
let dataset = match wrapper_after_load {
Some(wrapper) => dataset.with_object_store_wrappers([wrapper]),
None => dataset,
};
let uri = dataset.uri().to_string();
let dataset = DatasetConsistencyWrapper::new_latest(dataset, read_consistency_interval);
@@ -3323,6 +3378,12 @@ impl BaseTable for NativeTable {
Ok(result)
}
async fn refresh_column(&self, column: &str) -> Result<RefreshColumnResult> {
let result = refresh::execute_refresh_column(self, column).await?;
self.bump_freshness();
Ok(result)
}
async fn alter_columns(&self, alterations: &[ColumnAlteration]) -> Result<AlterColumnsResult> {
let result = schema_evolution::execute_alter_columns(self, alterations).await?;
self.bump_freshness();
@@ -3689,8 +3750,8 @@ pub struct FragmentSummaryStats {
#[cfg(test)]
#[allow(deprecated)]
mod tests {
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use arrow_array::{
@@ -3772,50 +3833,73 @@ mod tests {
);
}
#[tokio::test]
async fn test_open_not_found_when_empty_directory_exists() {
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
std::fs::create_dir(&dataset_path).unwrap();
/// Write a table and then break it, leaving the `<name>.lance` directory in place.
///
/// `remove_all` reproduces an interrupted drop + re-create (the directory is left
/// empty); otherwise only the manifests are removed, leaving the data files behind.
async fn write_then_corrupt_table(dir: &std::path::Path, remove_all: bool) -> String {
let dataset_path = dir.join("test.lance");
let uri = dataset_path.to_str().unwrap().to_string();
let err = NativeTable::open(dataset_path.to_str().unwrap())
.await
.unwrap_err();
let batch = make_test_batches();
let reader = RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
Dataset::write(reader, &uri, None).await.unwrap();
if remove_all {
for entry in std::fs::read_dir(&dataset_path).unwrap() {
let entry = entry.unwrap();
if entry.file_type().unwrap().is_dir() {
std::fs::remove_dir_all(entry.path()).unwrap();
} else {
std::fs::remove_file(entry.path()).unwrap();
}
}
assert_eq!(std::fs::read_dir(&dataset_path).unwrap().count(), 0);
} else {
let versions = dataset_path.join("_versions");
assert!(versions.is_dir(), "expected manifests under {versions:?}");
std::fs::remove_dir_all(&versions).unwrap();
assert!(std::fs::read_dir(&dataset_path).unwrap().count() > 0);
}
uri
}
#[tokio::test]
async fn test_open_corrupt_empty_dir() {
let tmp_dir = tempdir().unwrap();
let uri = write_then_corrupt_table(tmp_dir.path(), true).await;
let err = NativeTable::open(&uri).await.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "test"),
matches!(&err, Error::TableCorrupted { name, .. } if name == "test"),
"got {err:?}"
);
}
#[tokio::test]
async fn test_open_not_found_when_only_uncommitted_storage_exists() {
async fn test_open_corrupt_missing_manifest() {
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let data_dir = dataset_path.join("data");
std::fs::create_dir_all(&data_dir).unwrap();
std::fs::write(data_dir.join("orphan.lance"), b"uncommitted").unwrap();
let uri = write_then_corrupt_table(tmp_dir.path(), false).await;
let err = NativeTable::open(dataset_path.to_str().unwrap())
.await
.unwrap_err();
let err = NativeTable::open(&uri).await.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "test"),
matches!(&err, Error::TableCorrupted { name, .. } if name == "test"),
"got {err:?}"
);
}
/// Listing databases discover physical `*.lance` entries. That snapshot is not an
/// authoritative table-existence check: only a committed manifest makes a table
/// openable, and the entry could also be concurrently created or dropped.
/// A table listed by `table_names()` must not be reported as missing by
/// `open_table()`. See <https://github.com/lancedb/lancedb/issues/3127>.
#[tokio::test]
async fn test_table_names_may_include_uncommitted_storage() {
async fn test_open_table_corrupt_is_still_listed() {
let tmp_dir = tempdir().unwrap();
let db = connect(tmp_dir.path().to_str().unwrap())
.execute()
.await
.unwrap();
std::fs::create_dir(tmp_dir.path().join("test.lance")).unwrap();
write_then_corrupt_table(tmp_dir.path(), true).await;
assert_eq!(
db.table_names().execute().await.unwrap(),
@@ -3823,177 +3907,12 @@ mod tests {
);
let err = db.open_table("test").execute().await.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "test"),
"physical storage without a committed manifest is not a table: {err:?}"
);
}
#[derive(Debug)]
struct ParentListGuardStore {
inner: Arc<dyn object_store::ObjectStore>,
parent: object_store::path::Path,
parent_list_calls: Arc<AtomicUsize>,
}
impl std::fmt::Display for ParentListGuardStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("ParentListGuardStore")
}
}
#[async_trait::async_trait]
#[deny(clippy::missing_trait_methods)]
impl object_store::ObjectStore for ParentListGuardStore {
async fn put_opts(
&self,
location: &object_store::path::Path,
payload: object_store::PutPayload,
opts: object_store::PutOptions,
) -> object_store::Result<object_store::PutResult> {
self.inner.put_opts(location, payload, opts).await
}
async fn put_multipart_opts(
&self,
location: &object_store::path::Path,
opts: object_store::PutMultipartOptions,
) -> object_store::Result<Box<dyn object_store::MultipartUpload>> {
self.inner.put_multipart_opts(location, opts).await
}
async fn get_opts(
&self,
location: &object_store::path::Path,
options: object_store::GetOptions,
) -> object_store::Result<object_store::GetResult> {
self.inner.get_opts(location, options).await
}
async fn get_ranges(
&self,
location: &object_store::path::Path,
ranges: &[std::ops::Range<u64>],
) -> object_store::Result<Vec<bytes::Bytes>> {
self.inner.get_ranges(location, ranges).await
}
fn delete_stream(
&self,
locations: futures::stream::BoxStream<
'static,
object_store::Result<object_store::path::Path>,
>,
) -> futures::stream::BoxStream<'static, object_store::Result<object_store::path::Path>>
{
self.inner.delete_stream(locations)
}
fn list(
&self,
prefix: Option<&object_store::path::Path>,
) -> futures::stream::BoxStream<'static, object_store::Result<object_store::ObjectMeta>>
{
if prefix == Some(&self.parent) {
self.parent_list_calls.fetch_add(1, Ordering::Relaxed);
}
self.inner.list(prefix)
}
fn list_with_offset(
&self,
prefix: Option<&object_store::path::Path>,
offset: &object_store::path::Path,
) -> futures::stream::BoxStream<'static, object_store::Result<object_store::ObjectMeta>>
{
if prefix == Some(&self.parent) {
self.parent_list_calls.fetch_add(1, Ordering::Relaxed);
}
self.inner.list_with_offset(prefix, offset)
}
async fn list_with_delimiter(
&self,
prefix: Option<&object_store::path::Path>,
) -> object_store::Result<object_store::ListResult> {
if prefix == Some(&self.parent) {
self.parent_list_calls.fetch_add(1, Ordering::Relaxed);
}
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(
&self,
from: &object_store::path::Path,
to: &object_store::path::Path,
options: object_store::CopyOptions,
) -> object_store::Result<()> {
self.inner.copy_opts(from, to, options).await
}
async fn rename_opts(
&self,
from: &object_store::path::Path,
to: &object_store::path::Path,
options: object_store::RenameOptions,
) -> object_store::Result<()> {
self.inner.rename_opts(from, to, options).await
}
}
#[derive(Debug)]
struct ParentListGuardWrapper {
parent_list_calls: Arc<AtomicUsize>,
}
impl WrappingObjectStore for ParentListGuardWrapper {
fn wrap(
&self,
_store_prefix: &str,
inner: Arc<dyn object_store::ObjectStore>,
) -> Arc<dyn object_store::ObjectStore> {
Arc::new(ParentListGuardStore {
inner,
parent: object_store::path::Path::from("database"),
parent_list_calls: self.parent_list_calls.clone(),
})
}
}
#[tokio::test]
async fn test_open_missing_never_lists_database_parent() {
let parent_list_calls = Arc::new(AtomicUsize::new(0));
let params = ReadParams {
store_options: Some(ObjectStoreParams {
object_store_wrapper: Some(Arc::new(ParentListGuardWrapper {
parent_list_calls: parent_list_calls.clone(),
})),
..Default::default()
}),
..Default::default()
};
let err = NativeTable::open_with_params(
"memory:///database/missing.lance",
"missing",
Vec::new(),
None,
Some(params),
None,
None,
HashSet::new(),
None,
)
.await
.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "missing"),
matches!(&err, Error::TableCorrupted { name, .. } if name == "test"),
"got {err:?}"
);
assert_eq!(
parent_list_calls.load(Ordering::Relaxed),
0,
"opening one missing table must not enumerate sibling tables"
assert!(
err.to_string().contains("exists but could not be loaded"),
"got {err}"
);
}
@@ -4062,66 +3981,6 @@ mod tests {
}
}
#[derive(Debug)]
struct OrderedStoreWrapper {
name: &'static str,
order: Arc<Mutex<Vec<&'static str>>>,
}
impl WrappingObjectStore for OrderedStoreWrapper {
fn wrap(
&self,
_store_prefix: &str,
original: Arc<dyn object_store::ObjectStore>,
) -> Arc<dyn object_store::ObjectStore> {
self.order.lock().unwrap().push(self.name);
original
}
}
#[tokio::test]
async fn test_open_with_params_keeps_caller_store_wrapper_outermost() {
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let uri = dataset_path.to_str().unwrap();
let batch = make_test_batches();
let reader = RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
Dataset::write(reader, uri, None).await.unwrap();
let order = Arc::new(Mutex::new(Vec::new()));
let caller_wrapper = Arc::new(OrderedStoreWrapper {
name: "caller",
order: order.clone(),
});
let compatibility_wrapper = Arc::new(OrderedStoreWrapper {
name: "compatibility",
order: order.clone(),
});
let params = ReadParams {
store_options: Some(ObjectStoreParams {
object_store_wrapper: Some(caller_wrapper),
..Default::default()
}),
..Default::default()
};
NativeTable::open_with_params(
uri,
"test",
vec![],
Some(compatibility_wrapper),
Some(params),
None,
None,
HashSet::new(),
None,
)
.await
.unwrap();
assert_eq!(*order.lock().unwrap(), vec!["compatibility", "caller"]);
}
#[tokio::test]
async fn test_open_table_options() {
let tmp_dir = tempdir().unwrap();
+120 -22
View File
@@ -8,6 +8,7 @@ use std::sync::Arc;
use lance::dataset::NewColumnTransform;
use super::BaseTable;
use super::computed_columns;
use super::schema_evolution::AddColumnsResult;
use crate::{Error, Result};
@@ -15,6 +16,7 @@ use crate::{Error, Result};
pub struct AddColumnsBuilder {
parent: Arc<dyn BaseTable>,
transform: Option<NewColumnTransform>,
computed: Vec<(String, String)>,
read_columns: Option<Vec<String>>,
}
@@ -23,6 +25,7 @@ impl std::fmt::Debug for AddColumnsBuilder {
f.debug_struct("AddColumnsBuilder")
.field("parent", &self.parent)
.field("has_transform", &self.transform.is_some())
.field("computed", &self.computed)
.field("read_columns", &self.read_columns)
.finish()
}
@@ -33,19 +36,57 @@ impl AddColumnsBuilder {
Self {
parent,
transform: None,
computed: Vec::new(),
read_columns: None,
}
}
/// Set how the new columns' values are produced. Required.
/// Set how the new columns' values are produced.
pub fn transform(mut self, transform: NewColumnTransform) -> Self {
self.transform = Some(transform);
self
}
/// Add a column defined by `expression`, evaluated by a later refresh
/// rather than by this commit. Its type and inputs are derived from the
/// expression.
///
/// The column is committed with no values, so declaring one costs the same
/// on an empty table as on a large one. Rows get values from
/// [`Table::refresh_column`](super::Table::refresh_column), which fills
/// every fragment that has none -- including fragments appended since the
/// last refresh.
///
/// 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.
///
/// Local tables only: LanceDB Cloud and Enterprise reject a declaration
/// with `NotSupported`.
///
/// ```
/// # use lancedb::Table;
/// # async fn declare(table: &Table) -> Result<(), Box<dyn std::error::Error>> {
/// table
/// .add_columns()
/// .computed("doubled", "x * 2")
/// .execute()
/// .await?;
/// let filled = table.refresh_column("doubled").await?;
/// println!("filled {} rows", filled.rows_filled);
/// # Ok(())
/// # }
/// ```
pub fn computed(mut self, name: impl Into<String>, expression: impl Into<String>) -> Self {
self.computed.push((name.into(), expression.into()));
self
}
/// Limit which existing columns a [`NewColumnTransform::BatchUDF`] mapper
/// receives. Every other transform determines what it reads, so setting
/// this alongside one is an error rather than a silent no-op.
/// receives. Every other transform, and a computed column, determines what
/// it reads, so setting this alongside one is an error rather than a silent
/// no-op.
pub fn read_columns(mut self, columns: impl IntoIterator<Item = impl Into<String>>) -> Self {
self.read_columns = Some(columns.into_iter().map(Into::into).collect());
self
@@ -56,24 +97,43 @@ impl AddColumnsBuilder {
let Self {
parent,
transform,
computed,
read_columns,
} = self;
let Some(transform) = transform else {
return Err(Error::InvalidInput {
message: "add_columns requires a transform".into(),
});
};
if read_columns.is_some() && !matches!(transform, NewColumnTransform::BatchUDF(_)) {
return Err(Error::InvalidInput {
message: "read_columns applies only to a BatchUDF transform; \
every other transform determines what it reads"
match (transform, computed.is_empty()) {
(None, true) => Err(Error::InvalidInput {
message: "add_columns requires a transform or a computed column".into(),
}),
// The two commit through different transforms, so one call covering
// both would be two commits and could half-apply.
(Some(_), false) => Err(Error::InvalidInput {
message: "add_columns cannot mix a transform with computed columns; \
they cannot be added atomically in one call"
.into(),
});
}),
(Some(transform), true) => {
if read_columns.is_some() && !matches!(transform, NewColumnTransform::BatchUDF(_)) {
return Err(Error::InvalidInput {
message: "read_columns applies only to a BatchUDF transform; \
every other transform determines what it reads"
.into(),
});
}
parent.add_columns(transform, read_columns).await
}
(None, false) => {
if read_columns.is_some() {
return Err(Error::InvalidInput {
message: "read_columns applies only to a BatchUDF transform; \
a computed column's inputs come from its expression"
.into(),
});
}
let transform = computed_columns::declare(parent.schema().await?, &computed)?;
parent.add_columns(transform, None).await
}
}
parent.add_columns(transform, read_columns).await
}
}
@@ -85,8 +145,8 @@ mod tests {
use arrow_schema::{DataType, Field, Schema};
use lance::dataset::{BatchUDF, NewColumnTransform};
use crate::Table;
use crate::connect;
use crate::{Error, Table};
async fn table_with_two_columns(name: &str) -> Table {
let conn = connect("memory://").execute().await.unwrap();
@@ -98,10 +158,7 @@ mod tests {
async fn test_requires_a_transform() {
let table = table_with_two_columns("no_transform").await;
let err = table.add_columns().execute().await.unwrap_err();
assert!(
err.to_string().contains("requires a transform"),
"got: {err}"
);
assert!(matches!(err, Error::InvalidInput { .. }));
}
#[tokio::test]
@@ -117,7 +174,7 @@ mod tests {
.execute()
.await
.unwrap_err();
assert!(err.to_string().contains("BatchUDF"), "got: {err}");
assert!(matches!(err, Error::InvalidInput { .. }));
let schema = table.schema().await.unwrap();
assert!(
@@ -126,6 +183,47 @@ mod tests {
);
}
#[tokio::test]
async fn test_mixing_transform_and_computed_is_rejected() {
let table = table_with_two_columns("mixed_add").await;
let err = table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"eager".into(),
"x * 2".into(),
)]))
.computed("lazy", "x * 3")
.execute()
.await
.unwrap_err();
assert!(matches!(err, Error::InvalidInput { .. }));
let schema = table.schema().await.unwrap();
assert!(schema.field_with_name("eager").is_err());
assert!(schema.field_with_name("lazy").is_err());
}
#[tokio::test]
async fn test_read_columns_with_computed_is_rejected() {
let table = table_with_two_columns("read_cols_computed").await;
let err = table
.add_columns()
.computed("doubled", "x * 2")
.read_columns(["x"])
.execute()
.await
.unwrap_err();
assert!(matches!(err, Error::InvalidInput { .. }));
assert!(
table
.schema()
.await
.unwrap()
.field_with_name("doubled")
.is_err()
);
}
#[tokio::test]
async fn test_read_columns_limits_what_a_batch_udf_sees() {
let table = table_with_two_columns("read_cols_udf").await;
+705
View File
@@ -0,0 +1,705 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Computed columns.
//!
//! A computed column is defined by a rule rather than by values supplied at
//! write time. Declaring one commits the column carrying that rule in field
//! metadata but no data, so the cost does not scale with the table; a later
//! refresh fills the rows.
//!
//! The rule is tagged by kind ([`ComputedColumnKind`]) because kinds differ in
//! where the column's type and inputs come from. A SQL expression is
//! self-describing -- both are derived from the expression, so a caller writes
//! neither -- while a kind resolved through a registry cannot be typed without
//! consulting it. Only SQL exists today; the tag is what lets another kind be
//! added without a second reading of the same key.
//!
//! [`computed_columns`] and [`computed_column_from_field`] read declarations
//! back off a schema.
use std::collections::HashMap;
use std::sync::Arc;
use arrow_schema::{Field as ArrowField, Schema as ArrowSchema, SchemaRef};
use lance::dataset::NewColumnTransform;
use lance_datafusion::planner::Planner;
use crate::{Error, Result};
/// Field metadata key marking a column as computed. The value is `"true"`.
pub const COMPUTED_COLUMN_META_KEY: &str = "computed_column";
/// Field metadata key naming the kind of rule that defines the column.
pub const KIND_META_KEY: &str = "computed_column.kind";
/// Field metadata key holding the SQL expression that defines the column.
pub const EXPRESSION_META_KEY: &str = "computed_column.expression";
/// Field metadata key holding the column's inputs, as a JSON array of names.
pub const INPUTS_META_KEY: &str = "computed_column.inputs";
/// Value of [`KIND_META_KEY`] for a column defined by a SQL expression.
pub const SQL_KIND: &str = "sql";
/// The rule that defines a computed column's values.
///
/// Non-exhaustive: a kind added later is an additive change, and a caller that
/// only handles the kinds it knows keeps compiling.
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum ComputedColumnKind {
/// A SQL expression evaluated by DataFusion. It is the whole definition:
/// the column's type and its inputs are both derived from it.
Sql {
/// The expression.
expression: String,
},
/// A kind this version does not understand, written by a newer one.
///
/// Reported rather than hidden so a caller can tell a column it cannot
/// refresh apart from one that was never computed. Nothing produces this.
Unrecognized {
/// The kind as it was found in the metadata.
kind: String,
},
}
/// A computed column's declaration, as read back from field metadata.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ComputedColumn {
/// Name of the computed column.
pub name: String,
/// The rule that defines it.
pub kind: ComputedColumnKind,
/// Columns the rule reads, recorded at declaration time.
///
/// Outside the kind because every kind has inputs and the consumers that
/// use them -- refresh planning, dependency ordering -- do not care which
/// kind produced them. Where they come from does differ, and that is
/// settled at declaration: derived from a SQL expression, supplied by the
/// caller for a kind that cannot be parsed.
pub inputs: Vec<String>,
}
/// Build the field metadata recording a SQL binding.
fn computed_column_metadata(expression: &str, inputs: &[String]) -> HashMap<String, String> {
HashMap::from([
(COMPUTED_COLUMN_META_KEY.to_string(), "true".to_string()),
(KIND_META_KEY.to_string(), SQL_KIND.to_string()),
(EXPRESSION_META_KEY.to_string(), expression.to_string()),
(
INPUTS_META_KEY.to_string(),
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
),
])
}
/// Read a field's computed-column declaration, if it carries one.
///
/// A field flagged computed but carrying no kind, or a SQL one missing its
/// expression, is not a computed column here: without the rule there is
/// nothing to refresh from, so it is reported as absent rather than as a
/// half-formed declaration. An unrecognized kind is different -- the rule is
/// there and intact, this version just cannot act on it -- and comes back as
/// [`ComputedColumnKind::Unrecognized`].
pub fn computed_column_from_field(field: &ArrowField) -> Option<ComputedColumn> {
let metadata = field.metadata();
if metadata.get(COMPUTED_COLUMN_META_KEY).map(String::as_str) != Some("true") {
return None;
}
let kind = match metadata.get(KIND_META_KEY)?.as_str() {
SQL_KIND => ComputedColumnKind::Sql {
expression: metadata.get(EXPRESSION_META_KEY)?.clone(),
},
other => ComputedColumnKind::Unrecognized {
kind: other.to_string(),
},
};
let inputs = metadata
.get(INPUTS_META_KEY)
.and_then(|raw| serde_json::from_str::<Vec<String>>(raw).ok())
.unwrap_or_default();
Some(ComputedColumn {
name: field.name().clone(),
kind,
inputs,
})
}
/// Read every computed-column declaration carried by `schema`, in field order.
///
/// Introspection is a pure read of the schema the caller already holds, the
/// way a SQL catalog reports a generation expression as another column of
/// `information_schema.columns`.
pub fn computed_columns(schema: &ArrowSchema) -> Vec<ComputedColumn> {
schema
.fields()
.iter()
.filter_map(|field| computed_column_from_field(field))
.collect()
}
/// Reject a schema change to a column some declaration reads.
///
/// A binding is SQL text naming its inputs, so renaming, retyping or dropping
/// one leaves an expression that no longer resolves. Refusing the change keeps
/// a declaration that survived [`plan`] evaluable for as long as it exists.
///
/// Paths are compared at their root: a declaration reading `metadata` is
/// invalidated by a change to `metadata.age` just as surely.
pub(crate) fn ensure_not_an_input(schema: &ArrowSchema, paths: &[&str]) -> Result<()> {
let root = |path: &str| path.split('.').next().unwrap_or(path).to_string();
for declaration in computed_columns(schema) {
for path in paths {
// A declaration does not read itself, so it is free to be dropped
// or renamed along with its binding.
if declaration.name == root(path) {
continue;
}
if declaration
.inputs
.iter()
.any(|input| root(input) == root(path))
{
return Err(Error::InvalidInput {
message: format!(
"column '{}' is read by computed column '{}'; drop that column first",
path, declaration.name
),
});
}
}
}
Ok(())
}
/// Resolve `(name, expression)` pairs against `schema` into fields carrying
/// their bindings.
///
/// Everything that can be known statically is checked here rather than at
/// refresh time: that the expression parses, that every column it reads
/// exists, and that the target name is free. A declaration that survives this
/// is one a refresh can always act on.
pub(crate) fn plan(schema: SchemaRef, columns: &[(String, String)]) -> Result<Vec<ArrowField>> {
if columns.is_empty() {
return Err(Error::InvalidInput {
message: "at least one computed column is required".into(),
});
}
let planner = Planner::new(schema.clone());
let mut fields = Vec::with_capacity(columns.len());
let mut declared: Vec<&str> = Vec::with_capacity(columns.len());
for (name, expression) in columns {
if schema.field_with_name(name).is_ok() || declared.contains(&name.as_str()) {
return Err(Error::ColumnAlreadyExists { name: name.clone() });
}
let expr = planner
.parse_expr(expression)
.and_then(|expr| planner.optimize_expr(expr))
.map_err(|e| Error::InvalidExpression {
column: name.clone(),
message: e.to_string(),
})?;
let mut inputs = Planner::column_names_in_expr(&expr);
inputs.sort();
inputs.dedup();
// Resolved here rather than left to the planner so an unknown column
// names itself in the error instead of surfacing as a plan failure.
let mut indices = Vec::with_capacity(inputs.len());
for input in &inputs {
let index = schema
.index_of(input)
.map_err(|_| Error::InvalidExpression {
column: name.clone(),
message: format!("unknown column '{input}'"),
})?;
indices.push(index);
}
// Physical expressions address columns by position, so the planner
// that types the expression has to be built on the projected schema
// the refresh will actually read.
let read_schema =
Arc::new(
schema
.project(&indices)
.map_err(|e| Error::InvalidExpression {
column: name.clone(),
message: e.to_string(),
})?,
);
let physical = Planner::new(read_schema.clone())
.create_physical_expr(&expr)
.map_err(|e| Error::InvalidExpression {
column: name.clone(),
message: e.to_string(),
})?;
let data_type =
physical
.data_type(read_schema.as_ref())
.map_err(|e| Error::InvalidExpression {
column: name.clone(),
message: e.to_string(),
})?;
// Declared columns start entirely null, so nullability is a property
// of the declaration rather than of what the expression yields.
fields.push(
ArrowField::new(name, data_type, true)
.with_metadata(computed_column_metadata(expression, &inputs)),
);
declared.push(name);
}
Ok(fields)
}
/// Build the transform that declares `columns` against `schema`.
///
/// An all-null column is how a binding with no values yet is carried into a
/// commit; that it is spelled `AllNulls` is a detail of the commit, not of the
/// column, which is why this is internal and
/// [`AddColumnsBuilder::computed`](super::AddColumnsBuilder::computed) is the
/// public way in.
pub(crate) fn declare(
schema: SchemaRef,
columns: &[(String, String)],
) -> Result<NewColumnTransform> {
let fields = plan(schema, columns)?;
Ok(NewColumnTransform::AllNulls(Arc::new(ArrowSchema::new(
fields,
))))
}
/// Commit a declaration of a kind this version does not produce, the way a
/// newer lancedb would leave one behind. Shared with the refresh tests, which
/// need the same column to check that refresh refuses it.
#[cfg(test)]
pub(super) async fn add_foreign_kind(table: &crate::Table, name: &str, kind: &str) {
use arrow_schema::DataType;
let field = ArrowField::new(name, DataType::Int32, true).with_metadata(HashMap::from([
(COMPUTED_COLUMN_META_KEY.to_string(), "true".to_string()),
(KIND_META_KEY.to_string(), kind.to_string()),
(INPUTS_META_KEY.to_string(), r#"["x"]"#.to_string()),
]));
table
.add_columns()
.transform(NewColumnTransform::AllNulls(Arc::new(ArrowSchema::new(
vec![field],
))))
.execute()
.await
.unwrap();
}
#[cfg(test)]
mod tests {
use arrow_array::record_batch;
use arrow_schema::DataType;
use futures::TryStreamExt;
use lance::dataset::ColumnAlteration;
use super::*;
use crate::connect;
use crate::query::{ExecutableQuery, QueryBase, Select};
use crate::{Error, Table};
async fn table_with_ints(name: &str) -> Table {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(("x", Int32, [1, 2, 3])).unwrap();
conn.create_table(name, batch).execute().await.unwrap()
}
/// Declare `columns` the way a caller would: plan the expressions, then
/// add them through the ordinary column API.
async fn add_computed(table: &Table, columns: &[(String, String)]) -> Result<u64> {
let mut builder = table.add_columns();
for (name, expression) in columns {
builder = builder.computed(name, expression);
}
Ok(builder.execute().await?.version)
}
async fn declared(table: &Table) -> Vec<ComputedColumn> {
computed_columns(table.schema().await.unwrap().as_ref())
}
#[tokio::test]
async fn test_declare_infers_type_and_inputs() {
let table = table_with_ints("declare_infers").await;
let initial = table.version().await.unwrap();
let version = add_computed(&table, &[("doubled".into(), "x * 2".into())])
.await
.unwrap();
assert!(version > initial);
let schema = table.schema().await.unwrap();
let field = schema.field_with_name("doubled").unwrap();
assert_eq!(field.data_type(), &DataType::Int32);
assert!(field.is_nullable());
assert_eq!(
declared(&table).await,
vec![ComputedColumn {
name: "doubled".into(),
kind: ComputedColumnKind::Sql {
expression: "x * 2".into()
},
inputs: vec!["x".into()],
}]
);
}
/// The binding reaches the schema only if `AllNulls` carries per-field
/// metadata through the commit. The whole representation rests on it.
#[tokio::test]
async fn test_all_nulls_preserves_field_metadata() {
let table = table_with_ints("metadata_survives").await;
add_computed(&table, &[("doubled".into(), "x * 2".into())])
.await
.unwrap();
let schema = table.schema().await.unwrap();
let metadata = schema.field_with_name("doubled").unwrap().metadata();
assert_eq!(
metadata.get(COMPUTED_COLUMN_META_KEY).map(String::as_str),
Some("true")
);
assert_eq!(metadata.get(KIND_META_KEY).map(String::as_str), Some("sql"));
assert_eq!(
metadata.get(EXPRESSION_META_KEY).map(String::as_str),
Some("x * 2")
);
assert_eq!(
metadata.get(INPUTS_META_KEY).map(String::as_str),
Some(r#"["x"]"#)
);
}
#[tokio::test]
async fn test_declared_column_is_all_null() {
let table = table_with_ints("declare_is_null").await;
add_computed(&table, &[("doubled".into(), "x * 2".into())])
.await
.unwrap();
let batches = table
.query()
.select(Select::columns(&["doubled"]))
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
let total: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(total, 3);
for batch in &batches {
assert_eq!(batch["doubled"].null_count(), batch.num_rows());
}
}
#[tokio::test]
async fn test_unknown_column_fails_at_declare_time() {
let table = table_with_ints("unknown_input").await;
let err = add_computed(&table, &[("bad".into(), "missing + 1".into())])
.await
.unwrap_err();
assert!(matches!(err, Error::InvalidExpression { column, .. } if column == "bad"));
let schema = table.schema().await.unwrap();
assert!(schema.field_with_name("bad").is_err());
}
#[tokio::test]
async fn test_unparsable_expression_fails_at_declare_time() {
let table = table_with_ints("bad_syntax").await;
let err = add_computed(&table, &[("bad".into(), "x *".into())])
.await
.unwrap_err();
assert!(matches!(err, Error::InvalidExpression { column, .. } if column == "bad"));
assert!(
table
.schema()
.await
.unwrap()
.field_with_name("bad")
.is_err()
);
}
/// A user-defined function is an expression like any other; only its
/// resolution is missing. When a registry-aware planner exists this
/// becomes a supported declaration rather than a new API.
#[tokio::test]
async fn test_unregistered_function_is_rejected_for_now() {
let table = table_with_ints("udf_not_yet").await;
let err = add_computed(&table, &[("vec".into(), "embed(x)".into())])
.await
.unwrap_err();
assert!(matches!(err, Error::InvalidExpression { column, .. } if column == "vec"));
assert!(
table
.schema()
.await
.unwrap()
.field_with_name("vec")
.is_err()
);
}
#[tokio::test]
async fn test_existing_column_name_is_rejected() {
let table = table_with_ints("name_taken").await;
let err = add_computed(&table, &[("x".into(), "x * 2".into())])
.await
.unwrap_err();
assert!(matches!(err, Error::ColumnAlreadyExists { name } if name == "x"));
assert!(declared(&table).await.is_empty());
}
#[tokio::test]
async fn test_constant_expression_needs_no_inputs() {
let table = table_with_ints("constant").await;
add_computed(&table, &[("answer".into(), "42".into())])
.await
.unwrap();
let declared = declared(&table).await;
assert_eq!(declared.len(), 1);
assert!(declared[0].inputs.is_empty());
}
#[tokio::test]
async fn test_multiple_columns_in_one_commit() {
let table = table_with_ints("multi").await;
let initial = table.version().await.unwrap();
add_computed(
&table,
&[
("plus".into(), "x + 1".into()),
("squared".into(), "x * x".into()),
],
)
.await
.unwrap();
assert_eq!(table.version().await.unwrap(), initial + 1);
let declared = declared(&table).await;
assert_eq!(declared.len(), 2);
assert_eq!(declared[0].name, "plus");
assert_eq!(declared[1].name, "squared");
}
#[tokio::test]
async fn test_duplicate_declaration_in_one_call_is_rejected() {
let table = table_with_ints("dupe").await;
let err = add_computed(
&table,
&[
("dup".into(), "x + 1".into()),
("dup".into(), "x + 2".into()),
],
)
.await
.unwrap_err();
assert!(matches!(err, Error::ColumnAlreadyExists { name } if name == "dup"));
assert!(declared(&table).await.is_empty());
}
/// A column added by an ordinary transform is materialized, not bound, so
/// it carries no declaration to report.
#[tokio::test]
async fn test_ordinary_columns_are_not_reported_as_computed() {
let table = table_with_ints("plain").await;
assert!(declared(&table).await.is_empty());
table
.add_columns()
.transform(NewColumnTransform::SqlExpressions(vec![(
"eager".into(),
"x * 2".into(),
)]))
.execute()
.await
.unwrap();
assert!(declared(&table).await.is_empty());
}
/// Built-in functions type the column the same way an operator does.
#[tokio::test]
async fn test_builtin_function_inference() {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(("name", Utf8, ["ada", "grace"]), ("n", Int32, [-1, 2])).unwrap();
let table = conn
.create_table("builtins", batch)
.execute()
.await
.unwrap();
add_computed(
&table,
&[
("shout".into(), "upper(name)".into()),
("width".into(), "length(name)".into()),
("magnitude".into(), "abs(n)".into()),
],
)
.await
.unwrap();
let schema = table.schema().await.unwrap();
assert_eq!(
schema.field_with_name("shout").unwrap().data_type(),
&DataType::Utf8
);
assert_eq!(
schema.field_with_name("magnitude").unwrap().data_type(),
&DataType::Int32
);
// length() returns a width-dependent integer type; assert it is one
// rather than pinning which.
assert!(
schema
.field_with_name("width")
.unwrap()
.data_type()
.is_integer()
);
let declared = declared(&table).await;
assert_eq!(declared.len(), 3);
assert_eq!(declared[0].inputs, vec!["name".to_string()]);
assert_eq!(declared[2].inputs, vec!["n".to_string()]);
}
/// The reason the kind is tagged: a declaration written by a newer version
/// has to read back as a computed column this one cannot evaluate, not as
/// an ordinary column. Reported as absent it would be refreshable by
/// nothing and redeclarable over, silently.
#[tokio::test]
async fn test_unrecognized_kind_is_reported_rather_than_hidden() {
let table = table_with_ints("foreign_kind").await;
super::add_foreign_kind(&table, "embedding", "udf").await;
assert_eq!(
declared(&table).await,
vec![ComputedColumn {
name: "embedding".into(),
kind: ComputedColumnKind::Unrecognized { kind: "udf".into() },
inputs: vec!["x".into()],
}]
);
let err = add_computed(&table, &[("embedding".into(), "x * 2".into())])
.await
.unwrap_err();
assert!(matches!(err, Error::ColumnAlreadyExists { name } if name == "embedding"));
}
/// A kind is what makes a declaration readable at all, so the flag alone
/// is half-formed in the same way a missing expression is.
#[test]
fn test_flag_without_a_kind_is_not_a_declaration() {
let field =
ArrowField::new("half", DataType::Int32, true).with_metadata(HashMap::from([(
COMPUTED_COLUMN_META_KEY.to_string(),
"true".to_string(),
)]));
assert_eq!(computed_column_from_field(&field), None);
}
/// A SQL declaration is its expression; without one there is nothing to
/// refresh from.
#[test]
fn test_sql_kind_without_an_expression_is_not_a_declaration() {
let field = ArrowField::new("half", DataType::Int32, true).with_metadata(HashMap::from([
(COMPUTED_COLUMN_META_KEY.to_string(), "true".to_string()),
(KIND_META_KEY.to_string(), SQL_KIND.to_string()),
]));
assert_eq!(computed_column_from_field(&field), None);
}
#[tokio::test]
async fn test_inputs_are_deduplicated_and_sorted() {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(("b", Int32, [1, 2]), ("a", Int32, [3, 4])).unwrap();
let table = conn.create_table("dedupe", batch).execute().await.unwrap();
add_computed(&table, &[("total".into(), "b + a + b".into())])
.await
.unwrap();
assert_eq!(
declared(&table).await[0].inputs,
vec!["a".to_string(), "b".to_string()]
);
}
#[tokio::test]
async fn test_dropping_an_input_is_refused() {
let table = table_with_ints("drop_input").await;
add_computed(&table, &[("doubled".into(), "x * 2".into())])
.await
.unwrap();
let err = table.drop_columns(&["x"]).await.unwrap_err();
assert!(
matches!(&err, Error::InvalidInput { message } if message.contains("doubled")),
"{err:?}"
);
}
#[tokio::test]
async fn test_renaming_an_input_is_refused() {
let table = table_with_ints("rename_input").await;
add_computed(&table, &[("doubled".into(), "x * 2".into())])
.await
.unwrap();
let err = table
.alter_columns(&[ColumnAlteration::new("x".into()).rename("y".into())])
.await
.unwrap_err();
assert!(
matches!(&err, Error::InvalidInput { message } if message.contains("doubled")),
"{err:?}"
);
}
/// Nothing resolves against nullability, so it is not a rebinding.
#[tokio::test]
async fn test_altering_an_input_nullability_is_allowed() {
let table = table_with_ints("nullable_input").await;
add_computed(&table, &[("doubled".into(), "x * 2".into())])
.await
.unwrap();
table
.alter_columns(&[ColumnAlteration::new("x".into()).set_nullable(true)])
.await
.unwrap();
}
/// A declaration does not read itself, so it travels with its binding.
#[tokio::test]
async fn test_dropping_the_computed_column_is_allowed() {
let table = table_with_ints("drop_computed").await;
add_computed(&table, &[("doubled".into(), "x * 2".into())])
.await
.unwrap();
table.drop_columns(&["doubled"]).await.unwrap();
assert!(declared(&table).await.is_empty());
}
}
+523
View File
@@ -0,0 +1,523 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Filling computed columns.
//!
//! 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.
//!
//! Convergence comes from staging nothing when nothing would change, so an
//! expression yielding null settles after one pass rather than re-selecting
//! the same rows forever. Fragments that already cover the column and hold no
//! nulls are skipped without evaluating it at all.
use std::sync::Arc;
use arrow_array::RecordBatch;
use arrow_schema::Schema as ArrowSchema;
use futures::{TryStreamExt, stream};
use lance::Dataset;
use lance::dataset::WriteDestination;
use lance::dataset::fragment::FileFragment;
use lance::dataset::transaction::Operation;
use lance_core::ROW_ID;
use lance_core::datatypes::Schema as LanceSchema;
use serde::{Deserialize, Serialize};
use super::NativeTable;
use super::computed_columns::{ComputedColumnKind, computed_column_from_field};
use crate::{Error, Result};
/// Alias the expression is projected under, so its result and the column's
/// current values can be read side by side.
const COMPUTED_ALIAS: &str = "__lancedb_computed";
/// The result of refreshing a computed column.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize, Default)]
pub struct RefreshColumnResult {
/// Rows that had a value computed.
#[serde(default)]
pub rows_filled: u64,
/// The commit version associated with the operation.
#[serde(default)]
pub version: u64,
}
/// Internal implementation of the refresh logic.
pub(crate) async fn execute_refresh_column(
table: &NativeTable,
column: &str,
) -> Result<RefreshColumnResult> {
table.dataset.ensure_mutable()?;
let dataset = table.dataset.get().await?;
let expression = declared_expression(&dataset, column)?;
let field = dataset
.schema()
.field(column)
.ok_or_else(|| Error::ColumnNotFound {
name: column.to_string(),
})?;
// The dataset's own field, so the identity write_column checks against the
// manifest holds by construction.
let column_schema = LanceSchema {
fields: vec![field.clone()],
metadata: Default::default(),
};
let mut rows_filled = 0u64;
let mut replacements = Vec::new();
for fragment in fragments_to_consider(&dataset, column, field.id).await? {
let Some((filled, values)) =
fill_fragment(&dataset, &fragment, column, &expression).await?
else {
continue;
};
rows_filled += filled;
replacements.push(
fragment
.write_column(stream::iter(values.into_iter().map(Ok)), &column_schema)
.await?,
);
}
if replacements.is_empty() {
return Ok(RefreshColumnResult {
rows_filled: 0,
version: dataset.version().version,
});
}
let read_version = dataset.version().version;
let new_dataset = Dataset::commit(
WriteDestination::Dataset(dataset.clone()),
Operation::DataReplacement { replacements },
Some(read_version),
None,
None,
Arc::new(Default::default()),
false,
)
.await?;
let version = new_dataset.version().version;
table.dataset.update(new_dataset);
Ok(RefreshColumnResult {
rows_filled,
version,
})
}
/// The SQL expression `column` is declared with.
fn declared_expression(dataset: &Dataset, column: &str) -> Result<String> {
let schema = ArrowSchema::from(dataset.schema());
let field = schema
.field_with_name(column)
.map_err(|_| Error::ColumnNotFound {
name: column.to_string(),
})?;
let declaration =
computed_column_from_field(field).ok_or_else(|| Error::NotAComputedColumn {
name: column.to_string(),
})?;
match declaration.kind {
ComputedColumnKind::Sql { expression } => Ok(expression),
ComputedColumnKind::Unrecognized { kind } => Err(Error::NotSupported {
message: format!(
"computed column '{column}' is defined by '{kind}', which this version of \
lancedb cannot evaluate"
),
}),
}
}
/// Quote `name` as a lance SQL identifier.
///
/// Lance's dialect delimits with backticks, so a double-quoted name would
/// parse as a string literal rather than a column.
fn quote_identifier(name: &str) -> String {
format!("`{}`", name.replace('`', "``"))
}
/// Fragments that could hold a row needing a value.
///
/// A fragment whose data files do not carry the field cannot hold one that
/// does. One that carries it is asked, since a row rewrite -- an update, or a
/// compaction folding an unfilled fragment into a filled one -- can leave
/// nulls behind a covering file.
async fn fragments_to_consider(
dataset: &Dataset,
column: &str,
field_id: i32,
) -> Result<Vec<FileFragment>> {
let unfilled = format!("{} IS NULL", quote_identifier(column));
let mut considered = Vec::new();
for fragment in dataset.get_fragments() {
let covered = fragment
.metadata()
.files
.iter()
.any(|file| file.fields.contains(&field_id));
if !covered || fragment.count_rows(Some(unfilled.clone())).await? > 0 {
considered.push(fragment);
}
}
Ok(considered)
}
/// Compute one fragment's column, keeping every value it already holds.
///
/// `Ok(None)` when no live row gained a value, which is what keeps a refresh
/// from restaging a fragment whose expression yields null. 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 has to cover them.
async fn fill_fragment(
dataset: &Dataset,
fragment: &FileFragment,
column: &str,
expression: &str,
) -> Result<Option<(u64, Vec<RecordBatch>)>> {
let mut scanner = dataset.scan();
scanner
.with_fragments(vec![fragment.metadata().clone()])
.with_row_id()
.include_deleted_rows()
.project_with_transform(&[
(column, quote_identifier(column).as_str()),
(COMPUTED_ALIAS, expression),
])?;
let projected = Arc::new(ArrowSchema::new(vec![
ArrowSchema::from(dataset.schema())
.field_with_name(column)
.map_err(|_| Error::ColumnNotFound {
name: column.to_string(),
})?
.clone(),
]));
let missing = |name: &str| Error::Runtime {
message: format!("refreshing {column} produced no {name} column"),
};
let mut filled = 0u64;
let mut values = Vec::new();
let mut batches = scanner.try_into_stream().await?;
while let Some(batch) = batches.try_next().await? {
let existing = batch
.column_by_name(column)
.ok_or_else(|| missing(column))?;
let computed = batch
.column_by_name(COMPUTED_ALIAS)
.ok_or_else(|| missing("expression"))?;
let row_ids = batch
.column_by_name(ROW_ID)
.ok_or_else(|| missing(ROW_ID))?;
// A row is filled only if it gains a value: an expression yielding null
// leaves it as unfilled as it was, which is what lets a refresh settle.
// A deleted row has a null row id; its value is written but not counted.
let unfilled = arrow::compute::is_null(existing.as_ref())?;
filled += (0..unfilled.len())
.filter(|i| unfilled.value(*i) && row_ids.is_valid(*i) && computed.is_valid(*i))
.count() as u64;
let merged = arrow_select::zip::zip(&unfilled, computed, existing)?;
values.push(RecordBatch::try_new(projected.clone(), vec![merged])?);
}
Ok((filled > 0).then_some((filled, values)))
}
#[cfg(test)]
mod tests {
use arrow_array::{Int32Array, record_batch};
use futures::TryStreamExt;
use crate::connect;
use crate::query::{ExecutableQuery, QueryBase, Select};
use crate::{Error, Result, Table};
async fn table_with(name: &str, values: Vec<i32>) -> Table {
let conn = connect("memory://").execute().await.unwrap();
let batch = record_batch!(("x", Int32, values)).unwrap();
conn.create_table(name, batch).execute().await.unwrap()
}
async fn declare_doubled(table: &Table) -> Result<u64> {
Ok(table
.add_columns()
.computed("doubled", "x * 2")
.execute()
.await?
.version)
}
async fn read(table: &Table, column: &str) -> Vec<Option<i32>> {
let batches = table
.query()
.select(Select::columns(&[column]))
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
let mut values: Vec<Option<i32>> = batches
.iter()
.flat_map(|batch| {
batch[column]
.as_any()
.downcast_ref::<Int32Array>()
.unwrap()
.iter()
.collect::<Vec<_>>()
})
.collect();
values.sort();
values
}
async fn append(table: &Table, values: Vec<i32>) {
let batch = record_batch!(("x", Int32, values)).unwrap();
table.add(batch).execute().await.unwrap();
}
#[tokio::test]
async fn test_refresh_fills_a_declared_column() {
let table = table_with("refresh_fills", vec![1, 2, 3]).await;
let declared = declare_doubled(&table).await.unwrap();
assert_eq!(read(&table, "doubled").await, vec![None, None, None]);
let result = table.refresh_column("doubled").await.unwrap();
assert!(result.version > declared);
assert_eq!(result.rows_filled, 3);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(6)]
);
}
/// Values written after the last refresh must be reachable by another one.
#[tokio::test]
async fn test_refresh_fills_rows_appended_since_the_last_refresh() {
let table = table_with("refresh_appended", vec![1, 2]).await;
declare_doubled(&table).await.unwrap();
table.refresh_column("doubled").await.unwrap();
append(&table, vec![5, 6]).await;
assert_eq!(
read(&table, "doubled").await,
vec![None, None, Some(2), Some(4)]
);
let result = table.refresh_column("doubled").await.unwrap();
assert_eq!(result.rows_filled, 2);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(10), Some(12)]
);
}
#[tokio::test]
async fn test_refresh_with_nothing_to_fill() {
let table = table_with("refresh_noop", vec![1, 2, 3]).await;
declare_doubled(&table).await.unwrap();
table.refresh_column("doubled").await.unwrap();
let again = table.refresh_column("doubled").await.unwrap();
assert_eq!(again.rows_filled, 0);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(6)]
);
}
/// 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. 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;
let declared = table
.add_columns()
.computed("maybe", "nullif(x, x)")
.execute()
.await
.unwrap()
.version;
let first = table.refresh_column("maybe").await.unwrap();
assert_eq!(first.rows_filled, 0);
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, declared);
}
/// 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_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();
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
table.update().column("x", "3").execute().await.unwrap();
let again = table.refresh_column("doubled").await.unwrap();
assert_eq!(again.rows_filled, 0);
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
}
/// A row rewrite before the first refresh materializes the declared
/// column as null behind a covering data file. Those rows are still
/// unfilled and a later refresh has to reach them.
#[tokio::test]
async fn test_update_before_the_first_refresh() {
let table = table_with("refresh_update_first", vec![1]).await;
declare_doubled(&table).await.unwrap();
table.update().column("x", "3").execute().await.unwrap();
let result = table.refresh_column("doubled").await.unwrap();
assert_eq!(result.rows_filled, 1);
assert_eq!(read(&table, "doubled").await, vec![Some(6)]);
}
/// 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_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();
append(&table, vec![5]).await;
table
.update()
.column("x", "100")
.only_if("x = 1")
.execute()
.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, 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)]
);
}
/// 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_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();
append(&table, vec![5]).await;
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, 1);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(4), Some(10)]
);
}
#[tokio::test]
async fn test_refresh_leaves_deleted_rows_alone() {
let table = table_with("refresh_deleted", vec![1, 2, 3, 4]).await;
declare_doubled(&table).await.unwrap();
table.delete("x = 2").await.unwrap();
let result = table.refresh_column("doubled").await.unwrap();
assert_eq!(result.rows_filled, 3);
assert_eq!(
read(&table, "doubled").await,
vec![Some(2), Some(6), Some(8)]
);
}
#[tokio::test]
async fn test_refresh_a_constant_expression() {
let table = table_with("refresh_constant", vec![1, 2, 3]).await;
table
.add_columns()
.computed("answer", "42")
.execute()
.await
.unwrap();
let result = table.refresh_column("answer").await.unwrap();
assert_eq!(result.rows_filled, 3);
}
/// A name needing quotes reaches the evaluator intact: it is carried as a
/// projection alias, never spliced into SQL text.
#[tokio::test]
async fn test_refresh_a_column_whose_name_needs_quoting() {
let table = table_with("refresh_quoted", vec![1, 2, 3]).await;
table
.add_columns()
.computed("double value", "x * 2")
.execute()
.await
.unwrap();
let result = table.refresh_column("double value").await.unwrap();
assert_eq!(result.rows_filled, 3);
assert_eq!(
read(&table, "double value").await,
vec![Some(2), Some(4), Some(6)]
);
}
#[tokio::test]
async fn test_refresh_rejects_a_plain_column() {
let table = table_with("refresh_plain", vec![1, 2, 3]).await;
let err = table.refresh_column("x").await.unwrap_err();
assert!(matches!(err, Error::NotAComputedColumn { name } if name == "x"));
}
#[tokio::test]
async fn test_refresh_rejects_an_unknown_column() {
let table = table_with("refresh_missing", vec![1, 2, 3]).await;
let err = table.refresh_column("nope").await.unwrap_err();
assert!(matches!(err, Error::ColumnNotFound { name } if name == "nope"));
}
/// A declaration of a kind this version cannot evaluate is refused by
/// name, rather than mistaken for a plain column or fed to the SQL path.
#[tokio::test]
async fn test_refresh_rejects_a_kind_it_cannot_evaluate() {
let table = table_with("refresh_foreign", vec![1, 2, 3]).await;
super::super::computed_columns::add_foreign_kind(&table, "embedding", "udf").await;
let err = table.refresh_column("embedding").await.unwrap_err();
assert!(matches!(err, Error::NotSupported { message } if message.contains("udf")));
}
}
@@ -8,11 +8,13 @@
//! - [`alter_columns`](execute_alter_columns): Rename columns, change types, or modify nullability
//! - [`drop_columns`](execute_drop_columns): Remove columns from the table
use arrow_schema::Schema as ArrowSchema;
use lance::dataset::{ColumnAlteration, NewColumnTransform};
use serde::{Deserialize, Serialize};
use std::collections::HashMap;
use super::NativeTable;
use super::computed_columns;
use crate::Result;
/// The result of an add columns operation.
@@ -116,6 +118,14 @@ pub(crate) async fn execute_alter_columns(
) -> Result<AlterColumnsResult> {
table.dataset.ensure_mutable()?;
let mut dataset = (*table.dataset.get().await?).clone();
// Nullability is not part of what an expression resolves against, so only
// a rename or a retype can invalidate a binding.
let rebinding = alterations
.iter()
.filter(|alteration| alteration.rename.is_some() || alteration.data_type.is_some())
.map(|alteration| alteration.path.as_str())
.collect::<Vec<_>>();
computed_columns::ensure_not_an_input(&ArrowSchema::from(dataset.schema()), &rebinding)?;
dataset.alter_columns(alterations).await?;
let version = dataset.version().version;
table.dataset.update(dataset);
@@ -131,6 +141,7 @@ pub(crate) async fn execute_drop_columns(
) -> Result<DropColumnsResult> {
table.dataset.ensure_mutable()?;
let mut dataset = (*table.dataset.get().await?).clone();
computed_columns::ensure_not_an_input(&ArrowSchema::from(dataset.schema()), columns)?;
dataset.drop_columns(columns).await?;
let version = dataset.version().version;
table.dataset.update(dataset);
+7 -54
View File
@@ -14,7 +14,6 @@ use lance::arrow::json::JsonDataType;
use lance::dataset::{ReadParams, WriteParams};
use lance::index::vector::utils::infer_vector_dim;
use lance::io::{ObjectStoreParams, WrappingObjectStore};
use lance_io::object_store::ChainedWrappingObjectStore;
use std::pin::Pin;
use crate::error::{Error, Result};
@@ -38,13 +37,13 @@ impl PatchStoreParam for Option<ObjectStoreParams> {
wrapper: Arc<dyn WrappingObjectStore>,
) -> Result<Option<ObjectStoreParams>> {
let mut params = self.unwrap_or_default();
params.object_store_wrapper = Some(match params.object_store_wrapper.take() {
// The wrapper being patched in is connection-level compatibility
// behavior. Keep it closest to the target store so an existing
// caller wrapper remains outermost and can observe every operation.
Some(existing) => Arc::new(ChainedWrappingObjectStore::new(vec![wrapper, existing])),
None => wrapper,
});
if params.object_store_wrapper.is_some() {
return Err(Error::Other {
message: "can not patch param because object store is already set".into(),
source: None,
});
}
params.object_store_wrapper = Some(wrapper);
Ok(Some(params))
}
@@ -473,60 +472,14 @@ impl Stream for MaxBatchLengthStream {
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use arrow_array::Int32Array;
use arrow_schema::Field;
use datafusion_physical_plan::stream::RecordBatchStreamAdapter;
use futures::{StreamExt, stream};
use object_store::{ObjectStore, memory::InMemory};
use tokio::time::sleep;
use super::*;
#[derive(Debug)]
struct OrderedStoreWrapper {
name: &'static str,
order: Arc<Mutex<Vec<&'static str>>>,
}
impl WrappingObjectStore for OrderedStoreWrapper {
fn wrap(
&self,
_store_prefix: &str,
original: Arc<dyn ObjectStore>,
) -> Arc<dyn ObjectStore> {
self.order.lock().unwrap().push(self.name);
original
}
}
#[test]
fn test_patch_store_param_keeps_caller_wrapper_outermost() {
let order = Arc::new(Mutex::new(Vec::new()));
let params = Some(ObjectStoreParams {
object_store_wrapper: Some(Arc::new(OrderedStoreWrapper {
name: "caller",
order: order.clone(),
})),
..Default::default()
});
let params = params
.patch_with_store_wrapper(Arc::new(OrderedStoreWrapper {
name: "compatibility",
order: order.clone(),
}))
.unwrap()
.unwrap();
params
.object_store_wrapper
.unwrap()
.wrap("memory", Arc::new(InMemory::new()) as Arc<dyn ObjectStore>);
assert_eq!(*order.lock().unwrap(), vec!["compatibility", "caller"]);
}
#[test]
fn test_guess_default_column() {
let schema_no_vector = Schema::new(vec![