Compare commits

...

16 Commits

Author SHA1 Message Date
Lance Release 357b535405 Bump version: 0.38.0-beta.3 → 0.38.0-beta.4 2026-08-22 16:38:08 +00:00
Wyatt Alt 68749ecfa3 feat(nodejs): materialized view bindings (#3935)
Exposes materialized views to TypeScript: createMaterializedView,
openMaterializedView and listMaterializedViews on Connection, and a
MaterializedView handle carrying the parsed definition and
refresh({full, sourceVersion}), which returns the typed refresh result.
select accepts column names, [alias, expression] pairs, or a record of
the
same; the definition reads back off the stored schema, so a reopened
handle
needs no side channel. Remote connections surface the core's
not-supported
error up front.

The napi crate needed the same recursion-limit raise as the core crate:
the
refresh future's type graph overflows the default trait-recursion depth.


<sub>Stack created with <a
href="https://github.com/github/gh-stack">GitHub Stacks CLI</a> • <a
href="https://gh.io/stacks-feedback">Give Feedback 💬</a></sub>
2026-08-21 23:48:43 -07:00
Drew Gallardo e98d8ac685 feat!: rename branch merge to cherry_pick (#3986)
This PR is a **breaking** rename of #3686.

merge reads like git merge w/ three-way, replay history, combine two
lines of work. That is not this API.

This call takes one additive change on a branch and lands it on main.
New column, including a blob column. Main's existing columns are not
rewritten. If it cannot land, you get `status="failed"` and
`diff.errors`, not a merge conflict to resolve.

Cherry-pick is terminology that aligns more with that.

```python
table = db.open_table("images")
table.branches.create("exp")
exp = table.branches.checkout("exp")

exp.add_columns({"tag": "cast('draft' as string)"})

diff = table.branches.diff("exp")
preview = table.branches.cherry_pick("exp", dry_run=True)
result = table.branches.cherry_pick("exp")

if result["status"] == "cherryPicked":
    print("landed at", result["mainVersionAfter"])
elif result["status"] == "failed":
    print(result["diff"]["errors"])
```

### Behavior

- Remote / Enterprise only. Local still NotSupported.
- HTTP 409 is not an exception. It is Ok with status="failed" and
diff.errors (CherryPickError).
- Unknown error / status codes still parse as Unknown.
- Requests are not retried. 409 is final and carries the body.
- Endpoint is POST /v1/table/{id}/branches/cherry_pick/.
- merge_insert and Table.merge are unchanged.

### Testing
- `cargo test -p lancedb --features remote diff_branch`
- `cargo test -p lancedb --features remote cherry_pick`
- `pytest python/python/tests/test_remote_db.py -k cherry_pick`
- node `remote.test.ts` diffs / cherry-picks path
2026-08-21 23:37:12 -07:00
Wyatt Alt 851fa16b47 feat(python): materialized view bindings (#3933)
Exposes materialized views to Python in both the async and sync clients:
create_materialized_view / open_materialized_view /
list_materialized_views
on the connections, and MaterializedView / AsyncMaterializedView handles
carrying the parsed definition and refresh(full=, source_version=),
which
returns the typed refresh result. select accepts column names, (alias,
expression) pairs, or a dict of the same; the definition reads back off
the
stored schema, so a reopened handle needs no side channel. Remote
connections raise NotImplementedError up front rather than failing deep
in
a request, matching the computed-column convention.


<sub>Stack created with <a
href="https://github.com/github/gh-stack">GitHub Stacks CLI</a> • <a
href="https://gh.io/stacks-feedback">Give Feedback 💬</a></sub>
2026-08-21 23:08:41 -07:00
Wyatt Alt d04ac7ed20 test: differential refresh harness for materialized views (#3932)
Example tests pin behaviors; the refresh contract is a property: after
any
sequence of source mutations, a view maintained by default refreshes
equals
the definition evaluated against the source directly, and so does a
forced
rebuild. This drives every mutation sequence up to length three --
appends,
deletes, updates crossing the filter, compactions, unrelated column adds
--
over an identity and a filtered view shape, checking against an oracle
that
shares nothing with the refresh path: a plain column scan with the
filter
applied in Rust. The oracle runs after every step because a later
rebuild-forcing mutation silently heals an incremental error; end-state
checks miss exactly the transient bugs that matter. A length-four sweep
runs behind
ignore.

Named regressions additionally assert the refresh mode, which value
comparison cannot: a wrongly rebuilding classifier still matches the
oracle, so the append, unrelated-column and compaction cases pin that
the
incremental path actually ran.



<sub>Stack created with <a
href="https://github.com/github/gh-stack">GitHub Stacks CLI</a> • <a
href="https://gh.io/stacks-feedback">Give Feedback 💬</a></sub>
2026-08-21 23:04:04 -07:00
Wyatt Alt a578e9ff7f feat: refresh materialized views (#4010)
A declared view holds no rows; refresh computes them. It pins one source
version, brings the view to exactly the definition's result at that
version,
and records the version as a watermark in the view's schema metadata.

It is incremental when it can reconcile what changed: appended rows are
computed and appended, and rows the source deleted or updated are found
by
the lance delta and evicted by their __source_row_id provenance, the
updated
ones recomputed in the same commit. Compaction rearranges rows without
changing
them, so its outputs cost nothing -- which is what keeps routine
background
compaction from rebuilding the view. A vacuumed watermark, a
delta the transaction-log walk cannot classify, a Legacy-storage source,
or
more staged ids than a fixed cap all fall back to a rebuild; rebuilding
an
indexed view swaps every fragment in one Update, so readers never see it
unindexed or empty.

Concurrent refreshes serialize at commit -- each carries the
same sentinel row id in its inserted-rows filter, so the loser lands
nothing. On the append path the watermark moves in a follow-up commit,
so a
crash between the two re-appends those rows. Bumps lance
to v11.0.0-beta.19 for the delta reader.
2026-08-21 22:39:47 -07:00
LanceDB Robot 7801e2746a chore: update lance dependency to v11.0.0-beta.19 (#4025)
Updates the Lance dependencies and Java lance-core dependency to
v11.0.0-beta.19. No compatibility fixes were required; workspace clippy
with all features passes. Triggering tag:
https://github.com/lance-format/lance/releases/tag/v11.0.0-beta.19
2026-08-21 21:43:45 -07:00
lancedb-gatefixer[bot] 5468f3d490 fix(rust): reject bitmap indexes on JSON fields (#3895)
## Summary

- reject whole-document `lance.json` fields during native BITMAP index
preparation
- preserve BITMAP support for raw `LargeBinary` fields
- return guidance to use a JSON-path scalar index or FTS instead
- add regression coverage for the logical JSON type while retaining the
existing raw binary coverage

## Root cause

Native scalar-index validation resolved the complete Arrow field but
checked BITMAP compatibility only against its physical data type.
Because `lance.json` is stored as `LargeBinary`, it was incorrectly
accepted under the raw binary compatibility rule.

The fix reuses Lance’s `lance_arrow::json::is_json_field` helper before
physical type validation. Remote serialization is unchanged, so remote
clients continue to send the requested BITMAP type for server-side
validation.

## Validation

- `cargo fmt --all -- --check`
- `cargo test --quiet --features remote -p lancedb
test_create_bitmap_index -- --nocapture`
- `cargo check --quiet --features remote --tests --examples`
- `cargo clippy --quiet --features remote --tests --examples`
- `cargo test --quiet --features remote --tests`

Fixes #3889

<!-- lance-gatekeeper-fix:v1 agent=e097fc02a548edc0d0be2e18c65c03a3
generation=1 -->

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
2026-08-21 16:44:21 -07:00
lancedb-gatefixer[bot] c0df2c63b6 test(rust): cover fixed-size-list merge overflow (#3907)
## Summary

- add a merge-insert regression test whose fixed-size-list child count
crosses `u32::MAX`
- verify delete-by-source updates the matching row, deletes every other
row, and completes without an Arrow panic
- use a null child array so the boundary case avoids allocating a real
vector payload

## Root cause and fix

The affected Lance merge fallback carried the target payload through a
full outer hash join. Arrow's fixed-size-list take kernel uses `u32`
child indices, so taking a target row whose child offset crossed
`u32::MAX` wrapped the offset and produced child data shorter than the
parent array, triggering the reported `ArrayData::slice` assertion.

The projection-aware merge path in the Lance version now used by `main`
avoids materializing the target fixed-size-list payload in that join.
This regression test locks in that production behavior at the exact
child-index boundary.

## Validation

- `cargo fmt --all`
- `cargo test --quiet --features remote -p lancedb
test_merge_insert_fixed_size_list_above_u32_child_count`
- `cargo check --quiet --features remote --tests --examples`
- `cargo clippy --quiet --features remote --tests --examples`

Fixes #2874

<!-- lance-gatekeeper-fix:v1 agent=582e68bcad65739e189352cb3cbf144c
generation=3 -->

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
2026-08-21 16:39:13 -07:00
lancedb-gatefixer[bot] 9e8f1c1a6d fix(python): expose FTS build memory limits (#3796)
## Summary

- expose `memory_limit` and `num_workers` on the Python FTS
configuration for local builds
- forward both build-only settings to the Lance inverted-index builder
- add an end-to-end regression proving the configured memory budget
reaches the native build

## Root cause

LanceDB 0.26.1 pinned Lance 1.0.1. That Lance version used an FTS
partition-merge path whose retained data made memory grow with merge
progress on very large indexes. Upstream Lance
[#5754](https://github.com/lance-format/lance/pull/5754) changed
partition merging to stream its inputs, reducing peak memory by about
25%. Lance [#6174](https://github.com/lance-format/lance/pull/6174) then
removed the old merge phase, compressed posting lists during
construction, reduced indexing memory by about 60%, and introduced a
total build `memory_limit` for bounded workers.

Current `main` pins Lance 11.0.0-beta.3, which contains those
architectural fixes. This PR does not duplicate or claim the upstream
leak fix; it addresses the remaining Python API gap.

## This repair

LanceDB Python did not expose the native FTS builder resource controls.
`memory_limit` now sets the total local-build budget in MiB, divided
among effective workers, and `num_workers` controls build parallelism.
Both are build-only settings and do not affect remote builds or
persisted index configuration.

## Validation

- `cargo check --quiet --features remote --tests --examples`
- `cargo fmt --all`
- `uv run --project python --extra tests --extra dev ruff check .`
- `uv run --project python --extra tests --extra dev ruff format --check
python/python/lancedb/index.py python/python/tests/test_fts.py`
- `uv run --project python --extra tests pytest python/tests/test_fts.py
-q` (51 passed)

Fixes #2923

<!-- lance-gatekeeper-fix:v1 agent=a1ceedf74531e0212cb6f1ebf9390a26
generation=1 -->

---------

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
2026-08-21 16:31:34 -07:00
Wyatt Alt 01679e37fd feat: materialized view declarations on local tables (#3930)
A materialized view is a table whose contents are defined by a query
over
one source table and maintained by refresh rather than by writes.

The declaration half: create_materialized_view(name, source) resolves a
projected, filtered and limited definition against the source schema --
output types come from the DataFusion planner, never the caller -- and
commits an empty table carrying it as kind-tagged JSON in schema
metadata.
The tag lets a kind added later read back as a view this version cannot
refresh rather than as a plain table. Views open and list as ordinary
tables.

Sources must have stable row ids, checked here because the property
cannot
be enabled later: each view row records its source row in
__source_row_id,
and that provenance survives compactions, updates and deletes only when
row
ids are stable.

A view inherits the metadata describing its columns and none governing
how a
table is written, so blob markers carry through while declarations its
always-nullable fields would contradict are stripped. Embedding
configuration is rewritten to the view's column names, and dropped where
it
does not project both ends of a function.


<sub>Stack created with <a
href="https://github.com/github/gh-stack">GitHub Stacks CLI</a> • <a
href="https://gh.io/stacks-feedback">Give Feedback 💬</a></sub>
2026-08-21 16:17:03 -07:00
lancedb-gatefixer[bot] c7cb0b9afa docs(python): clarify threading on two-CPU containers (#3807)
## Summary

- document that current LanceDB releases use one compute worker without
warning on two-vCPU containers
- distinguish compute-worker tuning from storage I/O concurrency
- direct users of affected LanceDB 0.21.1 installations to upgrade and
link the current threading guidance

## Root cause

The Lance version bundled with LanceDB 0.21.1 warned whenever the
detected CPU count was less than or equal to its default two-core I/O
reservation. A two-vCPU deployment therefore emitted the warning on
every query even though falling back to one compute worker was the
intended behavior. Lance fixed that warning condition upstream in
lance-format/lance#3710, and LanceDB current main already pins a version
containing the runtime fix; the Python package documentation did not
explain the corrected behavior or the distinct thread controls.

## Validation

- `git diff --check`
- verified the linked Lance threading-model documentation returns HTTP
200

Fixes #2326

<!-- lance-gatekeeper-fix:v1 agent=9f141242416a6dbeb43be0e80404dd4d
generation=1 -->

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
2026-08-21 16:13:49 -07:00
lancedb-gatefixer[bot] a35f7044ee test(rust): cover Azure table URI separators on Windows (#3810)
## Summary
- add cross-platform regression coverage for Azure table URI
construction
- assert that az:// database paths always produce forward-slash blob
keys

## Root cause
ListingDatabase previously used the host filesystem Path join operation
for object-store URIs, which inserted a backslash on Windows. The URI
construction was corrected in #2575, but the original Azure report had
no regression coverage and remained open.

## Validation
- cargo fmt --all
- cargo test --quiet --features remote -p lancedb
test_table_uri_uses_forward_slashes_for_azure
- cargo check --quiet --features remote --tests --examples
- cargo clippy --quiet --features remote --tests --examples

Fixes #2283

<!-- lance-gatekeeper-fix:v1 agent=dd96adfd0fcf303c11e873300663d8f6
generation=1 -->

Co-authored-by: Gatefixer <313497061+lancedb-gatefixer[bot]@users.noreply.github.com>
2026-08-21 15:43:04 -07:00
Wyatt Alt 29822306d2 fix(python): skip the unrunnable FunctionVersion doctest example (#4014)
The example binds an undefined `function`; only its last line was
skipped, so the doctest suite fails on main and on every PR.
2026-08-21 14:20:03 -07:00
Jack Ye f39a7a4dd9 feat: support remote tables in the data loader (#3981)
`StreamingDataset`, `PermutationBuilder`, and `Permutation` now work
against a `RemoteTable` (LanceDB Cloud and Enterprise), which unblocks
benchmarking the loader against the enterprise cluster cache.

```python
db = lancedb.connect("db://my-db", api_key=..., host_override=...)
ds = StreamingDataset(db.open_table("training"), world_size=8, rank=r)
```

Rows are addressed by `_rowid` exactly as before —
`PermutationReader::load_batch` already built the same `_rowid IN (...)`
filter that `Table::take_row_ids` sends, so the loader's fetch was
always the take path. It just was never allowed to run.

### The guard

`PermutationBuilder.__init__` rejected anything without `_inner`, so a
`RemoteTable` raised `TypeError` before reaching the PyO3 layer — which
already unwraps one via `_table._inner`.

### A bounded schema lookup

`PermutationReader::output_schema` reads the schema off a query plan,
and building a plan on a remote table *executes* the query
(`create_plan` → `execute_query`). With no limit that is `k =
isize::MAX`, so asking a remote table for its output schema pulled the
whole table over HTTP and threw it away — once per assigned split, on
every epoch, since `StreamingDataset.__iter__` constructs a
`Permutation` per split.

One row rather than zero, deliberately: lance gates its limit node on
`self.limit.unwrap_or(0) > 0`, so `Some(0)` means *no limit*.

### Tables with an LSM write spec are refused

A permutation references rows by row id, and rows that have not been
flushed to the base table do not have one yet. The loader could read
around them, but they would then be missing from training with nothing
said about it, so the build refuses such a table up front instead of
half supporting it.

### Fallible identity construction

`PermutationReader::identity` resolved `inner_new` with `unwrap`. That
was near total against a local dataset, but construction counts the base
table — an HTTP round trip for a remote one — so a transient network or
auth failure became a panic across the PyO3 boundary.

### Tests

End-to-end `permutation_builder` and `StreamingDataset` runs against a
mock server, the former torch-free so it runs wherever the suite does,
plus a test that a build succeeds without an LSM write spec and is
refused once one is installed.
2026-08-21 13:45:39 -07:00
Xuanwo 1baada89ef feat(python): bind function versions to columns (#4012)
A registered `FunctionVersion` has an exact identity and grouped output
contract, but the Python SDK cannot currently bind it to table columns
without manually constructing wire models.

Calling a `FunctionVersion` with named `col(...)` references now returns
one immutable `FunctionApplication` pinned to that exact version. The
application preserves named-struct outputs as one sibling group, while
`rename(columns=...)` defines the result-field to table-column mapping
consumed by `Table.add_columns`. Derived expressions and incomplete or
unknown input names fail before declaration.
2026-08-22 02:01:45 +08:00
83 changed files with 8739 additions and 433 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.38.0-beta.3"
current_version = "0.38.0-beta.4"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
Generated
+42 -42
View File
@@ -3455,8 +3455,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
[[package]]
name = "fsst"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"rand 0.9.5",
@@ -4815,8 +4815,8 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
[[package]]
name = "lance"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arc-swap",
"arrow",
@@ -4888,8 +4888,8 @@ dependencies = [
[[package]]
name = "lance-arrow"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4911,7 +4911,7 @@ dependencies = [
[[package]]
name = "lance-arrow-scalar"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4925,7 +4925,7 @@ dependencies = [
[[package]]
name = "lance-arrow-stats"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -4934,8 +4934,8 @@ dependencies = [
[[package]]
name = "lance-bitpacking"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrayref",
"crunchy",
@@ -4945,8 +4945,8 @@ dependencies = [
[[package]]
name = "lance-core"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4983,8 +4983,8 @@ dependencies = [
[[package]]
name = "lance-datafusion"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow",
"arrow-array",
@@ -5013,8 +5013,8 @@ dependencies = [
[[package]]
name = "lance-datagen"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow",
"arrow-array",
@@ -5031,8 +5031,8 @@ dependencies = [
[[package]]
name = "lance-derive"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"proc-macro2",
"quote",
@@ -5041,8 +5041,8 @@ dependencies = [
[[package]]
name = "lance-encoding"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5075,8 +5075,8 @@ dependencies = [
[[package]]
name = "lance-file"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5107,8 +5107,8 @@ dependencies = [
[[package]]
name = "lance-index"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arc-swap",
"arrow",
@@ -5172,8 +5172,8 @@ dependencies = [
[[package]]
name = "lance-index-core"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5195,8 +5195,8 @@ dependencies = [
[[package]]
name = "lance-io"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow",
"arrow-array",
@@ -5232,8 +5232,8 @@ dependencies = [
[[package]]
name = "lance-linalg"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5247,8 +5247,8 @@ dependencies = [
[[package]]
name = "lance-namespace"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow",
"async-trait",
@@ -5260,8 +5260,8 @@ dependencies = [
[[package]]
name = "lance-namespace-impls"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow",
"arrow-ipc",
@@ -5314,8 +5314,8 @@ dependencies = [
[[package]]
name = "lance-select"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5329,8 +5329,8 @@ dependencies = [
[[package]]
name = "lance-table"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow",
"arrow-array",
@@ -5370,8 +5370,8 @@ dependencies = [
[[package]]
name = "lance-testing"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5384,8 +5384,8 @@ dependencies = [
[[package]]
name = "lance-tokenizer"
version = "11.0.0-beta.18"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.18#7b6e2d3586e9c1b99326313533aba2150557ed65"
version = "11.0.0-beta.19"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.19#3128c0024427cb5bf8c04d492893ae45e78b0511"
dependencies = [
"frostem",
"icu_segmenter",
+14 -14
View File
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
lance = { "version" = "=11.0.0-beta.18", default-features = false, "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=11.0.0-beta.18", default-features = false, "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=11.0.0-beta.18", default-features = false, "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=11.0.0-beta.18", "tag" = "v11.0.0-beta.18", "git" = "https://github.com/lance-format/lance.git" }
lance = { "version" = "=11.0.0-beta.19", default-features = false, "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=11.0.0-beta.19", default-features = false, "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=11.0.0-beta.19", default-features = false, "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=11.0.0-beta.19", "tag" = "v11.0.0-beta.19", "git" = "https://github.com/lance-format/lance.git" }
lancedb = { path = "rust/lancedb", default-features = false }
ahash = "0.8"
# Note that this one does not include pyarrow
+1 -1
View File
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
<dependency>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-core</artifactId>
<version>0.38.0-beta.3</version>
<version>0.38.0-beta.4</version>
</dependency>
```
+25 -25
View File
@@ -37,6 +37,31 @@ latest and stays writable.
***
### cherryPick()
```ts
cherryPick(fromBranch, dryRun): Promise<CherryPickResult>
```
Cherry-pick a branch onto main.
Set `dryRun` to `true` to preview. A failed cherry-pick resolves
with `status: "failed"` instead of throwing.
#### Parameters
* **fromBranch**: `string`
Branch to cherry-pick from.
* **dryRun**: `boolean` = `false`
When true, only preview. Defaults to false.
#### Returns
`Promise`&lt;[`CherryPickResult`](../interfaces/CherryPickResult.md)&gt;
***
### create()
```ts
@@ -112,28 +137,3 @@ List all branches, mapping name to branch metadata.
#### Returns
`Promise`&lt;`Record`&lt;`string`, [`BranchContents`](BranchContents.md)&gt;&gt;
***
### merge()
```ts
merge(fromBranch, dryRun): Promise<MergeBranchResult>
```
Merge a branch into main.
Set `dryRun` to `true` to preview the merge. A rejected merge resolves
with `status: "rejected"` instead of throwing.
#### Parameters
* **fromBranch**: `string`
Branch to merge from.
* **dryRun**: `boolean` = `false`
When true, only preview the merge. Defaults to false.
#### Returns
`Promise`&lt;[`MergeBranchResult`](../interfaces/MergeBranchResult.md)&gt;
+75 -5
View File
@@ -169,6 +169,45 @@ Creates a new empty Table
***
### createMaterializedView()
```ts
abstract createMaterializedView(
name,
source,
options?): Promise<MaterializedView>
```
Define a materialized view named `name` over the table `source`.
The view is created empty, with the query recorded in its schema
metadata; `view.refresh()` computes the rows. The view is a normal
table: it can be queried, indexed and searched, and it appears in
`tableNames`. The source table must have stable row ids (create it with
the `newTableEnableStableRowIds` storage option); they keep the view's
provenance valid across source compactions and cannot be enabled after
a table exists. Local databases only.
#### Parameters
* **name**: `string`
* **source**: `string`
* **options?**
* **options.limit?**: `number`
* **options.select?**: [`MaterializedViewSelect`](../type-aliases/MaterializedViewSelect.md)
* **options.where?**: `string`
#### Returns
`Promise`&lt;[`MaterializedView`](MaterializedView.md)&gt;
***
### createNamespace()
```ts
@@ -499,6 +538,22 @@ List server-side jobs across the database's tables.
***
### listMaterializedViews()
```ts
abstract listMaterializedViews(): Promise<string[]>
```
The names of the materialized views in this database.
Found by reading every table's schema, so this costs an open per table.
#### Returns
`Promise`&lt;`string`[]&gt;
***
### listNamespaces()
```ts
@@ -529,6 +584,26 @@ Child namespace names and
***
### openMaterializedView()
```ts
abstract openMaterializedView(name): Promise<MaterializedView>
```
Open the materialized view named `name`.
Rejects a table that exists but is not a materialized view.
#### Parameters
* **name**: `string`
#### Returns
`Promise`&lt;[`MaterializedView`](MaterializedView.md)&gt;
***
### openTable()
```ts
@@ -538,18 +613,13 @@ abstract openTable(
options?): Promise<Table>
```
Open a table in the database.
#### Parameters
* **name**: `string`
The name of the table
* **namespacePath?**: `string`[]
The namespace path of the table (defaults to root namespace)
* **options?**: `Partial`&lt;[`OpenTableOptions`](../interfaces/OpenTableOptions.md)&gt;
Additional options
#### Returns
+101
View File
@@ -0,0 +1,101 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / MaterializedView
# Class: MaterializedView
A handle on a materialized view: its table plus its definition.
Obtained from [Connection#createMaterializedView](Connection.md#creatematerializedview) or
[Connection#openMaterializedView](Connection.md#openmaterializedview). The view is a normal table --
queries, indexes and search all apply through [MaterializedView#table](MaterializedView.md#table)
-- whose contents are maintained by [MaterializedView#refresh](MaterializedView.md#refresh).
## Constructors
### new MaterializedView()
```ts
new MaterializedView(table): MaterializedView
```
#### Parameters
* **table**: [`Table`](Table.md)
#### Returns
[`MaterializedView`](MaterializedView.md)
## Accessors
### name
```ts
get name(): string
```
#### Returns
`string`
## Methods
### definition()
```ts
definition(): Promise<MaterializedViewDefinition>
```
The query that defines the view, read from its stored schema.
#### Returns
`Promise`&lt;[`MaterializedViewDefinition`](../interfaces/MaterializedViewDefinition.md)&gt;
***
### refresh()
```ts
refresh(options?): Promise<RefreshMaterializedViewResult>
```
Recompute the view from its source.
The refresh is incremental when the source's changes can be reconciled
into the view -- rows added, changed or removed since the last one --
and otherwise rebuilds. `full` forces a rebuild; `sourceVersion`
refreshes to that source version instead of the latest.
Concurrent refreshes of one view do not duplicate its rows. Two that
plan the same source rows conflict on commit, and the loser throws
rather than writing them a second time.
#### Parameters
* **options?**
* **options.full?**: `boolean`
* **options.sourceVersion?**: `number`
#### Returns
`Promise`&lt;[`RefreshMaterializedViewResult`](../interfaces/RefreshMaterializedViewResult.md)&gt;
***
### table()
```ts
table(): Table
```
The view, as the table it is.
#### Returns
[`Table`](Table.md)
+7 -3
View File
@@ -28,6 +28,7 @@
- [Job](classes/Job.md)
- [MakeArrowTableOptions](classes/MakeArrowTableOptions.md)
- [MatchQuery](classes/MatchQuery.md)
- [MaterializedView](classes/MaterializedView.md)
- [MergeInsertBuilder](classes/MergeInsertBuilder.md)
- [MultiMatchQuery](classes/MultiMatchQuery.md)
- [NativeJsHeaderProvider](classes/NativeJsHeaderProvider.md)
@@ -59,6 +60,9 @@
- [BranchIndexSummary](interfaces/BranchIndexSummary.md)
- [BranchRowCountSummary](interfaces/BranchRowCountSummary.md)
- [BucketStats](interfaces/BucketStats.md)
- [CherryPickError](interfaces/CherryPickError.md)
- [CherryPickPreview](interfaces/CherryPickPreview.md)
- [CherryPickResult](interfaces/CherryPickResult.md)
- [ClientConfig](interfaces/ClientConfig.md)
- [ColumnAlteration](interfaces/ColumnAlteration.md)
- [ColumnOrdering](interfaces/ColumnOrdering.md)
@@ -98,10 +102,8 @@
- [ListNamespacesResponse](interfaces/ListNamespacesResponse.md)
- [LsmStats](interfaces/LsmStats.md)
- [LsmWriteSpec](interfaces/LsmWriteSpec.md)
- [MaterializedViewDefinition](interfaces/MaterializedViewDefinition.md)
- [MemtableStats](interfaces/MemtableStats.md)
- [MergeBlocker](interfaces/MergeBlocker.md)
- [MergeBranchResult](interfaces/MergeBranchResult.md)
- [MergePreview](interfaces/MergePreview.md)
- [MergeResult](interfaces/MergeResult.md)
- [NativeOAuthConfig](interfaces/NativeOAuthConfig.md)
- [OAuthConfig](interfaces/OAuthConfig.md)
@@ -110,6 +112,7 @@
- [OptimizeStats](interfaces/OptimizeStats.md)
- [QueryExecutionOptions](interfaces/QueryExecutionOptions.md)
- [RefreshColumnResult](interfaces/RefreshColumnResult.md)
- [RefreshMaterializedViewResult](interfaces/RefreshMaterializedViewResult.md)
- [RemovalStats](interfaces/RemovalStats.md)
- [RenameTableOptions](interfaces/RenameTableOptions.md)
- [RestNamespaceConfig](interfaces/RestNamespaceConfig.md)
@@ -142,6 +145,7 @@
- [FieldLike](type-aliases/FieldLike.md)
- [IntoSql](type-aliases/IntoSql.md)
- [IntoVector](type-aliases/IntoVector.md)
- [MaterializedViewSelect](type-aliases/MaterializedViewSelect.md)
- [MultiVector](type-aliases/MultiVector.md)
- [RecordBatchLike](type-aliases/RecordBatchLike.md)
- [SchemaLike](type-aliases/SchemaLike.md)
+8 -16
View File
@@ -50,6 +50,14 @@ changedColumns: BranchColumnChange[];
***
### errors
```ts
errors: CherryPickError[];
```
***
### fromBranch
```ts
@@ -66,22 +74,6 @@ mainVersion: number;
***
### mergeBlockers
```ts
mergeBlockers: MergeBlocker[];
```
***
### mergeable
```ts
mergeable: boolean;
```
***
### parentVersion
```ts
@@ -2,11 +2,11 @@
***
[@lancedb/lancedb](../globals.md) / MergeBlocker
[@lancedb/lancedb](../globals.md) / CherryPickError
# Interface: MergeBlocker
# Interface: CherryPickError
A reason why a branch cannot currently be merged.
A reason why a cherry-pick cannot currently land.
## Properties
@@ -0,0 +1,17 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / CherryPickPreview
# Interface: CherryPickPreview
Changes that would be, or were, promoted by a cherry-pick.
## Properties
### promotedColumns
```ts
promotedColumns: string[];
```
@@ -2,11 +2,11 @@
***
[@lancedb/lancedb](../globals.md) / MergeBranchResult
[@lancedb/lancedb](../globals.md) / CherryPickResult
# Interface: MergeBranchResult
# Interface: CherryPickResult
Result of previewing or attempting a branch merge.
Result of previewing or attempting a cherry-pick.
## Properties
@@ -29,7 +29,7 @@ optional mainVersionAfter: number;
### preview
```ts
preview: MergePreview;
preview: CherryPickPreview;
```
***
@@ -38,9 +38,9 @@ preview: MergePreview;
```ts
status:
| "failed"
| "unknown"
| "rejected"
| "ready"
| "notImplemented"
| "merged";
| "cherryPicked";
```
@@ -0,0 +1,59 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / MaterializedViewDefinition
# Interface: MaterializedViewDefinition
The query that defines a materialized view.
## Properties
### filter?
```ts
optional filter: string;
```
SQL predicate selecting the source rows the view holds.
***
### inputs
```ts
inputs: string[];
```
Source columns the projections and filter read.
***
### limit?
```ts
optional limit: number;
```
Cap on the number of rows the view holds.
***
### projections
```ts
projections: [string, string][];
```
`[output column, SQL expression]` pairs, in view schema order.
***
### sourceTable
```ts
sourceTable: string;
```
Name of the source table, in the same database as the view.
-17
View File
@@ -1,17 +0,0 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / MergePreview
# Interface: MergePreview
Changes that would be, or were, promoted by a branch merge.
## Properties
### promotedColumns
```ts
promotedColumns: string[];
```
@@ -0,0 +1,41 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / RefreshMaterializedViewResult
# Interface: RefreshMaterializedViewResult
## Properties
### mode
```ts
mode: string;
```
How the view was brought up to date: "rebuild", "incremental" or "no_op".
***
### rowsWritten
```ts
rowsWritten: number;
```
***
### sourceVersion
```ts
sourceVersion: number;
```
***
### version
```ts
version: number;
```
@@ -0,0 +1,14 @@
[**@lancedb/lancedb**](../README.md) • **Docs**
***
[@lancedb/lancedb](../globals.md) / MaterializedViewSelect
# Type Alias: MaterializedViewSelect
```ts
type MaterializedViewSelect: (string | [string, string])[] | Record<string, string>;
```
The view's columns: column names, `[alias, SQL expression]` pairs, or a
record of the same. A bare name projects itself.
+10
View File
@@ -102,6 +102,12 @@ listing a storage directory.
::: lancedb.job.AsyncJob
## Materialized Views (Synchronous)
::: lancedb.materialized_view.MaterializedView
::: lancedb.materialized_view.MaterializedViewDefinition
## Expressions
Type-safe expression builder for filters and projections. Use these instead
@@ -295,6 +301,10 @@ Table hold your actual data as a collection of records / rows.
::: lancedb.table.AsyncBranches
## Materialized Views (Asynchronous)
::: lancedb.materialized_view.AsyncMaterializedView
## Indices (Asynchronous)
Indices can be created on a table to speed up queries. This section
+1 -1
View File
@@ -8,7 +8,7 @@
<parent>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.38.0-beta.3</version>
<version>0.38.0-beta.4</version>
<relativePath>../pom.xml</relativePath>
</parent>
+2 -2
View File
@@ -6,7 +6,7 @@
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.38.0-beta.3</version>
<version>0.38.0-beta.4</version>
<packaging>pom</packaging>
<name>${project.artifactId}</name>
<description>LanceDB Java SDK Parent POM</description>
@@ -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.18</lance-core.version>
<lance-core.version>11.0.0-beta.19</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>
+1 -1
View File
@@ -1,7 +1,7 @@
[package]
name = "lancedb-nodejs"
edition.workspace = true
version = "0.38.0-beta.3"
version = "0.38.0-beta.4"
publish = false
license.workspace = true
description.workspace = true
+48
View File
@@ -487,4 +487,52 @@ describe("embedding functions", () => {
expect(stringSchema3).toEqual(stringExpectedSchema);
},
);
test("parses one function writing several vector columns", async () => {
class MockEmbeddingFunction extends EmbeddingFunction<string> {
ndims() {
return 3;
}
embeddingDataType(): Float {
return new Float32();
}
async computeQueryEmbeddings(_data: string) {
return [1, 2, 3];
}
async computeSourceEmbeddings(data: string[]) {
return Array.from({ length: data.length }).fill([
1, 2, 3,
]) as number[][];
}
}
const registry = getRegistry();
registry.register("multi_output_mock")(MockEmbeddingFunction);
// A materialized view can project one source vector column under two
// names, so a table's configuration names the same function twice.
const parsed = await registry.parseFunctions(
new Map([
[
"embedding_functions",
JSON.stringify([
{
name: "multi_output_mock",
sourceColumn: "text",
vectorColumn: "vector_a",
model: {},
},
{
name: "multi_output_mock",
sourceColumn: "text",
vectorColumn: "vector_b",
model: {},
},
]),
],
]),
);
expect(
[...parsed.values()].map(({ vectorColumn }) => vectorColumn).sort(),
).toEqual(["vector_a", "vector_b"]);
});
});
+147
View File
@@ -0,0 +1,147 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
import * as tmp from "tmp";
import { Connection, connect } from "../lancedb";
import {
DEFINITION_META_KEY,
definitionFromMetadata,
} from "../lancedb/materialized_view";
describe("materialized views", () => {
let tmpDir: tmp.DirResult;
let db: Connection;
beforeEach(async () => {
tmpDir = tmp.dirSync({ unsafeCleanup: true });
db = await connect(tmpDir.name);
await db.createTable(
"people",
[
{ name: "ada", age: 36 },
{ name: "kid", age: 7 },
{ name: "grace", age: 85 },
],
{ storageOptions: { newTableEnableStableRowIds: "true" } },
);
});
afterEach(() => tmpDir.removeCallback());
it("rejects a stored limit a number cannot carry", () => {
const big = new Map([
[
DEFINITION_META_KEY,
'{"kind":"select","source_table":"people","limit":9007199254740993}',
],
]);
expect(() => definitionFromMetadata(big, "v")).toThrow(
/too large to represent exactly/,
);
const safe = new Map([
[
DEFINITION_META_KEY,
'{"kind":"select","source_table":"people","limit":42}',
],
]);
expect(definitionFromMetadata(safe, "v").limit).toBe(42);
});
it("creates, refreshes and queries a view", async () => {
const view = await db.createMaterializedView("adults", "people", {
select: ["name", ["shout", "upper(name)"]],
where: "age >= 18",
});
expect(view.name).toBe("adults");
expect(await view.table().countRows()).toBe(0);
const result = await view.refresh();
expect(result.mode).toBe("rebuild");
expect(Number(result.rowsWritten)).toBe(2);
const rows = await view.table().query().toArray();
expect(rows.map((r) => r.shout).sort()).toEqual(["ADA", "GRACE"]);
});
it("round-trips the definition", async () => {
await db.createMaterializedView("adults", "people", {
where: "age >= 18",
});
const view = await db.openMaterializedView("adults");
const definition = await view.definition();
expect(definition.sourceTable).toBe("people");
expect(definition.filter).toBe("age >= 18");
expect(definition.projections).toEqual([
["name", "`name`"],
["age", "`age`"],
]);
expect(definition.inputs).toEqual(["age", "name"]);
});
it("refreshes incrementally after an append", async () => {
const view = await db.createMaterializedView("copy", "people");
await view.refresh();
const people = await db.openTable("people");
await people.add([{ name: "alan", age: 41 }]);
const result = await view.refresh();
expect(result.mode).toBe("incremental");
expect(Number(result.rowsWritten)).toBe(1);
expect(await view.table().countRows()).toBe(4);
expect((await view.refresh()).mode).toBe("no_op");
});
it("lists views and rejects non-views", async () => {
await db.createMaterializedView("adults", "people", {
where: "age >= 18",
});
expect(await db.listMaterializedViews()).toEqual(["adults"]);
await expect(db.openMaterializedView("people")).rejects.toThrow(
"not a materialized view",
);
});
it("rejects an invalid expression at create time", async () => {
await expect(
db.createMaterializedView("bad", "people", {
select: [["x", "missing + 1"]],
}),
).rejects.toThrow("missing");
});
it("rejects invalid numeric options before creating anything", async () => {
for (const limit of [-5, 1.5, Infinity, NaN]) {
await expect(
db.createMaterializedView("bad", "people", { limit }),
).rejects.toThrow("non-negative integer");
}
expect(await db.listMaterializedViews()).toEqual([]);
const view = await db.createMaterializedView("copy", "people");
for (const sourceVersion of [-1, 1.5, Infinity, NaN]) {
await expect(view.refresh({ sourceVersion })).rejects.toThrow(
"non-negative integer",
);
}
});
it("quotes bare select names", async () => {
await db.createTable("odd_names", [{ "order item": "widget" }], {
storageOptions: { newTableEnableStableRowIds: "true" },
});
const view = await db.createMaterializedView("quoted", "odd_names", {
select: ["order item"],
});
const result = await view.refresh();
expect(Number(result.rowsWritten)).toBe(1);
});
it("requires stable row ids on the source", async () => {
await db.createTable("plain", [{ x: 1 }]);
await expect(db.createMaterializedView("v", "plain")).rejects.toThrow(
"stable row ids",
);
});
});
+31 -14
View File
@@ -75,6 +75,25 @@ async function withMockDatabase(
}
describe("remote connection", () => {
it("refuses materialized views before issuing any request", async () => {
const paths: string[] = [];
await withMockDatabase(
(req, res) => {
paths.push(req.url ?? "");
res.writeHead(404).end();
},
async (db) => {
await expect(db.openMaterializedView("secret_table")).rejects.toThrow(
/only on local databases/,
);
await expect(db.listMaterializedViews()).rejects.toThrow(
/only on local databases/,
);
expect(paths).toEqual([]);
},
);
});
it("should accept partial connection options", async () => {
await connect("db://test", {
apiKey: "fake",
@@ -311,7 +330,7 @@ describe("remote connection", () => {
expect(createIndexBody?.["custom_stop_words"]).toEqual(["the"]);
});
it("diffs and merges remote branches", async () => {
it("diffs and cherry-picks remote branches", async () => {
const sampleDiff = {
fromBranch: "exp",
parentVersion: 1,
@@ -333,10 +352,9 @@ describe("remote connection", () => {
changedColumns: [],
addedIndexes: [],
removedIndexes: [],
mergeable: true,
mergeBlockers: [],
errors: [],
};
const mergeBodies: Record<string, unknown>[] = [];
const cherryPickBodies: Record<string, unknown>[] = [];
await withMockDatabase(
(req, res) => {
@@ -366,17 +384,16 @@ describe("remote connection", () => {
.end(JSON.stringify(sampleDiff));
return;
}
if (path.endsWith("/branches/merge/")) {
mergeBodies.push(body);
if (path.endsWith("/branches/cherry_pick/")) {
cherryPickBodies.push(body);
const dryRun = body["dry_run"] === true;
const response = {
status: dryRun ? "ready" : "rejected",
status: dryRun ? "ready" : "failed",
diff: dryRun
? sampleDiff
: {
...sampleDiff,
mergeable: false,
mergeBlockers: [
errors: [
{ code: "baseMoved", message: "main has advanced" },
],
},
@@ -398,19 +415,19 @@ describe("remote connection", () => {
await expect(branches.diff("exp")).resolves.toEqual(sampleDiff);
const rejected = await branches.merge("exp");
expect(rejected.status).toBe("rejected");
expect(rejected.diff.mergeBlockers).toEqual([
const failed = await branches.cherryPick("exp");
expect(failed.status).toBe("failed");
expect(failed.diff.errors).toEqual([
{ code: "baseMoved", message: "main has advanced" },
]);
const preview = await branches.merge("exp", true);
const preview = await branches.cherryPick("exp", true);
expect(preview.status).toBe("ready");
expect(preview.preview.promotedColumns).toEqual(["tag"]);
},
);
expect(mergeBodies).toEqual([
expect(cherryPickBodies).toEqual([
// biome-ignore lint/style/useNamingConvention: snake_case mandated by the server wire format
{ from_branch: "exp", dry_run: false },
// biome-ignore lint/style/useNamingConvention: snake_case mandated by the server wire format
+1 -1
View File
@@ -2953,7 +2953,7 @@ describe("column name options", () => {
.limit(10)
.toArray();
expect(results2.length).toBe(10);
});
}, 30_000);
});
describe("when creating an empty table", () => {
+70
View File
@@ -16,6 +16,12 @@ import {
makeEmptyTable,
} from "./arrow";
import { EmbeddingFunctionConfig, getRegistry } from "./embedding/registry";
import {
MaterializedView,
MaterializedViewSelect,
normalizeSelect,
validateNonNegativeInteger,
} from "./materialized_view";
import { Connection as LanceDbConnection } from "./native";
import type {
CreateNamespaceResponse,
@@ -247,6 +253,41 @@ export abstract class Connection {
* @param {string[]} namespacePath - The namespace path of the table (defaults to root namespace)
* @param {Partial<OpenTableOptions>} options - Additional options
*/
/**
* Define a materialized view named `name` over the table `source`.
*
* The view is created empty, with the query recorded in its schema
* metadata; `view.refresh()` computes the rows. The view is a normal
* table: it can be queried, indexed and searched, and it appears in
* `tableNames`. The source table must have stable row ids (create it with
* the `newTableEnableStableRowIds` storage option); they keep the view's
* provenance valid across source compactions and cannot be enabled after
* a table exists. Local databases only.
*/
abstract createMaterializedView(
name: string,
source: string,
options?: {
select?: MaterializedViewSelect;
where?: string;
limit?: number;
},
): Promise<MaterializedView>;
/**
* Open the materialized view named `name`.
*
* Rejects a table that exists but is not a materialized view.
*/
abstract openMaterializedView(name: string): Promise<MaterializedView>;
/**
* The names of the materialized views in this database.
*
* Found by reading every table's schema, so this costs an open per table.
*/
abstract listMaterializedViews(): Promise<string[]>;
abstract openTable(
name: string,
namespacePath?: string[],
@@ -531,6 +572,35 @@ export class LocalConnection extends Connection {
);
}
async createMaterializedView(
name: string,
source: string,
options?: {
select?: MaterializedViewSelect;
where?: string;
limit?: number;
},
): Promise<MaterializedView> {
validateNonNegativeInteger(options?.limit, "limit");
const innerTable = await this.inner.createMaterializedView(
name,
source,
normalizeSelect(options?.select),
options?.where,
options?.limit,
);
return new MaterializedView(new LocalTable(innerTable));
}
async openMaterializedView(name: string): Promise<MaterializedView> {
const innerTable = await this.inner.openMaterializedView(name);
return new MaterializedView(new LocalTable(innerTable));
}
async listMaterializedViews(): Promise<string[]> {
return await this.inner.listMaterializedViews();
}
async openTable(
name: string,
namespacePath?: string[],
+9 -3
View File
@@ -21,6 +21,11 @@ import type { BaseTokenizer } from "./indices";
import type { FtsToken } from "./table";
// Re-export native header provider for use with connectWithHeaderProvider
export {
MaterializedView,
MaterializedViewDefinition,
MaterializedViewSelect,
} from "./materialized_view";
export { JsHeaderProvider as NativeJsHeaderProvider } from "./native.js";
// OpenTelemetry metrics bridge. Only the high-level entry point is public; the
@@ -51,6 +56,7 @@ export {
AddResult,
AddColumnsResult,
RefreshColumnResult,
RefreshMaterializedViewResult,
AlterColumnsResult,
UpdateFieldMetadataResult,
DeleteResult,
@@ -135,10 +141,10 @@ export {
BranchColumnChange,
BranchIndexSummary,
BranchRowCountSummary,
MergeBlocker,
CherryPickError,
BranchDiff,
MergePreview,
MergeBranchResult,
CherryPickPreview,
CherryPickResult,
AddDataOptions,
UpdateOptions,
OptimizeOptions,
+161
View File
@@ -0,0 +1,161 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
import { RefreshMaterializedViewResult } from "./native";
import { Table } from "./table";
/** Schema metadata key holding a materialized view's definition. */
export const DEFINITION_META_KEY = "mv.definition";
/** The query that defines a materialized view. */
export interface MaterializedViewDefinition {
/** Name of the source table, in the same database as the view. */
sourceTable: string;
/** `[output column, SQL expression]` pairs, in view schema order. */
projections: [string, string][];
/** SQL predicate selecting the source rows the view holds. */
filter?: string;
/** Cap on the number of rows the view holds. */
limit?: number;
/** Source columns the projections and filter read. */
inputs: string[];
}
/**
* The view's columns: column names, `[alias, SQL expression]` pairs, or a
* record of the same. A bare name projects itself.
*/
export type MaterializedViewSelect =
| (string | [string, string])[]
| Record<string, string>;
/**
* @internal Reject a numeric option N-API would otherwise silently coerce:
* `Infinity` reaches Rust as 0, `1.5` as 1.
*/
export function validateNonNegativeInteger(
value: number | undefined,
name: string,
): void {
if (value !== undefined && !(Number.isSafeInteger(value) && value >= 0)) {
throw new Error(`${name} must be a non-negative integer`);
}
}
/** @internal Quote a column name as a Lance SQL identifier (backticks). */
function quoteIdentifier(name: string): string {
return "`" + name.replace(/`/g, "``") + "`";
}
/**
* @internal Normalize a select argument into `[alias, expression]` pairs.
* A bare name projects itself and is quoted, so any valid column name works;
* pair and record entries are kept verbatim because their right side is an
* expression.
*/
export function normalizeSelect(
select?: MaterializedViewSelect,
): [string, string][] | undefined {
if (select === undefined) {
return undefined;
}
if (Array.isArray(select)) {
return select.map((item) =>
typeof item === "string" ? [item, quoteIdentifier(item)] : item,
);
}
return Object.entries(select);
}
/** @internal Parse a definition off a table's stored schema metadata. */
export function definitionFromMetadata(
metadata: Map<string, string>,
name: string,
): MaterializedViewDefinition {
const raw = metadata.get(DEFINITION_META_KEY);
if (raw === undefined) {
throw new Error(`Table '${name}' is not a materialized view`);
}
// biome-ignore lint/suspicious/noExplicitAny: raw JSON
const value: any = JSON.parse(raw);
if (value.kind !== "select") {
throw new Error(
`materialized view '${name}' is defined by '${value.kind}', which this ` +
"version of lancedb cannot refresh",
);
}
const limit = value.limit ?? undefined;
// JSON.parse rounds integers past 2^53; every exact u64 parses to a safe
// integer and every rounded one does not, so this rejects precisely the
// values a number cannot carry.
if (limit !== undefined && !Number.isSafeInteger(limit)) {
throw new Error(
`materialized view '${name}' has a stored limit too large to represent exactly`,
);
}
return {
sourceTable: value.source_table,
// biome-ignore lint/suspicious/noExplicitAny: raw JSON
projections: (value.projections ?? []).map((p: any) => [
p.output,
p.expression,
]),
filter: value.filter ?? undefined,
limit,
inputs: value.inputs ?? [],
};
}
/**
* A handle on a materialized view: its table plus its definition.
*
* Obtained from {@link Connection#createMaterializedView} or
* {@link Connection#openMaterializedView}. The view is a normal table --
* queries, indexes and search all apply through {@link MaterializedView#table}
* -- whose contents are maintained by {@link MaterializedView#refresh}.
*/
export class MaterializedView {
private readonly inner: Table;
constructor(table: Table) {
this.inner = table;
}
get name(): string {
return this.inner.name;
}
/** The view, as the table it is. */
table(): Table {
return this.inner;
}
/** The query that defines the view, read from its stored schema. */
async definition(): Promise<MaterializedViewDefinition> {
const schema = await this.inner.schema();
return definitionFromMetadata(schema.metadata, this.name);
}
/**
* Recompute the view from its source.
*
* The refresh is incremental when the source's changes can be reconciled
* into the view -- rows added, changed or removed since the last one --
* and otherwise rebuilds. `full` forces a rebuild; `sourceVersion`
* refreshes to that source version instead of the latest.
*
* Concurrent refreshes of one view do not duplicate its rows. Two that
* plan the same source rows conflict on commit, and the loser throws
* rather than writing them a second time.
*/
async refresh(options?: {
full?: boolean;
sourceVersion?: number;
}): Promise<RefreshMaterializedViewResult> {
validateNonNegativeInteger(options?.sourceVersion, "sourceVersion");
return await this.inner.refreshMaterializedView(
options?.full,
options?.sourceVersion,
);
}
}
+38 -19
View File
@@ -35,6 +35,7 @@ import {
Branches as NativeBranches,
OptimizeStats,
RefreshColumnResult,
RefreshMaterializedViewResult,
TableStatistics,
Tags,
UpdateFieldMetadataResult,
@@ -602,6 +603,18 @@ export abstract class Table {
*/
abstract refreshColumnAsync(column: string): Promise<Job>;
/**
* Recompute this table's contents from its materialized-view definition.
*
* Plumbing for {@link MaterializedView.refresh}, which is the way to call
* it: rejects tables that carry no view definition. Local tables only.
* @ignore
*/
abstract refreshMaterializedView(
full?: boolean,
sourceVersion?: number,
): Promise<RefreshMaterializedViewResult>;
/**
* Alter the name or nullability of columns.
* @param {ColumnAlteration[]} columnAlterations One or more alterations to
@@ -1264,6 +1277,13 @@ export class LocalTable extends Table {
return await this.inner.refreshColumnAsync(column);
}
async refreshMaterializedView(
full?: boolean,
sourceVersion?: number,
): Promise<RefreshMaterializedViewResult> {
return await this.inner.refreshMaterializedView(full, sourceVersion);
}
async alterColumns(
columnAlterations: ColumnAlteration[],
): Promise<AlterColumnsResult> {
@@ -1557,8 +1577,8 @@ export interface BranchRowCountSummary {
deltaAvailable: boolean;
}
/** A reason why a branch cannot currently be merged. */
export interface MergeBlocker {
/** A reason why a cherry-pick cannot currently land. */
export interface CherryPickError {
code: string;
message: string;
}
@@ -1578,20 +1598,19 @@ export interface BranchDiff {
changedColumns: BranchColumnChange[];
addedIndexes: BranchIndexSummary[];
removedIndexes: BranchIndexSummary[];
mergeable: boolean;
mergeBlockers: MergeBlocker[];
errors: CherryPickError[];
}
/** Changes that would be, or were, promoted by a branch merge. */
export interface MergePreview {
/** Changes that would be, or were, promoted by a cherry-pick. */
export interface CherryPickPreview {
promotedColumns: string[];
}
/** Result of previewing or attempting a branch merge. */
export interface MergeBranchResult {
status: "ready" | "rejected" | "notImplemented" | "merged" | "unknown";
/** Result of previewing or attempting a cherry-pick. */
export interface CherryPickResult {
status: "ready" | "failed" | "notImplemented" | "cherryPicked" | "unknown";
diff: BranchDiff;
preview: MergePreview;
preview: CherryPickPreview;
mainVersionAfter?: number;
}
@@ -1654,21 +1673,21 @@ export class Branches {
}
/**
* Merge a branch into main.
* Cherry-pick a branch onto main.
*
* Set `dryRun` to `true` to preview the merge. A rejected merge resolves
* with `status: "rejected"` instead of throwing.
* Set `dryRun` to `true` to preview. A failed cherry-pick resolves
* with `status: "failed"` instead of throwing.
*
* @param fromBranch Branch to merge from.
* @param dryRun When true, only preview the merge. Defaults to false.
* @param fromBranch Branch to cherry-pick from.
* @param dryRun When true, only preview. Defaults to false.
*/
async merge(
async cherryPick(
fromBranch: string,
dryRun: boolean = false,
): Promise<MergeBranchResult> {
return (await this.#inner.merge(
): Promise<CherryPickResult> {
return (await this.#inner.cherryPick(
fromBranch,
dryRun,
)) as unknown as MergeBranchResult;
)) as unknown as CherryPickResult;
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-darwin-arm64",
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"os": ["darwin"],
"cpu": ["arm64"],
"main": "lancedb.darwin-arm64.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-gnu",
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-musl",
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-gnu",
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-musl",
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-arm64-msvc",
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"os": [
"win32"
],
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-x64-msvc",
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"os": ["win32"],
"cpu": ["x64"],
"main": "lancedb.win32-x64-msvc.node",
+1 -1
View File
@@ -11,7 +11,7 @@
"ann"
],
"private": false,
"version": "0.38.0-beta.3",
"version": "0.38.0-beta.4",
"main": "dist/index.js",
"exports": {
".": "./dist/index.js",
+52
View File
@@ -266,6 +266,58 @@ impl Connection {
Ok(Table::new(tbl))
}
#[napi(catch_unwind)]
pub async fn create_materialized_view(
&self,
name: String,
source: String,
projections: Option<Vec<Vec<String>>>,
filter: Option<String>,
limit: Option<i64>,
) -> napi::Result<Table> {
let mut builder = self.get_inner()?.create_materialized_view(name, source);
if let Some(projections) = projections {
let mut pairs = Vec::with_capacity(projections.len());
for pair in projections {
let [output, expression]: [String; 2] = pair.try_into().map_err(|_| {
napi::Error::from_reason("each projection must be an [output, expression] pair")
})?;
pairs.push((output, expression));
}
builder = builder.select(pairs);
}
if let Some(filter) = filter {
builder = builder.only_if(filter);
}
if let Some(limit) = limit {
let limit = u64::try_from(limit)
.map_err(|_| napi::Error::from_reason("limit must be a non-negative integer"))?;
builder = builder.limit(limit);
}
let view = builder.execute().await.default_error()?;
Ok(Table::new(view.table().clone()))
}
#[napi(catch_unwind)]
pub async fn open_materialized_view(&self, name: String) -> napi::Result<Table> {
let view = self
.get_inner()?
.open_materialized_view(&name)
.await
.default_error()?;
Ok(Table::new(view.table().clone()))
}
#[napi(catch_unwind)]
pub async fn list_materialized_views(&self) -> napi::Result<Vec<String>> {
let views = self
.get_inner()?
.list_materialized_views()
.await
.default_error()?;
Ok(views.into_iter().map(|v| v.name).collect())
}
#[napi(catch_unwind)]
pub async fn open_table(
&self,
+4
View File
@@ -1,6 +1,10 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
// The materialized-view refresh future deepens the type graph past the
// default trait-recursion depth; same raise as the core crate applies.
#![recursion_limit = "256"]
use std::collections::HashMap;
use env_logger::Env;
+48 -3
View File
@@ -381,6 +381,26 @@ impl Table {
Ok(crate::job::Job::new(job))
}
#[napi(catch_unwind)]
pub async fn refresh_materialized_view(
&self,
full: Option<bool>,
source_version: Option<i64>,
) -> napi::Result<RefreshMaterializedViewResult> {
let view = lancedb::MaterializedView::from_table(self.inner_ref()?.clone())
.await
.default_error()?;
let mut builder = view.refresh().full(full.unwrap_or(false));
if let Some(version) = source_version {
let version = u64::try_from(version).map_err(|_| {
napi::Error::from_reason("sourceVersion must be a non-negative integer")
})?;
builder = builder.source_version(version);
}
let result = builder.execute().await.default_error()?;
Ok(result.into())
}
#[napi(catch_unwind)]
pub async fn add_columns_with_schema(
&self,
@@ -1387,6 +1407,31 @@ pub struct RefreshColumnResult {
pub version: i64,
}
#[napi(object)]
pub struct RefreshMaterializedViewResult {
/// How the view was brought up to date: "rebuild", "incremental" or "no_op".
pub mode: String,
pub rows_written: i64,
pub source_version: i64,
pub version: i64,
}
impl From<lancedb::RefreshMaterializedViewResult> for RefreshMaterializedViewResult {
fn from(value: lancedb::RefreshMaterializedViewResult) -> Self {
let mode = match value.mode {
lancedb::RefreshMode::Rebuild => "rebuild",
lancedb::RefreshMode::Incremental => "incremental",
lancedb::RefreshMode::NoOp => "no_op",
};
Self {
mode: mode.to_string(),
rows_written: value.rows_written as i64,
source_version: value.source_version as i64,
version: value.version as i64,
}
}
}
impl From<lancedb::table::RefreshColumnResult> for RefreshColumnResult {
fn from(value: lancedb::table::RefreshColumnResult) -> Self {
Self {
@@ -1605,18 +1650,18 @@ impl Branches {
}
#[napi(ts_return_type = "Promise<Record<string, unknown>>")]
pub async fn merge(
pub async fn cherry_pick(
&self,
from_branch: String,
dry_run: Option<bool>,
) -> napi::Result<serde_json::Value> {
let result = self
.inner
.merge_branch(&from_branch, dry_run.unwrap_or(false))
.cherry_pick(&from_branch, dry_run.unwrap_or(false))
.await
.default_error()?;
serde_json::to_value(result).map_err(|err| {
napi::Error::from_reason(format!("failed to serialize branch merge result: {err}"))
napi::Error::from_reason(format!("failed to serialize cherry-pick result: {err}"))
})
}
}
+5 -6
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb-python"
version = "0.38.0-beta.3"
version = "0.38.0-beta.4"
publish = false
edition.workspace = true
description = "Python bindings for LanceDB"
@@ -26,7 +26,9 @@ lance-namespace-impls.workspace = true
lance-io.workspace = true
env_logger.workspace = true
log.workspace = true
pyo3 = { version = "0.28", features = ["extension-module", "abi3-py310", "chrono"] }
# Maturin enables extension-module mode for Python builds. Keeping it out of
# Cargo features lets Rust unit tests link against libpython.
pyo3 = { version = "0.28", features = ["abi3-py310", "chrono"] }
chrono.workspace = true
pyo3-async-runtimes = { version = "0.28", features = [
"attributes",
@@ -41,10 +43,7 @@ tokio.workspace = true
libc = "0.2"
[build-dependencies]
pyo3-build-config = { version = "0.28", features = [
"extension-module",
"abi3-py310",
] }
pyo3-build-config = { version = "0.28", features = ["abi3-py310"] }
[features]
default = ["remote", "lancedb/aws", "lancedb/gcs", "lancedb/azure", "lancedb/dynamodb", "lancedb/oss", "lancedb/huggingface", "lancedb/cos", "lancedb/goosefs", "lancedb/metrics-otel"]
+19
View File
@@ -38,6 +38,25 @@ Stable releases are created about every 2 weeks. For the latest features and bug
pip install --pre --extra-index-url https://pypi.fury.io/lancedb/ lancedb
```
### Threading in CPU-limited containers
LanceDB uses separate pools for compute work and storage I/O. On a container with
two visible CPUs, current releases intentionally use one compute worker by default;
no manual configuration is needed. If every query logs an I/O core reservation
warning on a two-CPU container, upgrade from LanceDB 0.21.1 or earlier.
The two commonly tuned environment variables control different resources:
- `LANCE_CPU_THREADS` overrides the number of compute workers. One worker is the
appropriate setting for a two-CPU container when an explicit override is needed.
- `LANCE_IO_THREADS` controls concurrent storage operations, not reserved CPU
cores. Its default can be greater than the number of CPUs because I/O workers
spend much of their time waiting for storage.
Keep the defaults unless measurements show that the workload benefits from an
override. See the [Lance threading model](https://lance.org/guide/performance/#threading-model)
for the current defaults and tuning guidance.
## Usage
### Basic Example
+1 -1
View File
@@ -103,7 +103,7 @@ python-source = "python"
module-name = "lancedb._lancedb"
[build-system]
requires = ["maturin>=1.4"]
requires = ["maturin>=1.9.4"]
build-backend = "maturin"
[tool.ruff.lint]
+8
View File
@@ -32,6 +32,11 @@ from .functions import (
UdfDefinition as UdfDefinition,
udf as udf,
)
from .materialized_view import (
AsyncMaterializedView,
MaterializedView,
MaterializedViewDefinition,
)
from .table import AsyncTable, Table
from .types import BaseTokenizerType
from ._lancedb import Session
@@ -506,6 +511,9 @@ async def connect_async(
__all__ = [
"AsyncMaterializedView",
"MaterializedView",
"MaterializedViewDefinition",
"connect",
"connect_async",
"tokenize",
+19 -1
View File
@@ -197,6 +197,15 @@ class Connection(object):
cur_namespace_path: Optional[List[str]] = None,
new_namespace_path: Optional[List[str]] = None,
) -> None: ...
async def create_materialized_view(
self,
name: str,
source: str,
projections: Optional[List[Tuple[str, str]]] = None,
filter: Optional[str] = None,
limit: Optional[int] = None,
) -> Table: ...
async def list_materialized_views(self) -> List[str]: ...
async def drop_table(
self, name: str, namespace_path: Optional[List[str]] = None
) -> None: ...
@@ -355,6 +364,9 @@ class Table:
) -> AddColumnsResult: ...
async def refresh_column(self, column: str) -> RefreshColumnResult: ...
async def refresh_column_async(self, column: str) -> Job: ...
async def refresh_materialized_view(
self, full: bool = False, source_version: Optional[int] = None
) -> RefreshMaterializedViewResult: ...
async def add_columns_with_schema(self, schema: pa.Schema) -> AddColumnsResult: ...
async def alter_columns(
self, columns: list[dict[str, Any]]
@@ -420,7 +432,7 @@ class Branches:
async def checkout(self, name: str, version: Optional[int] = None) -> Table: ...
async def delete(self, name: str) -> None: ...
async def diff(self, from_branch: str) -> Dict[str, Any]: ...
async def merge(
async def cherry_pick(
self, from_branch: str, dry_run: bool = False
) -> Dict[str, Any]: ...
@@ -704,6 +716,12 @@ class RefreshColumnResult:
rows_filled: int
version: int
class RefreshMaterializedViewResult:
mode: str
rows_written: int
source_version: int
version: int
class AlterColumnsResult:
version: int
+166
View File
@@ -47,6 +47,12 @@ from . import __version__
from ._lancedb import connect as lancedb_connect # type: ignore
from .functions import FunctionVersion, UdfDefinition
from .job import AsyncJob, Job, _function_job
from .materialized_view import (
AsyncMaterializedView,
MaterializedView,
SelectArg,
normalize_select,
)
from .table import (
AsyncTable,
LanceTable,
@@ -510,6 +516,70 @@ class DBConnection(EnforceOverrides):
"""
raise NotImplementedError
def create_materialized_view(
self,
name: str,
source: str,
*,
select: SelectArg = None,
where: Optional[str] = None,
limit: Optional[int] = None,
) -> MaterializedView:
"""Define a materialized view named ``name`` over the table ``source``.
The view is created empty, with the query recorded in its schema
metadata; ``view.refresh()`` computes the rows. The view is a normal
table: it can be queried, indexed and searched, and it appears in
``table_names``. Local databases only.
The source table must have stable row ids (create it with the
``new_table_enable_stable_row_ids`` storage option): they keep the
view's provenance valid across source compactions, and cannot be
enabled after a table exists.
Parameters
----------
name: str
The name of the view.
source: str
The name of the source table, in this database.
select: list or dict, optional
The view's columns: column names, ``(alias, SQL expression)``
pairs, or a dict of the same. Omitting it selects every source
column, expanded against the source schema at creation time.
where: str, optional
SQL predicate; only matching source rows appear in the view.
limit: int, optional
Cap the view at this many rows, in materialization order.
Returns
-------
MaterializedView
"""
raise NotImplementedError(
"materialized views are not supported on this connection type"
)
def open_materialized_view(self, name: str) -> MaterializedView:
"""Open the materialized view named ``name``.
Raises ``ValueError`` if the table exists but is not a materialized
view.
"""
raise NotImplementedError(
"materialized views are not supported on this connection type"
)
def list_materialized_views(self) -> List[str]:
"""The names of the materialized views in this database.
Found by reading every table's schema, so this costs an open per
table.
"""
raise NotImplementedError(
"materialized views are not supported on this connection type"
)
def drop_table(self, name: str, namespace_path: Optional[List[str]] = None):
"""Drop a table from the database.
@@ -1136,6 +1206,58 @@ class LanceDBConnection(DBConnection):
tbl.checkout(version)
return tbl
@override
def create_materialized_view(
self,
name: str,
source: str,
*,
select: SelectArg = None,
where: Optional[str] = None,
limit: Optional[int] = None,
) -> MaterializedView:
"""Define a materialized view named ``name`` over the table ``source``.
See
[DBConnection.create_materialized_view][lancedb.DBConnection.create_materialized_view].
Examples
--------
>>> import lancedb
>>> db = lancedb.connect(
... "./.lancedb",
... storage_options={"new_table_enable_stable_row_ids": "true"},
... )
>>> data = [{"name": "ada", "age": 36}, {"name": "kid", "age": 7}]
>>> table = db.create_table("people", data)
>>> view = db.create_materialized_view(
... "adults",
... "people",
... select=["name", ("shout", "upper(name)")],
... where="age >= 18",
... )
>>> result = view.refresh()
>>> result.rows_written
1
"""
LOOP.run(
self._conn.create_materialized_view(
name, source, select=select, where=where, limit=limit
)
)
return MaterializedView(self.open_table(name))
@override
def open_materialized_view(self, name: str) -> MaterializedView:
"""Open the materialized view named ``name``."""
view = MaterializedView(self.open_table(name))
view.definition
return view
@override
def list_materialized_views(self) -> List[str]:
"""The names of the materialized views in this database."""
return LOOP.run(self._conn.list_materialized_views())
def clone_table(
self,
target_table_name: str,
@@ -1906,6 +2028,50 @@ class AsyncConnection(object):
await tbl.checkout(version)
return tbl
async def create_materialized_view(
self,
name: str,
source: str,
*,
select: SelectArg = None,
where: Optional[str] = None,
limit: Optional[int] = None,
) -> AsyncMaterializedView:
"""Define a materialized view named ``name`` over the table ``source``.
See
[DBConnection.create_materialized_view][lancedb.DBConnection.create_materialized_view].
"""
inner = await self._inner.create_materialized_view(
name,
source,
projections=normalize_select(select),
filter=where,
limit=limit,
)
return AsyncMaterializedView(AsyncTable(inner))
async def open_materialized_view(self, name: str) -> AsyncMaterializedView:
"""Open the materialized view named ``name``.
Raises ``ValueError`` if the table exists but is not a materialized
view.
"""
if self.uri.startswith("db://"):
raise NotImplementedError(
"materialized views are supported only on local databases"
)
view = AsyncMaterializedView(await self.open_table(name))
await view.definition()
return view
async def list_materialized_views(self) -> List[str]:
"""The names of the materialized views in this database.
Found by reading every table's schema, so this costs an open per
table.
"""
return await self._inner.list_materialized_views()
async def clone_table(
self,
target_table_name: str,
+3 -2
View File
@@ -85,8 +85,9 @@ class Expr:
# for dict keys / set membership.
__hash__ = None # type: ignore[assignment]
def __init__(self, inner: PyExpr) -> None:
def __init__(self, inner: PyExpr, *, column_path: str | None = None) -> None:
self._inner = inner
self._column_path = column_path
# ── comparisons ──────────────────────────────────────────────────────────
@@ -273,7 +274,7 @@ def col(name: str) -> Expr:
>>> col("age") > lit(18)
Expr((age > 18))
"""
return Expr(expr_col(name))
return Expr(expr_col(name), column_path=name)
def lit(value: Union[bool, int, float, str, bytes, date, datetime, Decimal]) -> Expr:
+67 -1
View File
@@ -20,6 +20,7 @@ import re
import sys
import textwrap
import types
import uuid
from collections.abc import Mapping
from datetime import date, datetime
from typing import (
@@ -265,6 +266,64 @@ class FunctionVersion(_RemoteValue):
required_secrets: tuple[str, ...] = ()
created_at: str
def __call__(self, **inputs: Any) -> FunctionApplication:
"""Bind this exact version to named table columns.
Every input must be a direct [lancedb.col][lancedb.expr.col]
reference. The returned application is immutable and retains a
named-struct output as one sibling group, so every row's sibling values
come from one logical Function evaluation. Map result fields to table
columns with
[FunctionApplication.rename][lancedb.functions.FunctionApplication.rename],
then pass the application to
[Table.add_columns][lancedb.table.Table.add_columns].
Examples
--------
>>> from lancedb import col
>>> application = function( # doctest: +SKIP
... title=col("title"),
... body=col("body"),
... ).rename(columns={
... "normalized_text": "search_text",
... "token_count": "search_token_count",
... })
>>> table.add_columns(application) # doctest: +SKIP
"""
from lancedb.expr import Expr
parameters = tuple(parameter.name for parameter in self.signature.inputs)
missing = [parameter for parameter in parameters if parameter not in inputs]
unknown = sorted(set(inputs) - set(parameters))
if missing or unknown:
details = []
if missing:
details.append(f"missing inputs: {missing!r}")
if unknown:
details.append(f"unknown inputs: {unknown!r}")
raise TypeError("invalid Function inputs (" + "; ".join(details) + ")")
bindings = []
for parameter in parameters:
value = inputs[parameter]
if not isinstance(value, Expr) or value._column_path is None:
raise TypeError(
f"Function input {parameter!r} must be a direct col(...) reference"
)
bindings.append(
ApplicationInput(
parameter=parameter,
kind="column",
value={"path": value._column_path},
)
)
return FunctionApplication(
function=FunctionVersionRef(name=self.name, version=self.version),
inputs=tuple(bindings),
output=self.signature.output,
group_id=f"fg_{uuid.uuid4().hex}",
)
class FunctionRegistrationRequest(_RemoteValue):
"""Stable remote registration envelope produced by :func:`udf`.
@@ -304,7 +363,14 @@ class ApplicationInput(_OpenRemoteValue):
class FunctionApplication(_OpenRemoteValue):
"""Immutable pre-declaration application of an exact Function version."""
"""Immutable pre-declaration application of an exact Function version.
A named-struct output remains one grouped application through table
declaration and execution.
[FunctionApplication.rename][lancedb.functions.FunctionApplication.rename]
records the result-field to table-column mapping without splitting sibling
outputs into separate UDF calls.
"""
function: FunctionVersionRef
inputs: tuple[ApplicationInput, ...]
+11
View File
@@ -163,6 +163,15 @@ class FTS:
The number of documents per compressed posting block. Supported values
are 128 and 256. A value of 256 uses the experimental FTS V3 format
and may introduce breaking changes.
memory_limit : int, optional
The total memory limit in MiB for the local FTS build stage. The limit
is divided evenly among indexing workers. This build-only setting is
not persisted with the index and does not apply to remote tables.
num_workers : int, optional
The number of workers for a local FTS build. By default Lance uses
roughly half of the available CPU cores. The effective value is
limited by the available compute capacity. This build-only setting is
not persisted with the index and does not apply to remote tables.
Notes
-----
@@ -185,6 +194,8 @@ class FTS:
prefix_only: bool = False
block_size: int = 128
custom_stop_words: Optional[List[str]] = None
memory_limit: Optional[int] = None
num_workers: Optional[int] = None
@dataclass
+178
View File
@@ -0,0 +1,178 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
"""Materialized views: tables defined by a query over a source table and
maintained by refresh. See ``DBConnection.create_materialized_view``."""
from __future__ import annotations
import json
from dataclasses import dataclass, field
from typing import TYPE_CHECKING, Dict, List, Optional, Sequence, Tuple, Union
from .background_loop import LOOP
if TYPE_CHECKING:
import pyarrow as pa
from ._lancedb import RefreshMaterializedViewResult
from .table import AsyncTable, LanceTable
DEFINITION_META_KEY = b"mv.definition"
SelectArg = Union[
str,
Sequence[Union[str, Tuple[str, str]]],
Dict[str, str],
None,
]
@dataclass
class MaterializedViewDefinition:
"""The query that defines a materialized view."""
source_table: str
"""Name of the source table, in the same database as the view."""
projections: List[Tuple[str, str]]
"""``(output column, SQL expression)`` pairs, in view schema order."""
filter: Optional[str] = None
"""SQL predicate selecting the source rows the view holds."""
limit: Optional[int] = None
"""Cap on the number of rows the view holds."""
inputs: List[str] = field(default_factory=list)
"""Source columns the projections and filter read."""
def _definition_from_schema(
schema: "pa.Schema", name: str
) -> MaterializedViewDefinition:
metadata = schema.metadata or {}
raw = metadata.get(DEFINITION_META_KEY)
if raw is None:
raise ValueError(f"Table '{name}' is not a materialized view")
value = json.loads(raw)
kind = value.get("kind")
if kind != "select":
raise NotImplementedError(
f"materialized view '{name}' is defined by '{kind}', which this "
"version of lancedb cannot refresh"
)
return MaterializedViewDefinition(
source_table=value["source_table"],
projections=[
(p["output"], p["expression"]) for p in value.get("projections", [])
],
filter=value.get("filter"),
limit=value.get("limit"),
inputs=value.get("inputs", []),
)
def _quote_identifier(name: str) -> str:
"""Quote a column name as a Lance SQL identifier (backticks)."""
escaped = name.replace("`", "``")
return f"`{escaped}`"
def normalize_select(select: SelectArg) -> Optional[List[Tuple[str, str]]]:
"""``select`` items may be a column name, an ``(alias, expression)`` pair,
or a dict of the same. A bare name projects itself and is quoted, so any
valid column name works; dict and pair entries are kept verbatim because
their right side is an expression.
A lone string is one column, not a sequence of its characters."""
if select is None:
return None
if isinstance(select, str):
select = [select]
if isinstance(select, dict):
return list(select.items())
normalized = []
for item in select:
if isinstance(item, str):
normalized.append((item, _quote_identifier(item)))
else:
alias, expression = item
normalized.append((alias, expression))
return normalized
class AsyncMaterializedView:
"""A handle on a materialized view: its table plus its definition.
Obtained from ``AsyncConnection.create_materialized_view`` or
``AsyncConnection.open_materialized_view``.
"""
def __init__(self, table: "AsyncTable"):
self._table = table
def __repr__(self) -> str:
return f"AsyncMaterializedView(name={self.name!r})"
@property
def name(self) -> str:
return self._table.name
@property
def table(self) -> "AsyncTable":
"""The view, as the table it is. Queries, indexes and search all
apply; writes are not blocked, but a rebuild replaces them."""
return self._table
async def definition(self) -> MaterializedViewDefinition:
"""The query that defines the view, read from its stored schema."""
return _definition_from_schema(await self._table.schema(), self.name)
async def refresh(
self, *, full: bool = False, source_version: Optional[int] = None
) -> "RefreshMaterializedViewResult":
"""Recompute the view from its source.
The refresh is incremental when the source's changes can be
reconciled into the view -- rows added, changed or removed since the
last one -- and otherwise rebuilds. ``full=True`` forces a rebuild;
``source_version`` refreshes to that source version instead of the
latest.
Concurrent refreshes of one view do not duplicate its rows. Two that
plan the same source rows conflict on commit, and the loser raises
rather than writing them a second time.
"""
return await self._table._inner.refresh_materialized_view(
full=full, source_version=source_version
)
class MaterializedView:
"""Synchronous variant of
[AsyncMaterializedView][lancedb.materialized_view.AsyncMaterializedView]."""
def __init__(self, table: "LanceTable"):
self._table = table
self._async = AsyncMaterializedView(table._table)
def __repr__(self) -> str:
return f"MaterializedView(name={self.name!r})"
@property
def name(self) -> str:
return self._table.name
@property
def table(self) -> "LanceTable":
"""The view, as the table it is."""
return self._table
@property
def definition(self) -> MaterializedViewDefinition:
"""The query that defines the view, read from its stored schema."""
return _definition_from_schema(self._table.schema, self.name)
def refresh(
self, *, full: bool = False, source_version: Optional[int] = None
) -> "RefreshMaterializedViewResult":
"""Recompute the view from its source. See
[AsyncMaterializedView.refresh][lancedb.materialized_view.AsyncMaterializedView.refresh]."""
return LOOP.run(self._async.refresh(full=full, source_version=source_version))
+68
View File
@@ -61,6 +61,11 @@ from lance_namespace import (
NamespaceExistsRequest,
TableExistsRequest,
)
from lancedb.materialized_view import (
AsyncMaterializedView,
MaterializedView,
SelectArg,
)
from lancedb.table import AsyncTable, LanceTable, Table
from lancedb.util import validate_table_name
from lancedb.common import DATA
@@ -619,6 +624,42 @@ class LanceNamespaceDBConnection(DBConnection):
tbl.checkout(version)
return tbl
@override
def create_materialized_view(
self,
name: str,
source: str,
*,
select: "SelectArg" = None,
where: Optional[str] = None,
limit: Optional[int] = None,
) -> "MaterializedView":
"""Define a materialized view over a table in the root namespace.
See
[DBConnection.create_materialized_view][lancedb.DBConnection.create_materialized_view].
"""
return MaterializedView(
self.open_table(
LOOP.run(
self._inner.create_materialized_view(
name, source, select=select, where=where, limit=limit
)
).name
)
)
@override
def open_materialized_view(self, name: str) -> "MaterializedView":
"""Open the materialized view named ``name``."""
view = MaterializedView(self.open_table(name))
view.definition
return view
@override
def list_materialized_views(self) -> List[str]:
"""The names of the materialized views in the root namespace."""
return LOOP.run(self._inner.list_materialized_views())
@override
def drop_table(self, name: str, namespace_path: Optional[List[str]] = None):
if namespace_path is None:
@@ -1141,6 +1182,33 @@ class AsyncLanceNamespaceDBConnection:
route_pushdown_to_rust=self._route_pushdown_to_rust,
)
async def create_materialized_view(
self,
name: str,
source: str,
*,
select: "SelectArg" = None,
where: Optional[str] = None,
limit: Optional[int] = None,
) -> "AsyncMaterializedView":
"""Define a materialized view over a table in the root namespace."""
view = await self._inner.create_materialized_view(
name, source, select=select, where=where, limit=limit
)
# Reopen through the namespace so the view's table carries the
# namespace client and pushdown configuration a bare inner table lacks.
return AsyncMaterializedView(await self.open_table(view.name))
async def open_materialized_view(self, name: str) -> "AsyncMaterializedView":
"""Open the materialized view named ``name``."""
view = AsyncMaterializedView(await self.open_table(name))
await view.definition()
return view
async def list_materialized_views(self) -> List[str]:
"""The names of the materialized views in the root namespace."""
return await self._inner.list_materialized_views()
async def drop_table(self, name: str, namespace_path: Optional[List[str]] = None):
"""Drop a table from the namespace."""
if namespace_path is None:
+6 -12
View File
@@ -41,21 +41,15 @@ class PermutationBuilder:
The permutation is stored in memory and will be lost when the program exits.
"""
def __init__(self, table: LanceTable):
def __init__(self, table: Table):
"""
Creates a new permutation builder for the given table.
By default, the permutation builder will create a single split that contains all
rows in the same order as the base table.
Tables with an LSM write spec are rejected: unflushed rows have no row id.
"""
if not hasattr(table, "_inner"):
raise TypeError(
f"PermutationBuilder requires a local LanceTable, "
f"got {type(table).__name__}. "
"The permutation API is not supported on remote tables. "
"Remote tables connect to LanceDB Cloud or Enterprise and do not have "
"direct access to the underlying Lance dataset needed for permutations."
)
self._async = async_permutation_builder(table)
def split_random(
@@ -231,7 +225,7 @@ class PermutationBuilder:
return LOOP.run(do_execute())
def permutation_builder(table: LanceTable) -> PermutationBuilder:
def permutation_builder(table: Table) -> PermutationBuilder:
return PermutationBuilder(table)
@@ -248,7 +242,7 @@ class Permutations:
Attributes
----------
base_table: LanceTable
base_table: Table
The base table that the permutations are based on.
permutation_table: LanceTable
The permutation table that defines the splits.
@@ -282,7 +276,7 @@ class Permutations:
{'train': 0, 'test': 1}
"""
def __init__(self, base_table: LanceTable, permutation_table: LanceTable):
def __init__(self, base_table: Table, permutation_table: LanceTable):
self.base_table = base_table
self.permutation_table = permutation_table
+27
View File
@@ -25,6 +25,7 @@ from ..common import DATA
from ..db import DBConnection, LOOP
from ..functions import FunctionVersion, UdfDefinition
from ..job import AsyncJob, Job
from ..materialized_view import MaterializedView, SelectArg
if TYPE_CHECKING:
from .._lancedb import JobDescription, JobInfo
@@ -648,6 +649,32 @@ class RemoteDBConnection(DBConnection):
namespace_path=namespace_path,
)
@override
def create_materialized_view(
self,
name: str,
source: str,
*,
select: SelectArg = None,
where: Optional[str] = None,
limit: Optional[int] = None,
) -> MaterializedView:
raise NotImplementedError(
"materialized views are supported only on local databases"
)
@override
def open_materialized_view(self, name: str) -> MaterializedView:
raise NotImplementedError(
"materialized views are supported only on local databases"
)
@override
def list_materialized_views(self) -> List[str]:
raise NotImplementedError(
"materialized views are supported only on local databases"
)
@override
def drop_table(self, name: str, namespace_path: Optional[List[str]] = None):
"""Drop a table from the database.
+12 -10
View File
@@ -6801,21 +6801,21 @@ class Branches:
"""Diff a branch against main."""
return LOOP.run(self._table.branches.diff(from_branch))
def merge(self, from_branch: str, dry_run: bool = False) -> Dict[str, Any]:
"""Merge a branch into main, or dry-run.
def cherry_pick(self, from_branch: str, dry_run: bool = False) -> Dict[str, Any]:
"""Cherry-pick a branch onto main, or dry-run.
Parameters
----------
from_branch: str
Branch to merge from.
Branch to cherry-pick from.
dry_run: bool, default False
When True, only preview. When False, attempt the merge.
When True, only preview. When False, attempt the cherry-pick.
Notes
-----
A rejected merge returns ``status="rejected"`` instead of raising.
A failed cherry-pick returns ``status="failed"`` instead of raising.
"""
return LOOP.run(self._table.branches.merge(from_branch, dry_run))
return LOOP.run(self._table.branches.cherry_pick(from_branch, dry_run))
def _wrap(
self, async_table: "AsyncTable", version: Optional[int] = None
@@ -6951,9 +6951,11 @@ class AsyncBranches:
"""Diff a branch against main."""
return await self._table.branches.diff(from_branch)
async def merge(self, from_branch: str, dry_run: bool = False) -> Dict[str, Any]:
"""Merge a branch into main, or dry-run.
async def cherry_pick(
self, from_branch: str, dry_run: bool = False
) -> Dict[str, Any]:
"""Cherry-pick a branch onto main, or dry-run.
A rejected merge returns ``status="rejected"`` instead of raising.
A failed cherry-pick returns ``status="failed"`` instead of raising.
"""
return await self._table.branches.merge(from_branch, dry_run)
return await self._table.branches.cherry_pick(from_branch, dry_run)
@@ -37,6 +37,11 @@ from unittest.mock import patch
import lancedb
import pyarrow as pa
import pytest
from utils import (
MockPermutationServer,
assert_server_safe_row_id_requests,
mock_remote_table,
)
torch = pytest.importorskip("torch")
streaming = pytest.importorskip("lancedb.streaming")
@@ -2118,3 +2123,27 @@ def test_doc_example_checkpoint(lance_table):
assert sorted(consumed + remaining_original) == list(range(NUM_ROWS)), (
"Consumed + remaining must cover every row exactly once"
)
# ---------------------------------------------------------------------------
# Remote tables (LanceDB Cloud / Enterprise)
# ---------------------------------------------------------------------------
def test_streaming_dataset_over_remote_table():
"""StreamingDataset reads a remote table, with server-safe requests.
Builds a permutation over a remote table, then fetches batches from it by row id.
"""
server = MockPermutationServer()
with mock_remote_table(server) as table:
ds = StreamingDataset(table, num_splits=2, shuffle_seed=SHUFFLE_SEED)
ids = [row["id"] for row in ds]
assert sorted(ids) == list(range(server.num_rows)), (
"Every row of the remote table must be yielded exactly once"
)
assert len(server.scans) == 1, "the permutation is built with one row-id scan"
assert server.takes, "rows must be fetched with row-id takes"
assert_server_safe_row_id_requests(server)
@@ -6,6 +6,7 @@ from pathlib import Path
import pytest
from lancedb import col
import lancedb.functions as functions
from lancedb.functions import (
FunctionApplication,
@@ -120,6 +121,83 @@ def test_function_version_identity_is_immutable_and_exact():
assert FunctionVersion(**changed) != version
def test_function_version_binds_named_columns_as_one_immutable_group():
version = FunctionVersion.from_json(
json.dumps(job_result("remote_function_job.json"))
)
application = version(text=col("documents.body"))
assert application.function.name == version.name
assert application.function.version == version.version
assert application.output is version.signature.output
assert application.group_id.startswith("fg_")
assert [
(value.parameter, value.kind, value.value["path"])
for value in application.inputs
] == [("text", "column", "documents.body")]
with pytest.raises((TypeError, ValueError)):
application.group_id = "fg_changed"
def test_function_version_binding_validates_names_and_direct_columns():
version = FunctionVersion.from_json(
json.dumps(job_result("remote_function_job.json"))
)
with pytest.raises(TypeError, match=r"missing inputs: \['text'\]"):
version()
with pytest.raises(TypeError, match=r"unknown inputs: \['body'\]"):
version(text=col("text"), body=col("body"))
with pytest.raises(TypeError, match="direct col"):
version(text=col("text").lower())
def test_function_version_keeps_named_struct_outputs_in_one_application():
value = job_result("remote_function_job.json")
value["name"] = "text_features"
value["version"] = "fv_grouped"
value["signature"] = {
"inputs": [
{"name": "title", "arrow_type": "utf8", "nullable": True},
{"name": "body", "arrow_type": "utf8", "nullable": True},
],
"output": {
"kind": "named_struct",
"fields": [
{
"name": "normalized_text",
"arrow_type": "utf8",
"nullable": False,
},
{
"name": "token_count",
"arrow_type": "int64",
"nullable": False,
},
],
},
}
version = FunctionVersion(**value)
application = version(body=col("body"), title=col("title")).rename(
columns={
"normalized_text": "search_text",
"token_count": "search_token_count",
}
)
assert [value.parameter for value in application.inputs] == ["title", "body"]
assert [field.name for field in application.output.fields] == [
"normalized_text",
"token_count",
]
assert dict(application.columns) == {
"normalized_text": "search_text",
"token_count": "search_token_count",
}
def test_unknown_fields_and_discriminators_are_forward_decodable():
value = job_result("remote_function_job.json")
value["future_version_metadata"] = {"retention_class": "catalog"}
+8
View File
@@ -245,6 +245,14 @@ def test_create_inverted_index_rejects_invalid_block_size(table):
table.create_index("text", config=FTS(block_size=129))
def test_create_inverted_index_respects_build_memory_limit(table):
with pytest.raises(ValueError, match="exceeds worker memory limit"):
table.create_index(
"text",
config=FTS(memory_limit=0, num_workers=1),
)
def test_custom_stop_words_list(table):
table.create_index(
"text",
@@ -0,0 +1,268 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
import lancedb
import pytest
from lancedb.materialized_view import MaterializedViewDefinition
STABLE_ROW_IDS = {"new_table_enable_stable_row_ids": "true"}
def make_db(tmp_path):
db = lancedb.connect(tmp_path, storage_options=STABLE_ROW_IDS)
db.create_table(
"people",
[
{"name": "ada", "age": 36},
{"name": "kid", "age": 7},
{"name": "grace", "age": 85},
],
)
return db
def test_create_refresh_and_query(tmp_path):
db = make_db(tmp_path)
view = db.create_materialized_view(
"adults",
"people",
select=["name", ("shout", "upper(name)")],
where="age >= 18",
)
assert view.name == "adults"
assert view.table.count_rows() == 0
result = view.refresh()
assert result.mode == "rebuild"
assert result.rows_written == 2
rows = view.table.search().to_list()
assert sorted(row["shout"] for row in rows) == ["ADA", "GRACE"]
def test_definition_round_trips(tmp_path):
db = make_db(tmp_path)
db.create_materialized_view("adults", "people", where="age >= 18")
view = db.open_materialized_view("adults")
assert view.definition == MaterializedViewDefinition(
source_table="people",
projections=[("name", "`name`"), ("age", "`age`")],
filter="age >= 18",
inputs=["age", "name"],
)
def test_incremental_refresh_after_append(tmp_path):
db = make_db(tmp_path)
view = db.create_materialized_view("copy", "people")
view.refresh()
db.open_table("people").add([{"name": "alan", "age": 41}])
result = view.refresh()
assert result.mode == "incremental"
assert result.rows_written == 1
assert view.table.count_rows() == 4
assert view.refresh().mode == "no_op"
def test_incremental_refresh_after_update(tmp_path):
db = make_db(tmp_path)
view = db.create_materialized_view("copy", "people")
view.refresh()
db.open_table("people").update(where="name = 'kid'", values={"age": 8})
result = view.refresh()
assert result.mode == "incremental"
assert result.rows_written == 1
rows = view.table.search().to_list()
assert sorted(row["age"] for row in rows) == [8, 36, 85]
def test_legacy_storage_source_update_rebuilds(tmp_path):
db = lancedb.connect(
tmp_path,
storage_options={**STABLE_ROW_IDS, "new_table_data_storage_version": "legacy"},
)
db.create_table("people", [{"name": "ada", "age": 36}, {"name": "kid", "age": 7}])
view = db.create_materialized_view("copy", "people")
view.refresh()
db.open_table("people").update(where="name = 'kid'", values={"age": 8})
result = view.refresh()
assert result.mode == "rebuild"
rows = view.table.search().to_list()
assert sorted(row["age"] for row in rows) == [8, 36]
def test_list_and_not_a_view(tmp_path):
db = make_db(tmp_path)
db.create_materialized_view("adults", "people", where="age >= 18")
assert db.list_materialized_views() == ["adults"]
with pytest.raises(ValueError, match="not a materialized view"):
db.open_materialized_view("people")
def test_invalid_expression_fails_at_create(tmp_path):
db = make_db(tmp_path)
with pytest.raises(Exception, match="missing"):
db.create_materialized_view("bad", "people", select=[("x", "missing + 1")])
assert "bad" not in db.list_tables().tables
@pytest.mark.asyncio
async def test_async_create_refresh_and_open(tmp_path):
db = await lancedb.connect_async(tmp_path, storage_options=STABLE_ROW_IDS)
await db.create_table("people", [{"name": "ada", "age": 36}])
view = await db.create_materialized_view(
"shouts", "people", select=[("shout", "upper(name)")]
)
result = await view.refresh()
assert result.mode == "rebuild"
assert result.rows_written == 1
reopened = await db.open_materialized_view("shouts")
definition = await reopened.definition()
assert definition.projections == [("shout", "upper(name)")]
assert await db.list_materialized_views() == ["shouts"]
@pytest.mark.asyncio
async def test_async_incremental(tmp_path):
db = await lancedb.connect_async(tmp_path, storage_options=STABLE_ROW_IDS)
await db.create_table("people", [{"name": "ada", "age": 36}])
view = await db.create_materialized_view("copy", "people")
await view.refresh()
table = await db.open_table("people")
await table.add([{"name": "alan", "age": 41}])
result = await view.refresh()
assert result.mode == "incremental"
assert result.rows_written == 1
def test_source_requires_stable_row_ids(tmp_path):
db = lancedb.connect(tmp_path)
db.create_table("plain", [{"x": 1}])
with pytest.raises(Exception, match="stable row ids"):
db.create_materialized_view("v", "plain")
def test_bare_select_names_are_quoted(tmp_path):
db = lancedb.connect(tmp_path, storage_options=STABLE_ROW_IDS)
db.create_table("odd_names", [{"order item": "widget", "select": 2}])
view = db.create_materialized_view(
"quoted", "odd_names", select=["order item", "select"]
)
result = view.refresh()
assert result.rows_written == 1
rows = view.table.search().to_list()
assert rows[0]["order item"] == "widget"
assert rows[0]["select"] == 2
@pytest.mark.asyncio
async def test_async_remote_is_refused_without_network():
db = await lancedb.connect_async(
"db://nowhere", api_key="sk_test", region="us-east-1"
)
with pytest.raises(NotImplementedError, match="local"):
await db.create_materialized_view("v", "src")
with pytest.raises(NotImplementedError, match="local"):
await db.open_materialized_view("v")
with pytest.raises(NotImplementedError, match="local"):
await db.list_materialized_views()
def test_scalar_select_is_one_column(tmp_path):
db = make_db(tmp_path)
view = db.create_materialized_view("just_name", "people", select="name")
view.refresh()
rows = view.table.search().to_list()
assert set(rows[0]) - {"__source_row_id"} == {"name"}
assert sorted(row["name"] for row in rows) == ["ada", "grace", "kid"]
@pytest.mark.asyncio
async def test_async_scalar_select_is_one_column(tmp_path):
db = await lancedb.connect_async(tmp_path, storage_options=STABLE_ROW_IDS)
await db.create_table("people", [{"name": "ada", "age": 36}])
view = await db.create_materialized_view("just_name", "people", select="name")
await view.refresh()
rows = await view.table.query().to_list()
assert set(rows[0]) - {"__source_row_id"} == {"name"}
def test_limit_above_i64_max_is_refused(tmp_path):
db = make_db(tmp_path)
with pytest.raises(ValueError, match="exceeds the maximum"):
db.create_materialized_view("too_big", "people", limit=2**63)
# The boundary is fine, and zero still means an empty view.
db.create_materialized_view("at_max", "people", limit=2**63 - 1)
empty = db.create_materialized_view("none", "people", limit=0)
empty.refresh()
assert empty.table.count_rows() == 0
def _namespace_db(tmp_path):
return lancedb.connect_namespace(
"dir",
{"root": str(tmp_path)},
storage_options=STABLE_ROW_IDS,
)
def test_namespace_connection_materialized_views(tmp_path):
db = _namespace_db(tmp_path)
db.create_table(
"people",
[{"name": "ada", "age": 36}, {"name": "kid", "age": 7}],
storage_options=STABLE_ROW_IDS,
)
view = db.create_materialized_view("adults", "people", where="age >= 18")
view.refresh()
assert view.table.count_rows() == 1
assert db.list_materialized_views() == ["adults"]
reopened = db.open_materialized_view("adults")
assert reopened.definition.source_table == "people"
with pytest.raises(ValueError, match="not a materialized view"):
db.open_materialized_view("people")
@pytest.mark.asyncio
async def test_async_namespace_connection_materialized_views(tmp_path):
db = lancedb.connect_namespace_async(
"dir",
{"root": str(tmp_path)},
storage_options=STABLE_ROW_IDS,
)
await db.create_table(
"people",
[{"name": "ada", "age": 36}, {"name": "kid", "age": 7}],
storage_options=STABLE_ROW_IDS,
)
view = await db.create_materialized_view("adults", "people", where="age >= 18")
await view.refresh()
assert await view.table.count_rows() == 1
assert await db.list_materialized_views() == ["adults"]
reopened = await db.open_materialized_view("adults")
assert (await reopened.definition()).source_table == "people"
# The view's table came through the namespace, not straight from the
# inner connection: a bare inner table carries no namespace context, so
# its pushdown routing differs from a table the namespace opened.
through_namespace = await db.open_table("adults")
for handle in (view.table, reopened.table):
assert (
handle._route_pushdown_to_rust == through_namespace._route_pushdown_to_rust
)
assert handle._namespace_path == through_namespace._namespace_path
+59
View File
@@ -8,6 +8,11 @@ import pytest
from lancedb import DBConnection, Table, connect
from lancedb.background_loop import LOOP
from lancedb.permutation import Permutation, Permutations, permutation_builder
from utils import (
MockPermutationServer,
assert_server_safe_row_id_requests,
mock_remote_table,
)
def test_split_random_ratios(mem_db):
@@ -1214,3 +1219,57 @@ def test_remove_rowid_after_select(some_permutation: Permutation):
perm_without_rowid = perm_with_rowid.remove_columns(["_rowid"])
assert "_rowid" not in perm_without_rowid.column_names
assert perm_without_rowid.column_names == ["id"]
def test_permutation_is_stable_when_remote_scan_order_varies():
"""Splits are assigned by scan position, and every rank builds its own
permutation, so two ranks seeing different scan orders must still agree."""
server = MockPermutationServer(num_rows=16, vary_scan_order=True)
def split_of_each_row(permutation_tbl):
# Sequential splits are assigned by position, so a reversed scan would put
# the last rows in split 0. Compare the mapping rather than the table order,
# which the split-id sort does not pin down.
rows = permutation_tbl.search(None).to_arrow().to_pydict()
return dict(zip(rows["row_id"], rows["split_id"]))
with mock_remote_table(server) as table:
first = split_of_each_row(
permutation_builder(table).split_sequential(fixed=2).execute()
)
second = split_of_each_row(
permutation_builder(table).split_sequential(fixed=2).execute()
)
assert server.scan_calls == 2, "both builds must have scanned"
assert first == second
assert first[0] == 0 and first[server.num_rows - 1] == 1, first
def test_permutation_over_remote_table():
"""The permutation API accepts a remote table, addressing rows by `_rowid` just
as `take_row_ids` does. Also pins the request shapes sent to the server.
"""
server = MockPermutationServer()
with mock_remote_table(server) as table:
permutation_tbl = permutation_builder(table).split_sequential(fixed=2).execute()
assert permutation_tbl.count_rows() == server.num_rows
permutation = Permutation.from_tables(table, permutation_tbl, 0)
assert permutation.num_rows == server.num_rows // 2
# Compare against the permutation's own order; the split-id sort is not stable.
rows = permutation_tbl.search(None).to_arrow().to_pydict()
split0 = [
row_id
for row_id, split in zip(rows["row_id"], rows["split_id"])
if not split
]
# The mock table's `id` equals its `_rowid`.
assert permutation.take_offsets([2, 0]) == [
{"id": split0[2]},
{"id": split0[0]},
]
assert_server_safe_row_id_requests(server)
+8 -9
View File
@@ -242,8 +242,8 @@ def test_remote_table_branches_sync():
table.branches.delete("exp")
def test_remote_table_branch_merge_defaults_to_execute():
merge_bodies = []
def test_remote_table_cherry_pick_defaults_to_execute():
cherry_pick_bodies = []
diff = {
"fromBranch": "exp",
"parentVersion": 1,
@@ -265,8 +265,7 @@ def test_remote_table_branch_merge_defaults_to_execute():
"changedColumns": [],
"addedIndexes": [],
"removedIndexes": [],
"mergeable": True,
"mergeBlockers": [],
"errors": [],
}
def handler(request):
@@ -276,11 +275,11 @@ def test_remote_table_branch_merge_defaults_to_execute():
else:
content_len = int(request.headers.get("Content-Length"))
request_body = json.loads(request.rfile.read(content_len))
merge_bodies.append(request_body)
cherry_pick_bodies.append(request_body)
dry_run = request_body["dry_run"]
status = 200 if dry_run else 409
body = {
"status": "ready" if dry_run else "rejected",
"status": "ready" if dry_run else "failed",
"diff": diff,
"preview": {"promotedColumns": []},
}
@@ -292,10 +291,10 @@ def test_remote_table_branch_merge_defaults_to_execute():
with mock_lancedb_connection(handler) as db:
branches = db.open_table("test").branches
assert branches.merge("exp")["status"] == "rejected"
assert branches.merge("exp", dry_run=True)["status"] == "ready"
assert branches.cherry_pick("exp")["status"] == "failed"
assert branches.cherry_pick("exp", dry_run=True)["status"] == "ready"
assert merge_bodies == [
assert cherry_pick_bodies == [
{"from_branch": "exp", "dry_run": False},
{"from_branch": "exp", "dry_run": True},
]
+206
View File
@@ -1,7 +1,17 @@
# SPDX-License-Identifier: Apache-2.0
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
import contextlib
import http.server
import json
import re
import threading
import lancedb
import pyarrow as pa
import pytest
ARROW_FILE_CONTENT_TYPE = "application/vnd.apache.arrow.file"
def exception_output(e_info: pytest.ExceptionInfo):
import traceback
@@ -9,3 +19,199 @@ def exception_output(e_info: pytest.ExceptionInfo):
# skip traceback part, since it's not worth checking in tests
lines = traceback.format_exception_only(e_info.type, e_info.value)
return "".join(lines).strip()
def parse_in_list(filter_sql: str) -> list[int]:
"""Pull the integers out of a `<col> IN (a, b, c)` predicate.
Scoped to the parenthesised list so a cast in the SQL adds no phantom values.
"""
match = re.search(r"\bIN\s*\(([^)]*)\)", filter_sql, re.IGNORECASE)
assert match is not None, f"expected an IN list, got: {filter_sql}"
return [int(m) for m in re.findall(r"-?\d+", match.group(1))]
def is_row_id_take(body) -> bool:
"""True when a query body fetches specific rows by row id."""
return "_rowid" in (body.get("filter") or "")
def arrow_file_bytes(table: pa.Table) -> bytes:
"""Serialize to the Arrow IPC *file* framing the /query/ route answers with."""
sink = pa.BufferOutputStream()
with pa.ipc.new_file(sink, table.schema) as writer:
writer.write_table(table)
return sink.getvalue().to_pybytes()
class MockPermutationServer:
"""A stand-in LanceDB server hosting one table whose ``id`` equals its ``_rowid``.
Records every ``/query/`` body so tests can assert on the request shapes sent to
the server, which is the part that has to stay compatible.
"""
def __init__(self, name="remote_data", num_rows=8, vary_scan_order=False):
self.name = name
self.num_rows = num_rows
self.query_bodies = []
# Stand in for a distributed scan that answers in no fixed order.
self.vary_scan_order = vary_scan_order
self.scan_calls = 0
def __call__(self, request):
path = request.path
if path == f"/v1/table/{self.name}/describe/":
return self._json(
request,
{
"version": 1,
"schema": {
"fields": [
{"name": "id", "type": {"type": "int64"}, "nullable": False}
]
},
},
)
if path == f"/v1/table/{self.name}/get_lsm_write_spec/":
self._read_body(request)
# Null spec: this table has no LSM write path.
return self._json(request, {"lsm_write_spec": None})
if path == f"/v1/table/{self.name}/count_rows/":
self._read_body(request)
return self._json(request, self.num_rows)
if path == f"/v1/table/{self.name}/query/":
return self._query(request, self._read_body(request))
# Drain first, so an unexpected route cannot desync a keep-alive connection.
self._read_body(request)
request.send_response(404)
request.end_headers()
@property
def scans(self):
"""Bodies of the permutation build scan: the row id column, nothing else."""
return [b for b in self.query_bodies if b.get("columns") == ["_rowid"]]
@property
def takes(self):
"""Bodies of the row-id takes the loader fetches batches with.
Keyed on `_rowid`, not "has a filter": the schema probe also has a predicate.
"""
return [b for b in self.query_bodies if is_row_id_take(b)]
@staticmethod
def _read_body(request):
content_len = int(request.headers.get("Content-Length") or 0)
return json.loads(request.rfile.read(content_len)) if content_len else {}
@staticmethod
def _json(request, payload):
request.send_response(200)
request.send_header("Content-Type", "application/json")
request.end_headers()
request.wfile.write(json.dumps(payload).encode())
@staticmethod
def _arrow(request, table):
body = arrow_file_bytes(table)
request.send_response(200)
request.send_header("Content-Type", ARROW_FILE_CONTENT_TYPE)
request.send_header("Content-Length", str(len(body)))
request.end_headers()
request.wfile.write(body)
def _query(self, request, body):
self.query_bodies.append(body)
if is_row_id_take(body):
# A row-id take. Answer ascending, so tests prove the client reorders.
row_ids = sorted(parse_in_list(body["filter"]))
return self._arrow(
request,
pa.table(
{
"id": pa.array(row_ids, pa.int64()),
"_rowid": pa.array(row_ids, pa.uint64()),
}
),
)
if body.get("columns") == ["_rowid"]:
# The permutation build scan: row ids and nothing else.
row_ids = list(range(self.num_rows))
if self.vary_scan_order and self.scan_calls % 2:
row_ids.reverse()
self.scan_calls += 1
return self._arrow(
request,
pa.table({"_rowid": pa.array(row_ids, pa.uint64())}),
)
# The schema probe: filtered to nothing, so it carries schema and no rows.
return self._arrow(request, pa.table({"id": pa.array([], pa.int64())}))
def _make_handler(serve):
class MockLanceDBHandler(http.server.BaseHTTPRequestHandler):
def do_GET(self):
serve(self)
def do_POST(self):
serve(self)
def log_message(self, *args):
pass # keep pytest output readable
return MockLanceDBHandler
@contextlib.contextmanager
def mock_remote_table(server):
"""Run ``server`` on a local port and yield an open remote table against it.
Threading: the loader fans out fetch threads a single-threaded server would
serialize, hiding the prefetch overlap under test.
"""
with http.server.ThreadingHTTPServer(
("localhost", 0), _make_handler(server)
) as srv:
thread = threading.Thread(target=srv.serve_forever)
thread.start()
try:
db = lancedb.connect(
"db://dev",
api_key="fake",
host_override=f"http://localhost:{srv.server_address[1]}",
client_config={"timeout_config": {"connect_timeout": 5}},
)
yield db.open_table(server.name)
finally:
srv.shutdown()
thread.join()
def assert_server_safe_row_id_requests(server):
"""Assert the loader fetched rows by row id and bounded everything else.
`.get`, not `[...]`, so a dropped field reads as the assertion, not a KeyError.
"""
for body in server.takes:
# The fetch needs the row id back to restore the requested order.
assert body.get("with_row_id") is True, body
assert "_rowid" in body["filter"], body
# Only the one-off permutation scan may scan the whole table; the schema probe is
# built once per split per epoch. `k == 0` counts as unbounded: lance reads a zero
# limit as "no limit".
def is_unbounded(body):
if is_row_id_take(body):
return False
k = body.get("k")
return k is None or k == 0 or k > server.num_rows
unbounded = [b for b in server.query_bodies if is_unbounded(b)]
assert unbounded == server.scans, (
f"only the permutation scan may be unbounded, got {unbounded}"
)
+34
View File
@@ -333,6 +333,40 @@ impl Connection {
})
}
#[pyo3(signature = (name, source, projections=None, filter=None, limit=None))]
pub fn create_materialized_view(
self_: PyRef<'_, Self>,
name: String,
source: String,
projections: Option<Vec<(String, String)>>,
filter: Option<String>,
limit: Option<u64>,
) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.get_inner()?.clone();
future_into_py(self_.py(), async move {
let mut builder = inner.create_materialized_view(name, source);
if let Some(projections) = projections {
builder = builder.select(projections);
}
if let Some(filter) = filter {
builder = builder.only_if(filter);
}
if let Some(limit) = limit {
builder = builder.limit(limit);
}
let view = builder.execute().await.infer_error()?;
Ok(Table::new(view.table().clone()))
})
}
pub fn list_materialized_views(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.get_inner()?.clone();
future_into_py(self_.py(), async move {
let views = inner.list_materialized_views().await.infer_error()?;
Ok(views.into_iter().map(|view| view.name).collect::<Vec<_>>())
})
}
#[pyo3(signature = (name, namespace_path=None))]
pub fn drop_table(
self_: PyRef<'_, Self>,
+57 -1
View File
@@ -42,7 +42,7 @@ pub fn extract_index_params(source: &Option<Bound<'_, PyAny>>) -> PyResult<Lance
"Fm" => Ok(LanceDbIndex::Fm(FmIndexBuilder::default())),
"FTS" => {
let params = source.extract::<FtsParams>()?;
let inner_opts = FtsIndexBuilder::default()
let mut inner_opts = FtsIndexBuilder::default()
.base_tokenizer(params.base_tokenizer)
.language(&params.language)
.map_err(|_| {
@@ -61,6 +61,12 @@ pub fn extract_index_params(source: &Option<Bound<'_, PyAny>>) -> PyResult<Lance
.ngram_max_length(params.ngram_max_length)
.ngram_prefix_only(params.prefix_only)
.custom_stop_words(params.custom_stop_words);
if let Some(memory_limit) = params.memory_limit {
inner_opts = inner_opts.memory_limit_mb(memory_limit);
}
if let Some(num_workers) = params.num_workers {
inner_opts = inner_opts.num_workers(num_workers);
}
let inner_opts = inner_opts
.block_size(params.block_size)
.map_err(|err| PyValueError::new_err(err.to_string()))?;
@@ -213,6 +219,8 @@ struct FtsParams {
ngram_max_length: u32,
prefix_only: bool,
block_size: usize,
memory_limit: Option<u64>,
num_workers: Option<usize>,
}
#[derive(FromPyObject)]
@@ -444,3 +452,51 @@ impl IndexConfig {
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use pyo3::types::{PyDict, PyDictMethods};
use serde_json::json;
#[test]
fn fts_build_controls_are_forwarded() {
Python::initialize();
Python::attach(|py| {
let locals = PyDict::new(py);
py.run(
c"class FTS:
with_position = True
base_tokenizer = 'simple'
language = 'English'
max_token_length = None
lower_case = True
stem = False
remove_stop_words = False
custom_stop_words = None
ascii_folding = False
ngram_min_length = 3
ngram_max_length = 3
prefix_only = False
block_size = 128
memory_limit = 2048
num_workers = 7
config = FTS()",
None,
Some(&locals),
)
.unwrap();
let config = locals.get_item("config").unwrap().unwrap();
let index = extract_index_params(&Some(config)).unwrap();
let LanceDbIndex::FTS(params) = index else {
panic!("expected FTS index parameters");
};
let training_json = params.to_training_json().unwrap();
assert_eq!(training_json.get("memory_limit"), Some(&json!(2048)));
assert_eq!(training_json.get("num_workers"), Some(&json!(7)));
});
}
}
+3 -2
View File
@@ -16,8 +16,8 @@ use query::{FTSQuery, HybridQuery, Query, VectorQuery};
use session::Session;
use table::{
AddColumnsResult, AddResult, AlterColumnsResult, DeleteResult, DropColumnsResult, FtsToken,
LsmWriteSpec, MergeResult, PyBlobFile, RefreshColumnResult, Table, UpdateFieldMetadataResult,
UpdateResult,
LsmWriteSpec, MergeResult, PyBlobFile, RefreshColumnResult, RefreshMaterializedViewResult,
Table, UpdateFieldMetadataResult, UpdateResult,
};
pub mod arrow;
@@ -60,6 +60,7 @@ pub fn _lancedb(_py: Python, m: &Bound<'_, PyModule>) -> PyResult<()> {
m.add_class::<RecordBatchStream>()?;
m.add_class::<AddColumnsResult>()?;
m.add_class::<RefreshColumnResult>()?;
m.add_class::<RefreshMaterializedViewResult>()?;
m.add_class::<AlterColumnsResult>()?;
m.add_class::<UpdateFieldMetadataResult>()?;
m.add_class::<AddResult>()?;
+3 -1
View File
@@ -268,7 +268,9 @@ impl PyPermutationReader {
.await
.infer_error()?
} else {
PermutationReader::identity(base_table).await
PermutationReader::identity(base_table)
.await
.infer_error()?
};
Ok(Self::from_reader(reader))
})
+57 -2
View File
@@ -441,6 +441,41 @@ impl From<lancedb::table::RefreshColumnResult> for RefreshColumnResult {
}
}
#[pyclass(get_all, from_py_object)]
#[derive(Clone, Debug)]
pub struct RefreshMaterializedViewResult {
pub mode: String,
pub rows_written: u64,
pub source_version: u64,
pub version: u64,
}
#[pymethods]
impl RefreshMaterializedViewResult {
pub fn __repr__(&self) -> String {
format!(
"RefreshMaterializedViewResult(mode={}, rows_written={}, source_version={}, version={})",
self.mode, self.rows_written, self.source_version, self.version
)
}
}
impl From<lancedb::RefreshMaterializedViewResult> for RefreshMaterializedViewResult {
fn from(result: lancedb::RefreshMaterializedViewResult) -> Self {
let mode = match result.mode {
lancedb::RefreshMode::Rebuild => "rebuild",
lancedb::RefreshMode::Incremental => "incremental",
lancedb::RefreshMode::NoOp => "no_op",
};
Self {
mode: mode.to_string(),
rows_written: result.rows_written,
source_version: result.source_version,
version: result.version,
}
}
}
#[pymethods]
impl AddColumnsResult {
pub fn __repr__(&self) -> String {
@@ -1588,6 +1623,26 @@ impl Table {
})
}
#[pyo3(signature = (full=false, source_version=None))]
pub fn refresh_materialized_view(
self_: PyRef<'_, Self>,
full: bool,
source_version: Option<u64>,
) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let view = lancedb::MaterializedView::from_table(inner)
.await
.infer_error()?;
let mut builder = view.refresh().full(full);
if let Some(version) = source_version {
builder = builder.source_version(version);
}
let result = builder.execute().await.infer_error()?;
Ok(RefreshMaterializedViewResult::from(result))
})
}
pub fn add_columns_with_schema(
self_: PyRef<'_, Self>,
schema: PyArrowType<Schema>,
@@ -1885,7 +1940,7 @@ impl Branches {
}
#[pyo3(signature = (from_branch, dry_run=false))]
pub fn merge(
pub fn cherry_pick(
self_: PyRef<'_, Self>,
from_branch: String,
dry_run: bool,
@@ -1893,7 +1948,7 @@ impl Branches {
let inner = self_.inner.clone();
future_into_py(self_.py(), async move {
let result = inner
.merge_branch(&from_branch, dry_run)
.cherry_pick(&from_branch, dry_run)
.await
.infer_error()?;
Python::attach(|py| struct_to_wire_py(py, &result))
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb"
version = "0.38.0-beta.3"
version = "0.38.0-beta.4"
edition.workspace = true
description = "LanceDB: A serverless, low-latency vector database for AI applications"
license.workspace = true
+1 -1
View File
@@ -41,7 +41,7 @@ use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
mod create_table;
fn merge_storage_options(
pub(crate) fn merge_storage_options(
store_params: &mut ObjectStoreParams,
pairs: impl IntoIterator<Item = (String, String)>,
) {
+255 -62
View File
@@ -765,60 +765,13 @@ impl ListingDatabase {
}
}
/// Extract storage option overrides from the request
fn extract_storage_overrides(
&self,
request: &CreateTableRequest,
) -> Result<(Option<LanceFileVersion>, Option<bool>, Option<bool>)> {
let storage_options = request
.write_options
.lance_write_params
.as_ref()
.and_then(|p| p.store_params.as_ref())
.and_then(|sp| sp.storage_options());
let storage_version_override = storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_STORAGE_VERSION))
.map(|s| s.parse::<LanceFileVersion>())
.transpose()?;
let v2_manifest_override = storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_V2_MANIFEST_PATHS))
.map(|s| s.parse::<bool>())
.transpose()
.map_err(|_| Error::InvalidInput {
message: "enable_v2_manifest_paths must be a boolean".to_string(),
})?;
let stable_row_ids_override = storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS))
.map(|s| s.parse::<bool>())
.transpose()
.map_err(|_| Error::InvalidInput {
message: "enable_stable_row_ids must be a boolean".to_string(),
})?;
Ok((
storage_version_override,
v2_manifest_override,
stable_row_ids_override,
))
}
/// Prepare write parameters for table creation
fn prepare_write_params(
&self,
request: &CreateTableRequest,
storage_version_override: Option<LanceFileVersion>,
v2_manifest_override: Option<bool>,
stable_row_ids_override: Option<bool>,
mut write_params: lance::dataset::WriteParams,
overrides: NewTableConfig,
) -> lance::dataset::WriteParams {
let mut write_params = request
.write_options
.lance_write_params
.clone()
.unwrap_or_default();
// Only modify the storage options if we actually have something to
// inherit. There is a difference between storage_options=None and
// storage_options=Some({}). Using storage_options=None will cause the
@@ -842,18 +795,21 @@ impl ListingDatabase {
store_params.storage_options_accessor = Some(Arc::new(accessor));
}
write_params.data_storage_version = storage_version_override
write_params.data_storage_version = overrides
.data_storage_version
.or(write_params.data_storage_version)
.or(self.new_table_config.data_storage_version);
if let Some(enable_v2_manifest_paths) =
v2_manifest_override.or(self.new_table_config.enable_v2_manifest_paths)
if let Some(enable_v2_manifest_paths) = overrides
.enable_v2_manifest_paths
.or(self.new_table_config.enable_v2_manifest_paths)
{
write_params.enable_v2_manifest_paths = enable_v2_manifest_paths;
}
let data_schema = request.data.arrow_schema();
if let Some(enable_stable_row_ids) = stable_row_ids_override
if let Some(enable_stable_row_ids) = overrides
.enable_stable_row_ids
.or(self.new_table_config.enable_stable_row_ids)
.or(has_blob_columns(&data_schema).then_some(true))
{
@@ -1048,15 +1004,13 @@ impl Database for ListingDatabase {
.clone()
.unwrap_or_else(|| self.table_uri(&request.name).unwrap());
let (storage_version_override, v2_manifest_override, stable_row_ids_override) =
self.extract_storage_overrides(&request)?;
let write_params = self.prepare_write_params(
&request,
storage_version_override,
v2_manifest_override,
stable_row_ids_override,
);
let mut write_params = request
.write_options
.lance_write_params
.clone()
.unwrap_or_default();
let overrides = take_request_creation_overrides(&mut write_params)?;
let write_params = self.prepare_write_params(&request, write_params, overrides);
let data_schema = request.data.arrow_schema();
@@ -1288,8 +1242,232 @@ impl Database for ListingDatabase {
}
}
/// Parse the request-level `new_table_*` creation keys into overrides and
/// strip them from the store options in one step: every create path that
/// honors them must also keep them out of the object store.
pub(crate) fn take_request_creation_overrides(
params: &mut lance::dataset::WriteParams,
) -> Result<NewTableConfig> {
let storage_options = params
.store_params
.as_ref()
.and_then(|sp| sp.storage_options());
let overrides = NewTableConfig {
data_storage_version: storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_STORAGE_VERSION))
.map(|s| s.parse::<LanceFileVersion>())
.transpose()?,
enable_v2_manifest_paths: storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_V2_MANIFEST_PATHS))
.map(|s| s.parse::<bool>())
.transpose()
.map_err(|_| Error::InvalidInput {
message: "enable_v2_manifest_paths must be a boolean".to_string(),
})?,
enable_stable_row_ids: storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS))
.map(|s| s.parse::<bool>())
.transpose()
.map_err(|_| Error::InvalidInput {
message: "enable_stable_row_ids must be a boolean".to_string(),
})?,
};
if let Some(store_params) = params.store_params.as_mut() {
strip_new_table_creation_keys(store_params);
}
Ok(overrides)
}
/// Strip the `new_table_*` creation keys from request store options: they are
/// creation config, not credentials, and left in place they fork a fresh
/// store connection for the request.
fn strip_new_table_creation_keys(store_params: &mut ObjectStoreParams) {
let mut options = store_params.storage_options().cloned().unwrap_or_default();
let mut removed = false;
for key in [
OPT_NEW_TABLE_STORAGE_VERSION,
OPT_NEW_TABLE_V2_MANIFEST_PATHS,
OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
] {
removed |= options.remove(key).is_some();
}
if !removed {
return;
}
let provider = store_params
.storage_options_accessor
.as_ref()
.and_then(|accessor| accessor.provider().cloned());
store_params.storage_options_accessor = match (options.is_empty(), provider) {
(true, None) => None,
(true, Some(provider)) => Some(Arc::new(StorageOptionsAccessor::with_provider(provider))),
(false, Some(provider)) => Some(Arc::new(
StorageOptionsAccessor::with_initial_and_provider(options, provider),
)),
(false, None) => Some(Arc::new(StorageOptionsAccessor::with_static_options(
options,
))),
};
}
#[cfg(test)]
mod tests {
#[tokio::test]
async fn request_level_creation_keys_do_not_fork_the_store() {
use crate::query::ExecutableQuery;
use futures::TryStreamExt;
let db = crate::connect("memory://").execute().await.unwrap();
let batch = arrow_array::record_batch!(("x", Int32, [1, 2])).unwrap();
let store_params = ObjectStoreParams {
storage_options_accessor: Some(Arc::new(StorageOptionsAccessor::with_static_options(
HashMap::from([(
OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS.to_string(),
"true".to_string(),
)]),
))),
..Default::default()
};
db.create_table("t", batch)
.write_options(crate::table::WriteOptions {
lance_write_params: Some(lance::dataset::WriteParams {
store_params: Some(store_params),
..Default::default()
}),
})
.execute()
.await
.unwrap();
let table = db.open_table("t").execute().await.unwrap();
let rows: usize = table
.query()
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap()
.iter()
.map(|b| b.num_rows())
.sum();
assert_eq!(rows, 2, "the table must live in the session's store");
}
mod strip_new_table_creation_keys {
use super::super::*;
#[derive(Debug)]
struct EmptyProvider;
#[async_trait::async_trait]
impl StorageOptionsProvider for EmptyProvider {
async fn fetch_storage_options(
&self,
) -> lance_core::Result<Option<HashMap<String, String>>> {
Ok(Some(HashMap::new()))
}
fn provider_id(&self) -> String {
"empty-test-provider".into()
}
}
fn params_with_static(options: &[(&str, &str)]) -> ObjectStoreParams {
ObjectStoreParams {
storage_options_accessor: Some(Arc::new(
StorageOptionsAccessor::with_static_options(
options
.iter()
.map(|(k, v)| (k.to_string(), v.to_string()))
.collect(),
),
)),
..Default::default()
}
}
#[test]
fn creation_keys_are_removed_and_store_keys_kept() {
let mut params = params_with_static(&[
("region", "us-west-2"),
(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS, "true"),
]);
strip_new_table_creation_keys(&mut params);
let options = params.storage_options().cloned().unwrap();
assert_eq!(options.get("region").map(String::as_str), Some("us-west-2"));
assert!(!options.contains_key(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS));
// Creation keys alone: no accessor survives to fork a store.
let mut params = params_with_static(&[(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS, "true")]);
strip_new_table_creation_keys(&mut params);
assert!(params.storage_options_accessor.is_none());
}
/// A provider must survive every shape of strip: untouched accessors
/// keep their identity, emptied ones still fetch, and residual
/// statics ride along.
#[test]
fn provider_accessors_survive_the_strip() {
let accessor = Arc::new(StorageOptionsAccessor::with_provider(Arc::new(
EmptyProvider,
)));
let mut params = ObjectStoreParams {
storage_options_accessor: Some(accessor.clone()),
..Default::default()
};
strip_new_table_creation_keys(&mut params);
assert!(Arc::ptr_eq(
params.storage_options_accessor.as_ref().unwrap(),
&accessor
));
let mut params = ObjectStoreParams {
storage_options_accessor: Some(Arc::new(
StorageOptionsAccessor::with_initial_and_provider(
HashMap::from([
("region".to_string(), "us-west-2".to_string()),
(
OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS.to_string(),
"true".to_string(),
),
]),
Arc::new(EmptyProvider),
),
)),
..Default::default()
};
strip_new_table_creation_keys(&mut params);
let accessor = params.storage_options_accessor.unwrap();
assert!(accessor.has_provider());
assert_eq!(
accessor
.initial_storage_options()
.and_then(|o| o.get("region").cloned())
.as_deref(),
Some("us-west-2")
);
// Emptied entirely: a first-fetch accessor, not one caching {}.
let mut params = ObjectStoreParams {
storage_options_accessor: Some(Arc::new(
StorageOptionsAccessor::with_initial_and_provider(
HashMap::from([(
OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS.to_string(),
"true".to_string(),
)]),
Arc::new(EmptyProvider),
),
)),
..Default::default()
};
strip_new_table_creation_keys(&mut params);
let accessor = params.storage_options_accessor.unwrap();
assert!(accessor.has_provider());
assert!(accessor.initial_storage_options().is_none());
}
}
use super::*;
use crate::Table;
use crate::arrow::{SendableRecordBatchStream, SimpleRecordBatchStream};
@@ -2569,6 +2747,21 @@ mod tests {
}
}
/// Regression test for https://github.com/lancedb/lancedb/issues/2283.
///
/// Object-store URIs must use `/` on every platform. In particular, joining
/// with `std::path::Path` used to insert a `\\` into Azure blob keys on
/// Windows.
#[tokio::test]
async fn test_table_uri_uses_forward_slashes_for_azure() {
let (_tempdir, mut db) = setup_database().await;
db.uri = "az://test/db/test".to_string();
let uri = db.table_uri("test").unwrap();
assert_eq!(uri, "az://test/db/test/test.lance");
}
/// Regression: connecting via a URL-style URI (which goes through
/// `url::Url::parse` and the `query_pairs_mut()` path) must not
/// append a trailing `?` to per-table URIs when the input URI has
+149 -53
View File
@@ -26,10 +26,7 @@ use lance_table::io::commit::external_manifest::ExternalManifestCommitHandler;
use crate::blob::{ensure_blob_storage_version, has_blob_columns};
use crate::connection::NamespaceClientPushdownOperation;
use crate::database::ReadConsistency;
use crate::database::listing::{
NewTableConfig, OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS, OPT_NEW_TABLE_STORAGE_VERSION,
OPT_NEW_TABLE_V2_MANIFEST_PATHS,
};
use crate::database::listing::{NewTableConfig, take_request_creation_overrides};
use crate::database::read_freshness::{
FreshnessBaselines, ReadFreshnessContextProvider, TableFreshness,
};
@@ -197,69 +194,28 @@ impl LanceNamespaceDatabase {
TableFreshness::new(self.freshness_baselines.clone(), key)
}
fn extract_storage_overrides(
&self,
request: &DbCreateTableRequest,
) -> Result<(
Option<lance_file::version::LanceFileVersion>,
Option<bool>,
Option<bool>,
)> {
let storage_options = request
.write_options
.lance_write_params
.as_ref()
.and_then(|p| p.store_params.as_ref())
.and_then(|sp| sp.storage_options());
let storage_version_override = storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_STORAGE_VERSION))
.map(|s| s.parse::<lance_file::version::LanceFileVersion>())
.transpose()?;
let v2_manifest_override = storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_V2_MANIFEST_PATHS))
.map(|s| s.parse::<bool>())
.transpose()
.map_err(|_| Error::InvalidInput {
message: "enable_v2_manifest_paths must be a boolean".to_string(),
})?;
let stable_row_ids_override = storage_options
.and_then(|opts| opts.get(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS))
.map(|s| s.parse::<bool>())
.transpose()
.map_err(|_| Error::InvalidInput {
message: "enable_stable_row_ids must be a boolean".to_string(),
})?;
Ok((
storage_version_override,
v2_manifest_override,
stable_row_ids_override,
))
}
fn apply_new_table_config(
&self,
params: &mut lance::dataset::WriteParams,
request: &DbCreateTableRequest,
) -> Result<()> {
let (storage_version_override, v2_manifest_override, stable_row_ids_override) =
self.extract_storage_overrides(request)?;
let overrides = take_request_creation_overrides(params)?;
params.data_storage_version = storage_version_override
params.data_storage_version = overrides
.data_storage_version
.or(params.data_storage_version)
.or(self.new_table_config.data_storage_version);
if let Some(enable_v2_manifest_paths) =
v2_manifest_override.or(self.new_table_config.enable_v2_manifest_paths)
if let Some(enable_v2_manifest_paths) = overrides
.enable_v2_manifest_paths
.or(self.new_table_config.enable_v2_manifest_paths)
{
params.enable_v2_manifest_paths = enable_v2_manifest_paths;
}
let data_schema = request.data.schema();
if let Some(enable_stable_row_ids) = stable_row_ids_override
if let Some(enable_stable_row_ids) = overrides
.enable_stable_row_ids
.or(self.new_table_config.enable_stable_row_ids)
.or(has_blob_columns(data_schema.as_ref()).then_some(true))
{
@@ -644,6 +600,146 @@ mod tests {
RecordBatch::try_new(schema, vec![Arc::new(id_array), Arc::new(name_array)]).unwrap()
}
/// The shared parse-and-sanitize boundary is wired into this path: the
/// request-level creation key must act as an override (the strip itself
/// is covered by the listing tests).
#[tokio::test]
async fn request_level_creation_keys_are_taken_as_overrides() {
use crate::database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS;
let tmp_dir = tempdir().unwrap();
let mut properties = HashMap::new();
properties.insert(
"root".to_string(),
tmp_dir.path().to_str().unwrap().to_string(),
);
let db = connect_namespace("dir", properties)
.execute()
.await
.unwrap();
let store_params = ObjectStoreParams {
storage_options_accessor: Some(Arc::new(StorageOptionsAccessor::with_static_options(
HashMap::from([(
OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS.to_string(),
"true".to_string(),
)]),
))),
..Default::default()
};
let table = db
.create_table("t", create_test_data())
.write_options(crate::table::WriteOptions {
lance_write_params: Some(lance::dataset::WriteParams {
store_params: Some(store_params),
..Default::default()
}),
})
.execute()
.await
.unwrap();
let native = table.as_native().unwrap();
assert!(
native
.dataset
.get()
.await
.unwrap()
.manifest
.uses_stable_row_ids(),
"the creation key must be honored as an override"
);
let table = db.open_table("t").execute().await.unwrap();
let rows: usize = table
.query()
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap()
.iter()
.map(|b| b.num_rows())
.sum();
assert_eq!(rows, 5);
}
/// Sanitation on this path: apply must strip the creation keys from the
/// store options while genuine options and the provider survive.
#[tokio::test]
async fn apply_new_table_config_sanitizes_request_store_options() {
use crate::database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS;
use lance_io::object_store::StorageOptionsProvider;
#[derive(Debug)]
struct EmptyProvider;
#[async_trait::async_trait]
impl StorageOptionsProvider for EmptyProvider {
async fn fetch_storage_options(
&self,
) -> lance_core::Result<Option<HashMap<String, String>>> {
Ok(Some(HashMap::new()))
}
fn provider_id(&self) -> String {
"empty-test-provider".into()
}
}
let tmp_dir = tempdir().unwrap();
let mut properties = HashMap::new();
properties.insert(
"root".to_string(),
tmp_dir.path().to_str().unwrap().to_string(),
);
let db = LanceNamespaceDatabase::connect_with_new_table_config(
"dir",
properties,
HashMap::new(),
None,
None,
HashSet::new(),
NewTableConfig::default(),
)
.await
.unwrap();
let request = DbCreateTableRequest::new("t".to_string(), Box::new(create_test_data()));
let mut params = lance::dataset::WriteParams {
store_params: Some(ObjectStoreParams {
storage_options_accessor: Some(Arc::new(
StorageOptionsAccessor::with_initial_and_provider(
HashMap::from([
("region".to_string(), "us-west-2".to_string()),
(
OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS.to_string(),
"true".to_string(),
),
]),
Arc::new(EmptyProvider),
),
)),
..Default::default()
}),
..Default::default()
};
db.apply_new_table_config(&mut params, &request).unwrap();
assert!(params.enable_stable_row_ids);
let store_params = params.store_params.unwrap();
let options = store_params.storage_options().cloned().unwrap();
assert!(!options.contains_key(OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS));
assert_eq!(options.get("region").map(String::as_str), Some("us-west-2"));
assert!(
store_params
.storage_options_accessor
.unwrap()
.has_provider()
);
}
#[tokio::test]
async fn test_namespace_connection_simple() {
// Test that namespace connections work with simple connect_namespace(impl_type, properties)
@@ -160,9 +160,10 @@ impl PermutationBuilder {
self
}
async fn sort_by_split_id(
async fn sort_by_column(
&self,
data: SendableRecordBatchStream,
column: &str,
) -> Result<SendableRecordBatchStream> {
let memory_limit = std::env::var("LANCEDB_PERM_BUILDER_MEMORY_LIMIT")
.unwrap_or_else(|_| DEFAULT_MEMORY_LIMIT.to_string())
@@ -188,25 +189,26 @@ impl PermutationBuilder {
let df = ctx
.read_one_shot(data.into_df_stream())
.map_err(|e| Error::Other {
message: format!("Failed to setup sort by split id: {}", e),
message: format!("Failed to setup sort by {}: {}", column, e),
source: Some(e.into()),
})?;
let df_stream = df
.sort_by(vec![col(SPLIT_ID_COLUMN)])
.sort_by(vec![col(column)])
.map_err(|e| Error::Other {
message: format!("Failed to plan sort by split id: {}", e),
message: format!("Failed to plan sort by {}: {}", column, e),
source: Some(e.into()),
})?
.execute_stream()
.await
.map_err(|e| Error::Other {
message: format!("Failed to sort by split id: {}", e),
message: format!("Failed to sort by {}: {}", column, e),
source: Some(e.into()),
})?;
let column = column.to_string();
let schema = df_stream.schema();
let stream = df_stream.map_err(|e| Error::Other {
message: format!("Failed to execute sort by split id: {}", e),
let stream = df_stream.map_err(move |e| Error::Other {
message: format!("Failed to execute sort by {}: {}", column, e),
source: Some(e.into()),
});
Ok(Box::pin(SimpleRecordBatchStream { schema, stream }))
@@ -238,7 +240,25 @@ impl PermutationBuilder {
/// Builds the permutation table and stores it in the given database.
pub async fn build(self) -> Result<Table> {
// First pass, apply filter and load row ids
// Unflushed rows have no row id, so a permutation cannot address them.
match self.base_table.base_table().get_lsm_write_spec().await {
Ok(Some(_)) => {
return Err(Error::NotSupported {
message: "the data loader does not support tables with an LSM write \
spec: rows that have not been flushed to the base table \
have no row id, so a permutation cannot reference them"
.to_string(),
});
}
Ok(None) => {}
// No LSM write path means no spec.
Err(Error::NotSupported { .. }) => {}
Err(err) => return Err(err),
}
// First pass, apply filter and load row ids. `Shuffler` permutes positions, so
// every rank must scan the rows in the same order to build the same permutation.
// TODO: pin the version resolved here; remote does not implement Lazy.
let mut rows = self.base_table.query().select(Select::columns(&[ROW_ID]));
if let Some(filter) = &self.config.filter {
@@ -263,6 +283,12 @@ impl PermutationBuilder {
// Apply splits
let rows = rows.execute().await?;
// Splits are assigned by position, so the scan has to arrive in a fixed order.
let rows = if self.base_table.base_table().scan_order_is_deterministic() {
rows
} else {
self.sort_by_column(rows, ROW_ID).await?
};
let split_data = splitter.apply(rows, num_rows).await?;
// Shuffle data if requested
@@ -284,7 +310,7 @@ impl PermutationBuilder {
needs_sort |= !matches!(self.config.shuffle_strategy, ShuffleStrategy::None);
let sorted = if needs_sort {
self.sort_by_split_id(shuffled).await?
self.sort_by_column(shuffled, SPLIT_ID_COLUMN).await?
} else {
shuffled
};
@@ -367,6 +393,22 @@ mod tests {
);
}
#[tokio::test]
async fn test_native_scan_order_is_deterministic() {
let temp_dir = tempfile::tempdir().unwrap();
let db = connect(temp_dir.path().to_str().unwrap())
.execute()
.await
.unwrap();
let data = lance_datagen::gen_batch()
.col("col_a", lance_datagen::array::step::<Int32Type>())
.into_ldb_stream(RowCount::from(10), BatchCount::from(1));
let table = db.create_table("t", data).execute().await.unwrap();
// Native tables skip the canonicalizing sort; remote does not.
assert!(table.base_table().scan_order_is_deterministic());
}
#[tokio::test]
async fn test_permutation_builder() {
let temp_dir = tempfile::tempdir().unwrap();
@@ -416,4 +458,48 @@ mod tests {
283
);
}
/// Rows that have not been flushed to the base table have no row id, so a
/// permutation cannot reference them. Reading the base table alone would drop
/// them from training without saying so, so the table is refused instead.
#[tokio::test]
async fn test_permutation_rejects_lsm_write_spec() {
use crate::table::LsmWriteSpec;
use arrow_array::{Int32Array, RecordBatchIterator};
use arrow_schema::{DataType, Field, Schema};
// MemWAL needs a real dataset directory and a non-nullable primary key.
let temp_dir = tempfile::tempdir().unwrap();
let db = connect(temp_dir.path().to_str().unwrap())
.execute()
.await
.unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("idx", DataType::Int32, false)]));
let batch = arrow_array::RecordBatch::try_new(
schema.clone(),
vec![Arc::new(Int32Array::from(vec![0, 1, 2, 3]))],
)
.unwrap();
let reader: Box<dyn arrow_array::RecordBatchReader + Send> =
Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema.clone()));
let table = db.create_table("tbl", reader).execute().await.unwrap();
// Without a spec the build succeeds.
PermutationBuilder::new(table.clone())
.build()
.await
.unwrap();
table.set_unenforced_primary_key(["idx"]).await.unwrap();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
let err = PermutationBuilder::new(table).build().await.unwrap_err();
assert!(
err.to_string().contains("LSM write spec"),
"expected the pre-check to refuse the table, got: {err}"
);
}
}
@@ -97,8 +97,10 @@ impl PermutationReader {
Self::inner_new(base_table, Some(permutation_table), split).await
}
pub async fn identity(base_table: Arc<dyn BaseTable>) -> Self {
Self::inner_new(base_table, None, 0).await.unwrap()
/// A reader over the base table in storage order, with no permutation.
/// Fallible because construction counts the base table.
pub async fn identity(base_table: Arc<dyn BaseTable>) -> Result<Self> {
Self::inner_new(base_table, None, 0).await
}
/// Validates the limit and offset and returns the number of rows that will be read
@@ -487,7 +489,13 @@ impl PermutationReader {
pub async fn output_schema(&self, selection: Select) -> Result<SchemaRef> {
let table = Table::from(self.base_table.clone());
table.query().select(selection).output_schema().await
// limit(1) because some table types execute the query to get its schema
table
.query()
.select(selection)
.limit(1)
.output_schema()
.await
}
pub fn count_rows(&self) -> u64 {
@@ -779,7 +787,9 @@ mod tests {
.into_mem_table("tbl", RowCount::from(10), BatchCount::from(1))
.await;
let reader = PermutationReader::identity(base_table.base_table().clone()).await;
let reader = PermutationReader::identity(base_table.base_table().clone())
.await
.unwrap();
// With no permutation table, take_offsets uses the base table directly
let offsets = vec![0, 2, 4, 6];
@@ -961,7 +971,9 @@ mod tests {
.into_mem_table("tbl", RowCount::from(10), BatchCount::from(1))
.await;
let reader = PermutationReader::identity(base_table.base_table().clone()).await;
let reader = PermutationReader::identity(base_table.base_table().clone())
.await
.unwrap();
let batch = reader.take_offsets(&[], Select::All).await.unwrap();
+2
View File
@@ -77,6 +77,8 @@ pub enum Error {
ColumnAlreadyExists { name: String },
#[snafu(display("Column '{name}' is not a computed column"))]
NotAComputedColumn { name: String },
#[snafu(display("Table '{name}' is not a materialized view"))]
NotAMaterializedView { name: String },
#[snafu(display("Invalid expression for column '{column}': {message}"))]
InvalidExpression { column: String, message: String },
+4
View File
@@ -186,6 +186,7 @@ pub mod index;
pub mod io;
pub mod ipc;
pub mod job;
pub mod materialized_view;
#[cfg(feature = "metrics-otel")]
pub mod metrics_otel;
#[cfg(feature = "polars")]
@@ -210,6 +211,9 @@ pub use function::FunctionVersion;
pub use job::Job;
use lance_index::vector::ApproxMode as LanceApproxMode;
use lance_linalg::distance::DistanceType as LanceDistanceType;
pub use materialized_view::{
MaterializedView, MaterializedViewDefinition, RefreshMaterializedViewResult, RefreshMode,
};
/// Re-export of the [`metrics`](https://docs.rs/metrics) crate facade. Enable
/// the `metrics` feature to publish LanceDB's internal metrics; install any
/// `metrics`-compatible recorder to collect them. See also [`metrics_otel`] for
File diff suppressed because it is too large Load Diff
@@ -0,0 +1,730 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Differential refresh testing.
//!
//! The refresh contract is a property: after any sequence of source
//! mutations, a view maintained by default (incremental-where-possible)
//! refreshes equals the definition evaluated against the source directly,
//! and so does a forced rebuild. The oracle is an independent read of the
//! source -- plain column scan, filter applied in Rust -- so it shares
//! nothing with the refresh path it checks.
//!
//! The oracle runs after every step, not just at the end: a later mutation
//! that forces a rebuild would silently heal an incremental error, and those
//! transient errors are exactly the bugs this exists to catch.
use arrow_array::{Float32Array, Int32Array, RecordBatch};
use arrow_schema::{DataType, Field as ArrowField, Schema as ArrowSchema};
use futures::{StreamExt, TryStreamExt};
use lance::dataset::NewColumnTransform;
use std::sync::Arc;
use super::MaterializedView;
use super::refresh::RefreshMode;
use crate::connect;
use crate::connection::Connection;
use crate::query::{ExecutableQuery, QueryBase, Select};
use crate::table::{CompactionOptions, OptimizeAction, Table};
/// One source mutation, one per correctness-relevant class.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum SrcOp {
/// Fresh non-colliding ids of both parities, so every other op has
/// view-resident rows to act on: the only op that should refresh
/// incrementally.
AppendNew,
/// Deletion in surviving fragments must break the pure-append check.
DeleteEven,
/// An in-place update; on the filtered shape it crosses the predicate,
/// so rows must leave the view.
UpdateOddScore,
/// Fragment rewrite/renumber must break the pure-append check.
Compact,
/// A column the view does not read must NOT force a rebuild.
AddColumn,
/// merge_insert commits an Update whose by-source arm deletes rows, so a
/// classifier that reads Update as "changed only" loses those deletions.
MergeDropLargest,
/// merge_insert that both changes existing rows and inserts new ones in
/// one transaction.
MergeUpsert,
}
const ALL_OPS: [SrcOp; 7] = [
SrcOp::AppendNew,
SrcOp::DeleteEven,
SrcOp::UpdateOddScore,
SrcOp::Compact,
SrcOp::AddColumn,
SrcOp::MergeDropLargest,
SrcOp::MergeUpsert,
];
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Shape {
/// SELECT id, score.
Identity,
/// SELECT id, score WHERE score > 50: additionally sensitive to rows
/// crossing the predicate.
Filtered,
/// SELECT id, score LIMIT 4. Which rows are held depends on the order
/// they were first materialized, so the oracle checks containment and
/// the cap rather than equality.
Limited,
}
impl Shape {
fn filter(&self) -> Option<&'static str> {
match self {
Self::Identity | Self::Limited => None,
Self::Filtered => Some("score > 50"),
}
}
fn matches(&self, score: f32) -> bool {
match self {
Self::Identity | Self::Limited => true,
Self::Filtered => score > 50.0,
}
}
fn limit(&self) -> Option<usize> {
match self {
Self::Limited => Some(4),
_ => None,
}
}
}
struct Case {
conn: Connection,
source: Table,
view: MaterializedView,
shape: Shape,
next_id: i32,
added_columns: u32,
}
fn rows_batch(ids: &[i32]) -> RecordBatch {
let scores: Vec<f32> = ids.iter().map(|id| (*id * 10) as f32).collect();
RecordBatch::try_new(
Arc::new(ArrowSchema::new(vec![
ArrowField::new("id", DataType::Int32, true),
ArrowField::new("score", DataType::Float32, true),
])),
vec![
Arc::new(Int32Array::from(ids.to_vec())),
Arc::new(Float32Array::from(scores)),
],
)
.unwrap()
}
fn merge_batch(ids: &[i32]) -> RecordBatch {
let scores: Vec<f32> = ids.iter().map(|id| (*id * 10 + 5) as f32).collect();
RecordBatch::try_new(
Arc::new(ArrowSchema::new(vec![
ArrowField::new("id", DataType::Int32, true),
ArrowField::new("score", DataType::Float32, true),
])),
vec![
Arc::new(Int32Array::from(ids.to_vec())),
Arc::new(Float32Array::from(scores)),
],
)
.unwrap()
}
impl Case {
async fn new(shape: Shape) -> Self {
let conn = connect("memory://").execute().await.unwrap();
let source = conn
.create_table("src", rows_batch(&[1, 2, 3, 4]))
.write_options(crate::materialized_view::tests::stable_row_ids())
.execute()
.await
.unwrap();
let mut builder = conn
.create_materialized_view("view", "src")
.select([("id", "id"), ("score", "score")]);
if let Some(filter) = shape.filter() {
builder = builder.only_if(filter);
}
if let Some(limit) = shape.limit() {
builder = builder.limit(limit as u64);
}
let view = builder.execute().await.unwrap();
Self {
conn,
source,
view,
shape,
next_id: 100,
added_columns: 0,
}
}
async fn apply(&mut self, op: SrcOp) {
match op {
SrcOp::AppendNew => {
// Mixed parity: the middle id is odd, so UpdateOddScore always
// has a filter-matching appended row to evict.
let ids = vec![self.next_id, self.next_id + 101, self.next_id + 202];
self.next_id += 303;
self.source.add(rows_batch(&ids)).execute().await.unwrap();
}
SrcOp::DeleteEven => {
self.source.delete("id % 2 = 0").await.unwrap();
}
SrcOp::UpdateOddScore => {
self.source
.update()
.column("score", "-1.0")
.only_if("id % 2 = 1")
.execute()
.await
.unwrap();
}
SrcOp::Compact => {
self.source
.optimize(OptimizeAction::Compact {
options: CompactionOptions::default(),
remap_options: None,
})
.await
.unwrap();
}
SrcOp::MergeDropLargest => {
let mut ids = self.source_ids().await;
ids.sort_unstable();
ids.pop();
if ids.is_empty() {
return;
}
let batch = rows_batch(&ids);
let reader =
arrow_array::RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
let mut merge = self.source.merge_insert(&["id"]);
merge.when_not_matched_by_source_delete(None);
merge.execute(Box::new(reader)).await.unwrap();
}
SrcOp::MergeUpsert => {
let mut ids = self.source_ids().await;
ids.sort_unstable();
// One row that exists (updated in place) and one that does not.
let existing = ids.first().copied().unwrap_or(self.next_id);
let fresh = self.next_id;
self.next_id += 1;
let batch = merge_batch(&[existing, fresh]);
let reader =
arrow_array::RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
let mut merge = self.source.merge_insert(&["id"]);
merge
.when_matched_update_all(None)
.when_not_matched_insert_all();
merge.execute(Box::new(reader)).await.unwrap();
}
SrcOp::AddColumn => {
self.added_columns += 1;
let field = ArrowField::new(
format!("extra_{}", self.added_columns),
DataType::Int32,
true,
);
self.source
.add_columns()
.transform(NewColumnTransform::AllNulls(Arc::new(ArrowSchema::new(
vec![field],
))))
.execute()
.await
.unwrap();
}
}
}
async fn source_ids(&self) -> Vec<i32> {
read_rows(
self.source
.query()
.select(Select::columns(&["id", "score"])),
)
.await
.into_iter()
.map(|(id, _)| id)
.collect()
}
/// The definition's result, read independently of the refresh path:
/// plain column scan, filter applied here, sorted.
async fn oracle(&self) -> Vec<(i32, i32)> {
let mut rows = read_rows(
self.source
.query()
.select(Select::columns(&["id", "score"])),
)
.await
.into_iter()
.filter(|(_, score)| self.shape.matches(*score as f32))
.collect::<Vec<_>>();
rows.sort_unstable();
rows
}
async fn view_rows(&self) -> Vec<(i32, i32)> {
let mut rows = read_rows(
self.view
.table()
.query()
.select(Select::columns(&["id", "score"])),
)
.await;
rows.sort_unstable();
rows
}
async fn check(&self, label: &str) -> Result<(), String> {
let expected = self.oracle().await;
let actual = self.view_rows().await;
let Some(cap) = self.shape.limit() else {
if expected != actual {
return Err(format!(
"{label}: view diverged from oracle\n expected: {expected:?}\n actual: {actual:?}"
));
}
return Ok(());
};
// A capped view holds some subset of the definition's result, never
// more than the cap, and never the same row twice.
if actual.len() > cap {
return Err(format!(
"{label}: view holds {} rows, over its cap of {cap}: {actual:?}",
actual.len()
));
}
let mut unique = actual.clone();
unique.dedup();
if unique.len() != actual.len() {
return Err(format!("{label}: view holds a row twice: {actual:?}"));
}
if let Some(stray) = actual.iter().find(|row| !expected.contains(row)) {
return Err(format!(
"{label}: view holds {stray:?}, which the definition does not select: {expected:?}"
));
}
// Below the cap the view must be complete, or a row was lost.
if actual.len() < cap.min(expected.len()) {
return Err(format!(
"{label}: view holds {} of {} selectable rows under a cap of {cap}: {actual:?}",
actual.len(),
expected.len()
));
}
Ok(())
}
}
async fn read_rows(query: impl ExecutableQuery) -> Vec<(i32, i32)> {
let batches = query
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
batches
.iter()
.flat_map(|batch| {
let ids = batch["id"].as_any().downcast_ref::<Int32Array>().unwrap();
let scores = batch["score"]
.as_any()
.downcast_ref::<Float32Array>()
.unwrap();
// Scores are integer-valued by construction; compare exactly.
(0..batch.num_rows())
.map(|i| (ids.value(i), scores.value(i) as i32))
.collect::<Vec<_>>()
})
.collect()
}
/// Drive one mutation sequence: refresh + oracle-check after every step,
/// then a forced rebuild checked against the same oracle.
async fn run_sequence(ops: &[SrcOp], shape: Shape) -> Result<(), String> {
let label = format!("{shape:?} {ops:?}");
let mut case = Case::new(shape).await;
case.view
.refresh()
.execute()
.await
.map_err(|e| format!("{label}: initial refresh failed: {e}"))?;
case.check(&format!("{label} (initial)")).await?;
for (step, op) in ops.iter().enumerate() {
case.apply(*op).await;
case.view
.refresh()
.execute()
.await
.map_err(|e| format!("{label}: refresh at step {step} failed: {e}"))?;
case.check(&format!("{label} (step {step}, {op:?})"))
.await?;
}
case.view
.refresh()
.full(true)
.execute()
.await
.map_err(|e| format!("{label}: final full refresh failed: {e}"))?;
case.check(&format!("{label} (final rebuild)")).await?;
// Silence the unused-connection lint without dropping it mid-case.
let _ = &case.conn;
Ok(())
}
/// Every op sequence up to `max_len`.
fn all_sequences(max_len: u32) -> Vec<Vec<SrcOp>> {
let mut sequences = Vec::new();
for len in 1..=max_len {
for mut index in 0..ALL_OPS.len().pow(len) {
let mut ops = Vec::with_capacity(len as usize);
for _ in 0..len {
ops.push(ALL_OPS[index % ALL_OPS.len()]);
index /= ALL_OPS.len();
}
sequences.push(ops);
}
}
sequences
}
async fn run_exhaustive(max_len: u32) {
let mut cases = Vec::new();
for shape in [Shape::Identity, Shape::Filtered, Shape::Limited] {
for ops in all_sequences(max_len) {
cases.push((ops, shape));
}
}
let failures: Vec<String> = futures::stream::iter(cases)
.map(|(ops, shape)| async move { run_sequence(&ops, shape).await.err() })
.buffer_unordered(8)
.filter_map(|failure| async move { failure })
.collect()
.await;
assert!(
failures.is_empty(),
"{} sequences diverged; first: {}",
failures.len(),
failures[0]
);
}
#[tokio::test(flavor = "multi_thread")]
async fn differential_exhaustive() {
run_exhaustive(3).await;
}
#[tokio::test(flavor = "multi_thread")]
#[ignore = "longer sweep; run manually"]
async fn differential_exhaustive_deep() {
run_exhaustive(4).await;
}
/// Named interleavings that double as repro handles. The mode assertions pin
/// the classifier, which value comparison alone cannot: a wrongly rebuilt
/// view still matches the oracle.
#[tokio::test(flavor = "multi_thread")]
async fn differential_named_regressions() {
// An append is the one op that must stay incremental.
let mut case = Case::new(Shape::Identity).await;
case.view.refresh().execute().await.unwrap();
case.apply(SrcOp::AppendNew).await;
let result = case.view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
case.check("append stays incremental").await.unwrap();
// A column the view does not read must not force a rebuild.
let mut case = Case::new(Shape::Identity).await;
case.view.refresh().execute().await.unwrap();
case.apply(SrcOp::AddColumn).await;
let result = case.view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 0);
// Compaction rearranges rows without changing them: the watermark
// advances and nothing rebuilds.
let mut case = Case::new(Shape::Identity).await;
case.view.refresh().execute().await.unwrap();
case.apply(SrcOp::AppendNew).await;
case.view.refresh().execute().await.unwrap();
case.apply(SrcOp::Compact).await;
let result = case.view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 0);
case.check("compaction alone").await.unwrap();
// Fragment bookkeeping stays coherent across the compaction: the next
// append is separable and computed alone.
case.apply(SrcOp::AppendNew).await;
let result = case.view.refresh().execute().await.unwrap();
assert_eq!(result.mode, RefreshMode::Incremental);
assert_eq!(result.rows_written, 3);
case.check("compact then append").await.unwrap();
// A row updated to no longer match the filter must leave the view --
// and the fixture must prove the eviction happened, not merely that the
// end state matches: an update that never touched a view-resident row
// would also "match".
let mut case = Case::new(Shape::Filtered).await;
case.apply(SrcOp::AppendNew).await;
case.view.refresh().execute().await.unwrap();
let before = case.view_rows().await.len();
case.apply(SrcOp::UpdateOddScore).await;
case.view.refresh().execute().await.unwrap();
let after = case.view_rows().await.len();
assert!(
after < before,
"no view-resident row was evicted ({before} -> {after}); the fixture \
no longer exercises the filtered-update transition"
);
case.check("update crosses the filter").await.unwrap();
}
// ---------------------------------------------------------------------------
// Concurrency
// ---------------------------------------------------------------------------
//
// The sequential cases above cannot observe a cross-process race: the
// per-view refresh lock is process-local, so a second refresh in this
// process queues behind the first. What is missing is not more op
// sequences but a second process. These cases add one, and assert the same
// property the harness always asserts -- the view holds each row once.
/// Rows the definition selects from the source: every id but the first,
/// read straight from the source, sharing nothing with the refresh path.
async fn concurrency_oracle(conn: &Connection) -> Vec<i32> {
let batches: Vec<RecordBatch> = conn
.open_table("src")
.execute()
.await
.unwrap()
.query()
.select(Select::columns(&["id"]))
.execute()
.await
.unwrap()
.try_collect()
.await
.unwrap();
let mut ids = Vec::new();
for batch in &batches {
let column = batch
.column(0)
.as_any()
.downcast_ref::<Int32Array>()
.unwrap();
for i in 0..batch.num_rows() {
if column.value(i) > 1 {
ids.push(column.value(i));
}
}
}
ids.sort_unstable();
ids
}
/// The view's ids, sorted.
async fn concurrency_view_ids(conn: &Connection) -> Vec<i32> {
let batches: Vec<RecordBatch> = conn
.open_table("mv")
.execute()
.await
.unwrap()
.query()
.select(Select::columns(&["id"]))
.execute()
.await
.unwrap()
.try_collect()
.await
.unwrap();
let mut ids = Vec::new();
for batch in &batches {
let column = batch
.column(0)
.as_any()
.downcast_ref::<Int32Array>()
.unwrap();
for i in 0..batch.num_rows() {
ids.push(column.value(i));
}
}
ids.sort_unstable();
ids
}
/// One refresh of the view at `MV_RACE_DIR`, in its own process.
///
/// Setup happens before the start barrier so warm-up does not stagger the
/// two processes. What makes the race certain rather than likely is the
/// second barrier inside `refresh()` itself, which holds every participant
/// between staging and commit.
#[tokio::test]
#[ignore = "spawned as a child process by the concurrency cases"]
async fn cross_process_refresh_child() {
let Ok(dir) = std::env::var("MV_RACE_DIR") else {
return;
};
let dir = std::path::PathBuf::from(dir);
let tag = std::env::var("MV_RACE_TAG").unwrap();
let conn = connect(dir.to_str().unwrap()).execute().await.unwrap();
let table = conn.open_table("mv").execute().await.unwrap();
let _ = table.schema().await.unwrap();
let _ = table.count_rows(None).await.unwrap();
let source = conn.open_table("src").execute().await.unwrap();
let _ = source.count_rows(None).await.unwrap();
let view = MaterializedView::from_table(table).await.unwrap();
std::fs::write(dir.join(format!("ready-{tag}")), b"1").unwrap();
while !dir.join("START").exists() {
std::thread::sleep(std::time::Duration::from_millis(2));
}
let outcome = match view.refresh().execute().await {
Ok(result) => format!("committed rows={}", result.rows_written),
Err(err) if is_commit_conflict(&err) => "conflicted".to_string(),
Err(err) => format!("failed {err}"),
};
std::fs::write(dir.join(format!("outcome-{tag}")), outcome).unwrap();
}
/// Whether a refresh lost its commit to a concurrent one, as opposed to
/// failing for any other reason.
fn is_commit_conflict(err: &crate::Error) -> bool {
let text = err.to_string();
text.contains("Retryable commit conflict") || text.contains("preempted by concurrent")
}
/// Two processes refreshing one view concurrently must leave the view
/// equal to the oracle: each selected row present exactly once.
///
/// Both plan the same incremental delta from one watermark. A refresh is
/// meant to land on the generation it planned or leave nothing behind, so
/// at most one of them may write.
#[tokio::test]
async fn concurrent_refreshes_hold_each_row_once() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().to_str().unwrap().to_string();
let conn = connect(&path).execute().await.unwrap();
conn.create_table("src", rows_batch(&[1, 2, 3, 4]))
.write_options(crate::materialized_view::tests::stable_row_ids())
.execute()
.await
.unwrap();
let view = conn
.create_materialized_view("mv", "src")
.select([("id", "id"), ("score", "score")])
.only_if("id > 1")
.execute()
.await
.unwrap();
// Seed the watermark so the racing refreshes are both incremental.
view.refresh().execute().await.unwrap();
// Large enough that a refresh is real work rather than a formality.
let ids: Vec<i32> = (100..200_100).collect();
conn.open_table("src")
.execute()
.await
.unwrap()
.add(rows_batch(&ids))
.execute()
.await
.unwrap();
let tags = ["a", "b"];
let exe = std::env::current_exe().unwrap();
let children: Vec<std::process::Child> = tags
.iter()
.map(|tag| {
std::process::Command::new(&exe)
.args([
"--exact",
"materialized_view::differential::cross_process_refresh_child",
"--ignored",
"--nocapture",
])
.env("MV_RACE_DIR", dir.path())
.env("MV_RACE_SYNC", dir.path())
.env("MV_RACE_PEERS", "2")
.env("MV_RACE_TAG", tag)
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::null())
.spawn()
.unwrap()
})
.collect();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(180);
while tags
.iter()
.any(|tag| !dir.path().join(format!("ready-{tag}")).exists())
{
assert!(
std::time::Instant::now() < deadline,
"children never became ready"
);
std::thread::sleep(std::time::Duration::from_millis(10));
}
std::fs::write(dir.path().join("START"), b"1").unwrap();
for (tag, mut child) in tags.iter().zip(children) {
let status = loop {
match child.try_wait().unwrap() {
Some(status) => break status,
None if std::time::Instant::now() >= deadline => {
child.kill().unwrap();
panic!("child {tag} never finished");
}
None => std::thread::sleep(std::time::Duration::from_millis(10)),
}
};
assert!(status.success(), "child {tag} exited {status}");
}
// Both refreshes reached the commit boundary before either committed --
// the in-refresh barrier guarantees it -- so exactly one may win.
let outcomes: Vec<String> = tags
.iter()
.map(|tag| {
std::fs::read_to_string(dir.path().join(format!("outcome-{tag}")))
.unwrap_or_else(|_| panic!("child {tag} recorded no outcome"))
})
.collect();
for tag in tags {
assert!(
dir.path().join(format!("planned-{tag}")).exists(),
"child {tag} never reached the commit boundary, so nothing was synchronized"
);
}
let committed = outcomes.iter().filter(|o| o.contains("committed")).count();
let conflicted = outcomes.iter().filter(|o| o.contains("conflicted")).count();
assert_eq!(
(committed, conflicted),
(1, 1),
"exactly one refresh may win the generation both planned: {outcomes:?}"
);
let expected = concurrency_oracle(&conn).await;
let actual = concurrency_view_ids(&conn).await;
assert_eq!(
actual.len(),
expected.len(),
"the view holds {} rows, the oracle {}: a losing refresh left rows behind",
actual.len(),
expected.len()
);
assert_eq!(actual, expected, "the view does not match the oracle");
}
File diff suppressed because it is too large Load Diff
+68 -35
View File
@@ -21,11 +21,11 @@ use crate::remote::job::RemoteJob;
use crate::table::AddColumnsResult;
use crate::table::AddResult;
use crate::table::BranchDiff;
use crate::table::CherryPickResult;
use crate::table::DeleteResult;
use crate::table::DropColumnsResult;
use crate::table::LsmStats;
use crate::table::LsmWriteSpec;
use crate::table::MergeBranchResult;
use crate::table::MergeResult;
use crate::table::Tags;
use crate::table::UpdateResult;
@@ -2031,7 +2031,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
async fn diff_branch(&self, from_branch: &str) -> Result<BranchDiff> {
if from_branch.trim().is_empty() {
return Err(Error::InvalidInput {
message: "from_branch must be a non-empty string".into(),
message: "Branch name cannot be empty.".into(),
});
}
let request = self
@@ -2058,20 +2058,23 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
})
}
async fn merge_branch(&self, from_branch: &str, dry_run: bool) -> Result<MergeBranchResult> {
async fn cherry_pick(&self, from_branch: &str, dry_run: bool) -> Result<CherryPickResult> {
if from_branch.trim().is_empty() {
return Err(Error::InvalidInput {
message: "from_branch must be a non-empty string".into(),
message: "Branch name cannot be empty.".into(),
});
}
let request = self
.client
.post(&format!("/v1/table/{}/branches/merge/", self.identifier))
.post(&format!(
"/v1/table/{}/branches/cherry_pick/",
self.identifier
))
.json(&serde_json::json!({
"from_branch": from_branch,
"dry_run": dry_run,
}));
// No retry. 409 rejected merge is final and carries a body.
// No retry. HTTP 409 is CherryPickStatus::Failed with a body, not a transport error.
let (request_id, response) = self.send(request, false).await?;
let status = response.status();
if status == StatusCode::NOT_FOUND {
@@ -2080,11 +2083,11 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
source: format!("branch '{}' does not exist", from_branch).into(),
});
}
// 200 and 409 both carry MergeBranchResult.
// 200 and 409 both carry CherryPickResult.
if status != StatusCode::OK && status != StatusCode::CONFLICT {
let body = response.text().await.unwrap_or_default();
return Err(Error::Http {
source: format!("unexpected status {status} from merge_branch: {body}").into(),
source: format!("unexpected status {status} from cherry_pick: {body}").into(),
request_id,
status_code: Some(status),
});
@@ -2092,7 +2095,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
let body = response.text().await.err_to_http(request_id.clone())?;
serde_json::from_str(&body).map_err(|err| Error::Http {
source: format!(
"Failed to parse merge_branch response: {}, body: {}",
"Failed to parse cherry_pick response: {}, body: {}",
err, body
)
.into(),
@@ -4303,6 +4306,19 @@ mod tests {
write_ipc_stream_uncompressed(&one_row_blob_batch(column))
}
#[tokio::test]
async fn test_remote_scan_order_is_not_deterministic() {
// A distributed scan answers in no fixed order, so callers that assign meaning
// to row position have to sort for themselves.
let table = Table::new_with_handler("my_table", |_| {
http::Response::builder()
.status(200)
.body(Vec::new())
.unwrap()
});
assert!(!table.base_table().scan_order_is_deterministic());
}
#[tokio::test]
async fn test_fetch_blobs_sends_the_checked_out_version() {
let ipc = one_row_blob_ipc_stream("image");
@@ -10397,6 +10413,20 @@ mod tests {
);
}
#[tokio::test]
async fn test_materialized_view_refused_without_a_request() {
// Materialized views are local-only. The table-level entry the
// bindings use must refuse a remote table before reading its schema,
// so the panicking handler is the assertion.
let table = Table::new_with_handler("my_table", |request| -> http::Response<String> {
panic!("unexpected request: {}", request.url().path())
});
let err = crate::MaterializedView::from_table(table)
.await
.unwrap_err();
assert!(matches!(err, Error::NotSupported { .. }), "got {err:?}");
}
#[tokio::test]
async fn test_create_branch_empty_name_rejected_client_side() {
use lance::dataset::refs::Ref;
@@ -10486,8 +10516,7 @@ mod tests {
"changedColumns":[],
"addedIndexes":[],
"removedIndexes":[],
"mergeable":true,
"mergeBlockers":[]
"errors":[]
}"#
}
@@ -10505,15 +10534,18 @@ mod tests {
});
let diff = table.diff_branch("exp").await.unwrap();
assert_eq!(diff.from_branch, "exp");
assert!(diff.mergeable);
assert!(diff.errors.is_empty());
assert_eq!(diff.added_columns.len(), 1);
assert_eq!(diff.added_columns[0].name, "tag");
}
#[tokio::test]
async fn test_merge_branch_dry_run() {
async fn test_cherry_pick_dry_run() {
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.url().path(), "/v1/table/my_table/branches/merge/");
assert_eq!(
request.url().path(),
"/v1/table/my_table/branches/cherry_pick/"
);
let body = request_body_json(&request);
assert_eq!(body["from_branch"], "exp");
assert_eq!(body["dry_run"], true);
@@ -10523,27 +10555,29 @@ mod tests {
);
http::Response::builder().status(200).body(resp).unwrap()
});
let result = table.merge_branch("exp", true).await.unwrap();
assert_eq!(result.status, crate::table::MergeBranchStatus::Ready);
let result = table.cherry_pick("exp", true).await.unwrap();
assert_eq!(result.status, crate::table::CherryPickStatus::Ready);
assert_eq!(result.preview.promoted_columns, vec!["tag".to_string()]);
assert!(result.main_version_after.is_none());
}
#[tokio::test]
async fn test_merge_branch_rejected_returns_ok_with_body() {
async fn test_cherry_pick_failed_returns_ok_with_body() {
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.url().path(), "/v1/table/my_table/branches/merge/");
assert_eq!(
request.url().path(),
"/v1/table/my_table/branches/cherry_pick/"
);
let body = request_body_json(&request);
assert_eq!(body["dry_run"], false);
let mut diff: serde_json::Value =
serde_json::from_str(sample_branch_diff_json()).unwrap();
diff["mergeable"] = serde_json::json!(false);
diff["mergeBlockers"] = serde_json::json!([{
diff["errors"] = serde_json::json!([{
"code": "baseMoved",
"message": "main has advanced"
}]);
let resp = serde_json::json!({
"status": "rejected",
"status": "failed",
"diff": diff,
"preview": { "promotedColumns": [] }
});
@@ -10552,24 +10586,23 @@ mod tests {
.body(resp.to_string())
.unwrap()
});
let result = table.merge_branch("exp", false).await.unwrap();
assert_eq!(result.status, crate::table::MergeBranchStatus::Rejected);
assert!(!result.diff.mergeable);
assert_eq!(result.diff.merge_blockers.len(), 1);
let result = table.cherry_pick("exp", false).await.unwrap();
assert_eq!(result.status, crate::table::CherryPickStatus::Failed);
assert!(!result.diff.errors.is_empty());
assert_eq!(result.diff.errors.len(), 1);
}
#[tokio::test]
async fn test_merge_branch_unknown_blocker_code_parses() {
async fn test_cherry_pick_unknown_error_code_parses() {
let table = Table::new_with_handler("my_table", |_| {
let mut diff: serde_json::Value =
serde_json::from_str(sample_branch_diff_json()).unwrap();
diff["mergeable"] = serde_json::json!(false);
diff["mergeBlockers"] = serde_json::json!([{
diff["errors"] = serde_json::json!([{
"code": "multipleCommits",
"message": "branch has more than one data commit"
}]);
let resp = serde_json::json!({
"status": "rejected",
"status": "failed",
"diff": diff,
"preview": { "operation": "append", "rowsAdded": 2 }
});
@@ -10578,24 +10611,24 @@ mod tests {
.body(resp.to_string())
.unwrap()
});
let result = table.merge_branch("exp", false).await.unwrap();
assert_eq!(result.status, crate::table::MergeBranchStatus::Rejected);
let result = table.cherry_pick("exp", false).await.unwrap();
assert_eq!(result.status, crate::table::CherryPickStatus::Failed);
assert_eq!(
result.diff.merge_blockers[0].code,
crate::table::MergeBlockerCode::Unknown
result.diff.errors[0].code,
crate::table::CherryPickErrorCode::Unknown
);
assert!(result.preview.promoted_columns.is_empty());
}
#[tokio::test]
async fn test_merge_branch_unexpected_2xx_is_error() {
async fn test_cherry_pick_unexpected_2xx_is_error() {
let table = Table::new_with_handler("my_table", |_| {
http::Response::builder()
.status(204)
.body(String::new())
.unwrap()
});
let err = table.merge_branch("exp", false).await.unwrap_err();
let err = table.cherry_pick("exp", false).await.unwrap_err();
match err {
Error::Http {
status_code: Some(code),
+31 -17
View File
@@ -66,8 +66,8 @@ use self::merge::MergeInsertBuilder;
pub mod add_columns;
mod add_data;
pub mod branch_merge;
pub mod checkpoint;
pub mod cherry_pick;
pub mod computed_columns;
mod create_index;
pub mod datafusion;
@@ -87,9 +87,9 @@ pub use add_columns::AddColumnsBuilder;
#[cfg(feature = "remote")]
pub(crate) use add_data::PreprocessingOutput;
pub use add_data::{AddDataBuilder, AddDataMode, AddResult, NaNVectorBehavior};
pub use branch_merge::{
BranchDiff, ColumnChange, ColumnSummary, IndexSummary, MergeBlocker, MergeBlockerCode,
MergeBranchResult, MergeBranchStatus, MergePreview, RowCountSummary,
pub use cherry_pick::{
BranchDiff, CherryPickError, CherryPickErrorCode, CherryPickPreview, CherryPickResult,
CherryPickStatus, ColumnChange, ColumnSummary, IndexSummary, RowCountSummary,
};
pub use chrono::Duration;
pub use computed_columns::{
@@ -789,6 +789,13 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
async fn checkout_tag(&self, tag: &str) -> Result<()>;
/// Checkout the latest version of the table.
async fn checkout_latest(&self) -> Result<()>;
/// Whether repeated identical scans return rows in the same order.
///
/// Callers that assign meaning to a row's position must order the results
/// themselves when this is false. Defaults to false so a table type opts in.
fn scan_order_is_deterministic(&self) -> bool {
false
}
/// Restore the table to the currently checked out version.
async fn restore(&self) -> Result<()>;
/// List the versions of the table.
@@ -825,14 +832,14 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
/// Diff a branch against main. Remote only.
async fn diff_branch(&self, _from_branch: &str) -> Result<BranchDiff> {
Err(Error::NotSupported {
message: "diff_branch is only supported on remote tables".into(),
message: "Branch diffs are only supported on Enterprise tables.".into(),
})
}
/// Merge a branch into main, or dry-run. Remote only.
/// HTTP 409 still returns [`Ok`] with [`MergeBranchStatus::Rejected`].
async fn merge_branch(&self, _from_branch: &str, _dry_run: bool) -> Result<MergeBranchResult> {
/// Cherry-pick a branch onto main, or dry-run. Remote only.
/// HTTP 409 still returns [`Ok`] with [`CherryPickStatus::Failed`].
async fn cherry_pick(&self, _from_branch: &str, _dry_run: bool) -> Result<CherryPickResult> {
Err(Error::NotSupported {
message: "merge_branch is only supported on remote tables".into(),
message: "Cherry-picking branches is only supported on Enterprise tables.".into(),
})
}
/// The branch this handle is scoped to, or `None` for `main`.
@@ -1059,6 +1066,11 @@ impl Table {
self.database.as_ref().unwrap()
}
/// The database this handle was opened through, when it was.
pub fn database_opt(&self) -> Option<&Arc<dyn Database>> {
self.database.as_ref()
}
pub fn embedding_registry(&self) -> &Arc<dyn EmbeddingRegistry> {
&self.embedding_registry
}
@@ -2251,14 +2263,10 @@ impl Table {
self.inner.diff_branch(from_branch).await
}
/// Merge a branch into main, or dry-run. Remote only.
/// HTTP 409 still returns [`Ok`] with [`MergeBranchStatus::Rejected`].
pub async fn merge_branch(
&self,
from_branch: &str,
dry_run: bool,
) -> Result<MergeBranchResult> {
self.inner.merge_branch(from_branch, dry_run).await
/// Cherry-pick a branch onto main, or dry-run. Remote only.
/// HTTP 409 still returns [`Ok`] with [`CherryPickStatus::Failed`].
pub async fn cherry_pick(&self, from_branch: &str, dry_run: bool) -> Result<CherryPickResult> {
self.inner.cherry_pick(from_branch, dry_run).await
}
/// The branch this handle is scoped to, or `None` for `main`.
@@ -2996,6 +3004,12 @@ impl BaseTable for NativeTable {
self
}
/// Lance scans fragments in order (`Scanner::ordered` defaults to true, and we
/// never clear it), so repeated identical scans agree.
fn scan_order_is_deterministic(&self) -> bool {
true
}
fn name(&self) -> &str {
self.name.as_str()
}
@@ -1,7 +1,7 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Types for remote branch diff / merge against main.
//! Types for remote branch diff / cherry-pick onto main.
use serde::{Deserialize, Serialize};
@@ -44,13 +44,13 @@ pub struct RowCountSummary {
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub enum MergeBlockerCode {
pub enum CherryPickErrorCode {
BaseMoved,
RowCountMismatch,
RowsChanged,
ColumnRemoved,
ColumnChanged,
NoMergeableChanges,
NothingToApply,
NoColumnChanges,
InputColumnDependency,
ParentNotMain,
@@ -60,8 +60,8 @@ pub enum MergeBlockerCode {
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct MergeBlocker {
pub code: MergeBlockerCode,
pub struct CherryPickError {
pub code: CherryPickErrorCode,
pub message: String,
}
@@ -81,34 +81,33 @@ pub struct BranchDiff {
pub changed_columns: Vec<ColumnChange>,
pub added_indexes: Vec<IndexSummary>,
pub removed_indexes: Vec<IndexSummary>,
pub mergeable: bool,
pub merge_blockers: Vec<MergeBlocker>,
pub errors: Vec<CherryPickError>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct MergePreview {
pub struct CherryPickPreview {
#[serde(default)]
pub promoted_columns: Vec<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub enum MergeBranchStatus {
pub enum CherryPickStatus {
Ready,
Rejected,
Failed,
NotImplemented,
Merged,
CherryPicked,
#[serde(other)]
Unknown,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct MergeBranchResult {
pub status: MergeBranchStatus,
pub struct CherryPickResult {
pub status: CherryPickStatus,
pub diff: BranchDiff,
pub preview: MergePreview,
pub preview: CherryPickPreview,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub main_version_after: Option<u64>,
}
+38
View File
@@ -12,6 +12,7 @@ use arrow_schema::{DataType, Field};
use lance::index::DatasetIndexExt;
use lance::index::vector::VectorIndexParams;
use lance::index::vector::utils::infer_vector_dim;
use lance_arrow::json::is_json_field;
use lance_index::IndexType;
use lance_index::scalar::{BuiltinIndexType, ScalarIndexParams};
use lance_index::vector::bq::RQBuildParams;
@@ -219,6 +220,14 @@ impl NativeTable {
)))
}
Index::Bitmap(_) => {
if is_json_field(field) {
return Err(Error::Schema {
message: format!(
"A BITMAP index cannot be created on the whole-document lance.json field `{}`. Create a JSON-path scalar index for structured equality or range predicates, or use FTS for document search",
field.name()
),
});
}
Self::validate_index_type(field, "Bitmap", supported_bitmap_data_type)?;
Ok(Box::new(ScalarIndexParams::for_builtin(
BuiltinIndexType::Bitmap,
@@ -1465,6 +1474,35 @@ mod tests {
assert_eq!(stats.distance_type, None);
}
#[tokio::test]
async fn test_create_bitmap_index_rejects_lance_json() {
let conn = connect("memory://").execute().await.unwrap();
let schema = Arc::new(Schema::new(vec![lance_arrow::json::json_field(
"metadata", true,
)]));
let table = conn
.create_empty_table("json_bitmap", schema)
.execute()
.await
.unwrap();
let err = table
.create_index(&["metadata"], Index::Bitmap(Default::default()))
.execute()
.await
.expect_err("a whole-document lance.json field must not support a bitmap index");
let message = err.to_string();
assert!(
message.contains("lance.json"),
"unexpected error: {message}"
);
assert!(
message.contains("JSON-path scalar index"),
"unexpected error: {message}"
);
assert!(message.contains("FTS"), "unexpected error: {message}");
}
#[tokio::test]
async fn test_create_label_list_index() {
let conn = connect("memory://").execute().await.unwrap();
+70 -1
View File
@@ -321,7 +321,8 @@ pub(crate) async fn execute_merge_insert(
mod tests {
use arrow_array::builder::FixedSizeBinaryBuilder;
use arrow_array::{
Int32Array, RecordBatch, RecordBatchIterator, RecordBatchReader, StringArray, UInt64Array,
FixedSizeListArray, Int32Array, NullArray, RecordBatch, RecordBatchIterator,
RecordBatchReader, StringArray, UInt32Array, UInt64Array,
};
use arrow_schema::{DataType, Field, Schema};
use std::sync::Arc;
@@ -529,6 +530,74 @@ mod tests {
assert_eq!(result.num_deleted_rows, 5);
assert_eq!(table.count_rows(None).await.unwrap(), 5);
}
#[tokio::test]
async fn test_merge_insert_fixed_size_list_above_u32_child_count() {
// Arrow's FixedSizeList take kernel uses u32 child indices. Previously,
// delete-by-source materialized the target payload in a full outer join,
// causing the final list below to overflow those indices and panic.
// A Null child keeps this boundary test small in memory.
const LIST_SIZE: i32 = 65_536;
const ROW_COUNT: usize = (u32::MAX as usize / LIST_SIZE as usize) + 1;
const BATCH_SIZE: usize = 8_192;
let item = Arc::new(Field::new("item", DataType::Null, true));
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::UInt32, false),
Field::new(
"vector",
DataType::FixedSizeList(item.clone(), LIST_SIZE),
false,
),
]));
let batch = |start: usize, len: usize| {
RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(UInt32Array::from_iter_values(
start as u32..(start + len) as u32,
)),
Arc::new(FixedSizeListArray::new(
item.clone(),
LIST_SIZE,
Arc::new(NullArray::new(len * LIST_SIZE as usize)),
None,
)),
],
)
.unwrap()
};
let target_batches = (0..ROW_COUNT)
.step_by(BATCH_SIZE)
.map(|start| {
let len = (ROW_COUNT - start).min(BATCH_SIZE);
Ok(batch(start, len))
})
.collect::<Vec<_>>();
let target_data: Box<dyn RecordBatchReader + Send> =
Box::new(RecordBatchIterator::new(target_batches, schema.clone()));
let conn = connect("memory://").execute().await.unwrap();
let table = conn
.create_table("fixed_size_list_overflow", target_data)
.execute()
.await
.unwrap();
let source = batch(ROW_COUNT - 1, 1);
let mut merge = table.merge_insert(&["id"]);
merge
.when_matched_update_all(None)
.when_not_matched_by_source_delete(None);
let result = merge
.execute(Box::new(RecordBatchIterator::new([Ok(source)], schema)))
.await
.unwrap();
assert_eq!(result.num_updated_rows, 1);
assert_eq!(result.num_deleted_rows, (ROW_COUNT - 1) as u64);
assert_eq!(table.count_rows(None).await.unwrap(), 1);
}
}
#[cfg(test)]
+7
View File
@@ -104,6 +104,13 @@ pub(crate) async fn set_lsm_write_spec(table: &NativeTable, spec: LsmWriteSpec)
.into(),
});
}
if crate::materialized_view::materialized_view_kind(&dataset.schema().metadata)?.is_some() {
return Err(Error::NotSupported {
message: "an LSM write spec cannot be installed on a materialized view: \
rows in un-compacted tiers are invisible to refresh"
.into(),
});
}
let mut builder = dataset.initialize_mem_wal();
let writer_config_defaults = match spec {
LsmWriteSpec::Bucket {
+2 -2
View File
@@ -193,7 +193,7 @@ fn declared_expression(dataset: &Dataset, column: &str) -> Result<String> {
///
/// 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 {
pub(crate) fn quote_identifier(name: &str) -> String {
format!("`{}`", name.replace('`', "``"))
}
@@ -597,7 +597,7 @@ mod tests {
/// A fragment spanning several scan batches exercises the streamed fill:
/// the probe buffers only until the first gained value and the rest flows
/// through write_column a batch at a time.
/// through write_columns a batch at a time.
#[tokio::test]
async fn test_refresh_streams_a_multi_batch_fragment() {
let values: Vec<i32> = (0..20_000).collect();