Compare commits

..
Author SHA1 Message Date
Lance Release 097f455ed5 Bump version: 0.39.0-beta.1 → 0.39.0-beta.2 2026-09-05 23:25:29 +00:00
26 changed files with 191 additions and 2020 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.39.0-beta.5"
current_version = "0.39.0-beta.2"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
Generated
+47 -49
View File
@@ -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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5323,23 +5323,21 @@ 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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
dependencies = [
"arrow",
"async-trait",
"bytes",
"lance-core",
"lance-namespace-reqwest-client",
"serde",
"serde_json",
"snafu 0.9.0",
]
[[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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
dependencies = [
"arrow",
"arrow-ipc",
@@ -5378,9 +5376,9 @@ dependencies = [
[[package]]
name = "lance-namespace-reqwest-client"
version = "0.12.0"
version = "0.11.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d8d23e54b1634d5bbb434f8dd33dc3c05f6e58d876a9a27b3b4aef58ddbe11af"
checksum = "1d06b1fbb5d41f93bc652b61e2872af92e8a6c5f6b4ce8839a8ecfa05365d359"
dependencies = [
"reqwest 0.12.28",
"serde",
@@ -5392,8 +5390,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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5407,8 +5405,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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
dependencies = [
"arrow",
"arrow-array",
@@ -5448,8 +5446,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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5462,8 +5460,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.11"
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.11#4a0e26895729feb86d0cb9c09d551bfd619c6472"
dependencies = [
"frostem",
"icu_segmenter",
@@ -5476,7 +5474,7 @@ dependencies = [
[[package]]
name = "lancedb"
version = "0.39.0-beta.4"
version = "0.39.0-beta.1"
dependencies = [
"ahash",
"anyhow",
@@ -5567,7 +5565,7 @@ dependencies = [
[[package]]
name = "lancedb-nodejs"
version = "0.39.0-beta.4"
version = "0.39.0-beta.1"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5592,7 +5590,7 @@ dependencies = [
[[package]]
name = "lancedb-python"
version = "0.39.0-beta.4"
version = "0.39.0-beta.1"
dependencies = [
"arrow",
"async-trait",
+14 -14
View File
@@ -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.11", default-features = false, "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=12.0.0-beta.11", default-features = false, "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=12.0.0-beta.11", default-features = false, "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=12.0.0-beta.11", "tag" = "v12.0.0-beta.11", "git" = "https://github.com/lance-format/lance.git" }
lancedb = { path = "rust/lancedb", default-features = false }
ahash = "0.8"
# Note that this one does not include pyarrow
+1 -1
View File
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
<dependency>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-core</artifactId>
<version>0.39.0-beta.5</version>
<version>0.39.0-beta.2</version>
</dependency>
```
+1 -1
View File
@@ -8,7 +8,7 @@
<parent>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.39.0-beta.5</version>
<version>0.39.0-beta.2</version>
<relativePath>../pom.xml</relativePath>
</parent>
+2 -2
View File
@@ -6,7 +6,7 @@
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.39.0-beta.5</version>
<version>0.39.0-beta.2</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.11</lance-core.version>
<spotless.skip>false</spotless.skip>
<spotless.version>2.30.0</spotless.version>
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
+1 -1
View File
@@ -1,7 +1,7 @@
[package]
name = "lancedb-nodejs"
edition.workspace = true
version = "0.39.0-beta.5"
version = "0.39.0-beta.2"
publish = false
license.workspace = true
description.workspace = true
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-darwin-arm64",
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"os": ["darwin"],
"cpu": ["arm64"],
"main": "lancedb.darwin-arm64.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-gnu",
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-musl",
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-gnu",
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-musl",
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-arm64-msvc",
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"os": [
"win32"
],
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-x64-msvc",
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"os": ["win32"],
"cpu": ["x64"],
"main": "lancedb.win32-x64-msvc.node",
+1 -1
View File
@@ -11,7 +11,7 @@
"ann"
],
"private": false,
"version": "0.39.0-beta.5",
"version": "0.39.0-beta.2",
"main": "dist/index.js",
"exports": {
".": "./dist/index.js",
+24 -24
View File
@@ -41,7 +41,7 @@ importers:
version: 3.7.0(@emnapi/core@1.10.0)(@emnapi/runtime@1.11.3)(@types/node@22.7.4)
'@opentelemetry/sdk-metrics':
specifier: ^2.10.0
version: 2.11.0(@opentelemetry/api@1.9.1)
version: 2.10.0(@opentelemetry/api@1.9.1)
'@types/axios':
specifier: ^0.14.0
version: 0.14.4
@@ -80,7 +80,7 @@ importers:
version: 0.2.7
ts-jest:
specifier: ^29.1.2
version: 29.4.12(@babel/core@7.29.7)(@jest/transform@29.7.0)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.29.7))(jest-util@29.7.0)(jest@29.7.0(@types/node@22.7.4))(typescript@5.5.4)
version: 29.4.9(@babel/core@7.29.7)(@jest/transform@29.7.0)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.29.7))(jest-util@29.7.0)(jest@29.7.0(@types/node@22.7.4))(typescript@5.5.4)
typedoc:
specifier: 0.26.4
version: 0.26.4(typescript@5.5.4)
@@ -1394,20 +1394,20 @@ packages:
resolution: {integrity: sha512-gLyJlPHPZYdAk1JENA9LeHejZe1Ti77/pTeFm/nMXmQH/HFZlcS/O2XJB+L8fkbrNSqhdtlvjBVjxwUYanNH5Q==}
engines: {node: '>=8.0.0'}
'@opentelemetry/core@2.11.0':
resolution: {integrity: sha512-7YP44XH0tV6+Mb54x2YGf84i7yi+31MBZlE8JwvozkxyTvXbSp10X7cI7YE49ChJ3shMJoBmCJF3+1QFBJctGA==}
'@opentelemetry/core@2.10.0':
resolution: {integrity: sha512-/wNZ8twnEQQA4HoHu22+vcsdru6pWPWxW+7w+FlxT6Id7PE/WIbZmVKkte+PF72e0F2dnImFeHD2syyE1Mw6MQ==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.0.0 <1.10.0'
'@opentelemetry/resources@2.11.0':
resolution: {integrity: sha512-Ie7+8q8MDF4FAEQCKVMTx3ReUvxiIAgIiiW3c9JdmP8+HMcDy20puT+AHjexnExgnbvBxjQ9fjkFDWrikJ2jQA==}
'@opentelemetry/resources@2.10.0':
resolution: {integrity: sha512-q6MMm2zhggzsHVNbabYwut+a6nbuQQe3URUoxaojM/8K1IBfwwPzvxIjNi2/lI1TFe+fMHMW9MWhrtDLEXEnkA==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.3.0 <1.10.0'
'@opentelemetry/sdk-metrics@2.11.0':
resolution: {integrity: sha512-7GXXcObyHyDUUSG+L+kJoquty01bzm7ivE7+SSgXXJcHuPzGviptxwARmI2c+bnnxjexGQbJnyNlN8HxBP/Y7A==}
'@opentelemetry/sdk-metrics@2.10.0':
resolution: {integrity: sha512-t6r1VSvXNtSDnPXU1FbZeetJb7yyovHmgu0wRSoftxtE0g2rSNhQZQUy69sRUCL+iioJpX8SN/S6wq6ZtvLySQ==}
engines: {node: ^18.19.0 || >=20.6.0}
peerDependencies:
'@opentelemetry/api': '>=1.9.0 <1.10.0'
@@ -1480,7 +1480,6 @@ packages:
'@smithy/core@3.24.1':
resolution: {integrity: sha512-3mT7o4qQyUWttYnVK3A0Z/u3Xha3E81tXn32Tz6vjZiUXhBrkEivpw1hBYfh84iFF9CSzkBU9Y1DJ3Q6RQ231g==}
engines: {node: '>=18.0.0'}
deprecated: Deprecated due to bug in browser bundling instructions https://github.com/smithy-lang/smithy-typescript/issues/2025
'@smithy/credential-provider-imds@4.3.1':
resolution: {integrity: sha512-0S/acwHnqX4WrjXzhdiDRxsG2s9SC0cpPIK9nZ1R6UOHd+j7uL28+4bHu22urbLk2TVw3fkp6na/+fkUt/pLNQ==}
@@ -3239,8 +3238,8 @@ packages:
peerDependencies:
typescript: '>=4.2.0'
ts-jest@29.4.12:
resolution: {integrity: sha512-Ov6ClY53Fflh6BGAnY2DlTq1hYDrTycz2PVTXBWFW2CU+9zrEqAp9fWdGXl42EXO5RLSFAcAZ2JFKbP+zBTFfw==}
ts-jest@29.4.9:
resolution: {integrity: sha512-LTb9496gYPMCqjeDLdPrKuXtncudeV1yRZnF4Wo5l3SFi0RYEnYRNgMrFIdg+FHvfzjCyQk1cLncWVqiSX+EvQ==}
engines: {node: ^14.15.0 || ^16.10.0 || ^18.0.0 || >=20.0.0}
hasBin: true
peerDependencies:
@@ -5111,22 +5110,22 @@ snapshots:
'@opentelemetry/api@1.9.1': {}
'@opentelemetry/core@2.11.0(@opentelemetry/api@1.9.1)':
'@opentelemetry/core@2.10.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/semantic-conventions': 1.43.0
'@opentelemetry/resources@2.11.0(@opentelemetry/api@1.9.1)':
'@opentelemetry/resources@2.10.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.11.0(@opentelemetry/api@1.9.1)
'@opentelemetry/core': 2.10.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions': 1.43.0
'@opentelemetry/sdk-metrics@2.11.0(@opentelemetry/api@1.9.1)':
'@opentelemetry/sdk-metrics@2.10.0(@opentelemetry/api@1.9.1)':
dependencies:
'@opentelemetry/api': 1.9.1
'@opentelemetry/core': 2.11.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources': 2.11.0(@opentelemetry/api@1.9.1)
'@opentelemetry/core': 2.10.0(@opentelemetry/api@1.9.1)
'@opentelemetry/resources': 2.10.0(@opentelemetry/api@1.9.1)
'@opentelemetry/semantic-conventions@1.43.0': {}
@@ -5575,7 +5574,7 @@ snapshots:
globby: 11.1.0
is-glob: 4.0.3
minimatch: 9.0.9
semver: 7.8.5
semver: 7.8.0
ts-api-utils: 1.4.3(typescript@5.5.4)
optionalDependencies:
typescript: 5.5.4
@@ -6427,7 +6426,7 @@ snapshots:
'@babel/parser': 7.29.3
'@istanbuljs/schema': 0.1.6
istanbul-lib-coverage: 3.2.2
semver: 7.8.5
semver: 7.8.0
transitivePeerDependencies:
- supports-color
@@ -6706,7 +6705,7 @@ snapshots:
jest-util: 29.7.0
natural-compare: 1.4.0
pretty-format: 29.7.0
semver: 7.8.5
semver: 7.8.0
transitivePeerDependencies:
- supports-color
@@ -6829,7 +6828,7 @@ snapshots:
make-dir@4.0.0:
dependencies:
semver: 7.8.5
semver: 7.8.0
make-error@1.3.6: {}
@@ -7163,7 +7162,8 @@ snapshots:
semver@7.8.0: {}
semver@7.8.5: {}
semver@7.8.5:
optional: true
sharp@0.35.4(@types/node@22.7.4):
dependencies:
@@ -7327,7 +7327,7 @@ snapshots:
dependencies:
typescript: 5.5.4
ts-jest@29.4.12(@babel/core@7.29.7)(@jest/transform@29.7.0)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.29.7))(jest-util@29.7.0)(jest@29.7.0(@types/node@22.7.4))(typescript@5.5.4):
ts-jest@29.4.9(@babel/core@7.29.7)(@jest/transform@29.7.0)(@jest/types@29.6.3)(babel-jest@29.7.0(@babel/core@7.29.7))(jest-util@29.7.0)(jest@29.7.0(@types/node@22.7.4))(typescript@5.5.4):
dependencies:
bs-logger: 0.2.6
fast-json-stable-stringify: 2.1.0
@@ -7336,7 +7336,7 @@ snapshots:
json5: 2.2.3
lodash.memoize: 4.1.2
make-error: 1.3.6
semver: 7.8.5
semver: 7.8.0
type-fest: 4.41.0
typescript: 5.5.4
yargs-parser: 21.1.1
+1 -1
View File
@@ -127,7 +127,7 @@ impl Job {
}
/// Serialise Arrow batches as a single IPC stream for the TypeScript layer.
fn batches_to_ipc_buffer(batches: &[RecordBatch]) -> napi::Result<Buffer> {
pub(crate) fn batches_to_ipc_buffer(batches: &[RecordBatch]) -> napi::Result<Buffer> {
let Some(first) = batches.first() else {
return Ok(Buffer::from(Vec::<u8>::new()));
};
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb-python"
version = "0.39.0-beta.5"
version = "0.39.0-beta.2"
publish = false
edition.workspace = true
description = "Python bindings for LanceDB"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb"
version = "0.39.0-beta.5"
version = "0.39.0-beta.2"
edition.workspace = true
description = "LanceDB: A serverless, low-latency vector database for AI applications"
license.workspace = true
-5
View File
@@ -3,7 +3,6 @@
//! 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};
@@ -305,10 +304,6 @@ 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
+17 -536
View File
@@ -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, new_null_array};
use arrow_schema::{FieldRef, Schema as ArrowSchema, SchemaRef};
use arrow_array::{RecordBatch, UInt64Array};
use arrow_schema::{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, UpdateMode};
use lance::dataset::transaction::{Operation, Transaction};
use lance::dataset::write::delete::DeleteBuilder;
use lance::dataset::write::merge_insert::inserted_rows::{
KeyExistenceFilter, KeyExistenceFilterBuilder, KeyValue,
@@ -51,9 +51,6 @@ 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};
@@ -170,52 +167,30 @@ pub(crate) async fn execute_refresh(
.map(|p| (p.output.clone(), p.expression.clone()))
.collect();
validate_inputs(&source_ds, definition)?;
let (replanned, planned_fields, _renames) = super::plan(
let (replanned, mut planned_fields, _renames) = super::plan(
source_schema,
&definition.source_table,
&definition.source_namespace,
Some(&projections),
&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 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
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
.fields()
.iter()
.filter(|f| computed_column_from_field(f).is_none())
.map(|f| (f.name().clone(), f.data_type().clone(), f.is_nullable()))
.collect();
// 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 {
if planned_shape != physical_shape {
return Err(Error::Schema {
message: format!(
"the stored definition of view '{}' does not produce this \
@@ -254,18 +229,11 @@ 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, except a fill of its
// computed columns, which rewrites nothing refresh certifies.
let recorded_view_version = metadata
// any other commit on the view since then is drift.
let view_intact = metadata
.get(VIEW_VERSION_META_KEY)
.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,
};
.and_then(|raw| raw.parse::<u64>().ok())
== Some(view_ds.version().version);
if !full && watermark == Some(source_version) && view_intact && recorded_ts == Some(source_ts) {
return Ok(RefreshMaterializedViewResult {
@@ -1122,69 +1090,6 @@ 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,
@@ -1253,10 +1158,6 @@ 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 {
@@ -2867,7 +2768,7 @@ mod tests {
let (conn, source) = db_with_source(vec![1]).await;
let prepared = crate::materialized_view::prepare_declaration(
&source,
Some(&[("x".into(), "x".into()), ("twice".into(), "x * 2".into())]),
&[("x".into(), "x".into()), ("twice".into(), "x * 2".into())],
None,
None,
)
@@ -3231,424 +3132,4 @@ 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
);
}
}
+1 -1
View File
@@ -8021,7 +8021,7 @@ mod tests {
match request.url().path() {
"/v1/table/my_table/backfill_column" => http::Response::builder()
.status(202)
.body(br#"{"job_id": "j-42"}"#.to_vec())
.body(r#"{"job_id": "j-42"}"#.as_bytes().to_vec())
.unwrap(),
"/v1/jobs/describe" => http::Response::builder()
.status(200)
+1 -2
View File
@@ -2804,7 +2804,7 @@ impl NativeTable {
namespace_client: Option<Arc<dyn LanceNamespace>>,
pushdown_operations: HashSet<NamespaceClientPushdownOperation>,
) -> Result<Self> {
let batches = computed_columns::admit_create_source(batches)?;
computed_columns::ensure_no_foreign_declarations(batches.arrow_schema().fields())?;
// Default params uses format v1.
let params = params.unwrap_or(WriteParams {
..Default::default()
@@ -2904,7 +2904,6 @@ 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());
+45 -322
View File
@@ -21,7 +21,6 @@
//! [`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;
@@ -760,33 +759,14 @@ fn function_output_field(name: &str, nullable: bool, raw: &str) -> Result<JsonAr
Ok(field)
}
/// Whether two fields describe the same Function output.
///
/// `compare_identity` covers the field's own name and nullability. Struct
/// children carry both as part of the declaration and compare with it on. List
/// children do not: Lance rewrites a list item's name and nullability when it
/// writes, so a stored `fixed_size_list<item: float not null>` comes back as
/// `fixed_size_list<item: float>` and never matches the declaration again.
/// Comparing those by type alone keeps this agreeing with the server, which
/// draws the same distinction and is what accepted the column when it was
/// declared.
fn function_output_field_matches(
expected: &ArrowField,
actual: &ArrowField,
compare_identity: bool,
) -> bool {
if compare_identity
&& (expected.name() != actual.name() || expected.is_nullable() != actual.is_nullable())
{
return false;
}
match (expected.is_blob_v2(), actual.is_blob_v2()) {
(false, false) => function_output_type_matches(expected.data_type(), actual.data_type()),
(true, true) => {
fn function_output_field_matches(expected: &ArrowField, actual: &ArrowField) -> bool {
expected.name() == actual.name()
&& expected.is_nullable() == actual.is_nullable()
&& if expected.is_blob_v2() {
has_supported_blob_v2_layout(expected) && has_supported_blob_v2_layout(actual)
} else {
function_output_type_matches(expected.data_type(), actual.data_type())
}
_ => false,
}
}
fn function_output_type_matches(expected: &DataType, actual: &DataType) -> bool {
@@ -799,19 +779,33 @@ fn function_output_type_matches(expected: &DataType, actual: &DataType) -> bool
&& expected
.iter()
.zip(actual)
.all(|(expected, actual)| function_output_field_matches(expected, actual, true))
.all(|(expected, actual)| function_output_field_matches(expected, actual))
}
(DataType::List(expected), DataType::List(actual))
| (DataType::LargeList(expected), DataType::LargeList(actual)) => {
function_output_field_matches(expected, actual, false)
function_output_field_matches(expected, actual)
}
(
DataType::FixedSizeList(expected, expected_size),
DataType::FixedSizeList(actual, actual_size),
) => expected_size == actual_size && function_output_field_matches(expected, actual, false),
) => expected_size == actual_size && function_output_field_matches(expected, actual),
(DataType::Map(expected, expected_sorted), DataType::Map(actual, actual_sorted)) => {
expected_sorted == actual_sorted
&& function_output_field_matches(expected, actual, true)
expected_sorted == actual_sorted && function_output_field_matches(expected, actual)
}
_ => false,
}
}
fn function_output_type_has_blob(data_type: &DataType) -> bool {
match data_type {
DataType::Struct(fields) => fields
.iter()
.any(|field| field.is_blob_v2() || function_output_type_has_blob(field.data_type())),
DataType::List(field)
| DataType::LargeList(field)
| DataType::FixedSizeList(field, _)
| DataType::Map(field, _) => {
field.is_blob_v2() || function_output_type_has_blob(field.data_type())
}
_ => false,
}
@@ -897,13 +891,16 @@ fn ensure_binding_matches_schema(schema: &ArrowSchema, binding: &FunctionBinding
binding.binding_id()
)));
}
let type_matches = if output.arrow_type == FUNCTION_BLOB_V2_TYPE {
has_supported_blob_v2_layout(field)
let (type_matches, has_semantic_blob) = if output.arrow_type == FUNCTION_BLOB_V2_TYPE {
(has_supported_blob_v2_layout(field), true)
} else {
let expected_type = parse_output_arrow_type(&output.arrow_type)?;
let expected_type = lance_namespace::schema::convert_json_arrow_type(&expected_type)
.map_err(|e| invalid_function(format!("invalid Function output type: {e}")))?;
function_output_type_matches(&expected_type, field.data_type())
(
function_output_type_matches(&expected_type, field.data_type()),
function_output_type_has_blob(&expected_type),
)
};
if !type_matches {
return Err(invalid_function(format!(
@@ -934,16 +931,19 @@ fn ensure_binding_matches_schema(schema: &ArrowSchema, binding: &FunctionBinding
binding.binding_id()
)));
}
// Rebuild from the declaration rather than from the stored field. The
// stored field carries Lance's write-time normalization, which would
// never round-trip back to the schema the binding recorded -- the same
// reason list children compare by type above. Whether the column on
// disk still matches is settled by that comparison, not here.
output_fields.push(function_output_field(
field.name(),
true,
&output.arrow_type,
)?);
if has_semantic_blob {
output_fields.push(function_output_field(
field.name(),
true,
&output.arrow_type,
)?);
} else {
let json = lance_namespace::schema::arrow_schema_to_json(&ArrowSchema::new(vec![
ArrowField::new(field.name().clone(), field.data_type().clone(), true),
]))
.map_err(|e| invalid_function(format!("invalid Function output schema: {e}")))?;
output_fields.push(json.fields.into_iter().next().unwrap());
}
}
if let Some(assignment) = binding.assignment() {
if binding
@@ -1339,106 +1339,6 @@ 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>(
@@ -1897,54 +1797,6 @@ 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
@@ -1963,69 +1815,6 @@ mod tests {
assert!(super::validate_declarations(schema, &declarations).is_err());
}
#[test]
fn list_children_match_by_type_but_struct_children_by_identity() {
use arrow_schema::Field as F;
// Lance rewrites a list item's name and nullability on write, so the
// stored field is no longer identical to what was declared. Comparing
// those by type keeps a table with a vector output usable.
let declared =
DataType::FixedSizeList(Arc::new(F::new("item", DataType::Float32, false)), 4);
let stored = DataType::FixedSizeList(Arc::new(F::new("item", DataType::Float32, true)), 4);
assert!(super::function_output_type_matches(&declared, &stored));
let renamed =
DataType::FixedSizeList(Arc::new(F::new("element", DataType::Float32, true)), 4);
assert!(super::function_output_type_matches(&declared, &renamed));
// The dimension is still part of the declaration.
let resized = DataType::FixedSizeList(Arc::new(F::new("item", DataType::Float32, true)), 8);
assert!(!super::function_output_type_matches(&declared, &resized));
// Struct children keep comparing by name and nullability.
let struct_declared =
DataType::Struct(vec![F::new("changed", DataType::Boolean, false)].into());
let struct_nullable =
DataType::Struct(vec![F::new("changed", DataType::Boolean, true)].into());
let struct_renamed =
DataType::Struct(vec![F::new("altered", DataType::Boolean, false)].into());
assert!(super::function_output_type_matches(
&struct_declared,
&struct_declared
));
assert!(!super::function_output_type_matches(
&struct_declared,
&struct_nullable
));
assert!(!super::function_output_type_matches(
&struct_declared,
&struct_renamed
));
// A list nested inside a struct gets the list rule.
let nested_declared = DataType::Struct(
vec![F::new(
"tokens",
DataType::List(Arc::new(F::new("item", DataType::Utf8, false))),
true,
)]
.into(),
);
let nested_stored = DataType::Struct(
vec![F::new(
"tokens",
DataType::List(Arc::new(F::new("item", DataType::Utf8, true))),
true,
)]
.into(),
);
assert!(super::function_output_type_matches(
&nested_declared,
&nested_stored
));
}
#[test]
fn output_arrow_type_grammar_matches_the_shared_golden() {
let golden: serde_json::Value = serde_json::from_str(include_str!(
@@ -2795,8 +2584,6 @@ 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();
@@ -2824,7 +2611,7 @@ mod tests {
.await
.unwrap_err();
assert!(
matches!(&err, Error::InvalidInput { message } if message.contains("computed column 'doubled'")),
matches!(&err, Error::InvalidInput { message } if message.contains("computed()")),
"{err:?}"
);
}
@@ -3502,70 +3289,6 @@ mod tests {
assert!(output_schema.field(0).is_blob_v2());
}
#[test]
fn binding_accepts_a_lance_normalized_list_child() {
// The whole guard, not just the type helper: this also reaches the
// output-schema comparison at the end of ensure_binding_matches_schema,
// which used to rebuild the schema from the stored field and so failed
// on exactly the same normalization.
let input = ArrowField::new("value", DataType::Int64, false);
let application = FunctionApplication::from_json(
&serde_json::json!({
"function": {"name": "embed", "version": "fv_embed"},
"inputs": [{
"parameter": "value",
"kind": "column",
"value": {"path": "value"}
}],
"output": {
"kind": "scalar",
"arrow_type": "fixed_size_list<float32, 4>",
"nullable": false
}
})
.to_string(),
)
.unwrap();
let plan = plan_function_application(
&ArrowSchema::new(vec![input.clone()]),
&application,
Some("embedding"),
)
.unwrap();
let binding = binding_from_plan(&plan);
// The declaration says the item is non-nullable; Lance rewrites it to
// nullable on write, so this is what the column looks like on disk.
let stored = DataType::FixedSizeList(
Arc::new(ArrowField::new("item", DataType::Float32, true)),
4,
);
let output = ArrowField::new("embedding", stored, true).with_metadata(
function_computed_column_metadata(binding.binding_id(), 0, &["value".into()]),
);
ensure_binding_matches_schema(&ArrowSchema::new(vec![input.clone(), output]), &binding)
.unwrap();
// A different element type is still a mismatch.
let wrong = ArrowField::new(
"embedding",
DataType::FixedSizeList(
Arc::new(ArrowField::new("item", DataType::Float64, true)),
4,
),
true,
)
.with_metadata(function_computed_column_metadata(
binding.binding_id(),
0,
&["value".into()],
));
assert!(
ensure_binding_matches_schema(&ArrowSchema::new(vec![input, wrong]), &binding).is_err()
);
}
#[test]
fn test_blob_scalar_binding_accepts_full_logical_layout() {
let input = crate::blob("image", false);
@@ -5,17 +5,12 @@
//!
//! [`super::cast::cast_to_table_schema`] calls [`coerce_blob_expr`].
use std::fmt;
use std::hash::{Hash, Hasher};
use std::sync::Arc;
use arrow_array::{Array, BooleanArray, RecordBatch};
use arrow_schema::{DataType, Field, FieldRef, Fields, Schema};
use arrow_select::nullif::nullif;
use arrow_schema::{DataType, Field, FieldRef, Fields};
use datafusion::functions::core::{get_field, named_struct};
use datafusion_common::ScalarValue;
use datafusion_common::config::ConfigOptions;
use datafusion_expr::ColumnarValue;
use datafusion_physical_expr::ScalarFunctionExpr;
use datafusion_physical_expr::expressions::{CastExpr, Literal};
use datafusion_physical_plan::PhysicalExpr;
@@ -138,102 +133,16 @@ pub(super) fn coerce_blob_expr(
ns_args.push(value);
}
let built: Arc<dyn PhysicalExpr> = Arc::new(ScalarFunctionExpr::new(
let expr: Arc<dyn PhysicalExpr> = Arc::new(ScalarFunctionExpr::new(
&format!("named_struct({})", table_field.name()),
named_struct(),
ns_args,
table_field.clone(),
config.clone(),
));
// `named_struct` always yields a valid struct, so a null input would land
// as a row that set neither `data` nor `uri` -- not an absent blob but a
// malformed one, which Lance rejects on write.
let expr: Arc<dyn PhysicalExpr> = Arc::new(AbsentBlobIsNull {
source: input_expr,
built,
field: table_field.clone(),
});
Ok((expr, table_field.clone()))
}
/// Carries the source column's nullity onto the struct built for it.
///
/// This is its own expression rather than a `CASE` because the projection
/// takes its output field from `return_field`, and the generic implementation
/// rebuilds a bare field -- which would drop the `lance.blob.v2` extension
/// metadata and stop the column being recognised as a blob at all.
#[derive(Debug, Clone)]
struct AbsentBlobIsNull {
source: Arc<dyn PhysicalExpr>,
built: Arc<dyn PhysicalExpr>,
field: FieldRef,
}
impl fmt::Display for AbsentBlobIsNull {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "absent_blob_is_null({}, {})", self.source, self.built)
}
}
impl PartialEq for AbsentBlobIsNull {
fn eq(&self, other: &Self) -> bool {
self.source.eq(&other.source) && self.built.eq(&other.built) && self.field == other.field
}
}
impl Eq for AbsentBlobIsNull {}
impl Hash for AbsentBlobIsNull {
fn hash<H: Hasher>(&self, state: &mut H) {
self.source.hash(state);
self.built.hash(state);
self.field.hash(state);
}
}
impl PhysicalExpr for AbsentBlobIsNull {
fn return_field(&self, _input_schema: &Schema) -> datafusion_common::Result<FieldRef> {
Ok(self.field.clone())
}
fn nullable(&self, _input_schema: &Schema) -> datafusion_common::Result<bool> {
Ok(true)
}
fn evaluate(&self, batch: &RecordBatch) -> datafusion_common::Result<ColumnarValue> {
let rows = batch.num_rows();
let built = self.built.evaluate(batch)?.into_array(rows)?;
let source = self.source.evaluate(batch)?.into_array(rows)?;
let Some(nulls) = source.logical_nulls() else {
return Ok(ColumnarValue::Array(built));
};
// `nullif` nulls the rows the mask marks true, which is where the
// source had no value.
let absent = BooleanArray::new(!nulls.inner(), None);
Ok(ColumnarValue::Array(nullif(built.as_ref(), &absent)?))
}
fn children(&self) -> Vec<&Arc<dyn PhysicalExpr>> {
vec![&self.source, &self.built]
}
fn with_new_children(
self: Arc<Self>,
children: Vec<Arc<dyn PhysicalExpr>>,
) -> datafusion_common::Result<Arc<dyn PhysicalExpr>> {
Ok(Arc::new(Self {
source: children[0].clone(),
built: children[1].clone(),
field: self.field.clone(),
}))
}
fn fmt_sql(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "{self}")
}
}
enum BlobInputShape<'a> {
Bytes,
String,
@@ -404,11 +313,6 @@ mod tests {
let data = image.column_by_name("data").unwrap();
assert!(!data.is_null(0));
assert!(data.is_null(1));
// The row itself has to be null, not merely a struct whose children
// are. A present-but-empty struct set neither `data` nor `uri`, which
// Lance rejects as malformed rather than reading as an absent blob.
assert!(!image.is_null(0));
assert!(image.is_null(1));
}
#[tokio::test]