mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-09 06:42:30 +00:00
Compare commits
5
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d3077b7641 | ||
|
|
3e3878b223 | ||
|
|
19fb665c76 | ||
|
|
1f95398c34 | ||
|
|
0111a72dc3 |
+1
-1
@@ -1,5 +1,5 @@
|
||||
[tool.bumpversion]
|
||||
current_version = "0.39.0-beta.4"
|
||||
current_version = "0.39.0-beta.6"
|
||||
parse = """(?x)
|
||||
(?P<major>0|[1-9]\\d*)\\.
|
||||
(?P<minor>0|[1-9]\\d*)\\.
|
||||
|
||||
Generated
+45
-45
@@ -3526,8 +3526,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
|
||||
|
||||
[[package]]
|
||||
name = "fsst"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"rand 0.9.5",
|
||||
@@ -4886,8 +4886,8 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
|
||||
|
||||
[[package]]
|
||||
name = "lance"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -4959,8 +4959,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-arrow"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4982,7 +4982,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-scalar"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4996,7 +4996,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-stats"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5005,8 +5005,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-bitpacking"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrayref",
|
||||
"crunchy",
|
||||
@@ -5016,8 +5016,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-core"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5054,8 +5054,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-datafusion"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5085,8 +5085,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-datagen"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5103,8 +5103,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-derive"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
@@ -5113,8 +5113,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-encoding"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5147,8 +5147,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-file"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5179,8 +5179,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-index"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -5244,8 +5244,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-index-core"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5267,8 +5267,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-io"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5308,8 +5308,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-linalg"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5323,8 +5323,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-namespace"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -5338,8 +5338,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-namespace-impls"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-ipc",
|
||||
@@ -5392,8 +5392,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-select"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5407,8 +5407,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-table"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5448,8 +5448,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-testing"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5462,8 +5462,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-tokenizer"
|
||||
version = "12.0.0-beta.14"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.14#a8a101774a1c9647065cc60137094feadbe55296"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"frostem",
|
||||
"icu_segmenter",
|
||||
@@ -5476,7 +5476,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb"
|
||||
version = "0.39.0-beta.3"
|
||||
version = "0.39.0-beta.5"
|
||||
dependencies = [
|
||||
"ahash",
|
||||
"anyhow",
|
||||
@@ -5567,7 +5567,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb-nodejs"
|
||||
version = "0.39.0-beta.3"
|
||||
version = "0.39.0-beta.5"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5592,7 +5592,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb-python"
|
||||
version = "0.39.0-beta.3"
|
||||
version = "0.39.0-beta.5"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
|
||||
+14
-14
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
lance = { "version" = "=12.0.0-beta.14", default-features = false, "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=12.0.0-beta.14", default-features = false, "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=12.0.0-beta.14", default-features = false, "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=12.0.0-beta.14", "tag" = "v12.0.0-beta.14", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "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
|
||||
|
||||
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
|
||||
<dependency>
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-core</artifactId>
|
||||
<version>0.39.0-beta.4</version>
|
||||
<version>0.39.0-beta.6</version>
|
||||
</dependency>
|
||||
```
|
||||
|
||||
|
||||
@@ -8,7 +8,7 @@
|
||||
<parent>
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-parent</artifactId>
|
||||
<version>0.39.0-beta.4</version>
|
||||
<version>0.39.0-beta.6</version>
|
||||
<relativePath>../pom.xml</relativePath>
|
||||
</parent>
|
||||
|
||||
|
||||
+2
-2
@@ -6,7 +6,7 @@
|
||||
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-parent</artifactId>
|
||||
<version>0.39.0-beta.4</version>
|
||||
<version>0.39.0-beta.6</version>
|
||||
<packaging>pom</packaging>
|
||||
<name>${project.artifactId}</name>
|
||||
<description>LanceDB Java SDK Parent POM</description>
|
||||
@@ -28,7 +28,7 @@
|
||||
<properties>
|
||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
||||
<arrow.version>15.0.0</arrow.version>
|
||||
<lance-core.version>12.0.0-beta.14</lance-core.version>
|
||||
<lance-core.version>12.0.0-beta.15</lance-core.version>
|
||||
<spotless.skip>false</spotless.skip>
|
||||
<spotless.version>2.30.0</spotless.version>
|
||||
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
|
||||
|
||||
+1
-1
@@ -1,7 +1,7 @@
|
||||
[package]
|
||||
name = "lancedb-nodejs"
|
||||
edition.workspace = true
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
publish = false
|
||||
license.workspace = true
|
||||
description.workspace = true
|
||||
|
||||
@@ -281,7 +281,7 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
|
||||
numIndices: 0,
|
||||
numRows: 3,
|
||||
// Full on-disk size of the two data files, footers and metadata included.
|
||||
totalBytes: 684,
|
||||
totalBytes: 550,
|
||||
});
|
||||
|
||||
// Index files count toward totalBytes too (only deletion files and
|
||||
@@ -289,7 +289,7 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
|
||||
await table.createIndex("id", { config: Index.btree() });
|
||||
const statsWithIndex = await table.stats();
|
||||
expect(statsWithIndex.numIndices).toBe(1);
|
||||
expect(statsWithIndex.totalBytes).toBeGreaterThan(684);
|
||||
expect(statsWithIndex.totalBytes).toBeGreaterThan(550);
|
||||
});
|
||||
|
||||
it("should overwrite data if asked", async () => {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-darwin-arm64",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["darwin"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.darwin-arm64.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-arm64-gnu",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.linux-arm64-gnu.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-arm64-musl",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.linux-arm64-musl.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-x64-gnu",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.linux-x64-gnu.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-x64-musl",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["linux"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.linux-x64-musl.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-win32-arm64-msvc",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": [
|
||||
"win32"
|
||||
],
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-win32-x64-msvc",
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"os": ["win32"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.win32-x64-msvc.node",
|
||||
|
||||
+1
-1
@@ -11,7 +11,7 @@
|
||||
"ann"
|
||||
],
|
||||
"private": false,
|
||||
"version": "0.39.0-beta.4",
|
||||
"version": "0.39.0-beta.6",
|
||||
"main": "dist/index.js",
|
||||
"exports": {
|
||||
".": "./dist/index.js",
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "lancedb-python"
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
publish = false
|
||||
edition.workspace = true
|
||||
description = "Python bindings for LanceDB"
|
||||
|
||||
@@ -297,7 +297,10 @@ def test_blob_v2_projection_sources_use_typed_column_name():
|
||||
|
||||
|
||||
def _legacy_v1_table(name):
|
||||
db = lancedb.connect("memory:///")
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
db = lancedb.connect(
|
||||
"memory:///", storage_options={"new_table_data_storage_version": "2.1"}
|
||||
)
|
||||
schema = pa.schema(
|
||||
[
|
||||
pa.field("id", pa.int64()),
|
||||
|
||||
@@ -193,7 +193,13 @@ class TestNamespaceConnection:
|
||||
),
|
||||
)
|
||||
|
||||
table = db.create_table("blob_table", data, namespace_path=["test_ns"])
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
table = db.create_table(
|
||||
"blob_table",
|
||||
data,
|
||||
namespace_path=["test_ns"],
|
||||
storage_options={"new_table_data_storage_version": "2.1"},
|
||||
)
|
||||
df = table.to_pandas(blob_mode="lazy").sort_values("id")
|
||||
|
||||
blob = df["blob"].iloc[0]
|
||||
|
||||
@@ -40,6 +40,10 @@ from utils import exception_output
|
||||
from importlib.util import find_spec
|
||||
|
||||
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
LEGACY_BLOB_STORAGE_OPTIONS = {"new_table_data_storage_version": "2.1"}
|
||||
|
||||
|
||||
def _blob_query_data():
|
||||
return pa.table(
|
||||
{
|
||||
@@ -119,13 +123,17 @@ def _assert_blob_bytes_projection(df):
|
||||
|
||||
def _blob_query_table(db, name, blob_schema):
|
||||
if blob_schema == "v1":
|
||||
return db.create_table(name, _blob_query_data())
|
||||
return db.create_table(
|
||||
name, _blob_query_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return _create_blob_v2_query_table(db, name)
|
||||
|
||||
|
||||
async def _blob_query_table_async(db, name, blob_schema):
|
||||
if blob_schema == "v1":
|
||||
return await db.create_table(name, _blob_query_data())
|
||||
return await db.create_table(
|
||||
name, _blob_query_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return await _create_blob_v2_query_table_async(db, name)
|
||||
|
||||
|
||||
@@ -275,7 +283,9 @@ async def test_query_to_pandas_kwargs(table, table_async):
|
||||
def test_plain_scan_query_to_pandas_blob_modes(tmp_db, blob_mode):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
f"test_query_to_pandas_blob_{blob_mode}", _blob_query_data()
|
||||
f"test_query_to_pandas_blob_{blob_mode}",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
df = (
|
||||
@@ -322,7 +332,9 @@ def test_plain_scan_query_to_pandas_blob_mode_does_not_collect_arrow(
|
||||
):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
"test_query_to_pandas_blob_no_arrow_collect", _blob_query_data()
|
||||
"test_query_to_pandas_blob_no_arrow_collect",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
query = table.search().where("id = 1").select(["id", "blob"])
|
||||
|
||||
@@ -347,7 +359,9 @@ def test_plain_scan_query_to_pandas_blob_descriptions_flatten_uses_scanner(
|
||||
):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
"test_query_to_pandas_blob_desc_flatten", _blob_query_data()
|
||||
"test_query_to_pandas_blob_desc_flatten",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
query = table.search().where("id = 1").select(["id", "blob"])
|
||||
|
||||
@@ -365,7 +379,11 @@ def test_plain_scan_query_to_pandas_blob_descriptions_flatten_uses_scanner(
|
||||
def test_plain_scan_query_to_pandas_scanner_state(tmp_db):
|
||||
pytest.importorskip("lance")
|
||||
data = _blob_query_data()
|
||||
table = tmp_db.create_table("test_query_to_pandas_scanner_state", data.slice(0, 2))
|
||||
table = tmp_db.create_table(
|
||||
"test_query_to_pandas_scanner_state",
|
||||
data.slice(0, 2),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
table.add(data.slice(2, 2))
|
||||
|
||||
fragments = table.to_lance().get_fragments()
|
||||
@@ -400,7 +418,9 @@ def test_plain_scan_query_to_pandas_scanner_state(tmp_db):
|
||||
async def test_async_plain_scan_query_to_pandas_blob_projection(tmp_db_async):
|
||||
pytest.importorskip("lance")
|
||||
table = await tmp_db_async.create_table(
|
||||
"test_async_query_to_pandas_blob_projection", _blob_query_data()
|
||||
"test_async_query_to_pandas_blob_projection",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
lazy_df = await (
|
||||
@@ -452,7 +472,9 @@ async def test_async_plain_scan_query_to_pandas_blob_mode_does_not_collect_arrow
|
||||
):
|
||||
pytest.importorskip("lance")
|
||||
table = await tmp_db_async.create_table(
|
||||
"test_async_query_to_pandas_blob_no_arrow_collect", _blob_query_data()
|
||||
"test_async_query_to_pandas_blob_no_arrow_collect",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
query = table.query().where("id = 1").select(["id", "blob"])
|
||||
|
||||
@@ -474,7 +496,11 @@ async def test_async_plain_scan_query_to_pandas_blob_mode_does_not_collect_arrow
|
||||
|
||||
def test_vector_query_to_pandas_blob_mode_requires_native_path(tmp_db):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table("test_vector_query_blob_mode", _blob_query_data())
|
||||
table = tmp_db.create_table(
|
||||
"test_vector_query_blob_mode",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
with pytest.raises(RuntimeError, match="Lance native pandas conversion"):
|
||||
table.search([1.0, 0.0]).select(["blob", "vector"]).limit(1).to_pandas(
|
||||
@@ -485,7 +511,9 @@ def test_vector_query_to_pandas_blob_mode_requires_native_path(tmp_db):
|
||||
def test_vector_query_to_pandas_blob_descriptions_requires_plain_scan(tmp_db):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(
|
||||
"test_vector_query_blob_descriptions", _blob_query_data()
|
||||
"test_vector_query_blob_descriptions",
|
||||
_blob_query_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
with pytest.raises(RuntimeError, match="plain scan query"):
|
||||
|
||||
@@ -64,15 +64,23 @@ async def _blob_v2_table_async(db: AsyncConnection, name: str):
|
||||
return table
|
||||
|
||||
|
||||
# Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
LEGACY_BLOB_STORAGE_OPTIONS = {"new_table_data_storage_version": "2.1"}
|
||||
|
||||
|
||||
def _blob_table(db: DBConnection, name: str, blob_schema: str):
|
||||
if blob_schema == "v1":
|
||||
return db.create_table(name, data=_blob_test_data())
|
||||
return db.create_table(
|
||||
name, data=_blob_test_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return _blob_v2_table(db, name)
|
||||
|
||||
|
||||
async def _blob_table_async(db: AsyncConnection, name: str, blob_schema: str):
|
||||
if blob_schema == "v1":
|
||||
return await db.create_table(name, data=_blob_test_data())
|
||||
return await db.create_table(
|
||||
name, data=_blob_test_data(), storage_options=LEGACY_BLOB_STORAGE_OPTIONS
|
||||
)
|
||||
return await _blob_v2_table_async(db, name)
|
||||
|
||||
|
||||
@@ -147,7 +155,11 @@ def test_table_to_pandas_invalid_blob_mode_non_blob_table(tmp_db: DBConnection):
|
||||
@pytest.mark.parametrize("blob_mode", ["lazy", "bytes", "descriptions"])
|
||||
def test_table_to_pandas_blob_modes(tmp_db: DBConnection, blob_mode):
|
||||
pytest.importorskip("lance")
|
||||
table = tmp_db.create_table(f"test_to_pandas_blob_{blob_mode}", _blob_test_data())
|
||||
table = tmp_db.create_table(
|
||||
f"test_to_pandas_blob_{blob_mode}",
|
||||
_blob_test_data(),
|
||||
storage_options=LEGACY_BLOB_STORAGE_OPTIONS,
|
||||
)
|
||||
|
||||
df = table.to_pandas(blob_mode=blob_mode)
|
||||
|
||||
@@ -3959,7 +3971,7 @@ def test_stats(mem_db: DBConnection):
|
||||
print(f"{stats=}")
|
||||
assert stats == {
|
||||
# Full on-disk size of the data file, footer and metadata included.
|
||||
"total_bytes": 633,
|
||||
"total_bytes": 637,
|
||||
"num_rows": 2,
|
||||
"num_indices": 0,
|
||||
"fragment_stats": {
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "lancedb"
|
||||
version = "0.39.0-beta.4"
|
||||
version = "0.39.0-beta.6"
|
||||
edition.workspace = true
|
||||
description = "LanceDB: A serverless, low-latency vector database for AI applications"
|
||||
license.workspace = true
|
||||
|
||||
@@ -532,10 +532,11 @@ mod tests {
|
||||
fn storage_version_bumps_to_v2_2() {
|
||||
let mut params = WriteParams::default();
|
||||
ensure_blob_storage_version(&blob_schema(), &mut params);
|
||||
assert_eq!(
|
||||
params.data_storage_version.unwrap().resolve(),
|
||||
ConcreteFileVersion::V2_2
|
||||
);
|
||||
let resolved = params
|
||||
.data_storage_version
|
||||
.unwrap_or(LanceFileVersion::Stable)
|
||||
.resolve();
|
||||
assert_eq!(resolved, ConcreteFileVersion::V2_2);
|
||||
assert!(!params.enable_stable_row_ids);
|
||||
}
|
||||
|
||||
@@ -547,10 +548,11 @@ mod tests {
|
||||
};
|
||||
ensure_blob_storage_version(&blob_schema(), &mut params);
|
||||
assert!(params.enable_stable_row_ids);
|
||||
assert_eq!(
|
||||
params.data_storage_version.unwrap().resolve(),
|
||||
ConcreteFileVersion::V2_2
|
||||
);
|
||||
let resolved = params
|
||||
.data_storage_version
|
||||
.unwrap_or(LanceFileVersion::Stable)
|
||||
.resolve();
|
||||
assert_eq!(resolved, ConcreteFileVersion::V2_2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
|
||||
//! Namespace-based database implementation that delegates table management to lance-namespace
|
||||
|
||||
use lance_datafusion::utils::StreamingWriteSource;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::sync::{Arc, Mutex};
|
||||
|
||||
@@ -304,6 +305,10 @@ impl Database for LanceNamespaceDatabase {
|
||||
}
|
||||
|
||||
async fn create_table(&self, request: DbCreateTableRequest) -> Result<Arc<dyn BaseTable>> {
|
||||
// Refuse a bad declaration before the namespace records a table.
|
||||
crate::table::computed_columns::ensure_declarations_are_planned(
|
||||
&request.data.arrow_schema(),
|
||||
)?;
|
||||
let mut table_id = request.namespace_path.clone();
|
||||
table_id.push(request.name.clone());
|
||||
let mut existing_table = None;
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -24,8 +24,8 @@ use std::time::{SystemTime, UNIX_EPOCH};
|
||||
|
||||
use arrow_array::cast::AsArray;
|
||||
use arrow_array::types::UInt64Type;
|
||||
use arrow_array::{RecordBatch, UInt64Array};
|
||||
use arrow_schema::{Schema as ArrowSchema, SchemaRef};
|
||||
use arrow_array::{RecordBatch, UInt64Array, new_null_array};
|
||||
use arrow_schema::{FieldRef, Schema as ArrowSchema, SchemaRef};
|
||||
use datafusion::common::ScalarValue;
|
||||
use datafusion::error::DataFusionError;
|
||||
use datafusion::physical_plan::SendableRecordBatchStream;
|
||||
@@ -34,7 +34,7 @@ use datafusion::prelude::{col, lit};
|
||||
use futures::{StreamExt, TryStreamExt};
|
||||
use lance::Dataset;
|
||||
use lance::dataset::mem_wal::DatasetMemWalExt;
|
||||
use lance::dataset::transaction::{Operation, Transaction};
|
||||
use lance::dataset::transaction::{Operation, Transaction, UpdateMode};
|
||||
use lance::dataset::write::delete::DeleteBuilder;
|
||||
use lance::dataset::write::merge_insert::inserted_rows::{
|
||||
KeyExistenceFilter, KeyExistenceFilterBuilder, KeyValue,
|
||||
@@ -51,6 +51,9 @@ use super::{
|
||||
definition_to_metadata,
|
||||
};
|
||||
use crate::database::OpenTableRequest;
|
||||
use crate::table::computed_columns::{
|
||||
computed_column_from_field, computed_columns, ensure_declarations_are_planned,
|
||||
};
|
||||
use crate::table::{NativeTable, NativeTableExt, Table};
|
||||
use crate::{Error, Result};
|
||||
|
||||
@@ -167,30 +170,52 @@ pub(crate) async fn execute_refresh(
|
||||
.map(|p| (p.output.clone(), p.expression.clone()))
|
||||
.collect();
|
||||
validate_inputs(&source_ds, definition)?;
|
||||
let (replanned, mut planned_fields, _renames) = super::plan(
|
||||
let (replanned, planned_fields, _renames) = super::plan(
|
||||
source_schema,
|
||||
&definition.source_table,
|
||||
&definition.source_namespace,
|
||||
&projections,
|
||||
Some(&projections),
|
||||
definition.filter.as_deref(),
|
||||
definition.limit,
|
||||
)?;
|
||||
let mut planned_fields = planned_fields;
|
||||
planned_fields.push(arrow_schema::Field::new(
|
||||
SOURCE_ROW_ID_COLUMN,
|
||||
arrow_schema::DataType::UInt64,
|
||||
false,
|
||||
));
|
||||
// A computed column is not planned from the source: refresh writes it
|
||||
// NULL and its declaration's owner fills it. Its declaration must still
|
||||
// be complete, and it must be able to hold NULL.
|
||||
let physical = ArrowSchema::from(view_ds.schema());
|
||||
let planned_shape: Vec<_> = planned_fields
|
||||
.iter()
|
||||
.map(|f| (f.name().clone(), f.data_type().clone(), f.is_nullable()))
|
||||
.collect();
|
||||
let physical_shape: Vec<_> = physical
|
||||
let mut computed = computed_columns(&physical).into_iter().map(|c| c.name);
|
||||
if let Some(name) = computed.by_ref().find(|name| {
|
||||
physical
|
||||
.field_with_name(name)
|
||||
.is_ok_and(|f| !f.is_nullable())
|
||||
}) {
|
||||
return Err(Error::Schema {
|
||||
message: format!(
|
||||
"computed column '{name}' of view '{}' cannot hold NULL; recreate the view",
|
||||
view.name()
|
||||
),
|
||||
});
|
||||
}
|
||||
ensure_declarations_are_planned(&physical)?;
|
||||
let physical_planned: Vec<&FieldRef> = physical
|
||||
.fields()
|
||||
.iter()
|
||||
.map(|f| (f.name().clone(), f.data_type().clone(), f.is_nullable()))
|
||||
.filter(|f| computed_column_from_field(f).is_none())
|
||||
.collect();
|
||||
if planned_shape != physical_shape {
|
||||
// A projected column that became nullable at the source still fits the
|
||||
// view's nullable field; the reverse would not.
|
||||
let matches = planned_fields.len() == physical_planned.len()
|
||||
&& planned_fields.iter().zip(&physical_planned).all(|(e, p)| {
|
||||
e.name() == p.name()
|
||||
&& e.data_type() == p.data_type()
|
||||
&& (p.is_nullable() || !e.is_nullable())
|
||||
});
|
||||
if !matches {
|
||||
return Err(Error::Schema {
|
||||
message: format!(
|
||||
"the stored definition of view '{}' does not produce this \
|
||||
@@ -229,11 +254,18 @@ pub(crate) async fn execute_refresh(
|
||||
.get(SOURCE_VERSION_TS_META_KEY)
|
||||
.and_then(|raw| raw.parse().ok());
|
||||
// The watermark speaks only for the view state its refresh left behind;
|
||||
// any other commit on the view since then is drift.
|
||||
let view_intact = metadata
|
||||
// any other commit on the view since then is drift, except a fill of its
|
||||
// computed columns, which rewrites nothing refresh certifies.
|
||||
let recorded_view_version = metadata
|
||||
.get(VIEW_VERSION_META_KEY)
|
||||
.and_then(|raw| raw.parse::<u64>().ok())
|
||||
== Some(view_ds.version().version);
|
||||
.and_then(|raw| raw.parse::<u64>().ok());
|
||||
let view_intact = match recorded_view_version {
|
||||
Some(recorded) if recorded == view_ds.version().version => true,
|
||||
Some(recorded) if recorded < view_ds.version().version => {
|
||||
only_computed_rewrites_since(&view_ds, recorded).await?
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
|
||||
if !full && watermark == Some(source_version) && view_intact && recorded_ts == Some(source_ts) {
|
||||
return Ok(RefreshMaterializedViewResult {
|
||||
@@ -1090,6 +1122,69 @@ struct RowScope {
|
||||
limit: Option<u64>,
|
||||
}
|
||||
|
||||
/// Whether every commit on the view after `recorded` is a fill of its
|
||||
/// computed columns: a column rewrite or data replacement touching only
|
||||
/// those fields and neither adding nor removing rows. A version whose
|
||||
/// transaction cannot be read is not proven, so it counts as drift.
|
||||
async fn only_computed_rewrites_since(view_ds: &Dataset, recorded: u64) -> Result<bool> {
|
||||
// A fill may write any field under a computed column, so the whole
|
||||
// subtree counts, not only the root.
|
||||
let physical = ArrowSchema::from(view_ds.schema());
|
||||
fn subtree(field: &lance_core::datatypes::Field, ids: &mut Vec<u32>) {
|
||||
ids.push(field.id as u32);
|
||||
for child in &field.children {
|
||||
subtree(child, ids);
|
||||
}
|
||||
}
|
||||
let mut computed_fields = Vec::new();
|
||||
for column in computed_columns(&physical) {
|
||||
if let Some(field) = view_ds.schema().field(&column.name) {
|
||||
subtree(field, &mut computed_fields);
|
||||
}
|
||||
}
|
||||
if computed_fields.is_empty() {
|
||||
return Ok(false);
|
||||
}
|
||||
for version in recorded + 1..=view_ds.version().version {
|
||||
let Some(transaction) = view_ds.read_transaction_by_version(version).await? else {
|
||||
return Ok(false);
|
||||
};
|
||||
let fill = match &transaction.operation {
|
||||
Operation::Update {
|
||||
removed_fragment_ids,
|
||||
new_fragments,
|
||||
fields_modified,
|
||||
update_mode: Some(UpdateMode::RewriteColumns),
|
||||
..
|
||||
} => {
|
||||
removed_fragment_ids.is_empty()
|
||||
&& new_fragments.is_empty()
|
||||
&& !fields_modified.is_empty()
|
||||
&& fields_modified
|
||||
.iter()
|
||||
.all(|field| computed_fields.contains(field))
|
||||
}
|
||||
// What `refresh_column` commits for a SQL declaration.
|
||||
Operation::DataReplacement { replacements } => {
|
||||
!replacements.is_empty()
|
||||
&& replacements.iter().all(|group| {
|
||||
!group.1.fields.is_empty()
|
||||
&& group
|
||||
.1
|
||||
.fields
|
||||
.iter()
|
||||
.all(|field| computed_fields.contains(&(*field as u32)))
|
||||
})
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
if !fill {
|
||||
return Ok(false);
|
||||
}
|
||||
}
|
||||
Ok(true)
|
||||
}
|
||||
|
||||
async fn compute_stream(
|
||||
source: &Dataset,
|
||||
definition: &MaterializedViewDefinition,
|
||||
@@ -1158,6 +1253,10 @@ async fn compute_stream(
|
||||
let batch = batch.map_err(|e| DataFusionError::External(Box::new(e)))?;
|
||||
let mut columns = Vec::with_capacity(out_schema.fields().len());
|
||||
for field in out_schema.fields() {
|
||||
if computed_column_from_field(field).is_some() {
|
||||
columns.push(new_null_array(field.data_type(), batch.num_rows()));
|
||||
continue;
|
||||
}
|
||||
let name = if field.name() == SOURCE_ROW_ID_COLUMN {
|
||||
ROW_ID
|
||||
} else {
|
||||
@@ -2768,7 +2867,7 @@ mod tests {
|
||||
let (conn, source) = db_with_source(vec![1]).await;
|
||||
let prepared = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
&[("x".into(), "x".into()), ("twice".into(), "x * 2".into())],
|
||||
Some(&[("x".into(), "x".into()), ("twice".into(), "x * 2".into())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
@@ -3132,4 +3231,424 @@ mod tests {
|
||||
let err = view.refresh().execute().await.unwrap_err();
|
||||
assert!(err.to_string().contains("source table 'src'"), "{err}");
|
||||
}
|
||||
|
||||
/// A view with a computed column, declared over `people` and refreshed.
|
||||
async fn refreshed_computed_view(conn: &Connection) -> MaterializedView {
|
||||
use crate::materialized_view::tests::{computed_field, people, test_binding};
|
||||
let source = people(conn).await;
|
||||
let view = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[
|
||||
("id".to_string(), "id".to_string()),
|
||||
("name".to_string(), "name".to_string()),
|
||||
]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.with_computed_columns(
|
||||
vec![(2, computed_field("emb", "fb_1", "name"))],
|
||||
&[test_binding("fb_1", "name", "emb")],
|
||||
)
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
let result = view.refresh().execute().await.unwrap();
|
||||
assert_eq!(result.mode, RefreshMode::Rebuild);
|
||||
view
|
||||
}
|
||||
|
||||
async fn unfilled(view: &MaterializedView) -> usize {
|
||||
view.table()
|
||||
.count_rows(Some("emb IS NULL".to_string()))
|
||||
.await
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
async fn append_people(conn: &Connection, ids: Vec<i32>, names: Vec<&str>) {
|
||||
let batch = record_batch!(("id", Int32, ids), ("name", Utf8, names)).unwrap();
|
||||
conn.open_table("people")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap()
|
||||
.add(batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// Commit the fill job's shape on the view: a column rewrite of
|
||||
/// `fields`, touching no rows. The data is left as it is; what matters
|
||||
/// here is how the next refresh classifies the commit.
|
||||
async fn commit_column_rewrite(view: &MaterializedView, fields: &[&str]) {
|
||||
let native = view.table().as_native().unwrap();
|
||||
native.dataset.reload().await.unwrap();
|
||||
let dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
let fields_modified = fields
|
||||
.iter()
|
||||
.map(|name| dataset.schema().field(name).unwrap().id as u32)
|
||||
.collect();
|
||||
let updated_fragments = dataset
|
||||
.get_fragments()
|
||||
.iter()
|
||||
.map(|fragment| fragment.metadata().clone())
|
||||
.collect();
|
||||
let operation = Operation::Update {
|
||||
removed_fragment_ids: Vec::new(),
|
||||
updated_fragments,
|
||||
new_fragments: Vec::new(),
|
||||
fields_modified,
|
||||
compacted_sstables: Vec::new(),
|
||||
fields_for_preserving_frag_bitmap: Vec::new(),
|
||||
update_mode: Some(UpdateMode::RewriteColumns),
|
||||
inserted_rows_filter: None,
|
||||
updated_fragment_offsets: None,
|
||||
};
|
||||
let read_version = dataset.version().version;
|
||||
CommitBuilder::new(WriteDestination::Dataset(Arc::new(dataset)))
|
||||
.execute(Transaction::new(read_version, operation, None))
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// Refresh never computes a computed column: every row it writes, on a
|
||||
/// rebuild, an append and a rewrite, carries NULL there, and the
|
||||
/// declaration survives all three.
|
||||
#[tokio::test]
|
||||
async fn test_computed_columns_are_written_null_and_kept() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
assert_eq!(unfilled(&view).await, 3);
|
||||
|
||||
append_people(&conn, vec![4, 5], vec!["d", "e"]).await;
|
||||
let result = view.refresh().execute().await.unwrap();
|
||||
assert_eq!(result.mode, RefreshMode::Incremental);
|
||||
assert_eq!(unfilled(&view).await, 5);
|
||||
|
||||
conn.open_table("people")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap()
|
||||
.update()
|
||||
.column("name", "'z'")
|
||||
.only_if("id = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
view.refresh().execute().await.unwrap();
|
||||
assert_eq!(unfilled(&view).await, 5);
|
||||
assert_eq!(read(view.table(), "id").await, vec![1, 2, 3, 4, 5]);
|
||||
|
||||
let schema = view.table().schema().await.unwrap();
|
||||
assert!(
|
||||
crate::table::computed_columns::function_bindings(&schema)
|
||||
.unwrap()
|
||||
.iter()
|
||||
.any(|b| b.binding_id() == "fb_1"),
|
||||
"the binding envelope was lost"
|
||||
);
|
||||
assert!(
|
||||
computed_column_from_field(schema.field_with_name("emb").unwrap()).is_some(),
|
||||
"the declaration was lost"
|
||||
);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
}
|
||||
|
||||
/// The fill job's commit rewrites only computed columns. It is the one
|
||||
/// commit on a view that is not drift: the next refresh carries on from
|
||||
/// its watermark instead of rebuilding, which would null what the fill
|
||||
/// just wrote.
|
||||
#[tokio::test]
|
||||
async fn test_a_computed_column_fill_is_not_drift() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
|
||||
commit_column_rewrite(&view, &["emb"]).await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
|
||||
commit_column_rewrite(&view, &["emb"]).await;
|
||||
append_people(&conn, vec![4], vec!["d"]).await;
|
||||
let result = view.refresh().execute().await.unwrap();
|
||||
assert_eq!(result.mode, RefreshMode::Incremental);
|
||||
assert_eq!(result.rows_written, 1);
|
||||
assert_eq!(read(view.table(), "id").await, vec![1, 2, 3, 4]);
|
||||
}
|
||||
|
||||
/// A column rewrite that reaches a projected column is drift like any
|
||||
/// other write: refresh certifies those columns and must recompute them.
|
||||
#[tokio::test]
|
||||
async fn test_a_rewrite_of_a_projected_column_is_drift() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
|
||||
commit_column_rewrite(&view, &["emb", "name"]).await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
}
|
||||
|
||||
/// The declaration contract is checked before any refresh mutation: a
|
||||
/// missing binding envelope and a column that lost its declaration both
|
||||
/// fail closed.
|
||||
#[tokio::test]
|
||||
async fn test_a_broken_declaration_is_refused_before_refresh() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
let native = view.table().as_native().unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.update_schema_metadata(vec![(
|
||||
crate::table::computed_columns::FUNCTION_BINDINGS_META_KEY.to_string(),
|
||||
None,
|
||||
)])
|
||||
.await
|
||||
.unwrap();
|
||||
let err = view.refresh().execute().await.unwrap_err().to_string();
|
||||
assert!(err.contains("references missing binding 'fb_1'"), "{err}");
|
||||
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
let native = view.table().as_native().unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.replace_field_metadata(vec![(
|
||||
dataset.schema().field("emb").unwrap().id as u32,
|
||||
HashMap::new(),
|
||||
)])
|
||||
.await
|
||||
.unwrap();
|
||||
let err = view.refresh().execute().await.unwrap_err().to_string();
|
||||
assert!(err.contains("does not match binding 'fb_1'"), "{err}");
|
||||
}
|
||||
|
||||
/// An input the view does not project is materialized on every refresh
|
||||
/// path, before the provenance column, with the source's values.
|
||||
#[tokio::test]
|
||||
async fn test_internal_inputs_are_materialized_and_refreshed() {
|
||||
use crate::materialized_view::tests::{computed_field, strict_people, test_binding};
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let source = strict_people(&conn).await;
|
||||
let mut prepared = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[("id".to_string(), "id".to_string())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
let input = prepared.input_column("name").unwrap();
|
||||
let view = prepared
|
||||
.with_computed_columns(
|
||||
vec![(1, computed_field("emb", "fb_1", &input))],
|
||||
&[test_binding("fb_1", &input, "emb")],
|
||||
)
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
let names: Vec<String> = view
|
||||
.table()
|
||||
.schema()
|
||||
.await
|
||||
.unwrap()
|
||||
.fields()
|
||||
.iter()
|
||||
.map(|f| f.name().clone())
|
||||
.collect();
|
||||
assert_eq!(names, ["id", "emb", "__input_name", SOURCE_ROW_ID_COLUMN]);
|
||||
|
||||
let unfilled_inputs = || async {
|
||||
view.table()
|
||||
.count_rows(Some("__input_name IS NULL".to_string()))
|
||||
.await
|
||||
.unwrap()
|
||||
};
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
assert_eq!(view.table().count_rows(None).await.unwrap(), 3);
|
||||
assert_eq!(unfilled_inputs().await, 0);
|
||||
|
||||
let more = arrow_array::RecordBatch::try_new(
|
||||
source.schema().await.unwrap(),
|
||||
vec![
|
||||
Arc::new(Int32Array::from(vec![4])),
|
||||
Arc::new(arrow_array::StringArray::from(vec!["d"])),
|
||||
],
|
||||
)
|
||||
.unwrap();
|
||||
source.add(more).execute().await.unwrap();
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Incremental
|
||||
);
|
||||
assert_eq!(unfilled_inputs().await, 0);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.count_rows(Some("__input_name = 'd'".to_string()))
|
||||
.await
|
||||
.unwrap(),
|
||||
1
|
||||
);
|
||||
|
||||
source
|
||||
.update()
|
||||
.column("name", "'z'")
|
||||
.only_if("id = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
view.refresh().execute().await.unwrap();
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.count_rows(Some("__input_name = 'z'".to_string()))
|
||||
.await
|
||||
.unwrap(),
|
||||
1
|
||||
);
|
||||
assert_eq!(
|
||||
unfilled(&view).await,
|
||||
4,
|
||||
"rewritten and new rows are unfilled"
|
||||
);
|
||||
}
|
||||
|
||||
/// A SQL declaration is filled by `refresh_column` on the view, which
|
||||
/// commits a data replacement; the next refresh continues from its
|
||||
/// watermark and keeps what the fill wrote, and only rows the view added
|
||||
/// since come back unfilled.
|
||||
#[tokio::test]
|
||||
async fn test_a_sql_fill_is_not_drift() {
|
||||
use crate::materialized_view::tests::{people, sql_field};
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let source = people(&conn).await;
|
||||
let view = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[("id".to_string(), "id".to_string())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.with_computed_columns(
|
||||
vec![(
|
||||
1,
|
||||
sql_field("next", arrow_schema::DataType::Int32, "id + 1", r#"["id"]"#),
|
||||
)],
|
||||
&[],
|
||||
)
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
let filled = || async {
|
||||
view.table()
|
||||
.count_rows(Some("next = id + 1".to_string()))
|
||||
.await
|
||||
.unwrap()
|
||||
};
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.refresh_column("next")
|
||||
.await
|
||||
.unwrap()
|
||||
.rows_filled,
|
||||
3
|
||||
);
|
||||
assert_eq!(filled().await, 3);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
assert_eq!(filled().await, 3);
|
||||
|
||||
append_people(&conn, vec![4], vec!["d"]).await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Incremental
|
||||
);
|
||||
assert_eq!(filled().await, 3);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.refresh_column("next")
|
||||
.await
|
||||
.unwrap()
|
||||
.rows_filled,
|
||||
1
|
||||
);
|
||||
assert_eq!(filled().await, 4);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
}
|
||||
|
||||
/// A fill of a nested computed column writes its child fields; that is
|
||||
/// still a fill, not drift.
|
||||
#[tokio::test]
|
||||
async fn test_a_nested_sql_fill_is_not_drift() {
|
||||
use crate::materialized_view::tests::{people, sql_field};
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let source = people(&conn).await;
|
||||
let payload = sql_field(
|
||||
"payload",
|
||||
arrow_schema::DataType::Struct(
|
||||
vec![arrow_schema::Field::new(
|
||||
"value",
|
||||
arrow_schema::DataType::Utf8,
|
||||
true,
|
||||
)]
|
||||
.into(),
|
||||
),
|
||||
"named_struct('value', name)",
|
||||
r#"["name"]"#,
|
||||
);
|
||||
let view = crate::materialized_view::prepare_declaration(
|
||||
&source,
|
||||
Some(&[("name".to_string(), "name".to_string())]),
|
||||
None,
|
||||
None,
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.with_computed_columns(vec![(1, payload)], &[])
|
||||
.unwrap()
|
||||
.create("v")
|
||||
.await
|
||||
.unwrap();
|
||||
view.refresh().execute().await.unwrap();
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.refresh_column("payload")
|
||||
.await
|
||||
.unwrap()
|
||||
.rows_filled,
|
||||
3
|
||||
);
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
assert_eq!(
|
||||
view.table()
|
||||
.count_rows(Some("payload.value = name".to_string()))
|
||||
.await
|
||||
.unwrap(),
|
||||
3
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -2804,7 +2804,7 @@ impl NativeTable {
|
||||
namespace_client: Option<Arc<dyn LanceNamespace>>,
|
||||
pushdown_operations: HashSet<NamespaceClientPushdownOperation>,
|
||||
) -> Result<Self> {
|
||||
computed_columns::ensure_no_foreign_declarations(batches.arrow_schema().fields())?;
|
||||
let batches = computed_columns::admit_create_source(batches)?;
|
||||
// Default params uses format v1.
|
||||
let params = params.unwrap_or(WriteParams {
|
||||
..Default::default()
|
||||
@@ -2904,6 +2904,7 @@ impl NativeTable {
|
||||
pushdown_operations: HashSet<NamespaceClientPushdownOperation>,
|
||||
session: Option<Arc<lance::session::Session>>,
|
||||
) -> Result<Self> {
|
||||
let batches = computed_columns::admit_create_source(batches)?;
|
||||
// Build table_id from namespace + name for the storage options provider
|
||||
let mut table_id = namespace.clone();
|
||||
table_id.push(name.to_string());
|
||||
@@ -5677,7 +5678,7 @@ mod tests {
|
||||
TableStatistics {
|
||||
num_rows: 250,
|
||||
num_indices: 0,
|
||||
total_bytes: 8925,
|
||||
total_bytes: 8969,
|
||||
fragment_stats: FragmentStatistics {
|
||||
num_fragments: 11,
|
||||
num_small_fragments: 11,
|
||||
|
||||
@@ -21,6 +21,7 @@
|
||||
//! [`computed_columns`] and [`computed_column_from_field`] read declarations
|
||||
//! back off a schema.
|
||||
|
||||
use futures::StreamExt;
|
||||
use std::collections::{BTreeSet, HashMap, HashSet};
|
||||
use std::sync::Arc;
|
||||
|
||||
@@ -1338,6 +1339,106 @@ pub(crate) fn ensure_batch_writes_no_computed_values(
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Validate every computed-column declaration `schema` carries against the
|
||||
/// schema itself: every field with declaration metadata is a complete
|
||||
/// declaration, a SQL declaration re-plans to the field it declares, a
|
||||
/// Function declaration satisfies the binding contract, and no declaration
|
||||
/// reads another computed column. What passes here is what `refresh_column`
|
||||
/// can execute.
|
||||
pub(crate) fn ensure_declarations_are_planned(schema: &ArrowSchema) -> Result<()> {
|
||||
let invalid = |message: String| Error::InvalidInput { message };
|
||||
// A field with any declaration key is a declaration; a partial one is
|
||||
// not "no declaration", it is a broken one.
|
||||
for field in schema.fields() {
|
||||
if field.metadata().keys().any(|k| is_declaration_key(k))
|
||||
&& computed_column_from_field(field).is_none()
|
||||
{
|
||||
return Err(invalid(format!(
|
||||
"field '{}' carries an incomplete computed-column declaration",
|
||||
field.name()
|
||||
)));
|
||||
}
|
||||
}
|
||||
let declared: HashSet<String> = computed_columns(schema)
|
||||
.into_iter()
|
||||
.map(|c| c.name)
|
||||
.collect();
|
||||
for column in computed_columns(schema) {
|
||||
let field = schema.field_with_name(&column.name)?;
|
||||
if !field.is_nullable() {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' must be nullable until a refresh fills it",
|
||||
column.name
|
||||
)));
|
||||
}
|
||||
match &column.kind {
|
||||
ComputedColumnKind::Sql { expression } => {
|
||||
let others: Vec<ArrowField> = schema
|
||||
.fields()
|
||||
.iter()
|
||||
.filter(|f| f.name() != &column.name)
|
||||
.map(|f| f.as_ref().clone())
|
||||
.collect();
|
||||
let bound = bind(Arc::new(ArrowSchema::new(others)), &column.name, expression)?;
|
||||
if let Some(input) = bound.roots.iter().find(|r| declared.contains(*r)) {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' reads computed column '{input}'",
|
||||
column.name
|
||||
)));
|
||||
}
|
||||
if &bound.data_type != field.data_type() {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' is declared as {} but its expression yields {}",
|
||||
column.name,
|
||||
field.data_type(),
|
||||
bound.data_type
|
||||
)));
|
||||
}
|
||||
let mut declared_inputs = column.inputs.clone();
|
||||
declared_inputs.sort();
|
||||
if declared_inputs != bound.inputs {
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' declares inputs {:?} but its expression reads {:?}",
|
||||
column.name, declared_inputs, bound.inputs
|
||||
)));
|
||||
}
|
||||
}
|
||||
ComputedColumnKind::Function { binding_id, .. } => {
|
||||
// The binding validator resolves each input's leaf; the
|
||||
// no-computed-input rule is about the root it hangs from.
|
||||
let bindings = function_bindings(schema)?;
|
||||
let Some(binding) = bindings.iter().find(|b| b.binding_id() == binding_id) else {
|
||||
continue; // reported by the binding validator below
|
||||
};
|
||||
// Roots come from the canonical path parser: a quoted
|
||||
// top-level name may itself contain a dot.
|
||||
if let Some(input) = binding
|
||||
.inputs()
|
||||
.iter()
|
||||
.filter_map(|input| resolve_field_path(schema, &input.field_path).ok())
|
||||
.map(|resolved| resolved.root.name().as_str())
|
||||
.find(|r| declared.contains(*r))
|
||||
{
|
||||
return Err(invalid(format!(
|
||||
"computed column '{}' reads computed column '{input}'",
|
||||
column.name
|
||||
)));
|
||||
}
|
||||
}
|
||||
ComputedColumnKind::Unrecognized { kind } => {
|
||||
return Err(Error::NotSupported {
|
||||
message: format!(
|
||||
"computed column '{}' is defined by '{kind}', which this version \
|
||||
of lancedb cannot fill",
|
||||
column.name
|
||||
),
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
ensure_supported_function_metadata(schema)
|
||||
}
|
||||
|
||||
/// Reject fields carrying declaration metadata that did not come through
|
||||
/// [`plan`]. One authority for creation, overwrite and raw transforms.
|
||||
pub(crate) fn ensure_no_foreign_declarations<'a>(
|
||||
@@ -1796,6 +1897,54 @@ pub(super) async fn add_foreign_kind(table: &crate::Table, name: &str, kind: &st
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// Admit a table's initial data: every declaration it carries is validated,
|
||||
/// and the stream refuses any batch with values in a computed column, whose
|
||||
/// values come from refresh alone. One boundary for every way a table is
|
||||
/// created.
|
||||
pub(crate) fn admit_create_source<S: lance_datafusion::utils::StreamingWriteSource>(
|
||||
batches: S,
|
||||
) -> Result<UnfilledDeclarations<S>> {
|
||||
let schema = batches.arrow_schema();
|
||||
ensure_declarations_are_planned(&schema)?;
|
||||
let declared = computed_columns(&schema)
|
||||
.into_iter()
|
||||
.map(|c| c.name)
|
||||
.collect();
|
||||
Ok(UnfilledDeclarations {
|
||||
inner: batches,
|
||||
declared,
|
||||
})
|
||||
}
|
||||
|
||||
/// A write source whose computed columns must arrive unfilled.
|
||||
pub(crate) struct UnfilledDeclarations<S> {
|
||||
inner: S,
|
||||
declared: Vec<String>,
|
||||
}
|
||||
|
||||
impl<S: lance_datafusion::utils::StreamingWriteSource> lance_datafusion::utils::StreamingWriteSource
|
||||
for UnfilledDeclarations<S>
|
||||
{
|
||||
fn arrow_schema(&self) -> SchemaRef {
|
||||
self.inner.arrow_schema()
|
||||
}
|
||||
|
||||
fn into_stream(self) -> datafusion_physical_plan::SendableRecordBatchStream {
|
||||
if self.declared.is_empty() {
|
||||
return self.inner.into_stream();
|
||||
}
|
||||
let schema = self.inner.arrow_schema();
|
||||
let declared = self.declared;
|
||||
let stream = self.inner.into_stream().map(move |batch| {
|
||||
let batch = batch?;
|
||||
ensure_batch_writes_no_computed_values(&declared, &batch)
|
||||
.map_err(|e| datafusion_common::DataFusionError::External(Box::new(e)))?;
|
||||
Ok(batch)
|
||||
});
|
||||
Box::pin(datafusion_physical_plan::stream::RecordBatchStreamAdapter::new(schema, stream))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
/// The gate's reproducer: the validator applies the same schema-level
|
||||
@@ -2646,6 +2795,8 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// A create carries a declaration only if it re-plans completely; this
|
||||
/// one lacks its inputs and is refused before its forged value matters.
|
||||
#[tokio::test]
|
||||
async fn test_create_table_cannot_inject_a_declaration() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
@@ -2673,7 +2824,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(
|
||||
matches!(&err, Error::InvalidInput { message } if message.contains("computed()")),
|
||||
matches!(&err, Error::InvalidInput { message } if message.contains("computed column 'doubled'")),
|
||||
"{err:?}"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -19,6 +19,7 @@ use lancedb::{
|
||||
connect, connect_namespace,
|
||||
database::listing::{
|
||||
ListingDatabaseOptions, NewTableConfig, OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
|
||||
OPT_NEW_TABLE_STORAGE_VERSION,
|
||||
},
|
||||
query::{ExecutableQuery, QueryBase},
|
||||
table::{AddDataMode, CompactionOptions, OptimizeAction, OptimizeStats, WriteOptions},
|
||||
@@ -146,7 +147,10 @@ async fn non_blob_table_keeps_default_format_and_row_id_setting() -> Result<()>
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
|
||||
let table = db.create_empty_table("t", schema).execute().await?;
|
||||
|
||||
assert!(!supports_blob_v2(storage_format_version(&table).await));
|
||||
assert_eq!(
|
||||
storage_format_version(&table).await,
|
||||
LanceFileVersion::Stable.resolve()
|
||||
);
|
||||
assert!(!uses_stable_row_ids(&table).await);
|
||||
Ok(())
|
||||
}
|
||||
@@ -809,7 +813,11 @@ async fn fetch_blobs_rejects_unknown_column() -> Result<()> {
|
||||
#[tokio::test]
|
||||
async fn fetch_blobs_rejects_legacy_v1_blob_column() -> Result<()> {
|
||||
let tmp = tempdir().unwrap();
|
||||
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
|
||||
// Legacy v1 blob columns are only writable at file version <= 2.1.
|
||||
let db = connect(tmp.path().to_str().unwrap())
|
||||
.storage_options([(OPT_NEW_TABLE_STORAGE_VERSION, "2.1")])
|
||||
.execute()
|
||||
.await?;
|
||||
let legacy = Field::new("image", DataType::LargeBinary, true).with_metadata(
|
||||
std::collections::HashMap::from([("lance-encoding:blob".to_string(), "true".to_string())]),
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user