mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-04 20:48:50 +00:00
Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3d562a7c28 | |||
| 7357d63e87 | |||
| 624a75edf7 | |||
| c7ea91f3ea | |||
| 8e24dd3828 |
Generated
+30
-30
@@ -3422,7 +3422,7 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
|
||||
[[package]]
|
||||
name = "fsst"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"rand 0.9.5",
|
||||
@@ -4778,7 +4778,7 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
|
||||
[[package]]
|
||||
name = "lance"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -4853,7 +4853,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4875,7 +4875,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-scalar"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4889,7 +4889,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-stats"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -4899,7 +4899,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-bitpacking"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrayref",
|
||||
"crunchy",
|
||||
@@ -4910,7 +4910,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-core"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4951,7 +4951,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-datafusion"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -4982,7 +4982,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-datagen"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5000,7 +5000,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-derive"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
@@ -5010,7 +5010,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-encoding"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5046,7 +5046,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-file"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5078,7 +5078,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-index"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -5146,7 +5146,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-index-core"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5169,7 +5169,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-io"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5207,7 +5207,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-linalg"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5224,7 +5224,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-namespace"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -5237,7 +5237,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-namespace-impls"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-ipc",
|
||||
@@ -5292,7 +5292,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-select"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5308,7 +5308,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-table"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5348,7 +5348,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-testing"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5362,7 +5362,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-tokenizer"
|
||||
version = "10.1.0-beta.1"
|
||||
source = "git+https://github.com/lance-format/lance.git?rev=b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564#b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v10.1.0-beta.1#68f4d4c1d0c4871b067557c61fc405078f1ab3b7"
|
||||
dependencies = [
|
||||
"icu_segmenter",
|
||||
"jieba-rs",
|
||||
@@ -7583,7 +7583,7 @@ version = "0.14.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "343d3bd7056eda839b03204e68deff7d1b13aba7af2b2fd16890697274262ee7"
|
||||
dependencies = [
|
||||
"heck 0.5.0",
|
||||
"heck 0.4.1",
|
||||
"itertools 0.14.0",
|
||||
"log",
|
||||
"multimap",
|
||||
@@ -8498,9 +8498,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rkyv"
|
||||
version = "0.8.16"
|
||||
version = "0.8.17"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "73389e0c99e664f919275ab5b5b0471391fe9a8de61e1dff9b1eaf56a90f16e3"
|
||||
checksum = "815cc8a37159a463064825246cadb07961e25cd9885908606f6d08a98d8f8874"
|
||||
dependencies = [
|
||||
"bytecheck",
|
||||
"bytes",
|
||||
@@ -8517,9 +8517,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "rkyv_derive"
|
||||
version = "0.8.16"
|
||||
version = "0.8.17"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "5d2ed0b54125315fb36bd021e82d314d1c126548f871634b483f46b31d13cac6"
|
||||
checksum = "c0ed1a78a1b19d184b0daa629dd9a024573173ec7d485b287cb369fb3607cc1c"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
@@ -9277,7 +9277,7 @@ version = "0.8.9"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "c1c97747dbf44bb1ca44a561ece23508e99cb592e862f22222dcf42f51d1e451"
|
||||
dependencies = [
|
||||
"heck 0.5.0",
|
||||
"heck 0.4.1",
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.117",
|
||||
@@ -9289,7 +9289,7 @@ version = "0.9.0"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "54254b8531cafa275c5e096f62d48c81435d1015405a91198ddb11e967301d40"
|
||||
dependencies = [
|
||||
"heck 0.5.0",
|
||||
"heck 0.4.1",
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
"syn 2.0.117",
|
||||
@@ -9730,7 +9730,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "32497e9a4c7b38532efcdebeef879707aa9f794296a4f0244f6f69e9bc8574bd"
|
||||
dependencies = [
|
||||
"fastrand",
|
||||
"getrandom 0.4.2",
|
||||
"getrandom 0.3.4",
|
||||
"once_cell",
|
||||
"rustix",
|
||||
"windows-sys 0.61.2",
|
||||
|
||||
+14
-17
@@ -13,23 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
# TEMPORARY: `ObjectStore::read_dir_stream` is not in a lance release yet, so these point at
|
||||
# lance-format/lance#8120 cherry-picked onto the v10.1.0-beta.1 tag. Put the tag back once that
|
||||
# PR has merged and shipped in a release.
|
||||
lance = { "version" = "=10.1.0-beta.1", default-features = false, "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=10.1.0-beta.1", default-features = false, "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=10.1.0-beta.1", default-features = false, "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=10.1.0-beta.1", "rev" = "b7ef2bae3c7aeb85da637a8c3bc6643ab7d81564", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
ahash = "0.8"
|
||||
# Note that this one does not include pyarrow
|
||||
arrow = { version = "58.0.0", optional = false }
|
||||
|
||||
+6
-2
@@ -339,7 +339,9 @@ impl Table {
|
||||
let transforms = NewColumnTransform::SqlExpressions(transforms);
|
||||
let res = self
|
||||
.inner_ref()?
|
||||
.add_columns(transforms, None)
|
||||
.add_columns()
|
||||
.transform(transforms)
|
||||
.execute()
|
||||
.await
|
||||
.default_error()?;
|
||||
Ok(res.into())
|
||||
@@ -356,7 +358,9 @@ impl Table {
|
||||
let transforms = NewColumnTransform::AllNulls(schema);
|
||||
let res = self
|
||||
.inner_ref()?
|
||||
.add_columns(transforms, None)
|
||||
.add_columns()
|
||||
.transform(transforms)
|
||||
.execute()
|
||||
.await
|
||||
.default_error()?;
|
||||
Ok(res.into())
|
||||
|
||||
@@ -707,6 +707,9 @@ class LanceDBConnection(DBConnection):
|
||||
self._namespace_client_properties = namespace_client_properties
|
||||
if _inner is not None:
|
||||
self._conn = _inner
|
||||
# Native-derived wrappers resolve this in their async reconstruction
|
||||
# path so construction never synchronously re-enters LOOP.
|
||||
self._read_consistency_interval = read_consistency_interval
|
||||
self._cached_namespace_client = None
|
||||
return
|
||||
|
||||
@@ -756,11 +759,14 @@ class LanceDBConnection(DBConnection):
|
||||
# storage_options. Also, this class really shouldn't be holding any state
|
||||
# beyond _conn.
|
||||
self._conn = AsyncConnection(LOOP.run(do_connect()))
|
||||
# Keep property access synchronous so debugger introspection cannot wait on
|
||||
# the background loop while that thread is suspended at a breakpoint.
|
||||
self._read_consistency_interval = read_consistency_interval
|
||||
self._cached_namespace_client: Optional[LanceNamespace] = None
|
||||
|
||||
@property
|
||||
def read_consistency_interval(self) -> Optional[timedelta]:
|
||||
return LOOP.run(self._conn.get_read_consistency_interval())
|
||||
return self._read_consistency_interval
|
||||
|
||||
@property
|
||||
def session(self) -> Optional[Session]:
|
||||
@@ -771,8 +777,16 @@ class LanceDBConnection(DBConnection):
|
||||
return self._conn.uri
|
||||
|
||||
@classmethod
|
||||
def from_inner(cls, inner: LanceDbConnection):
|
||||
return cls(None, _inner=inner)
|
||||
def from_inner(
|
||||
cls,
|
||||
inner: LanceDbConnection,
|
||||
read_consistency_interval: Optional[timedelta],
|
||||
):
|
||||
return cls(
|
||||
None,
|
||||
read_consistency_interval=read_consistency_interval,
|
||||
_inner=inner,
|
||||
)
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return f"{self.__class__.__name__}(uri={self._conn.uri!r})"
|
||||
|
||||
@@ -226,7 +226,7 @@ class PermutationBuilder:
|
||||
|
||||
async def do_execute():
|
||||
inner_tbl = await self._async.execute()
|
||||
return LanceTable.from_inner(inner_tbl)
|
||||
return await LanceTable.from_inner(inner_tbl)
|
||||
|
||||
return LOOP.run(do_execute())
|
||||
|
||||
|
||||
+134
-37
@@ -2182,11 +2182,15 @@ class LanceTable(Table):
|
||||
return self.name
|
||||
|
||||
@classmethod
|
||||
def from_inner(cls, tbl: LanceDBTable):
|
||||
from .db import LanceDBConnection
|
||||
async def from_inner(cls, tbl: LanceDBTable):
|
||||
from .db import AsyncConnection, LanceDBConnection
|
||||
|
||||
async_tbl = AsyncTable(tbl)
|
||||
conn = LanceDBConnection.from_inner(tbl.database())
|
||||
inner_conn = tbl.database()
|
||||
read_consistency_interval = await AsyncConnection(
|
||||
inner_conn
|
||||
).get_read_consistency_interval()
|
||||
conn = LanceDBConnection.from_inner(inner_conn, read_consistency_interval)
|
||||
return cls(
|
||||
conn,
|
||||
async_tbl.name,
|
||||
@@ -4031,7 +4035,10 @@ def _handle_bad_vectors(
|
||||
for vector_column in vector_columns:
|
||||
dim = vector_column["expected_dim"]
|
||||
if target_schema is not None and dim is None:
|
||||
dim = _infer_vector_dim(batch[vector_column["name"]])
|
||||
dim = _infer_vector_column_dim(
|
||||
batch[vector_column["name"]],
|
||||
vector_column["is_multivector"],
|
||||
)
|
||||
pending_dims.append(vector_column)
|
||||
batch = _handle_bad_vector_column(
|
||||
batch,
|
||||
@@ -4040,11 +4047,13 @@ def _handle_bad_vectors(
|
||||
fill_value=fill_value,
|
||||
expected_dim=dim,
|
||||
expected_value_type=vector_column["expected_value_type"],
|
||||
is_multivector=vector_column["is_multivector"],
|
||||
)
|
||||
for vector_column in pending_dims:
|
||||
if vector_column["expected_dim"] is None:
|
||||
vector_column["expected_dim"] = _infer_vector_dim(
|
||||
batch[vector_column["name"]]
|
||||
vector_column["expected_dim"] = _infer_vector_column_dim(
|
||||
batch[vector_column["name"]],
|
||||
vector_column["is_multivector"],
|
||||
)
|
||||
if batch.schema.equals(output_schema, check_metadata=True):
|
||||
yield batch
|
||||
@@ -4070,22 +4079,30 @@ def _find_vector_columns(
|
||||
if target_schema is None:
|
||||
vector_columns = []
|
||||
for field in reader_schema:
|
||||
named_vector_col = (
|
||||
_is_list_like(field.type)
|
||||
and pa.types.is_floating(field.type.value_type)
|
||||
and field.name == VECTOR_COLUMN_NAME
|
||||
is_multivector = _is_multivector_type(field.type)
|
||||
is_fixed_multivector = is_multivector and pa.types.is_fixed_size_list(
|
||||
field.type.value_type
|
||||
)
|
||||
named_vector_col = (
|
||||
_is_float_vector_type(field.type) or is_multivector
|
||||
) and field.name == VECTOR_COLUMN_NAME
|
||||
likely_vector_col = (
|
||||
pa.types.is_fixed_size_list(field.type)
|
||||
and pa.types.is_floating(field.type.value_type)
|
||||
and (field.type.list_size >= 10)
|
||||
)
|
||||
if named_vector_col or likely_vector_col:
|
||||
if named_vector_col or likely_vector_col or is_fixed_multivector:
|
||||
vector_type = field.type.value_type if is_multivector else field.type
|
||||
vector_columns.append(
|
||||
{
|
||||
"name": field.name,
|
||||
"expected_dim": None,
|
||||
"expected_value_type": None,
|
||||
"expected_dim": (
|
||||
vector_type.list_size
|
||||
if pa.types.is_fixed_size_list(vector_type)
|
||||
else None
|
||||
),
|
||||
"expected_value_type": vector_type.value_type,
|
||||
"is_multivector": is_multivector,
|
||||
}
|
||||
)
|
||||
return vector_columns
|
||||
@@ -4099,9 +4116,8 @@ def _find_vector_columns(
|
||||
for field in target_schema:
|
||||
if field.name not in reader_column_names:
|
||||
continue
|
||||
if not _is_list_like(field.type) or not pa.types.is_floating(
|
||||
field.type.value_type
|
||||
):
|
||||
is_multivector = _is_multivector_type(field.type)
|
||||
if not _is_float_vector_type(field.type) and not is_multivector:
|
||||
continue
|
||||
|
||||
reader_field = reader_schema.field(field.name)
|
||||
@@ -4116,16 +4132,18 @@ def _find_vector_columns(
|
||||
and reader_field.type.list_size >= 10
|
||||
)
|
||||
|
||||
if named_vector_col or typed_fixed_vector_col:
|
||||
if named_vector_col or typed_fixed_vector_col or is_multivector:
|
||||
vector_type = field.type.value_type if is_multivector else field.type
|
||||
vector_columns.append(
|
||||
{
|
||||
"name": field.name,
|
||||
"expected_dim": (
|
||||
field.type.list_size
|
||||
if pa.types.is_fixed_size_list(field.type)
|
||||
vector_type.list_size
|
||||
if pa.types.is_fixed_size_list(vector_type)
|
||||
else None
|
||||
),
|
||||
"expected_value_type": field.type.value_type,
|
||||
"expected_value_type": vector_type.value_type,
|
||||
"is_multivector": is_multivector,
|
||||
}
|
||||
)
|
||||
|
||||
@@ -4176,6 +4194,7 @@ def _handle_bad_vector_column(
|
||||
fill_value: float = 0.0,
|
||||
expected_dim: Optional[int] = None,
|
||||
expected_value_type: Optional[pa.DataType] = None,
|
||||
is_multivector: bool = False,
|
||||
) -> pa.RecordBatch:
|
||||
"""
|
||||
Ensure that the vector column exists and has type fixed_size_list(float)
|
||||
@@ -4196,6 +4215,8 @@ def _handle_bad_vector_column(
|
||||
vec_arr = data[vector_column_name]
|
||||
if not _is_list_like(vec_arr.type):
|
||||
return data
|
||||
if is_multivector and not _is_multivector_type(vec_arr.type):
|
||||
return data
|
||||
|
||||
if (
|
||||
expected_dim is not None
|
||||
@@ -4212,12 +4233,21 @@ def _handle_bad_vector_column(
|
||||
vec_arr = pa.array(vec_arr.to_pylist(), type=pa.list_(expected_value_type))
|
||||
data = data.set_column(position, vector_column_name, vec_arr)
|
||||
|
||||
if pa.types.is_floating(vec_arr.type.value_type):
|
||||
if is_multivector or pa.types.is_floating(vec_arr.type.value_type):
|
||||
has_nan = has_nan_values(vec_arr)
|
||||
else:
|
||||
has_nan = pa.array([False] * len(vec_arr))
|
||||
|
||||
if expected_dim is not None:
|
||||
if is_multivector:
|
||||
dim = (
|
||||
expected_dim
|
||||
if expected_dim is not None
|
||||
else _infer_vector_column_dim(vec_arr, True)
|
||||
)
|
||||
if dim is None:
|
||||
return data
|
||||
has_wrong_dim = _multivector_has_wrong_dim(vec_arr, dim)
|
||||
elif expected_dim is not None:
|
||||
dim = expected_dim
|
||||
elif pa.types.is_fixed_size_list(vec_arr.type):
|
||||
dim = vec_arr.type.list_size
|
||||
@@ -4226,15 +4256,16 @@ def _handle_bad_vector_column(
|
||||
if dim is None:
|
||||
return data
|
||||
|
||||
is_null = pc.is_null(vec_arr)
|
||||
# pc.list_value_length returns null for null list entries, so
|
||||
# pc.not_equal(null, dim) also returns null. Use or_kleene so that
|
||||
# True OR null = True (Kleene three-valued logic), ensuring null vectors
|
||||
# are counted as wrong-dim.
|
||||
has_wrong_dim = pc.or_kleene(
|
||||
is_null,
|
||||
pc.not_equal(pc.list_value_length(vec_arr), dim),
|
||||
)
|
||||
if not is_multivector:
|
||||
is_null = pc.is_null(vec_arr)
|
||||
# pc.list_value_length returns null for null list entries, so
|
||||
# pc.not_equal(null, dim) also returns null. Use or_kleene so that
|
||||
# True OR null = True (Kleene three-valued logic), ensuring null vectors
|
||||
# are counted as wrong-dim.
|
||||
has_wrong_dim = pc.or_kleene(
|
||||
is_null,
|
||||
pc.not_equal(pc.list_value_length(vec_arr), dim),
|
||||
)
|
||||
|
||||
has_bad_vectors = pc.any(has_nan).as_py() or pc.any(has_wrong_dim).as_py()
|
||||
|
||||
@@ -4269,7 +4300,10 @@ def _handle_bad_vector_column(
|
||||
raise ValueError(
|
||||
"`fill_value` must not be None if `on_bad_vectors` is 'fill'"
|
||||
)
|
||||
vec_arr = _fill_bad_vector_values(vec_arr, dim, fill_value)
|
||||
if is_multivector:
|
||||
vec_arr = _fill_bad_multivector_values(vec_arr, dim, fill_value)
|
||||
else:
|
||||
vec_arr = _fill_bad_vector_values(vec_arr, dim, fill_value)
|
||||
else:
|
||||
raise ValueError(f"Invalid value for on_bad_vectors: {on_bad_vectors}")
|
||||
|
||||
@@ -4321,17 +4355,60 @@ def _fill_bad_vector_values(
|
||||
return filled.cast(arr.type)
|
||||
|
||||
|
||||
def _fill_bad_multivector_values(
|
||||
arr: Union[pa.Array, pa.ChunkedArray], dim: int, fill_value: float
|
||||
) -> pa.Array:
|
||||
if not isinstance(arr, pa.ChunkedArray):
|
||||
arr = pa.chunked_array([arr])
|
||||
arr = arr.combine_chunks()
|
||||
|
||||
filled_vectors = _fill_bad_vector_values(arr.values, dim, fill_value)
|
||||
parent_nulls = pc.is_null(arr)
|
||||
if pa.types.is_large_list(arr.type):
|
||||
filled = pa.LargeListArray.from_arrays(
|
||||
arr.offsets, filled_vectors, mask=parent_nulls
|
||||
)
|
||||
else:
|
||||
filled = pa.ListArray.from_arrays(
|
||||
arr.offsets, filled_vectors, mask=parent_nulls
|
||||
)
|
||||
return filled.cast(arr.type)
|
||||
|
||||
|
||||
def _multivector_has_wrong_dim(
|
||||
arr: Union[pa.Array, pa.ChunkedArray], dim: int
|
||||
) -> pa.BooleanArray:
|
||||
if isinstance(arr, pa.ChunkedArray):
|
||||
results = [_multivector_has_wrong_dim(chunk, dim) for chunk in arr.chunks]
|
||||
return pa.concat_arrays(results) if results else pa.array([], type=pa.bool_())
|
||||
|
||||
vectors = arr.flatten()
|
||||
vector_is_wrong = pc.or_kleene(
|
||||
pc.is_null(vectors),
|
||||
pc.not_equal(pc.list_value_length(vectors), dim),
|
||||
)
|
||||
parent_indices = pc.list_parent_indices(arr)
|
||||
wrong_parent_indices = pc.unique(pc.filter(parent_indices, vector_is_wrong))
|
||||
indices = pa.array(range(len(arr)), type=pa.uint32())
|
||||
return pc.or_(pc.is_null(arr), pc.is_in(indices, wrong_parent_indices))
|
||||
|
||||
|
||||
def has_nan_values(arr: Union[pa.ListArray, pa.ChunkedArray]) -> pa.BooleanArray:
|
||||
if isinstance(arr, pa.ChunkedArray):
|
||||
values = pa.chunked_array([chunk.flatten() for chunk in arr.chunks])
|
||||
else:
|
||||
values = arr.flatten()
|
||||
if pa.types.is_float16(values.type):
|
||||
results = [has_nan_values(chunk) for chunk in arr.chunks]
|
||||
return pa.concat_arrays(results) if results else pa.array([], type=pa.bool_())
|
||||
|
||||
values = arr.flatten()
|
||||
if _is_list_like(values.type):
|
||||
values_has_nan = has_nan_values(values)
|
||||
elif pa.types.is_float16(values.type):
|
||||
# is_nan isn't yet implemented for f16, so we cast to f32
|
||||
# https://github.com/apache/arrow/issues/45083
|
||||
values_has_nan = pc.is_nan(values.cast(pa.float32()))
|
||||
else:
|
||||
elif pa.types.is_floating(values.type):
|
||||
values_has_nan = pc.is_nan(values)
|
||||
else:
|
||||
return pa.array([False] * len(arr))
|
||||
values_indices = pc.list_parent_indices(arr)
|
||||
has_nan_indices = pc.unique(pc.filter(values_indices, values_has_nan))
|
||||
indices = pa.array(range(len(arr)), type=pa.uint32())
|
||||
@@ -4346,6 +4423,16 @@ def _is_list_like(data_type: pa.DataType) -> bool:
|
||||
)
|
||||
|
||||
|
||||
def _is_float_vector_type(data_type: pa.DataType) -> bool:
|
||||
return _is_list_like(data_type) and pa.types.is_floating(data_type.value_type)
|
||||
|
||||
|
||||
def _is_multivector_type(data_type: pa.DataType) -> bool:
|
||||
return (
|
||||
pa.types.is_list(data_type) or pa.types.is_large_list(data_type)
|
||||
) and _is_float_vector_type(data_type.value_type)
|
||||
|
||||
|
||||
def _merge_metadata(*metadata_dicts: Optional[dict]) -> dict:
|
||||
merged = {}
|
||||
for metadata in metadata_dicts:
|
||||
@@ -4437,6 +4524,16 @@ def _infer_vector_dim(arr: Union[pa.Array, pa.ChunkedArray]) -> Optional[int]:
|
||||
return pc.mode(lengths)[0].as_py()["mode"]
|
||||
|
||||
|
||||
def _infer_vector_column_dim(
|
||||
arr: Union[pa.Array, pa.ChunkedArray], is_multivector: bool
|
||||
) -> Optional[int]:
|
||||
if not is_multivector:
|
||||
return _infer_vector_dim(arr)
|
||||
if isinstance(arr, pa.ChunkedArray):
|
||||
arr = arr.combine_chunks()
|
||||
return _infer_vector_dim(arr.flatten())
|
||||
|
||||
|
||||
def _validate_schema(schema: pa.Schema):
|
||||
"""
|
||||
Make sure the metadata is valid utf8
|
||||
|
||||
@@ -77,6 +77,23 @@ def test_sync_repr_does_not_use_background_loop(tmp_path, monkeypatch):
|
||||
assert repr(table) == f"LanceTable(name='test', _conn={db!r})"
|
||||
|
||||
|
||||
def test_read_consistency_interval_does_not_use_background_loop(tmp_path, monkeypatch):
|
||||
from lancedb.background_loop import LOOP
|
||||
from lancedb.db import LanceDBConnection
|
||||
|
||||
consistency_interval = timedelta(seconds=5)
|
||||
db = lancedb.connect(tmp_path, read_consistency_interval=consistency_interval)
|
||||
db_from_inner = LanceDBConnection.from_inner(db._inner, consistency_interval)
|
||||
|
||||
def fail_run(*args, **kwargs):
|
||||
raise AssertionError("properties should not use the Python background loop")
|
||||
|
||||
monkeypatch.setattr(LOOP, "run", fail_run)
|
||||
|
||||
assert db.read_consistency_interval == consistency_interval
|
||||
assert db_from_inner.read_consistency_interval == consistency_interval
|
||||
|
||||
|
||||
def test_ingest_pd(tmp_path):
|
||||
db = lancedb.connect(tmp_path)
|
||||
|
||||
|
||||
@@ -6,6 +6,7 @@ import math
|
||||
import pytest
|
||||
|
||||
from lancedb import DBConnection, Table, connect
|
||||
from lancedb.background_loop import LOOP
|
||||
from lancedb.permutation import Permutation, Permutations, permutation_builder
|
||||
|
||||
|
||||
@@ -31,6 +32,25 @@ def test_split_random_ratios(mem_db):
|
||||
assert 65 <= split_1_count <= 75 # ~70% ± tolerance
|
||||
|
||||
|
||||
def test_execute_does_not_reenter_background_loop(tmp_path, monkeypatch):
|
||||
import threading
|
||||
|
||||
db = connect(tmp_path)
|
||||
tbl = db.create_table("test_table", pa.table({"x": range(10)}))
|
||||
original_run = LOOP.run
|
||||
|
||||
def fail_on_reentry(future):
|
||||
assert threading.current_thread() is not LOOP.thread
|
||||
return original_run(future)
|
||||
|
||||
monkeypatch.setattr(LOOP, "run", fail_on_reentry)
|
||||
|
||||
permutation_tbl = permutation_builder(tbl).execute()
|
||||
|
||||
assert permutation_tbl.count_rows() == 10
|
||||
assert permutation_tbl._conn.read_consistency_interval is None
|
||||
|
||||
|
||||
def test_split_random_counts(mem_db):
|
||||
"""Test random splitting with absolute counts."""
|
||||
tbl = mem_db.create_table(
|
||||
|
||||
@@ -6,6 +6,7 @@ import os
|
||||
import sys
|
||||
import threading
|
||||
import warnings
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from datetime import date, datetime, timedelta
|
||||
from time import sleep
|
||||
from typing import List
|
||||
@@ -1709,6 +1710,19 @@ def test_create_with_nans(mem_db: DBConnection):
|
||||
assert np.allclose(filled_vectors[22.0], np.array([5.0, 0.0]))
|
||||
|
||||
|
||||
def test_create_with_nans_in_multivectors(mem_db: DBConnection):
|
||||
multivector_type = pa.list_(pa.list_(pa.float32(), 128))
|
||||
schema = pa.schema(
|
||||
[pa.field("filename", pa.string()), pa.field("vector", multivector_type)]
|
||||
)
|
||||
vector = [0.1] * 128
|
||||
vector[-1] = np.nan
|
||||
data = [{"filename": "img1.jpg", "vector": [vector]}]
|
||||
|
||||
with pytest.raises(RuntimeError, match="Vector column 'vector' has NaNs"):
|
||||
mem_db.create_table("nan_multivector", data=data, schema=schema)
|
||||
|
||||
|
||||
def test_add_with_nans(mem_db: DBConnection):
|
||||
schema = pa.schema(
|
||||
[
|
||||
@@ -2124,6 +2138,27 @@ def test_delete(mem_db: DBConnection):
|
||||
assert table.to_arrow()["id"].to_pylist() == [1]
|
||||
|
||||
|
||||
def test_concurrent_deletes_are_thread_safe(mem_db: DBConnection):
|
||||
num_workers = 8
|
||||
table = mem_db.create_table(
|
||||
"my_table", data=[{"id": row_id} for row_id in range(num_workers)]
|
||||
)
|
||||
barrier = threading.Barrier(num_workers)
|
||||
|
||||
def delete(row_id: int):
|
||||
barrier.wait()
|
||||
return table.delete(f"id = {row_id}")
|
||||
|
||||
with ThreadPoolExecutor(max_workers=num_workers) as pool:
|
||||
results = list(pool.map(delete, range(num_workers)))
|
||||
|
||||
assert all(result.num_deleted_rows == 1 for result in results)
|
||||
assert sorted(result.version for result in results) == list(
|
||||
range(2, num_workers + 2)
|
||||
)
|
||||
assert table.count_rows() == 0
|
||||
|
||||
|
||||
def test_delete_expr(mem_db: DBConnection):
|
||||
table = mem_db.create_table(
|
||||
"my_table",
|
||||
|
||||
@@ -400,6 +400,83 @@ def test_handle_bad_vectors_nan(on_bad_vectors):
|
||||
assert output["vector"].combine_chunks() == expected
|
||||
|
||||
|
||||
@pytest.mark.parametrize("on_bad_vectors", ["error", "drop", "fill", "null"])
|
||||
def test_handle_bad_multivectors_nan(on_bad_vectors):
|
||||
multivector_type = pa.list_(pa.list_(pa.float32(), 2))
|
||||
vectors = pa.array(
|
||||
[
|
||||
[[1.0, float("nan")], [2.0, 3.0]],
|
||||
[[4.0, 5.0]],
|
||||
],
|
||||
type=multivector_type,
|
||||
)
|
||||
data = pa.table({"vector": vectors})
|
||||
|
||||
if on_bad_vectors == "error":
|
||||
with pytest.raises(ValueError, match="Vector column 'vector' has NaNs"):
|
||||
_handle_bad_vectors(data.to_reader()).read_all()
|
||||
return
|
||||
|
||||
output = _handle_bad_vectors(
|
||||
data.to_reader(),
|
||||
on_bad_vectors=on_bad_vectors,
|
||||
fill_value=42.0,
|
||||
).read_all()
|
||||
|
||||
if on_bad_vectors == "drop":
|
||||
expected = pa.array([[[4.0, 5.0]]], type=multivector_type)
|
||||
elif on_bad_vectors == "fill":
|
||||
expected = pa.array(
|
||||
[[[1.0, 42.0], [2.0, 3.0]], [[4.0, 5.0]]],
|
||||
type=multivector_type,
|
||||
)
|
||||
else:
|
||||
expected = pa.array([None, [[4.0, 5.0]]], type=multivector_type)
|
||||
|
||||
assert output["vector"].combine_chunks() == expected
|
||||
|
||||
|
||||
@pytest.mark.parametrize("on_bad_vectors", ["error", "drop", "fill", "null"])
|
||||
def test_handle_bad_variable_multivectors(on_bad_vectors):
|
||||
target_type = pa.list_(pa.list_(pa.float32(), 2))
|
||||
vectors = pa.array(
|
||||
[
|
||||
[[1.0, float("nan")], [2.0, 3.0]],
|
||||
[[4.0]],
|
||||
[[5.0, 6.0]],
|
||||
]
|
||||
)
|
||||
data = pa.table({"vector": vectors})
|
||||
|
||||
if on_bad_vectors == "error":
|
||||
with pytest.raises(ValueError, match="variable length vectors"):
|
||||
_handle_bad_vectors(
|
||||
data.to_reader(),
|
||||
target_schema=pa.schema({"vector": target_type}),
|
||||
).read_all()
|
||||
return
|
||||
|
||||
output = _handle_bad_vectors(
|
||||
data.to_reader(),
|
||||
on_bad_vectors=on_bad_vectors,
|
||||
fill_value=42.0,
|
||||
target_schema=pa.schema({"vector": target_type}),
|
||||
).read_all()
|
||||
|
||||
if on_bad_vectors == "drop":
|
||||
expected = [[[5.0, 6.0]]]
|
||||
elif on_bad_vectors == "fill":
|
||||
expected = [
|
||||
[[1.0, 42.0], [2.0, 3.0]],
|
||||
[[4.0, 42.0]],
|
||||
[[5.0, 6.0]],
|
||||
]
|
||||
else:
|
||||
expected = [None, None, [[5.0, 6.0]]]
|
||||
|
||||
assert output["vector"].combine_chunks().to_pylist() == expected
|
||||
|
||||
|
||||
def test_handle_bad_vectors_noop():
|
||||
# ChunkedArray should be preserved as-is
|
||||
vector = pa.chunked_array(
|
||||
|
||||
+15
-2
@@ -745,6 +745,9 @@ impl Table {
|
||||
|
||||
#[allow(private_interfaces)]
|
||||
pub fn delete(self_: PyRef<'_, Self>, condition: PredicateArg) -> PyResult<Bound<'_, PyAny>> {
|
||||
// Do not hold the Python borrow across the await. The cloned Rust table
|
||||
// handle is thread-safe and allows deletes on the same Python table to
|
||||
// run concurrently without PyO3 reporting "Already borrowed".
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
let result = match &condition {
|
||||
@@ -1375,7 +1378,12 @@ impl Table {
|
||||
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
let result = inner.add_columns(definitions, None).await.infer_error()?;
|
||||
let result = inner
|
||||
.add_columns()
|
||||
.transform(definitions)
|
||||
.execute()
|
||||
.await
|
||||
.infer_error()?;
|
||||
Ok(AddColumnsResult::from(result))
|
||||
})
|
||||
}
|
||||
@@ -1389,7 +1397,12 @@ impl Table {
|
||||
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
let result = inner.add_columns(transform, None).await.infer_error()?;
|
||||
let result = inner
|
||||
.add_columns()
|
||||
.transform(transform)
|
||||
.execute()
|
||||
.await
|
||||
.infer_error()?;
|
||||
Ok(AddColumnsResult::from(result))
|
||||
})
|
||||
}
|
||||
|
||||
@@ -8,15 +8,12 @@ use std::fs::create_dir_all;
|
||||
use std::path::Path;
|
||||
use std::{collections::HashMap, sync::Arc};
|
||||
|
||||
use futures::TryStreamExt;
|
||||
use lance::dataset::refs::Ref;
|
||||
use lance::dataset::{ReadParams, WriteMode, builder::DatasetBuilder};
|
||||
use lance::io::{ObjectStore, ObjectStoreParams, WrappingObjectStore};
|
||||
use lance_datafusion::utils::StreamingWriteSource;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lance_io::object_store::{
|
||||
DirCursor, ReadDirOptions, StorageOptionsAccessor, StorageOptionsProvider,
|
||||
};
|
||||
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
|
||||
use lance_table::io::commit::commit_handler_from_url;
|
||||
use object_store::local::LocalFileSystem;
|
||||
use snafu::ResultExt;
|
||||
@@ -721,54 +718,6 @@ impl ListingDatabase {
|
||||
self.namespace_database.clone()
|
||||
}
|
||||
|
||||
/// List up to `limit` table names, resuming after the table named `start_after`.
|
||||
///
|
||||
/// The cursor and the page size go into the object store's list request rather than
|
||||
/// being applied to a full listing, so the cost of a page is set by the size of the
|
||||
/// page and not by the size of the database. Stores with no paginated list API fall
|
||||
/// back to a full listing, which is what this did for every store before.
|
||||
///
|
||||
/// Names come back in the order the store lists the directories in, which is by key:
|
||||
/// `foo-bar` precedes `foo`, because the `-` of `foo-bar.lance` sorts below the `.` of
|
||||
/// `foo.lance`. Pagination has to follow the order the cursor is pushed down in, so
|
||||
/// that is the order both listing methods report and the order `start_after` resumes
|
||||
/// in. It matches sorting by name except between a name and one that extends it.
|
||||
async fn list_table_dirs(
|
||||
&self,
|
||||
start_after: Option<&str>,
|
||||
limit: Option<usize>,
|
||||
) -> Result<Vec<String>> {
|
||||
let dir_suffix = format!(".{}", LANCE_EXTENSION);
|
||||
let options = ReadDirOptions {
|
||||
// An empty name means "from the start": that is how comparing names against it
|
||||
// behaved, and how a client looping on a page token spells its first request.
|
||||
// Built into a cursor it would instead sit after every name below `.lance`.
|
||||
resume_from: start_after
|
||||
.filter(|name| !name.is_empty())
|
||||
.map(|name| DirCursor::after_directory(format!("{name}{dir_suffix}"))),
|
||||
page_size: limit,
|
||||
};
|
||||
let mut entries = self
|
||||
.object_store
|
||||
.read_dir_stream(self.base_path.clone(), options);
|
||||
|
||||
let mut names = Vec::new();
|
||||
while limit.is_none_or(|limit| names.len() < limit) {
|
||||
let Some(entry) = entries.try_next().await? else {
|
||||
break;
|
||||
};
|
||||
// A table is the directory `<name>.lance`; anything else under the database
|
||||
// prefix belongs to something other than a table.
|
||||
if !entry.is_dir() {
|
||||
continue;
|
||||
}
|
||||
if let Some(name) = entry.name.strip_suffix(&dir_suffix) {
|
||||
names.push(name.to_string());
|
||||
}
|
||||
}
|
||||
Ok(names)
|
||||
}
|
||||
|
||||
async fn drop_tables(&self, names: Vec<String>) -> Result<()> {
|
||||
let object_store_params = ObjectStoreParams {
|
||||
storage_options_accessor: if self.storage_options.is_empty() {
|
||||
@@ -1010,37 +959,80 @@ impl Database for ListingDatabase {
|
||||
if !request.namespace_path.is_empty() {
|
||||
return self.namespace_database().table_names(request).await;
|
||||
}
|
||||
self.list_table_dirs(
|
||||
request.start_after.as_deref(),
|
||||
request.limit.map(|limit| limit as usize),
|
||||
)
|
||||
.await
|
||||
let mut f = self
|
||||
.object_store
|
||||
.read_dir(self.base_path.clone())
|
||||
.await?
|
||||
.iter()
|
||||
.map(Path::new)
|
||||
.filter(|path| {
|
||||
let is_lance = path
|
||||
.extension()
|
||||
.and_then(|e| e.to_str())
|
||||
.map(|e| e == LANCE_EXTENSION);
|
||||
is_lance.unwrap_or(false)
|
||||
})
|
||||
.filter_map(|p| p.file_stem().and_then(|s| s.to_str().map(String::from)))
|
||||
.collect::<Vec<String>>();
|
||||
f.sort();
|
||||
if let Some(start_after) = request.start_after {
|
||||
let index = f
|
||||
.iter()
|
||||
.position(|name| name.as_str() > start_after.as_str())
|
||||
.unwrap_or(f.len());
|
||||
f.drain(0..index);
|
||||
}
|
||||
if let Some(limit) = request.limit {
|
||||
f.truncate(limit as usize);
|
||||
}
|
||||
Ok(f)
|
||||
}
|
||||
|
||||
async fn list_tables(&self, request: ListTablesRequest) -> Result<ListTablesResponse> {
|
||||
if request.id.as_ref().map(|v| !v.is_empty()).unwrap_or(false) {
|
||||
return self.namespace_database().list_tables(request).await;
|
||||
}
|
||||
let limit = request.limit.map(|limit| limit as usize);
|
||||
// Reading one past the page is how we learn whether another page follows, without
|
||||
// a second request. The extra name is dropped before the response goes out.
|
||||
let mut tables = self
|
||||
.list_table_dirs(
|
||||
request.page_token.as_deref(),
|
||||
limit.map(|limit| limit.saturating_add(1)),
|
||||
)
|
||||
.await?;
|
||||
let mut f = self
|
||||
.object_store
|
||||
.read_dir(self.base_path.clone())
|
||||
.await?
|
||||
.iter()
|
||||
.map(Path::new)
|
||||
.filter(|path| {
|
||||
let is_lance = path
|
||||
.extension()
|
||||
.and_then(|e| e.to_str())
|
||||
.map(|e| e == LANCE_EXTENSION);
|
||||
is_lance.unwrap_or(false)
|
||||
})
|
||||
.filter_map(|p| p.file_stem().and_then(|s| s.to_str().map(String::from)))
|
||||
.collect::<Vec<String>>();
|
||||
f.sort();
|
||||
|
||||
let next_page_token = match limit {
|
||||
Some(limit) if tables.len() > limit => {
|
||||
tables.truncate(limit);
|
||||
tables.last().cloned()
|
||||
// Handle pagination with page_token
|
||||
if let Some(ref page_token) = request.page_token {
|
||||
let index = f
|
||||
.iter()
|
||||
.position(|name| name.as_str() > page_token.as_str())
|
||||
.unwrap_or(f.len());
|
||||
f.drain(0..index);
|
||||
}
|
||||
|
||||
// Determine if there's a next page
|
||||
let next_page_token = if let Some(limit) = request.limit {
|
||||
if f.len() > limit as usize {
|
||||
let token = f[limit as usize].clone();
|
||||
f.truncate(limit as usize);
|
||||
Some(token)
|
||||
} else {
|
||||
None
|
||||
}
|
||||
_ => None,
|
||||
} else {
|
||||
None
|
||||
};
|
||||
|
||||
Ok(ListTablesResponse {
|
||||
tables,
|
||||
tables: f,
|
||||
page_token: next_page_token,
|
||||
})
|
||||
}
|
||||
@@ -1330,130 +1322,6 @@ mod tests {
|
||||
(tempdir, db)
|
||||
}
|
||||
|
||||
async fn create_tables(db: &ListingDatabase, names: &[&str]) {
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
|
||||
for name in names {
|
||||
db.create_table(CreateTableRequest {
|
||||
name: name.to_string(),
|
||||
namespace_path: vec![],
|
||||
data: Box::new(RecordBatch::new_empty(schema.clone())) as Box<dyn Scannable>,
|
||||
mode: CreateTableMode::Create,
|
||||
write_options: Default::default(),
|
||||
location: None,
|
||||
namespace_client: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
/// Paging with the returned token has to visit every table exactly once. The token is
|
||||
/// the last name of the page, which is what `page_token` resumes after.
|
||||
#[tokio::test]
|
||||
async fn test_list_tables_pages_over_every_table_once() {
|
||||
let (_tempdir, db) = setup_database().await;
|
||||
create_tables(&db, &["a", "b", "c", "d", "e"]).await;
|
||||
|
||||
let mut seen = Vec::new();
|
||||
let mut page_token = None;
|
||||
loop {
|
||||
let page = db
|
||||
.list_tables(ListTablesRequest {
|
||||
limit: Some(2),
|
||||
page_token,
|
||||
..Default::default()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
seen.extend(page.tables);
|
||||
match page.page_token {
|
||||
Some(token) => page_token = Some(token),
|
||||
None => break,
|
||||
}
|
||||
}
|
||||
|
||||
assert_eq!(seen, vec!["a", "b", "c", "d", "e"]);
|
||||
}
|
||||
|
||||
/// The last page reports no token, so a caller paging by token knows to stop without
|
||||
/// asking for an empty page.
|
||||
#[tokio::test]
|
||||
async fn test_list_tables_exhausted_page_has_no_token() {
|
||||
let (_tempdir, db) = setup_database().await;
|
||||
create_tables(&db, &["a", "b"]).await;
|
||||
|
||||
let page = db
|
||||
.list_tables(ListTablesRequest {
|
||||
limit: Some(2),
|
||||
..Default::default()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(page.tables, vec!["a", "b"]);
|
||||
assert_eq!(page.page_token, None);
|
||||
}
|
||||
|
||||
/// Listing follows the order the object store lists directories in, so a name that
|
||||
/// extends another comes first: `-` sorts below the `.` of `.lance`. Pagination pushes
|
||||
/// its cursor into the list request, so it cannot report a different order than the
|
||||
/// one it resumes in.
|
||||
#[tokio::test]
|
||||
async fn test_listing_order_follows_the_store_not_the_table_name() {
|
||||
let (_tempdir, db) = setup_database().await;
|
||||
create_tables(&db, &["users", "users-archive", "users.old"]).await;
|
||||
|
||||
#[allow(deprecated)]
|
||||
let names = db.table_names(TableNamesRequest::default()).await.unwrap();
|
||||
assert_eq!(names, vec!["users-archive", "users", "users.old"]);
|
||||
|
||||
// Resuming after a name skips everything the store lists before it, which is what
|
||||
// paging by the previous page's last name relies on.
|
||||
#[allow(deprecated)]
|
||||
let after = db
|
||||
.table_names(TableNamesRequest {
|
||||
start_after: Some("users-archive".to_string()),
|
||||
..Default::default()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(after, vec!["users", "users.old"]);
|
||||
}
|
||||
|
||||
/// An empty `start_after` means "from the start". A name that sorts below `.lance` is
|
||||
/// what disappears if it is treated as a cursor instead.
|
||||
#[tokio::test]
|
||||
async fn test_empty_start_after_lists_from_the_start() {
|
||||
let (_tempdir, db) = setup_database().await;
|
||||
create_tables(&db, &["-dash", "alpha"]).await;
|
||||
|
||||
#[allow(deprecated)]
|
||||
let names = db
|
||||
.table_names(TableNamesRequest {
|
||||
start_after: Some(String::new()),
|
||||
..Default::default()
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(names, vec!["-dash", "alpha"]);
|
||||
}
|
||||
|
||||
/// Only directories named `<name>.lance` are tables; loose files and other directories
|
||||
/// under the database prefix are not.
|
||||
#[tokio::test]
|
||||
async fn test_listing_ignores_non_table_children() {
|
||||
let (tempdir, db) = setup_database().await;
|
||||
create_tables(&db, &["real"]).await;
|
||||
std::fs::write(tempdir.path().join("loose.lance"), b"not a table").unwrap();
|
||||
create_dir_all(tempdir.path().join("scratch")).unwrap();
|
||||
|
||||
#[allow(deprecated)]
|
||||
let names = db.table_names(TableNamesRequest::default()).await.unwrap();
|
||||
|
||||
assert_eq!(names, vec!["real"]);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_listing_database_root_ops_do_not_create_manifest() {
|
||||
let tempdir = tempdir().unwrap();
|
||||
|
||||
@@ -3089,10 +3089,12 @@ mod tests {
|
||||
Box::pin(table.delete("false").map_ok(|_| ())),
|
||||
Box::pin(
|
||||
table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![("x".into(), "y".into())]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"x".into(),
|
||||
"y".into(),
|
||||
)]))
|
||||
.execute()
|
||||
.map_ok(|_| ()),
|
||||
),
|
||||
Box::pin(async {
|
||||
@@ -6388,13 +6390,12 @@ mod tests {
|
||||
});
|
||||
|
||||
let result = table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![
|
||||
("b".into(), "a + 1".into()),
|
||||
("x".into(), "cast(NULL as int32)".into()),
|
||||
]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![
|
||||
("b".into(), "a + 1".into()),
|
||||
("x".into(), "cast(NULL as int32)".into()),
|
||||
]))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -7119,10 +7120,12 @@ mod tests {
|
||||
}
|
||||
"add_columns" => {
|
||||
let _ = table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![("c".into(), "a + 1".into())]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"c".into(),
|
||||
"a + 1".into(),
|
||||
)]))
|
||||
.execute()
|
||||
.await;
|
||||
}
|
||||
"drop_columns" => {
|
||||
@@ -9880,10 +9883,12 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
branch
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![("b".into(), "a + 1".into())]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"b".into(),
|
||||
"a + 1".into(),
|
||||
)]))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
branch
|
||||
|
||||
@@ -65,6 +65,7 @@ use crate::utils::{PatchReadParam, PatchWriteParam, resolve_arrow_field_path};
|
||||
use self::dataset::DatasetConsistencyWrapper;
|
||||
use self::merge::MergeInsertBuilder;
|
||||
|
||||
pub mod add_columns;
|
||||
mod add_data;
|
||||
pub mod branch_merge;
|
||||
mod create_index;
|
||||
@@ -79,6 +80,7 @@ pub mod schema_evolution;
|
||||
pub mod update;
|
||||
pub mod write_progress;
|
||||
use crate::index::waiter::wait_for_index;
|
||||
pub use add_columns::AddColumnsBuilder;
|
||||
#[cfg(feature = "remote")]
|
||||
pub(crate) use add_data::PreprocessingOutput;
|
||||
pub use add_data::{AddDataBuilder, AddDataMode, AddResult, NaNVectorBehavior};
|
||||
@@ -1620,12 +1622,8 @@ impl Table {
|
||||
}
|
||||
|
||||
/// Add new columns to the table, providing values to fill in.
|
||||
pub async fn add_columns(
|
||||
&self,
|
||||
transforms: NewColumnTransform,
|
||||
read_columns: Option<Vec<String>>,
|
||||
) -> Result<AddColumnsResult> {
|
||||
self.inner.add_columns(transforms, read_columns).await
|
||||
pub fn add_columns(&self) -> AddColumnsBuilder {
|
||||
AddColumnsBuilder::new(self.inner.clone())
|
||||
}
|
||||
|
||||
/// Change a column's name or nullability.
|
||||
|
||||
@@ -0,0 +1,161 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! Builder for adding columns to a table.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use lance::dataset::NewColumnTransform;
|
||||
|
||||
use super::BaseTable;
|
||||
use super::schema_evolution::AddColumnsResult;
|
||||
use crate::{Error, Result};
|
||||
|
||||
/// Adds columns to a table. See [`Table::add_columns`](super::Table::add_columns).
|
||||
pub struct AddColumnsBuilder {
|
||||
parent: Arc<dyn BaseTable>,
|
||||
transform: Option<NewColumnTransform>,
|
||||
read_columns: Option<Vec<String>>,
|
||||
}
|
||||
|
||||
impl std::fmt::Debug for AddColumnsBuilder {
|
||||
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||
f.debug_struct("AddColumnsBuilder")
|
||||
.field("parent", &self.parent)
|
||||
.field("has_transform", &self.transform.is_some())
|
||||
.field("read_columns", &self.read_columns)
|
||||
.finish()
|
||||
}
|
||||
}
|
||||
|
||||
impl AddColumnsBuilder {
|
||||
pub(crate) fn new(parent: Arc<dyn BaseTable>) -> Self {
|
||||
Self {
|
||||
parent,
|
||||
transform: None,
|
||||
read_columns: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Set how the new columns' values are produced. Required.
|
||||
pub fn transform(mut self, transform: NewColumnTransform) -> Self {
|
||||
self.transform = Some(transform);
|
||||
self
|
||||
}
|
||||
|
||||
/// Limit which existing columns a [`NewColumnTransform::BatchUDF`] mapper
|
||||
/// receives. Every other transform determines what it reads, so setting
|
||||
/// this alongside one is an error rather than a silent no-op.
|
||||
pub fn read_columns(mut self, columns: impl IntoIterator<Item = impl Into<String>>) -> Self {
|
||||
self.read_columns = Some(columns.into_iter().map(Into::into).collect());
|
||||
self
|
||||
}
|
||||
|
||||
/// Add the columns.
|
||||
pub async fn execute(self) -> Result<AddColumnsResult> {
|
||||
let Self {
|
||||
parent,
|
||||
transform,
|
||||
read_columns,
|
||||
} = self;
|
||||
|
||||
let Some(transform) = transform else {
|
||||
return Err(Error::InvalidInput {
|
||||
message: "add_columns requires a transform".into(),
|
||||
});
|
||||
};
|
||||
|
||||
if read_columns.is_some() && !matches!(transform, NewColumnTransform::BatchUDF(_)) {
|
||||
return Err(Error::InvalidInput {
|
||||
message: "read_columns applies only to a BatchUDF transform; \
|
||||
every other transform determines what it reads"
|
||||
.into(),
|
||||
});
|
||||
}
|
||||
|
||||
parent.add_columns(transform, read_columns).await
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::sync::Arc;
|
||||
|
||||
use arrow_array::{Int32Array, RecordBatch, record_batch};
|
||||
use arrow_schema::{DataType, Field, Schema};
|
||||
use lance::dataset::{BatchUDF, NewColumnTransform};
|
||||
|
||||
use crate::Table;
|
||||
use crate::connect;
|
||||
|
||||
async fn table_with_two_columns(name: &str) -> Table {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let batch = record_batch!(("x", Int32, [1, 2, 3]), ("y", Int32, [10, 20, 30])).unwrap();
|
||||
conn.create_table(name, batch).execute().await.unwrap()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_requires_a_transform() {
|
||||
let table = table_with_two_columns("no_transform").await;
|
||||
let err = table.add_columns().execute().await.unwrap_err();
|
||||
assert!(
|
||||
err.to_string().contains("requires a transform"),
|
||||
"got: {err}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_read_columns_with_sql_expressions_is_rejected() {
|
||||
let table = table_with_two_columns("read_cols_sql").await;
|
||||
let err = table
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"doubled".into(),
|
||||
"x * 2".into(),
|
||||
)]))
|
||||
.read_columns(["x"])
|
||||
.execute()
|
||||
.await
|
||||
.unwrap_err();
|
||||
assert!(err.to_string().contains("BatchUDF"), "got: {err}");
|
||||
|
||||
let schema = table.schema().await.unwrap();
|
||||
assert!(
|
||||
schema.field_with_name("doubled").is_err(),
|
||||
"a rejected call must not commit"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_read_columns_limits_what_a_batch_udf_sees() {
|
||||
let table = table_with_two_columns("read_cols_udf").await;
|
||||
|
||||
let output_schema = Arc::new(Schema::new(vec![Field::new("sum", DataType::Int32, true)]));
|
||||
let mapper_schema = output_schema.clone();
|
||||
let udf = BatchUDF {
|
||||
mapper: Box::new(move |batch: &RecordBatch| {
|
||||
assert!(batch.column_by_name("x").is_some());
|
||||
assert!(batch.column_by_name("y").is_none(), "y was not requested");
|
||||
let x = batch["x"].as_any().downcast_ref::<Int32Array>().unwrap();
|
||||
let doubled: Int32Array = x.iter().map(|v| v.map(|v| v * 2)).collect();
|
||||
Ok(RecordBatch::try_new(
|
||||
mapper_schema.clone(),
|
||||
vec![Arc::new(doubled)],
|
||||
)?)
|
||||
}),
|
||||
output_schema,
|
||||
result_checkpoint: None,
|
||||
};
|
||||
|
||||
table
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::BatchUDF(udf))
|
||||
.read_columns(["x"])
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let schema = table.schema().await.unwrap();
|
||||
assert!(schema.field_with_name("sum").is_ok());
|
||||
}
|
||||
}
|
||||
@@ -258,6 +258,7 @@ mod tests {
|
||||
FixedSizeListArray, Float32Array, Int32Array, LargeStringArray, ListArray, RecordBatch,
|
||||
RecordBatchIterator, record_batch,
|
||||
};
|
||||
use arrow_buffer::OffsetBuffer;
|
||||
use arrow_schema::{ArrowError, DataType, Field, Schema};
|
||||
use futures::TryStreamExt;
|
||||
use lance::dataset::{WriteMode, WriteParams};
|
||||
@@ -576,10 +577,12 @@ mod tests {
|
||||
|
||||
// Add a new physical column AFTER the embedding column.
|
||||
table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![("score".into(), "42.0".into())]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"score".into(),
|
||||
"42.0".into(),
|
||||
)]))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -683,7 +686,9 @@ mod tests {
|
||||
true,
|
||||
)]));
|
||||
table
|
||||
.add_columns(NewColumnTransform::AllNulls(nested_schema), None)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::AllNulls(nested_schema))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -881,6 +886,63 @@ mod tests {
|
||||
assert_eq!(row_count, 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_add_rejects_nan_multivectors() {
|
||||
let vector_type =
|
||||
DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Float32, true)), 4);
|
||||
let schema = Arc::new(Schema::new(vec![Field::new(
|
||||
"embedding",
|
||||
DataType::List(Arc::new(Field::new("item", vector_type.clone(), true))),
|
||||
false,
|
||||
)]));
|
||||
|
||||
let db = connect("memory://").execute().await.unwrap();
|
||||
let table = db
|
||||
.create_empty_table("nan_multivector_test", schema.clone())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let vectors = FixedSizeListArray::try_new(
|
||||
Arc::new(Field::new("item", DataType::Float32, true)),
|
||||
4,
|
||||
Arc::new(Float32Array::from(vec![
|
||||
0.1,
|
||||
0.2,
|
||||
0.3,
|
||||
0.4,
|
||||
0.5,
|
||||
f32::NAN,
|
||||
0.7,
|
||||
0.8,
|
||||
])),
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
let multivectors = ListArray::try_new(
|
||||
Arc::new(Field::new("item", vector_type, true)),
|
||||
OffsetBuffer::from_lengths([2]),
|
||||
Arc::new(vectors),
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
let batch = RecordBatch::try_new(schema, vec![Arc::new(multivectors)]).unwrap();
|
||||
|
||||
let err = table.add(batch.clone()).execute().await.unwrap_err();
|
||||
assert!(
|
||||
err.to_string().contains("NaN"),
|
||||
"Expected error mentioning NaN values, but got: {err:?}"
|
||||
);
|
||||
|
||||
table
|
||||
.add(batch)
|
||||
.on_nan_vectors(NaNVectorBehavior::Keep)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(table.count_rows(None).await.unwrap(), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_add_subschema() {
|
||||
let data = record_batch!(("id", Int64, [4, 5]), ("text", Utf8, ["foo", "bar"])).unwrap();
|
||||
|
||||
@@ -5,7 +5,7 @@
|
||||
|
||||
use std::sync::{Arc, LazyLock};
|
||||
|
||||
use arrow_array::{Array, FixedSizeListArray};
|
||||
use arrow_array::{Array, FixedSizeListArray, ListArray};
|
||||
use arrow_schema::{DataType, Field, FieldRef};
|
||||
use datafusion_common::config::ConfigOptions;
|
||||
use datafusion_expr::{ColumnarValue, ScalarFunctionArgs, ScalarUDFImpl, Signature, Volatility};
|
||||
@@ -19,16 +19,20 @@ use crate::{Error, Result};
|
||||
static REJECT_NAN_UDF: LazyLock<Arc<datafusion_expr::ScalarUDF>> =
|
||||
LazyLock::new(|| Arc::new(datafusion_expr::ScalarUDF::from(RejectNanUdf::new())));
|
||||
|
||||
/// Returns true if the field is a vector column: FixedSizeList<Float16/32/64>.
|
||||
/// Returns true if the field is a vector or multivector column.
|
||||
fn is_vector_field(field: &Field) -> bool {
|
||||
if let DataType::FixedSizeList(child, _) = field.data_type() {
|
||||
matches!(
|
||||
child.data_type(),
|
||||
DataType::Float16 | DataType::Float32 | DataType::Float64
|
||||
)
|
||||
} else {
|
||||
false
|
||||
fn is_vector_data_type(data_type: &DataType) -> bool {
|
||||
match data_type {
|
||||
DataType::FixedSizeList(child, _) => matches!(
|
||||
child.data_type(),
|
||||
DataType::Float16 | DataType::Float32 | DataType::Float64
|
||||
),
|
||||
DataType::List(child) => is_vector_data_type(child.data_type()),
|
||||
_ => false,
|
||||
}
|
||||
}
|
||||
|
||||
is_vector_data_type(field.data_type())
|
||||
}
|
||||
|
||||
/// Wraps the input plan with a projection that checks vector columns for NaN values.
|
||||
@@ -69,8 +73,8 @@ pub fn reject_nan_vectors(input: Arc<dyn ExecutionPlan>) -> Result<Arc<dyn Execu
|
||||
Ok(Arc::new(projection))
|
||||
}
|
||||
|
||||
/// A scalar UDF that passes through FixedSizeList arrays unchanged, but errors
|
||||
/// if any float values in the list are NaN.
|
||||
/// A scalar UDF that passes through vector arrays unchanged, but errors if any
|
||||
/// float values in the vector or multivector are NaN.
|
||||
#[derive(Debug, Hash, PartialEq, Eq)]
|
||||
struct RejectNanUdf {
|
||||
signature: Signature,
|
||||
@@ -113,44 +117,54 @@ impl ScalarUDFImpl for RejectNanUdf {
|
||||
}
|
||||
|
||||
fn check_no_nans(array: &dyn Array) -> datafusion_common::Result<()> {
|
||||
let fsl = array
|
||||
.as_any()
|
||||
.downcast_ref::<FixedSizeListArray>()
|
||||
.ok_or_else(|| {
|
||||
datafusion_common::DataFusionError::Internal(
|
||||
"reject_nan expected FixedSizeList".to_string(),
|
||||
)
|
||||
})?;
|
||||
|
||||
// Only inspect elements that are both in a valid parent row and non-null
|
||||
// themselves. Values backing null parent rows or null child elements may
|
||||
// contain garbage (including NaN) per the Arrow spec.
|
||||
let has_nan = (0..fsl.len()).filter(|i| fsl.is_valid(*i)).any(|i| {
|
||||
let row = fsl.value(i);
|
||||
match row.data_type() {
|
||||
DataType::Float16 => row
|
||||
fn contains_nan(array: &dyn Array) -> datafusion_common::Result<bool> {
|
||||
match array.data_type() {
|
||||
DataType::Float16 => Ok(array
|
||||
.as_any()
|
||||
.downcast_ref::<arrow_array::Float16Array>()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.any(|v| v.is_some_and(|v| v.is_nan())),
|
||||
DataType::Float32 => row
|
||||
.any(|v| v.is_some_and(|v| v.is_nan()))),
|
||||
DataType::Float32 => Ok(array
|
||||
.as_any()
|
||||
.downcast_ref::<arrow_array::Float32Array>()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.any(|v| v.is_some_and(|v| v.is_nan())),
|
||||
DataType::Float64 => row
|
||||
.any(|v| v.is_some_and(|v| v.is_nan()))),
|
||||
DataType::Float64 => Ok(array
|
||||
.as_any()
|
||||
.downcast_ref::<arrow_array::Float64Array>()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.any(|v| v.is_some_and(|v| v.is_nan())),
|
||||
_ => false,
|
||||
.any(|v| v.is_some_and(|v| v.is_nan()))),
|
||||
DataType::FixedSizeList(_, _) => {
|
||||
let lists = array.as_any().downcast_ref::<FixedSizeListArray>().unwrap();
|
||||
for i in (0..lists.len()).filter(|i| lists.is_valid(*i)) {
|
||||
if contains_nan(lists.value(i).as_ref())? {
|
||||
return Ok(true);
|
||||
}
|
||||
}
|
||||
Ok(false)
|
||||
}
|
||||
DataType::List(_) => {
|
||||
let lists = array.as_any().downcast_ref::<ListArray>().unwrap();
|
||||
for i in (0..lists.len()).filter(|i| lists.is_valid(*i)) {
|
||||
if contains_nan(lists.value(i).as_ref())? {
|
||||
return Ok(true);
|
||||
}
|
||||
}
|
||||
Ok(false)
|
||||
}
|
||||
data_type => Err(datafusion_common::DataFusionError::Internal(format!(
|
||||
"reject_nan expected a vector or multivector, got {data_type}"
|
||||
))),
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
if has_nan {
|
||||
// Only inspect elements that are both in a valid parent row and non-null
|
||||
// themselves. Values backing null parent rows or null child elements may
|
||||
// contain garbage (including NaN) per the Arrow spec.
|
||||
if contains_nan(array)? {
|
||||
return Err(datafusion_common::DataFusionError::ArrowError(
|
||||
Box::new(arrow_schema::ArrowError::ComputeError(
|
||||
"Vector column contains NaN values".to_string(),
|
||||
@@ -166,6 +180,7 @@ fn check_no_nans(array: &dyn Array) -> datafusion_common::Result<()> {
|
||||
mod tests {
|
||||
use super::*;
|
||||
use arrow_array::Float32Array;
|
||||
use arrow_buffer::OffsetBuffer;
|
||||
|
||||
#[test]
|
||||
fn test_passes_clean_vectors() {
|
||||
@@ -191,6 +206,46 @@ mod tests {
|
||||
assert!(check_no_nans(&fsl).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_rejects_nan_multivectors() {
|
||||
let vectors = FixedSizeListArray::try_new(
|
||||
Arc::new(Field::new("item", DataType::Float32, true)),
|
||||
2,
|
||||
Arc::new(Float32Array::from(vec![1.0, 2.0, 3.0, f32::NAN])),
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
let multivectors = ListArray::try_new(
|
||||
Arc::new(Field::new("item", vectors.data_type().clone(), true)),
|
||||
OffsetBuffer::from_lengths([2]),
|
||||
Arc::new(vectors),
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert!(check_no_nans(&multivectors).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_skips_null_multivector_rows() {
|
||||
let vectors = FixedSizeListArray::try_new(
|
||||
Arc::new(Field::new("item", DataType::Float32, true)),
|
||||
2,
|
||||
Arc::new(Float32Array::from(vec![f32::NAN, f32::NAN, 1.0, 2.0])),
|
||||
None,
|
||||
)
|
||||
.unwrap();
|
||||
let multivectors = ListArray::try_new(
|
||||
Arc::new(Field::new("item", vectors.data_type().clone(), true)),
|
||||
OffsetBuffer::from_lengths([1, 1]),
|
||||
Arc::new(vectors),
|
||||
Some(vec![false, true].into()),
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
assert!(check_no_nans(&multivectors).is_ok());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn test_skips_null_rows() {
|
||||
// Values backing null rows may contain NaN per the Arrow spec.
|
||||
@@ -254,6 +309,15 @@ mod tests {
|
||||
DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Float64, true)), 4),
|
||||
false,
|
||||
)));
|
||||
assert!(is_vector_field(&Field::new(
|
||||
"v",
|
||||
DataType::List(Arc::new(Field::new(
|
||||
"item",
|
||||
DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Float32, true)), 4,),
|
||||
true,
|
||||
))),
|
||||
false,
|
||||
)));
|
||||
assert!(!is_vector_field(&Field::new("id", DataType::Int32, false)));
|
||||
assert!(!is_vector_field(&Field::new(
|
||||
"v",
|
||||
|
||||
@@ -193,10 +193,12 @@ mod tests {
|
||||
|
||||
// Add a computed column
|
||||
let result = table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![("doubled".into(), "id * 2".into())]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"doubled".into(),
|
||||
"id * 2".into(),
|
||||
)]))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -251,13 +253,12 @@ mod tests {
|
||||
|
||||
// Add multiple columns at once
|
||||
table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![
|
||||
("y".into(), "x + 1".into()),
|
||||
("z".into(), "x * x".into()),
|
||||
]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![
|
||||
("y".into(), "x + 1".into()),
|
||||
("z".into(), "x * x".into()),
|
||||
]))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -283,10 +284,12 @@ mod tests {
|
||||
|
||||
// Add a column with a constant value
|
||||
table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![("constant".into(), "42".into())]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"constant".into(),
|
||||
"42".into(),
|
||||
)]))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
@@ -659,10 +662,12 @@ mod tests {
|
||||
|
||||
// Add column increments version
|
||||
let add_result = table
|
||||
.add_columns(
|
||||
NewColumnTransform::SqlExpressions(vec![("c".into(), "a + b".into())]),
|
||||
None,
|
||||
)
|
||||
.add_columns()
|
||||
.transform(NewColumnTransform::SqlExpressions(vec![(
|
||||
"c".into(),
|
||||
"a + b".into(),
|
||||
)]))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
assert!(add_result.version > v1);
|
||||
|
||||
@@ -9,14 +9,17 @@ use arrow_array::{
|
||||
};
|
||||
use arrow_schema::{DataType, Field, Fields, Schema};
|
||||
use futures::TryStreamExt;
|
||||
use lance::Dataset;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lancedb::{
|
||||
Connection, Error, Result, Table,
|
||||
blob::{BlobRangeRequest, blob},
|
||||
connect, connect_namespace,
|
||||
database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
|
||||
database::listing::{
|
||||
ListingDatabaseOptions, NewTableConfig, OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
|
||||
},
|
||||
query::{ExecutableQuery, QueryBase},
|
||||
table::{AddDataMode, CompactionOptions, OptimizeAction},
|
||||
table::{AddDataMode, CompactionOptions, OptimizeAction, OptimizeStats},
|
||||
};
|
||||
use tempfile::tempdir;
|
||||
|
||||
@@ -1075,3 +1078,252 @@ async fn fetch_blob_files_aligns_across_fragments_with_nulls_and_dups() -> Resul
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Rows exercising the null/empty interleavings from
|
||||
/// <https://github.com/lancedb/lancedb/issues/3744>: a payload, a null, a valid
|
||||
/// empty value, then payloads whose descriptors a fragment rewrite used to zero.
|
||||
fn null_empty_input_batch() -> RecordBatch {
|
||||
let owned = [
|
||||
Some(dedicated_blob_bytes(1)),
|
||||
None,
|
||||
Some(Vec::new()),
|
||||
Some(dedicated_blob_bytes(4)),
|
||||
Some(dedicated_blob_bytes(5)),
|
||||
Some(dedicated_blob_bytes(6)),
|
||||
];
|
||||
let payloads: Vec<Option<&[u8]>> = owned.iter().map(|payload| payload.as_deref()).collect();
|
||||
binary_input_batch(&[1, 2, 3, 4, 5, 6], &payloads)
|
||||
}
|
||||
|
||||
/// One `(id, Some((payload length, first byte)))` per live row, or `(id, None)`
|
||||
/// for a null blob. Comparing lengths and first bytes keeps failure output
|
||||
/// readable where comparing whole payloads would not.
|
||||
type BlobSummary = Vec<(i64, Option<(usize, Option<u8>)>)>;
|
||||
|
||||
/// The rows [`null_empty_input_batch`] leaves behind after `id IN (1, 4)` is
|
||||
/// deleted: a null, a valid empty value, and the two payloads that follow them.
|
||||
fn expected_null_empty_survivors() -> BlobSummary {
|
||||
vec![
|
||||
(2, None),
|
||||
(3, Some((0, None))),
|
||||
(5, Some((DEDICATED_BLOB_LEN, Some(5)))),
|
||||
(6, Some((DEDICATED_BLOB_LEN, Some(6)))),
|
||||
]
|
||||
}
|
||||
|
||||
/// `optimize()` only rewrites a fragment when lance's compaction planner selects
|
||||
/// it — here because the delete pushes the fragment past
|
||||
/// `materialize_deletions_threshold` (0.1 by default; these tests delete 2 of 6
|
||||
/// rows). Without this check, a planner or threshold change upstream would leave
|
||||
/// both regression tests green while no rewrite happened at all.
|
||||
fn assert_compacted(stats: &OptimizeStats) {
|
||||
let metrics = stats
|
||||
.compaction
|
||||
.as_ref()
|
||||
.expect("OptimizeAction::All runs compaction");
|
||||
assert!(
|
||||
metrics.fragments_removed >= 1,
|
||||
"optimize() rewrote no fragment, so this test proves nothing: {metrics:?}"
|
||||
);
|
||||
}
|
||||
|
||||
fn summarize(rows: &[(i64, Option<Vec<u8>>)]) -> BlobSummary {
|
||||
rows.iter()
|
||||
.map(|(id, payload)| {
|
||||
(
|
||||
*id,
|
||||
payload
|
||||
.as_ref()
|
||||
.map(|bytes| (bytes.len(), bytes.first().copied())),
|
||||
)
|
||||
})
|
||||
.collect()
|
||||
}
|
||||
|
||||
async fn sorted_id_rowid(table: &Table) -> Result<Vec<(i64, u64)>> {
|
||||
let mut pairs = collect_id_rowid(table).await?;
|
||||
pairs.sort_by_key(|(id, _)| *id);
|
||||
Ok(pairs)
|
||||
}
|
||||
|
||||
/// `{position, size}` descriptors of a legacy v1 blob column, keyed by `id`.
|
||||
async fn v1_blob_descriptors(table: &Table) -> Result<Vec<(i64, Option<(u64, u64)>)>> {
|
||||
let batches = table
|
||||
.query()
|
||||
.execute()
|
||||
.await?
|
||||
.try_collect::<Vec<_>>()
|
||||
.await?;
|
||||
let batch = arrow_select::concat::concat_batches(&batches[0].schema(), &batches).unwrap();
|
||||
let ids = batch
|
||||
.column_by_name("id")
|
||||
.unwrap()
|
||||
.as_any()
|
||||
.downcast_ref::<Int64Array>()
|
||||
.unwrap();
|
||||
let descriptors = batch
|
||||
.column_by_name("image")
|
||||
.unwrap()
|
||||
.as_any()
|
||||
.downcast_ref::<StructArray>()
|
||||
.expect("v1 blob column reads back as a descriptor struct");
|
||||
let position = descriptors
|
||||
.column_by_name("position")
|
||||
.unwrap()
|
||||
.as_any()
|
||||
.downcast_ref::<UInt64Array>()
|
||||
.unwrap();
|
||||
let size = descriptors
|
||||
.column_by_name("size")
|
||||
.unwrap()
|
||||
.as_any()
|
||||
.downcast_ref::<UInt64Array>()
|
||||
.unwrap();
|
||||
let mut rows: Vec<(i64, Option<(u64, u64)>)> = (0..batch.num_rows())
|
||||
.map(|row| {
|
||||
let descriptor =
|
||||
(!descriptors.is_null(row)).then(|| (position.value(row), size.value(row)));
|
||||
(ids.value(row), descriptor)
|
||||
})
|
||||
.collect();
|
||||
rows.sort_by_key(|(id, _)| *id);
|
||||
Ok(rows)
|
||||
}
|
||||
|
||||
/// Payload bytes of every live row of a legacy v1 blob column, keyed by `id`.
|
||||
/// [`Table::fetch_blobs`] rejects v1 columns, so read them through lance.
|
||||
async fn v1_blob_payloads(dataset_uri: &str, table: &Table) -> Result<Vec<(i64, Option<Vec<u8>>)>> {
|
||||
let pairs = sorted_id_rowid(table).await?;
|
||||
let row_ids: Vec<u64> = pairs.iter().map(|(_, row_id)| *row_id).collect();
|
||||
let dataset = Arc::new(Dataset::open(dataset_uri).await?);
|
||||
let files = dataset.take_blobs(&row_ids, "image").await?;
|
||||
assert_eq!(
|
||||
files.len(),
|
||||
pairs.len(),
|
||||
"take_blobs returned {} handles for {} live rows",
|
||||
files.len(),
|
||||
pairs.len()
|
||||
);
|
||||
let mut rows = Vec::with_capacity(pairs.len());
|
||||
for ((id, _), file) in pairs.iter().zip(files) {
|
||||
let payload = match file {
|
||||
Some(file) => Some(file.read().await?.to_vec()),
|
||||
None => None,
|
||||
};
|
||||
rows.push((*id, payload));
|
||||
}
|
||||
Ok(rows)
|
||||
}
|
||||
|
||||
/// Length and first byte of every live blob v2 value, keyed by `id`.
|
||||
async fn blob_v2_values(table: &Table) -> Result<BlobSummary> {
|
||||
let pairs = sorted_id_rowid(table).await?;
|
||||
let row_ids: Vec<u64> = pairs.iter().map(|(_, row_id)| *row_id).collect();
|
||||
let bytes = table.fetch_blobs("image", &row_ids).await?;
|
||||
Ok(pairs
|
||||
.iter()
|
||||
.enumerate()
|
||||
.map(|(slot, (id, _))| {
|
||||
let value = (!bytes.is_null(slot))
|
||||
.then(|| (bytes.value(slot).len(), bytes.value(slot).first().copied()));
|
||||
(*id, value)
|
||||
})
|
||||
.collect())
|
||||
}
|
||||
|
||||
/// Regression test for [#3744]: on storage 2.0 (legacy v1 descriptors),
|
||||
/// compaction rewrote every payload following a null or empty value in the same
|
||||
/// fragment as `{position: 0, size: 0}`, so the payload bytes read back as `b""`
|
||||
/// and the new fragment no longer referenced them at all.
|
||||
///
|
||||
/// [#3744]: https://github.com/lancedb/lancedb/issues/3744
|
||||
#[tokio::test]
|
||||
async fn optimize_preserves_v1_blob_payloads_with_null_and_empty() -> Result<()> {
|
||||
let tmp = tempdir().unwrap();
|
||||
let db_uri = tmp.path().to_str().unwrap().to_string();
|
||||
let db = connect(&db_uri)
|
||||
.database_options(&ListingDatabaseOptions {
|
||||
new_table_config: NewTableConfig {
|
||||
data_storage_version: Some(LanceFileVersion::V2_0),
|
||||
..Default::default()
|
||||
},
|
||||
..Default::default()
|
||||
})
|
||||
.execute()
|
||||
.await?;
|
||||
let legacy = Field::new("image", DataType::LargeBinary, true).with_metadata(
|
||||
std::collections::HashMap::from([("lance-encoding:blob".to_string(), "true".to_string())]),
|
||||
);
|
||||
let schema = Arc::new(Schema::new(vec![
|
||||
Field::new("id", DataType::Int64, false),
|
||||
legacy,
|
||||
]));
|
||||
let table = db.create_empty_table("t", schema).execute().await?;
|
||||
table.add(null_empty_input_batch()).execute().await?;
|
||||
assert_eq!(
|
||||
storage_format_version(&table).await,
|
||||
LanceFileVersion::V2_0.resolve(),
|
||||
"v1 blob descriptors only exist below storage 2.2"
|
||||
);
|
||||
let dataset_uri = table.uri().await?;
|
||||
|
||||
// Any rewrite triggers it; deleting rows is the shape from the issue.
|
||||
table.delete("id IN (1, 4)").await?;
|
||||
let descriptors_before = v1_blob_descriptors(&table).await?;
|
||||
let before = v1_blob_payloads(&dataset_uri, &table).await?;
|
||||
assert_eq!(
|
||||
summarize(&before),
|
||||
expected_null_empty_survivors(),
|
||||
"test setup no longer produces the null/empty/payload mix"
|
||||
);
|
||||
|
||||
let stats = table.optimize(OptimizeAction::All).await?;
|
||||
assert_compacted(&stats);
|
||||
|
||||
let descriptors_after = v1_blob_descriptors(&table).await?;
|
||||
let after = v1_blob_payloads(&dataset_uri, &table).await?;
|
||||
assert_eq!(
|
||||
summarize(&after),
|
||||
summarize(&before),
|
||||
"optimize() lost blob payloads; descriptors before={descriptors_before:?} after={descriptors_after:?}"
|
||||
);
|
||||
assert!(after == before, "optimize() changed blob payload bytes");
|
||||
Ok(())
|
||||
}
|
||||
|
||||
/// Regression test for the blob v2 half of [#3744]: compaction rewrote a valid
|
||||
/// empty value as null, destroying the null-vs-empty distinction.
|
||||
///
|
||||
/// [#3744]: https://github.com/lancedb/lancedb/issues/3744
|
||||
#[tokio::test]
|
||||
async fn optimize_preserves_blob_v2_null_and_empty_distinction() -> Result<()> {
|
||||
let tmp = tempdir().unwrap();
|
||||
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
|
||||
let table = db
|
||||
.create_empty_table("t", blob_table_schema())
|
||||
.execute()
|
||||
.await?;
|
||||
table.add(null_empty_input_batch()).execute().await?;
|
||||
assert!(
|
||||
storage_format_version(&table).await >= LanceFileVersion::V2_2,
|
||||
"blob v2 columns require storage >= 2.2"
|
||||
);
|
||||
|
||||
table.delete("id IN (1, 4)").await?;
|
||||
let before = blob_v2_values(&table).await?;
|
||||
assert_eq!(
|
||||
before,
|
||||
expected_null_empty_survivors(),
|
||||
"test setup no longer produces the null/empty/payload mix"
|
||||
);
|
||||
|
||||
let stats = table.optimize(OptimizeAction::All).await?;
|
||||
assert_compacted(&stats);
|
||||
|
||||
assert_eq!(
|
||||
blob_v2_values(&table).await?,
|
||||
before,
|
||||
"optimize() changed blob v2 values"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user