mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-08 14:29:03 +00:00
Compare commits
12
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f36d51d385 | ||
|
|
1d93d8ee7e | ||
|
|
e98d8ac685 | ||
|
|
851fa16b47 | ||
|
|
d04ac7ed20 | ||
|
|
a578e9ff7f | ||
|
|
7801e2746a | ||
|
|
5468f3d490 | ||
|
|
c0df2c63b6 | ||
|
|
9e8f1c1a6d | ||
|
|
01679e37fd | ||
|
|
c7cb0b9afa |
Generated
+42
-42
@@ -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
@@ -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
|
||||
|
||||
@@ -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`<[`CherryPickResult`](../interfaces/CherryPickResult.md)>
|
||||
|
||||
***
|
||||
|
||||
### create()
|
||||
|
||||
```ts
|
||||
@@ -112,28 +137,3 @@ List all branches, mapping name to branch metadata.
|
||||
#### Returns
|
||||
|
||||
`Promise`<`Record`<`string`, [`BranchContents`](BranchContents.md)>>
|
||||
|
||||
***
|
||||
|
||||
### 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`<[`MergeBranchResult`](../interfaces/MergeBranchResult.md)>
|
||||
|
||||
@@ -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`<[`MaterializedView`](MaterializedView.md)>
|
||||
|
||||
***
|
||||
|
||||
### 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`<`string`[]>
|
||||
|
||||
***
|
||||
|
||||
### 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`<[`MaterializedView`](MaterializedView.md)>
|
||||
|
||||
***
|
||||
|
||||
### 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`<[`OpenTableOptions`](../interfaces/OpenTableOptions.md)>
|
||||
Additional options
|
||||
|
||||
#### Returns
|
||||
|
||||
|
||||
@@ -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`<[`MaterializedViewDefinition`](../interfaces/MaterializedViewDefinition.md)>
|
||||
|
||||
***
|
||||
|
||||
### 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`<[`RefreshMaterializedViewResult`](../interfaces/RefreshMaterializedViewResult.md)>
|
||||
|
||||
***
|
||||
|
||||
### table()
|
||||
|
||||
```ts
|
||||
table(): Table
|
||||
```
|
||||
|
||||
The view, as the table it is.
|
||||
|
||||
#### Returns
|
||||
|
||||
[`Table`](Table.md)
|
||||
@@ -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)
|
||||
|
||||
@@ -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[];
|
||||
```
|
||||
+6
-6
@@ -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.
|
||||
@@ -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.
|
||||
@@ -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
@@ -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>
|
||||
|
||||
@@ -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"]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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",
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -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
|
||||
|
||||
@@ -2953,7 +2953,7 @@ describe("column name options", () => {
|
||||
.limit(10)
|
||||
.toArray();
|
||||
expect(results2.length).toBe(10);
|
||||
});
|
||||
}, 30_000);
|
||||
});
|
||||
|
||||
describe("when creating an empty table", () => {
|
||||
|
||||
@@ -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[],
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
@@ -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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
@@ -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}"))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
+4
-5
@@ -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"]
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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))
|
||||
@@ -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:
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
@@ -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},
|
||||
]
|
||||
|
||||
@@ -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
@@ -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(¶ms.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
@@ -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>()?;
|
||||
|
||||
+57
-2
@@ -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))
|
||||
|
||||
@@ -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)>,
|
||||
) {
|
||||
|
||||
@@ -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};
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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 },
|
||||
|
||||
|
||||
@@ -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
@@ -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(),
|
||||
@@ -10410,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;
|
||||
@@ -10499,8 +10516,7 @@ mod tests {
|
||||
"changedColumns":[],
|
||||
"addedIndexes":[],
|
||||
"removedIndexes":[],
|
||||
"mergeable":true,
|
||||
"mergeBlockers":[]
|
||||
"errors":[]
|
||||
}"#
|
||||
}
|
||||
|
||||
@@ -10518,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);
|
||||
@@ -10536,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": [] }
|
||||
});
|
||||
@@ -10565,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 }
|
||||
});
|
||||
@@ -10591,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),
|
||||
|
||||
+18
-17
@@ -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::{
|
||||
@@ -832,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`.
|
||||
@@ -1066,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
|
||||
}
|
||||
@@ -2258,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`.
|
||||
|
||||
@@ -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>,
|
||||
}
|
||||
@@ -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();
|
||||
|
||||
@@ -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)]
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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();
|
||||
|
||||
Reference in New Issue
Block a user