mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-11 15:52:17 +00:00
Compare commits
1
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d3077b7641 |
@@ -1,20 +0,0 @@
|
||||
name: Typo checker
|
||||
on:
|
||||
push:
|
||||
branches:
|
||||
- main
|
||||
pull_request:
|
||||
|
||||
permissions:
|
||||
contents: read
|
||||
|
||||
jobs:
|
||||
run:
|
||||
name: Spell Check with Typos
|
||||
runs-on: ubuntu-latest
|
||||
steps:
|
||||
- name: Check out code
|
||||
uses: actions/checkout@v6
|
||||
|
||||
- name: Check spelling of the entire repository
|
||||
uses: crate-ci/typos@6802cc60d4e7f78b9d5454f6cf3935c042d5e1e3 # v1.26.0
|
||||
@@ -10,10 +10,6 @@ repos:
|
||||
rev: v0.9.9
|
||||
hooks:
|
||||
- id: ruff
|
||||
- repo: https://github.com/crate-ci/typos
|
||||
rev: v1.26.0
|
||||
hooks:
|
||||
- id: typos
|
||||
# - repo: https://github.com/RobertCraigie/pyright-python
|
||||
# rev: v1.1.395
|
||||
# hooks:
|
||||
|
||||
-19
@@ -1,19 +0,0 @@
|
||||
[default]
|
||||
extend-ignore-re = ["(?Rm)^.*(#|//)\\s*spellchecker:disable-line$"]
|
||||
|
||||
[default.extend-words]
|
||||
# Azure Kubernetes Service, mentioned in rust/lancedb/src/remote/oauth.rs.
|
||||
AKS = "AKS"
|
||||
# RabitQ is the name of a vector quantization algorithm, not a typo of "Rabbit".
|
||||
Rabit = "Rabit"
|
||||
# `VarBuilder::from_mmaped_safetensors` is the real (if oddly-spelled) name of
|
||||
# the candle-core API we call in rust/lancedb/src/embeddings/sentence_transformers.rs.
|
||||
mmaped = "mmaped"
|
||||
# `WriteableBuffer` is the real name of a type from Python's `_typeshed` stubs,
|
||||
# used in python/python/lancedb/_blob.py.
|
||||
Writeable = "Writeable"
|
||||
|
||||
[files]
|
||||
extend-exclude = [
|
||||
"*_THIRD_PARTY_LICENSES.*",
|
||||
]
|
||||
Generated
+49
-50
@@ -3526,8 +3526,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
|
||||
|
||||
[[package]]
|
||||
name = "fsst"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"rand 0.9.5",
|
||||
@@ -4886,8 +4886,8 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
|
||||
|
||||
[[package]]
|
||||
name = "lance"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -4959,8 +4959,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-arrow"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4982,7 +4982,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-scalar"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -4996,7 +4996,7 @@ dependencies = [
|
||||
[[package]]
|
||||
name = "lance-arrow-stats"
|
||||
version = "58.0.0"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5005,8 +5005,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-bitpacking"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrayref",
|
||||
"crunchy",
|
||||
@@ -5016,8 +5016,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-core"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5054,8 +5054,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-datafusion"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5085,8 +5085,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-datagen"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5103,8 +5103,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-derive"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"proc-macro2",
|
||||
"quote",
|
||||
@@ -5113,8 +5113,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-encoding"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5147,8 +5147,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-file"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-arith",
|
||||
"arrow-array",
|
||||
@@ -5179,8 +5179,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-index"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arc-swap",
|
||||
"arrow",
|
||||
@@ -5244,8 +5244,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-index-core"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5267,8 +5267,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-io"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5308,8 +5308,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-linalg"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5323,8 +5323,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-namespace"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -5338,8 +5338,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-namespace-impls"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-ipc",
|
||||
@@ -5392,8 +5392,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-select"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5407,8 +5407,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-table"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"arrow-array",
|
||||
@@ -5448,8 +5448,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-testing"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-schema",
|
||||
@@ -5462,8 +5462,8 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lance-tokenizer"
|
||||
version = "12.0.0-beta.16"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.16#f7df098f5860cf9b3ac5339e7a641430cfbc6351"
|
||||
version = "12.0.0-beta.15"
|
||||
source = "git+https://github.com/lance-format/lance.git?tag=v12.0.0-beta.15#271a155ba505fd8e1094c095d4ce356707e93866"
|
||||
dependencies = [
|
||||
"frostem",
|
||||
"icu_segmenter",
|
||||
@@ -5476,7 +5476,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb"
|
||||
version = "0.39.0-beta.6"
|
||||
version = "0.39.0-beta.5"
|
||||
dependencies = [
|
||||
"ahash",
|
||||
"anyhow",
|
||||
@@ -5553,7 +5553,6 @@ dependencies = [
|
||||
"serde_json",
|
||||
"serde_with",
|
||||
"serial_test",
|
||||
"sha2 0.10.9",
|
||||
"snafu 0.8.9",
|
||||
"tempfile",
|
||||
"test-log",
|
||||
@@ -5568,7 +5567,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb-nodejs"
|
||||
version = "0.39.0-beta.6"
|
||||
version = "0.39.0-beta.5"
|
||||
dependencies = [
|
||||
"arrow-array",
|
||||
"arrow-buffer",
|
||||
@@ -5593,7 +5592,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb-python"
|
||||
version = "0.39.0-beta.6"
|
||||
version = "0.39.0-beta.5"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
@@ -7704,9 +7703,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "prost"
|
||||
version = "0.14.4"
|
||||
version = "0.14.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "528ac67416ff8646872a3c02cad9cc4ee5dc9f9540c9b10771855c95cb2e5ae1"
|
||||
checksum = "d2ea70524a2f82d518bce41317d0fae74151505651af45faf1ffbd6fd33f0568"
|
||||
dependencies = [
|
||||
"bytes",
|
||||
"prost-derive",
|
||||
@@ -7733,9 +7732,9 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "prost-derive"
|
||||
version = "0.14.4"
|
||||
version = "0.14.3"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf"
|
||||
checksum = "27c6023962132f4b30eb4c172c91ce92d933da334c59c23cddee82358ddafb0b"
|
||||
dependencies = [
|
||||
"anyhow",
|
||||
"itertools 0.14.0",
|
||||
|
||||
+14
-14
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
lance = { "version" = "=12.0.0-beta.16", default-features = false, "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=12.0.0-beta.16", default-features = false, "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=12.0.0-beta.16", default-features = false, "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=12.0.0-beta.16", "tag" = "v12.0.0-beta.16", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=12.0.0-beta.15", default-features = false, "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=12.0.0-beta.15", "tag" = "v12.0.0-beta.15", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lancedb = { path = "rust/lancedb", default-features = false }
|
||||
ahash = "0.8"
|
||||
# Note that this one does not include pyarrow
|
||||
|
||||
+1
-1
@@ -155,7 +155,7 @@ paths:
|
||||
vector:
|
||||
type: FixedSizeList
|
||||
description: |
|
||||
The targeted vector to search for. Required.
|
||||
The targetted vector to search for. Required.
|
||||
vector_column:
|
||||
type: string
|
||||
description: |
|
||||
|
||||
@@ -141,7 +141,7 @@ Currently this causes multiple copies of the row to be created
|
||||
but that behavior is subject to change.
|
||||
|
||||
An optional condition may be specified. If it is, then only
|
||||
matched rows that satisfy the condition will be updated. Any
|
||||
matched rows that satisfy the condtion will be updated. Any
|
||||
rows that do not satisfy the condition will be left as they
|
||||
are. Failing to satisfy the condition does not cause a
|
||||
"matched row" to become a "not matched" row.
|
||||
|
||||
@@ -74,10 +74,10 @@ now: the column is committed with no values, and rows get them from
|
||||
[Table#refreshColumn](Table.md#refreshcolumn). Declaring one therefore costs the same on a
|
||||
large table as on an empty one.
|
||||
|
||||
A refresh also recomputes the rows whose inputs changed since they were
|
||||
computed, so a mutated input is reflected by the next refresh. While a
|
||||
declaration reads a column, that column cannot be renamed, retyped or
|
||||
dropped.
|
||||
A refresh does not revisit rows it has already filled, so mutating an
|
||||
input leaves the value computed at fill time; recomputing means dropping
|
||||
the column and declaring it again. While a declaration reads a column,
|
||||
that column cannot be renamed, retyped or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
server, and the refresh runs as a server job -- see
|
||||
@@ -854,10 +854,10 @@ abstract refreshColumn(column): Promise<RefreshColumnResult>
|
||||
|
||||
Fill the rows of a computed column that hold no value yet.
|
||||
|
||||
Rows appended since the last refresh are filled by the next one, and
|
||||
rows whose inputs changed since they were computed are recomputed;
|
||||
everything else is left as it is. Local tables only: a remote refresh
|
||||
runs as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
|
||||
Rows appended since the last refresh are filled by the next one; rows
|
||||
already filled are left as they are, so the call is idempotent and does
|
||||
not observe a mutated input. Local tables only: a remote refresh runs
|
||||
as a server job, through [Table#refreshColumnAsync](Table.md#refreshcolumnasync).
|
||||
|
||||
#### Parameters
|
||||
|
||||
@@ -1266,7 +1266,7 @@ value is 0")
|
||||
Note: if your condition is something like "some_id_column == 7" and
|
||||
you are updating many rows (with different ids) then you will get
|
||||
better performance with a single [`merge_insert`] call instead of
|
||||
repeatedly calling this method.
|
||||
repeatedly calilng this method.
|
||||
|
||||
##### Parameters
|
||||
|
||||
|
||||
@@ -118,7 +118,7 @@ Number of sub-vectors of PQ.
|
||||
This value controls how much the vector is compressed during the quantization step.
|
||||
The more sub vectors there are the less the vector is compressed. The default is
|
||||
the dimension of the vector divided by 16. If the dimension is not evenly divisible
|
||||
by 16 we use the dimension divided by 8.
|
||||
by 16 we use the dimension divded by 8.
|
||||
|
||||
The above two cases are highly preferred. Having 8 or 16 values per subvector allows
|
||||
us to use efficient SIMD instructions.
|
||||
|
||||
@@ -16,7 +16,7 @@ optional config: Index;
|
||||
|
||||
Advanced index configuration
|
||||
|
||||
This option allows you to specify a specific index to create and also
|
||||
This option allows you to specify a specfic index to create and also
|
||||
allows you to pass in configuration for training the index.
|
||||
|
||||
See the static methods on Index for details on the various index types.
|
||||
|
||||
@@ -112,7 +112,7 @@ Number of sub-vectors of PQ.
|
||||
This value controls how much the vector is compressed during the quantization step.
|
||||
The more sub vectors there are the less the vector is compressed. The default is
|
||||
the dimension of the vector divided by 16. If the dimension is not evenly divisible
|
||||
by 16 we use the dimension divided by 8.
|
||||
by 16 we use the dimension divded by 8.
|
||||
|
||||
The above two cases are highly preferred. Having 8 or 16 values per subvector allows
|
||||
us to use efficient SIMD instructions.
|
||||
|
||||
+1
-1
@@ -28,7 +28,7 @@
|
||||
<properties>
|
||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
||||
<arrow.version>15.0.0</arrow.version>
|
||||
<lance-core.version>12.0.0-beta.16</lance-core.version>
|
||||
<lance-core.version>12.0.0-beta.15</lance-core.version>
|
||||
<spotless.skip>false</spotless.skip>
|
||||
<spotless.version>2.30.0</spotless.version>
|
||||
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
|
||||
|
||||
@@ -3252,7 +3252,7 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
|
||||
const db = await connect(tmpDir.name);
|
||||
const data = [
|
||||
{ text: "fa", vector: [0.1, 0.2, 0.3] },
|
||||
{ text: "fo", vector: [0.4, 0.5, 0.6] }, // spellchecker:disable-line
|
||||
{ text: "fo", vector: [0.4, 0.5, 0.6] },
|
||||
{ text: "fob", vector: [0.4, 0.5, 0.6] },
|
||||
{ text: "focus", vector: [0.4, 0.5, 0.6] },
|
||||
{ text: "foo", vector: [0.4, 0.5, 0.6] },
|
||||
@@ -3277,7 +3277,7 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
|
||||
const resultSet = new Set(fuzzyResults.map((r) => r.text));
|
||||
expect(resultSet.has("foo")).toBe(true);
|
||||
expect(resultSet.has("fob")).toBe(true);
|
||||
expect(resultSet.has("fo")).toBe(true); // spellchecker:disable-line
|
||||
expect(resultSet.has("fo")).toBe(true);
|
||||
expect(resultSet.has("food")).toBe(true);
|
||||
|
||||
const prefixResults = await table
|
||||
|
||||
@@ -600,7 +600,7 @@ function makeVector(
|
||||
}
|
||||
if (values.length === 0) {
|
||||
throw Error(
|
||||
"makeVector requires at least one value or the type must be specified",
|
||||
"makeVector requires at least one value or the type must be specfied",
|
||||
);
|
||||
}
|
||||
const sampleValue = values.find((val) => val !== null && val !== undefined);
|
||||
@@ -858,7 +858,7 @@ async function applyEmbeddings<T>(
|
||||
* customized by the `embeddingDataType` property of the embedding function.
|
||||
*
|
||||
* If a schema is provided in `makeTableOptions` then it should include the
|
||||
* embedding columns. If no schema is provided then embedding columns will
|
||||
* embedding columns. If no schema is provded then embedding columns will
|
||||
* be placed at the end of the table, after all of the input columns.
|
||||
*/
|
||||
export async function convertToTable(
|
||||
|
||||
@@ -26,7 +26,7 @@ export interface IvfPqOptions {
|
||||
* This value controls how much the vector is compressed during the quantization step.
|
||||
* The more sub vectors there are the less the vector is compressed. The default is
|
||||
* the dimension of the vector divided by 16. If the dimension is not evenly divisible
|
||||
* by 16 we use the dimension divided by 8.
|
||||
* by 16 we use the dimension divded by 8.
|
||||
*
|
||||
* The above two cases are highly preferred. Having 8 or 16 values per subvector allows
|
||||
* us to use efficient SIMD instructions.
|
||||
@@ -228,7 +228,7 @@ export interface HnswPqOptions {
|
||||
* This value controls how much the vector is compressed during the quantization step.
|
||||
* The more sub vectors there are the less the vector is compressed. The default is
|
||||
* the dimension of the vector divided by 16. If the dimension is not evenly divisible
|
||||
* by 16 we use the dimension divided by 8.
|
||||
* by 16 we use the dimension divded by 8.
|
||||
*
|
||||
* The above two cases are highly preferred. Having 8 or 16 values per subvector allows
|
||||
* us to use efficient SIMD instructions.
|
||||
@@ -825,7 +825,7 @@ export interface IndexOptions {
|
||||
/**
|
||||
* Advanced index configuration
|
||||
*
|
||||
* This option allows you to specify a specific index to create and also
|
||||
* This option allows you to specify a specfic index to create and also
|
||||
* allows you to pass in configuration for training the index.
|
||||
*
|
||||
* See the static methods on Index for details on the various index types.
|
||||
|
||||
@@ -27,7 +27,7 @@ export class MergeInsertBuilder {
|
||||
* but that behavior is subject to change.
|
||||
*
|
||||
* An optional condition may be specified. If it is, then only
|
||||
* matched rows that satisfy the condition will be updated. Any
|
||||
* matched rows that satisfy the condtion will be updated. Any
|
||||
* rows that do not satisfy the condition will be left as they
|
||||
* are. Failing to satisfy the condition does not cause a
|
||||
* "matched row" to become a "not matched" row.
|
||||
|
||||
@@ -3,7 +3,7 @@
|
||||
|
||||
// The utilities in this file help sanitize data from the user's arrow
|
||||
// library into the types expected by vectordb's arrow library. Node
|
||||
// generally allows for multiple versions of the same library (and sometimes
|
||||
// generally allows for mulitple versions of the same library (and sometimes
|
||||
// even multiple copies of the same version) to be installed at the same
|
||||
// time. However, arrow-js uses instanceof which expected that the input
|
||||
// comes from the exact same library instance. This is not always the case
|
||||
|
||||
@@ -313,7 +313,7 @@ export abstract class Table {
|
||||
* Note: if your condition is something like "some_id_column == 7" and
|
||||
* you are updating many rows (with different ids) then you will get
|
||||
* better performance with a single [`merge_insert`] call instead of
|
||||
* repeatedly calling this method.
|
||||
* repeatedly calilng this method.
|
||||
* @param {Map<string, string> | Record<string, string>} updates - the
|
||||
* columns to update
|
||||
* @returns {Promise<UpdateResult>} A promise that resolves to an object
|
||||
@@ -542,10 +542,10 @@ export abstract class Table {
|
||||
* {@link Table#refreshColumn}. Declaring one therefore costs the same on a
|
||||
* large table as on an empty one.
|
||||
*
|
||||
* A refresh also recomputes the rows whose inputs changed since they were
|
||||
* computed, so a mutated input is reflected by the next refresh. While a
|
||||
* declaration reads a column, that column cannot be renamed, retyped or
|
||||
* dropped.
|
||||
* A refresh does not revisit rows it has already filled, so mutating an
|
||||
* input leaves the value computed at fill time; recomputing means dropping
|
||||
* the column and declaring it again. While a declaration reads a column,
|
||||
* that column cannot be renamed, retyped or dropped.
|
||||
*
|
||||
* On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
* server, and the refresh runs as a server job -- see
|
||||
@@ -576,10 +576,10 @@ export abstract class Table {
|
||||
/**
|
||||
* Fill the rows of a computed column that hold no value yet.
|
||||
*
|
||||
* Rows appended since the last refresh are filled by the next one, and
|
||||
* rows whose inputs changed since they were computed are recomputed;
|
||||
* everything else is left as it is. Local tables only: a remote refresh
|
||||
* runs as a server job, through {@link Table#refreshColumnAsync}.
|
||||
* Rows appended since the last refresh are filled by the next one; rows
|
||||
* already filled are left as they are, so the call is idempotent and does
|
||||
* not observe a mutated input. Local tables only: a remote refresh runs
|
||||
* as a server job, through {@link Table#refreshColumnAsync}.
|
||||
* @param {string} column The name of the computed column to fill.
|
||||
* @returns {Promise<RefreshColumnResult>} A promise that resolves to the
|
||||
* number of rows filled and the new version number of the table.
|
||||
|
||||
@@ -21,7 +21,7 @@ class GteEmbeddings(TextEmbeddingFunction):
|
||||
An embedding function that uses GTE-LARGE MLX format(for Apple silicon devices only)
|
||||
as well as the standard cpu/gpu version from: https://huggingface.co/thenlper/gte-large.
|
||||
|
||||
For Apple users, you will need the mlx package installed, which can be done with:
|
||||
For Apple users, you will need the mlx package insalled, which can be done with:
|
||||
pip install mlx
|
||||
|
||||
Parameters
|
||||
|
||||
@@ -60,7 +60,7 @@ class InstructorEmbeddingFunction(TextEmbeddingFunction):
|
||||
|
||||
import lancedb
|
||||
from lancedb.pydantic import LanceModel, Vector
|
||||
from lancedb.embeddings import get_registry, InstructorEmbeddingFunction
|
||||
from lancedb.embeddings import get_registry, InstuctorEmbeddingFunction
|
||||
|
||||
instructor = get_registry().get("instructor").create(
|
||||
source_instruction="represent the document for retrieval",
|
||||
|
||||
@@ -751,7 +751,7 @@ class IvfPq:
|
||||
This value controls how much the vector is compressed during the
|
||||
quantization step. The more sub vectors there are the less the vector is
|
||||
compressed. The default is the dimension of the vector divided by 16. If
|
||||
the dimension is not evenly divisible by 16 we use the dimension divided by
|
||||
the dimension is not evenly divisible by 16 we use the dimension divded by
|
||||
8.
|
||||
|
||||
The above two cases are highly preferred. Having 8 or 16 values per
|
||||
|
||||
@@ -78,10 +78,6 @@ if TYPE_CHECKING:
|
||||
T = TypeVar("T", bound="LanceModel")
|
||||
AnalyzePlanDistributedMetrics = Literal["aggregate", "per_worker", "full"]
|
||||
|
||||
# Number of rows a hybrid query returns when no limit was set on it. This
|
||||
# mirrors the default the Rust query builder applies to its sub-queries.
|
||||
DEFAULT_HYBRID_LIMIT = 10
|
||||
|
||||
|
||||
@runtime_checkable
|
||||
class _LanceScanner(Protocol):
|
||||
@@ -863,7 +859,7 @@ class Query(pydantic.BaseModel):
|
||||
return query
|
||||
|
||||
# This tells pydantic to allow custom types (needed for the `vector` query since
|
||||
# pa.Array wouldn't be allowed otherwise)
|
||||
# pa.Array wouln't be allowed otherwise)
|
||||
model_config = pydantic.ConfigDict(arbitrary_types_allowed=True)
|
||||
|
||||
|
||||
@@ -3897,54 +3893,14 @@ class AsyncHybridQuery(AsyncStandardQuery, AsyncVectorQueryBase):
|
||||
|
||||
return self
|
||||
|
||||
def _create_child_queries(
|
||||
self,
|
||||
) -> Tuple["AsyncFTSQuery", "AsyncVectorQuery", int, int]:
|
||||
"""Build the sub-queries that make up this hybrid query.
|
||||
|
||||
Execution, `explain_plan` and `analyze_plan` all go through here so that
|
||||
the plans that are reported are the plans that actually run.
|
||||
|
||||
Returns the two sub-queries along with the effective limit and offset of
|
||||
the hybrid query itself.
|
||||
"""
|
||||
fts_query = AsyncFTSQuery(self._inner.to_fts_query(), self._table)
|
||||
vec_query = AsyncVectorQuery(self._inner.to_vector_query(), self._table)
|
||||
|
||||
fts_req = fts_query._inner.to_query_request()
|
||||
vec_req = vec_query._inner.to_query_request()
|
||||
|
||||
# Only one of the two sub-queries carries the limit when it was never
|
||||
# set explicitly: nearest_to()/nearest_to_text() build the sibling query
|
||||
# from scratch, and that is where the default gets filled in. Which one
|
||||
# that is depends on the order the hybrid query was built in, so look at
|
||||
# both rather than at a single side.
|
||||
limit = fts_req.limit if fts_req.limit is not None else vec_req.limit
|
||||
if limit is None:
|
||||
limit = DEFAULT_HYBRID_LIMIT
|
||||
offset = fts_req.offset or vec_req.offset or 0
|
||||
|
||||
fts_query.with_row_id()
|
||||
vec_query.with_row_id()
|
||||
|
||||
# offset() pushes the offset down into both sub-queries, which would make
|
||||
# each of them skip its own first `offset` rows. The window has to be
|
||||
# taken out of the combined, reranked results instead, so fetch the
|
||||
# skipped prefix here too and slice it off afterwards.
|
||||
fts_query.limit(limit + offset)
|
||||
vec_query.limit(limit + offset)
|
||||
fts_query.offset(0)
|
||||
vec_query.offset(0)
|
||||
|
||||
return fts_query, vec_query, limit, offset
|
||||
|
||||
async def to_batches(
|
||||
self,
|
||||
*,
|
||||
max_batch_length: Optional[int] = None,
|
||||
timeout: Optional[timedelta] = None,
|
||||
) -> AsyncRecordBatchReader:
|
||||
fts_query, vec_query, limit, offset = self._create_child_queries()
|
||||
fts_query = AsyncFTSQuery(self._inner.to_fts_query(), self._table)
|
||||
vec_query = AsyncVectorQuery(self._inner.to_vector_query(), self._table)
|
||||
|
||||
req = fts_query._inner.to_query_request()
|
||||
blob_auto_row_id = False
|
||||
@@ -3964,6 +3920,9 @@ class AsyncHybridQuery(AsyncStandardQuery, AsyncVectorQueryBase):
|
||||
self._blob_auto_row_id = blob_auto_row_id
|
||||
self._blob_paths = blob_paths
|
||||
|
||||
fts_query.with_row_id()
|
||||
vec_query.with_row_id()
|
||||
|
||||
fts_results, vector_results = await asyncio.gather(
|
||||
fts_query.to_arrow(timeout=timeout),
|
||||
vec_query.to_arrow(timeout=timeout),
|
||||
@@ -3975,9 +3934,8 @@ class AsyncHybridQuery(AsyncStandardQuery, AsyncVectorQueryBase):
|
||||
norm=self._norm,
|
||||
fts_query=fts_query.get_query(),
|
||||
reranker=self._reranker,
|
||||
limit=limit,
|
||||
limit=self._inner.get_limit(),
|
||||
with_row_ids=True,
|
||||
offset=offset,
|
||||
)
|
||||
if (
|
||||
not self._user_requested_row_id()
|
||||
@@ -4006,14 +3964,14 @@ class AsyncHybridQuery(AsyncStandardQuery, AsyncVectorQueryBase):
|
||||
... print(plan)
|
||||
>>> asyncio.run(doctest_example()) # doctest: +ELLIPSIS, +NORMALIZE_WHITESPACE
|
||||
RRFReranker(K=60)
|
||||
ProjectionExec: expr=[vector@0 as vector, text@3 as text, _distance@2 as _distance, _rowid@1 as _rowid]
|
||||
ProjectionExec: expr=[vector@0 as vector, text@3 as text, _distance@2 as _distance]
|
||||
LanceRead: uri=..., projection=[text], source=stream(_rowid)
|
||||
GlobalLimitExec: skip=0, fetch=10
|
||||
FilterExec: _distance@2 IS NOT NULL
|
||||
SortExec: TopK(fetch=10), expr=[_distance@2 ASC NULLS LAST, _rowid@1 ASC NULLS LAST], preserve_partitioning=[false]
|
||||
KNNVectorDistance: metric=l2
|
||||
LanceRead: uri=..., projection=[vector], ...
|
||||
ProjectionExec: expr=[vector@2 as vector, text@3 as text, _score@1 as _score, _rowid@0 as _rowid]
|
||||
ProjectionExec: expr=[vector@2 as vector, text@3 as text, _score@1 as _score]
|
||||
LanceRead: uri=..., projection=[vector, text], source=stream(_rowid)
|
||||
GlobalLimitExec: skip=0, fetch=10
|
||||
MatchQuery: column=text, query=[hello]
|
||||
@@ -4028,9 +3986,8 @@ class AsyncHybridQuery(AsyncStandardQuery, AsyncVectorQueryBase):
|
||||
plan : str
|
||||
""" # noqa: E501
|
||||
|
||||
fts_query, vec_query, _, _ = self._create_child_queries()
|
||||
vector_plan = await vec_query.explain_plan(verbose)
|
||||
fts_plan = await fts_query.explain_plan(verbose)
|
||||
vector_plan = await self._inner.to_vector_query().explain_plan(verbose)
|
||||
fts_plan = await self._inner.to_fts_query().explain_plan(verbose)
|
||||
# Indent sub-plans under the reranker
|
||||
indented_vector = "\n".join(" " + line for line in vector_plan.splitlines())
|
||||
indented_fts = "\n".join(" " + line for line in fts_plan.splitlines())
|
||||
@@ -4057,12 +4014,14 @@ class AsyncHybridQuery(AsyncStandardQuery, AsyncVectorQueryBase):
|
||||
-------
|
||||
plan : str
|
||||
"""
|
||||
fts_query, vec_query, _, _ = self._create_child_queries()
|
||||
|
||||
results = ["Vector Search Query:"]
|
||||
results.append(await vec_query.analyze_plan(distributed_metrics))
|
||||
results.append(
|
||||
await self._inner.to_vector_query().analyze_plan(distributed_metrics)
|
||||
)
|
||||
results.append("FTS Search Query:")
|
||||
results.append(await fts_query.analyze_plan(distributed_metrics))
|
||||
results.append(
|
||||
await self._inner.to_fts_query().analyze_plan(distributed_metrics)
|
||||
)
|
||||
|
||||
return "\n".join(results)
|
||||
|
||||
|
||||
@@ -720,7 +720,7 @@ class RemoteTable(Table):
|
||||
Parameters
|
||||
----------
|
||||
query: list/np.ndarray/str/PIL.Image.Image, default None
|
||||
The targeted vector to search for.
|
||||
The targetted vector to search for.
|
||||
|
||||
- *default None*.
|
||||
Acceptable types are: list, np.ndarray, PIL.Image.Image
|
||||
|
||||
@@ -175,7 +175,7 @@ class Reranker(ABC):
|
||||
if the results haven't been executed yet or the results in arrow format.
|
||||
query : str or None,
|
||||
The input query. Some rerankers might not need the query to rerank.
|
||||
In that case, it can be set to None explicitly. This is intended to
|
||||
In that case, it can be set to None explicitly. This is inteded to
|
||||
be handled by the reranker implementations.
|
||||
deduplicate : bool, optional
|
||||
Whether to deduplicate the results based on the `_rowid` column,
|
||||
|
||||
@@ -1619,7 +1619,7 @@ class Table(ABC):
|
||||
Parameters
|
||||
----------
|
||||
query: list/np.ndarray/str/PIL.Image.Image, default None
|
||||
The targeted vector to search for.
|
||||
The targetted vector to search for.
|
||||
|
||||
- *default None*.
|
||||
Acceptable types are: list, np.ndarray, PIL.Image.Image
|
||||
@@ -2188,10 +2188,10 @@ class Table(ABC):
|
||||
Declaring one therefore costs the same on a large table as on an
|
||||
empty one.
|
||||
|
||||
A refresh also recomputes the rows whose inputs changed since they
|
||||
were computed, so a mutated input is reflected by the next refresh.
|
||||
While a declaration reads a column, that column cannot be renamed,
|
||||
retyped or dropped.
|
||||
A refresh does not revisit rows it has already filled, so mutating
|
||||
an input leaves the value computed at fill time; recomputing means
|
||||
dropping the column and declaring it again. While a declaration
|
||||
reads a column, that column cannot be renamed, retyped or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
server, and the refresh runs as a server job -- see
|
||||
@@ -2211,7 +2211,7 @@ class Table(ABC):
|
||||
>>> table.add_columns(computed={"doubled": "x * 2"})
|
||||
AddColumnsResult(version=2)
|
||||
>>> table.refresh_column("doubled")
|
||||
RefreshColumnResult(rows_filled=2, version=4)
|
||||
RefreshColumnResult(rows_filled=2, version=3)
|
||||
>>> table.to_arrow().sort_by("x").to_pandas()
|
||||
x doubled
|
||||
0 1 2
|
||||
@@ -2225,8 +2225,8 @@ class Table(ABC):
|
||||
|
||||
Declared with ``add_columns(computed=...)``, a column starts empty and
|
||||
gets its values here. Rows appended since the last refresh are filled
|
||||
by the next one, and rows whose inputs changed since they were computed
|
||||
are recomputed; everything else is left as it is.
|
||||
by the next one; rows already filled are left as they are, so the call
|
||||
is idempotent and does not observe a mutated input.
|
||||
|
||||
Local tables only: a remote refresh runs as a server job, through
|
||||
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
@@ -3841,7 +3841,7 @@ class LanceTable(Table):
|
||||
Parameters
|
||||
----------
|
||||
query: list/np.ndarray/str/PIL.Image.Image, default None
|
||||
The targeted vector to search for.
|
||||
The targetted vector to search for.
|
||||
|
||||
- *default None*.
|
||||
Acceptable types are: list, np.ndarray, PIL.Image.Image
|
||||
@@ -4318,14 +4318,13 @@ class LanceTable(Table):
|
||||
return LOOP.run(self._table.add_columns(transforms, computed=computed))
|
||||
|
||||
def refresh_column(self, column: str) -> "RefreshColumnResult":
|
||||
"""Fill a computed column's unfilled rows and recompute those whose
|
||||
inputs changed. See
|
||||
"""Fill a computed column's unfilled rows. See
|
||||
[`AsyncTable.refresh_column`][lancedb.AsyncTable.refresh_column]."""
|
||||
return LOOP.run(self._table.refresh_column(column))
|
||||
|
||||
def refresh_column_async(self, column: str) -> Job[RefreshColumnJobResult]:
|
||||
"""Fill a computed column's unfilled rows and recompute those whose
|
||||
inputs changed, returning a handle to the refresh job. See
|
||||
"""Fill a computed column's unfilled rows, returning a handle to the
|
||||
refresh job. See
|
||||
[`Table.refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
"""
|
||||
return Job(LOOP.run(self._table.refresh_column_async(column)))
|
||||
@@ -5639,7 +5638,7 @@ class AsyncTable:
|
||||
if fill_value is None:
|
||||
fill_value = 0.0
|
||||
|
||||
# _sanitize_data is an old code path, but we will use it until the
|
||||
# _santitize_data is an old code path, but we will use it until the
|
||||
# new code path is ready.
|
||||
if mode == "overwrite":
|
||||
# For overwrite, apply the same preprocessing as create_table
|
||||
@@ -5815,7 +5814,7 @@ class AsyncTable:
|
||||
Parameters
|
||||
----------
|
||||
query: list/np.ndarray/str/PIL.Image.Image, default None
|
||||
The targeted vector to search for.
|
||||
The targetted vector to search for.
|
||||
|
||||
- *default None*.
|
||||
Acceptable types are: list, np.ndarray, PIL.Image.Image
|
||||
@@ -6313,10 +6312,10 @@ class AsyncTable:
|
||||
them from
|
||||
[`refresh_column`][lancedb.table.AsyncTable.refresh_column].
|
||||
|
||||
A refresh also recomputes the rows whose inputs changed since they
|
||||
were computed, so a mutated input is reflected by the next refresh.
|
||||
While a declaration reads a column, that column cannot be renamed,
|
||||
retyped or dropped.
|
||||
A refresh does not revisit rows it has already filled, so mutating
|
||||
an input leaves the value computed at fill time. While a
|
||||
declaration reads a column, that column cannot be renamed, retyped
|
||||
or dropped.
|
||||
|
||||
On LanceDB Cloud and Enterprise the expression is planned by
|
||||
the server. Cannot be combined with ``transforms``.
|
||||
@@ -6378,8 +6377,8 @@ class AsyncTable:
|
||||
|
||||
Declared with ``add_columns(computed=...)``, a column starts empty and
|
||||
gets its values here. Rows appended since the last refresh are filled
|
||||
by the next one, and rows whose inputs changed since they were computed
|
||||
are recomputed; everything else is left as it is.
|
||||
by the next one; rows already filled are left as they are, so the call
|
||||
is idempotent and does not observe a mutated input.
|
||||
|
||||
Local tables only: a remote refresh runs as a server job, through
|
||||
[`refresh_column_async`][lancedb.table.Table.refresh_column_async].
|
||||
|
||||
@@ -327,8 +327,8 @@ def test_embedding_function_with_pandas(tmp_path):
|
||||
) -> List[np.array]:
|
||||
return [np.random.randn(self.ndims()).tolist() for _ in range(len(texts))]
|
||||
|
||||
registry = get_registry()
|
||||
func = registry.get("mock-embedding").create()
|
||||
registery = get_registry()
|
||||
func = registery.get("mock-embedding").create()
|
||||
|
||||
class TestSchema(LanceModel):
|
||||
text: str = func.SourceField()
|
||||
@@ -394,9 +394,9 @@ def test_multiple_embeddings_for_pandas(tmp_path):
|
||||
) -> List[np.array]:
|
||||
return [np.random.randn(self.ndims()).tolist() for _ in range(len(texts))]
|
||||
|
||||
registry = get_registry()
|
||||
func1 = registry.get("mock-embedding").create()
|
||||
func2 = registry.get("mock-embedding2").create()
|
||||
registery = get_registry()
|
||||
func1 = registery.get("mock-embedding").create()
|
||||
func2 = registery.get("mock-embedding2").create()
|
||||
|
||||
class TestSchema(LanceModel):
|
||||
text: str = func1.SourceField()
|
||||
|
||||
@@ -1011,13 +1011,8 @@ def test_fts_ngram(mem_db: DBConnection):
|
||||
assert set(r["text"] for r in results) == {"lance database", "lance is cool"}
|
||||
|
||||
results = (
|
||||
table.search(
|
||||
"nce", # spellchecker:disable-line
|
||||
query_type="fts",
|
||||
)
|
||||
.limit(10)
|
||||
.to_list()
|
||||
)
|
||||
table.search("nce", query_type="fts").limit(10).to_list()
|
||||
) # spellchecker:disable-line
|
||||
assert len(results) == 2
|
||||
assert set(r["text"] for r in results) == {"lance database", "lance is cool"}
|
||||
|
||||
@@ -1039,13 +1034,8 @@ def test_fts_ngram(mem_db: DBConnection):
|
||||
assert set(r["text"] for r in results) == {"lance database", "lance is cool"}
|
||||
|
||||
results = (
|
||||
table.search(
|
||||
"nce", # spellchecker:disable-line
|
||||
query_type="fts",
|
||||
)
|
||||
.limit(10)
|
||||
.to_list()
|
||||
)
|
||||
table.search("nce", query_type="fts").limit(10).to_list()
|
||||
) # spellchecker:disable-line
|
||||
assert len(results) == 0
|
||||
|
||||
results = table.search("la", query_type="fts").limit(10).to_list()
|
||||
|
||||
@@ -203,93 +203,6 @@ async def test_async_hybrid_query_default_limit(table: AsyncTable):
|
||||
assert texts.count("a") == 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_hybrid_query_offset(table: AsyncTable):
|
||||
# The offset window of a hybrid query must be a suffix of the same query
|
||||
# run without an offset. Skipping the first rows of each sub-query instead
|
||||
# of the first rows of the fused result silently changes which rows land in
|
||||
# the window.
|
||||
full = await (
|
||||
table.query()
|
||||
.nearest_to([0.0, 0.4])
|
||||
.nearest_to_text("dog")
|
||||
.limit(4)
|
||||
.with_row_id()
|
||||
.to_arrow()
|
||||
)
|
||||
assert len(full) == 4
|
||||
|
||||
second_page = await (
|
||||
table.query()
|
||||
.nearest_to([0.0, 0.4])
|
||||
.nearest_to_text("dog")
|
||||
.offset(2)
|
||||
.limit(2)
|
||||
.with_row_id()
|
||||
.to_arrow()
|
||||
)
|
||||
assert second_page["_rowid"].to_pylist() == full["_rowid"].to_pylist()[2:]
|
||||
|
||||
first_page = await (
|
||||
table.query()
|
||||
.nearest_to([0.0, 0.4])
|
||||
.nearest_to_text("dog")
|
||||
.limit(2)
|
||||
.with_row_id()
|
||||
.to_arrow()
|
||||
)
|
||||
# Paging through the result must visit every row exactly once: no row
|
||||
# repeated from the previous page and none dropped between the two.
|
||||
paged = first_page["_rowid"].to_pylist() + second_page["_rowid"].to_pylist()
|
||||
assert sorted(paged) == sorted(full["_rowid"].to_pylist())
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_hybrid_query_fts_first_default_limit(table: AsyncTable):
|
||||
# nearest_to() and nearest_to_text() build their new sibling sub-query from
|
||||
# scratch, and that is the sub-query the default limit ends up on. So the
|
||||
# side that carries the limit depends on the order the hybrid query was
|
||||
# built in, and looking at only one side loses the limit for half the ways
|
||||
# a hybrid query can be written. Without a limit the combined results are
|
||||
# not truncated at all and the whole union of both candidate lists is
|
||||
# returned.
|
||||
await table.add([{"text": "dog", "vector": [50.0 + i, 50.0]} for i in range(10)])
|
||||
|
||||
result = await (
|
||||
table.query().nearest_to_text("dog").nearest_to([0.1, 0.1]).to_arrow()
|
||||
)
|
||||
assert len(result) == 10
|
||||
|
||||
offset_result = await (
|
||||
table.query().nearest_to_text("dog").nearest_to([0.1, 0.1]).offset(2).to_arrow()
|
||||
)
|
||||
assert len(offset_result) == 10
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_hybrid_query_explain_plan_matches_execution(table: AsyncTable):
|
||||
# Paging rewrites the sub-queries: each one fetches limit + offset rows with
|
||||
# no offset of its own, and the window is sliced out after fusion. The plans
|
||||
# have to be built from those rewritten sub-queries, otherwise explain_plan
|
||||
# and analyze_plan describe a query that is never run.
|
||||
query = (
|
||||
table.query().nearest_to([0.0, 0.4]).nearest_to_text("dog").offset(2).limit(2)
|
||||
)
|
||||
await query.to_arrow()
|
||||
|
||||
plan = await query.explain_plan()
|
||||
assert [
|
||||
line.strip() for line in plan.splitlines() if "GlobalLimitExec" in line
|
||||
] == [
|
||||
"GlobalLimitExec: skip=0, fetch=4",
|
||||
"GlobalLimitExec: skip=0, fetch=4",
|
||||
]
|
||||
|
||||
analyzed = await query.analyze_plan()
|
||||
assert analyzed.count("skip=0, fetch=4") == 2
|
||||
assert "skip=2" not in analyzed
|
||||
|
||||
|
||||
def test_hybrid_query_offset(sync_table: Table):
|
||||
# The offset window of a hybrid query must be a suffix of the same query
|
||||
# run without an offset -- it must not be silently ignored.
|
||||
|
||||
@@ -81,7 +81,7 @@ def get_test_table(tmp_path):
|
||||
"but his son was mortal",
|
||||
"there hasn't been a good battlefield game since 2142",
|
||||
"I wish they would make another one",
|
||||
"campaigns are not as good as they used to be",
|
||||
"campains are not as good as they used to be",
|
||||
"Multiplayer and open world games have destroyed the single player experience",
|
||||
"Maybe the future is console games",
|
||||
"I don't know",
|
||||
|
||||
@@ -3354,7 +3354,7 @@ def test_empty_query(mem_db: DBConnection):
|
||||
# None is the same as default
|
||||
df = table.search().select(["id"]).limit(None).to_arrow()
|
||||
assert df.num_rows == 100
|
||||
# invalid limist is the same as None, which is the same as default
|
||||
# invalid limist is the same as None, wihch is the same as default
|
||||
df = table.search().select(["id"]).limit(-1).to_arrow()
|
||||
assert df.num_rows == 100
|
||||
# valid limit should work
|
||||
@@ -4183,14 +4183,13 @@ def test_refresh_column_async_returns_job(tmp_path):
|
||||
assert result.rows_failed == 0
|
||||
assert result.rows_remaining == 0
|
||||
assert result.source_version == 2
|
||||
# The fill lands at 3; the stamp recording its inputs is published at 4.
|
||||
assert result.published_version == 4
|
||||
assert result.published_version == 3
|
||||
assert job.status() == "finished"
|
||||
assert sorted(table.to_arrow()["doubled"].to_pylist()) == [2, 4]
|
||||
|
||||
no_op = table.refresh_column_async("doubled").wait()
|
||||
assert no_op.rows_assigned == 0
|
||||
assert no_op.source_version == 4
|
||||
assert no_op.source_version == 3
|
||||
assert no_op.published_version is None
|
||||
|
||||
# Bad input raises at the call, not through the job.
|
||||
@@ -4209,6 +4208,6 @@ async def test_refresh_column_async_job_async_table(tmp_path):
|
||||
assert isinstance(result, lancedb.RefreshColumnResult)
|
||||
assert result.rows_assigned == 1
|
||||
assert result.source_version == 2
|
||||
assert result.published_version == 4
|
||||
assert result.published_version == 3
|
||||
assert await job.status() == "finished"
|
||||
assert (await table.to_arrow())["tripled"].to_pylist() == [9]
|
||||
|
||||
+1
-1
@@ -334,7 +334,7 @@ pub struct PyQueryRequest {
|
||||
pub column: Option<String>,
|
||||
pub query_vector: Option<PyQueryVectors>,
|
||||
pub minimum_nprobes: Option<usize>,
|
||||
// None means user did not set it and default should be used (currently 20)
|
||||
// None means user did not set it and default shoud be used (currenty 20)
|
||||
// Some(0) means user set it to None and there is no limit
|
||||
pub maximum_nprobes: Option<usize>,
|
||||
pub lower_bound: Option<f32>,
|
||||
|
||||
@@ -95,14 +95,13 @@ candle-transformers = { version = "0.9.1", optional = true }
|
||||
candle-nn = { version = "0.9.1", optional = true }
|
||||
tokenizers = { version = "0.19.1", optional = true }
|
||||
semver = { workspace = true }
|
||||
roaring = "0.11.4"
|
||||
sha2 = "0.10"
|
||||
|
||||
[dev-dependencies]
|
||||
anyhow = "1"
|
||||
lance-testing = { workspace = true }
|
||||
tempfile = { workspace = true }
|
||||
random_word = { version = "0.4.3", features = ["en"] }
|
||||
roaring = "0.11.4"
|
||||
tokio = { workspace = true, features = ["io-util", "macros", "net", "test-util"] }
|
||||
uuid = { workspace = true }
|
||||
walkdir = "2"
|
||||
|
||||
@@ -163,7 +163,7 @@ pub struct PolarsDataFrameRecordBatchReader {
|
||||
impl PolarsDataFrameRecordBatchReader {
|
||||
/// Creates a new `PolarsDataFrameRecordBatchReader` from a given Polars DataFrame.
|
||||
/// If the input dataframe does not have aligned chunks, this function undergoes
|
||||
/// the costly operation of reallocating each series as a single contiguous chunk.
|
||||
/// the costly operation of reallocating each series as a single contigous chunk.
|
||||
pub fn new(mut df: DataFrame) -> Result<Self> {
|
||||
df.align_chunks();
|
||||
let arrow_schema =
|
||||
|
||||
@@ -827,7 +827,7 @@ impl Connection {
|
||||
pub struct ConnectRequest {
|
||||
/// Database URI
|
||||
///
|
||||
/// ### Accepted URI formats
|
||||
/// ### Accpeted URI formats
|
||||
///
|
||||
/// - `/path/to/database` - local database on file system.
|
||||
/// - `s3://bucket/path/to/database` or `gs://bucket/path/to/database` - database on cloud object store
|
||||
|
||||
@@ -512,7 +512,7 @@ impl ListingDatabase {
|
||||
// iter thru the query params and extract the commit store param
|
||||
let mut engine = None;
|
||||
let mut mirrored_store = None;
|
||||
let mut filtered_queries = vec![];
|
||||
let mut filtered_querys = vec![];
|
||||
|
||||
// WARNING: specifying engine is NOT a publicly supported feature in lancedb yet
|
||||
// THE API WILL CHANGE
|
||||
@@ -528,13 +528,13 @@ impl ListingDatabase {
|
||||
mirrored_store = Some(value.to_string());
|
||||
} else {
|
||||
// to owned so we can modify the url
|
||||
filtered_queries.push((key.to_string(), value.to_string()));
|
||||
filtered_querys.push((key.to_string(), value.to_string()));
|
||||
}
|
||||
}
|
||||
|
||||
// Filter out the commit store query param -- it's a lancedb param
|
||||
url.query_pairs_mut().clear();
|
||||
url.query_pairs_mut().extend_pairs(filtered_queries);
|
||||
url.query_pairs_mut().extend_pairs(filtered_querys);
|
||||
// Take a copy of the query string so we can propagate it to lance.
|
||||
// `query_pairs_mut()` leaves the URL with `Some("")` even when no
|
||||
// pairs survive (or none existed in the first place), so an empty
|
||||
@@ -896,11 +896,11 @@ impl Database for ListingDatabase {
|
||||
}
|
||||
|
||||
async fn read_consistency(&self) -> Result<ReadConsistency> {
|
||||
if let Some(interval) = self.read_consistency_interval {
|
||||
if interval.is_zero() {
|
||||
if let Some(read_consistency_inverval) = self.read_consistency_interval {
|
||||
if read_consistency_inverval.is_zero() {
|
||||
Ok(ReadConsistency::Strong)
|
||||
} else {
|
||||
Ok(ReadConsistency::Eventual(interval))
|
||||
Ok(ReadConsistency::Eventual(read_consistency_inverval))
|
||||
}
|
||||
} else {
|
||||
Ok(ReadConsistency::Manual)
|
||||
@@ -3043,15 +3043,15 @@ mod tests {
|
||||
/// across platforms — see the `file://` test below).
|
||||
fn capture_query_like_connect(input_uri: &str) -> Option<String> {
|
||||
let mut url = url::Url::parse(input_uri).unwrap();
|
||||
let mut filtered_queries = Vec::new();
|
||||
let mut filtered_querys = Vec::new();
|
||||
for (key, value) in url.query_pairs() {
|
||||
if key == ENGINE || key == MIRRORED_STORE {
|
||||
continue;
|
||||
}
|
||||
filtered_queries.push((key.to_string(), value.to_string()));
|
||||
filtered_querys.push((key.to_string(), value.to_string()));
|
||||
}
|
||||
url.query_pairs_mut().clear();
|
||||
url.query_pairs_mut().extend_pairs(filtered_queries);
|
||||
url.query_pairs_mut().extend_pairs(filtered_querys);
|
||||
url.query().filter(|q| !q.is_empty()).map(|s| s.to_string())
|
||||
}
|
||||
|
||||
|
||||
@@ -251,11 +251,11 @@ impl Database for LanceNamespaceDatabase {
|
||||
}
|
||||
|
||||
async fn read_consistency(&self) -> Result<ReadConsistency> {
|
||||
if let Some(interval) = self.read_consistency_interval {
|
||||
if interval.is_zero() {
|
||||
if let Some(read_consistency_inverval) = self.read_consistency_interval {
|
||||
if read_consistency_inverval.is_zero() {
|
||||
Ok(ReadConsistency::Strong)
|
||||
} else {
|
||||
Ok(ReadConsistency::Eventual(interval))
|
||||
Ok(ReadConsistency::Eventual(read_consistency_inverval))
|
||||
}
|
||||
} else {
|
||||
Ok(ReadConsistency::Manual)
|
||||
|
||||
@@ -125,7 +125,7 @@ macro_rules! impl_pq_params_setter {
|
||||
/// This value controls how much the vector is compressed during the quantization step.
|
||||
/// The more sub vectors there are the less the vector is compressed. The default is
|
||||
/// the dimension of the vector divided by 16. If the dimension is not evenly divisible
|
||||
/// by 16 we use the dimension divided by 8.
|
||||
/// by 16 we use the dimension divded by 8.
|
||||
///
|
||||
/// The above two cases are highly preferred. Having 8 or 16 values per subvector allows
|
||||
/// us to use efficient SIMD instructions.
|
||||
|
||||
@@ -1124,9 +1124,8 @@ struct RowScope {
|
||||
|
||||
/// 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, or the freshness
|
||||
/// stamp a fill leaves on them. A version whose transaction cannot be read
|
||||
/// is not proven, so it counts as drift.
|
||||
/// 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.
|
||||
@@ -1177,19 +1176,6 @@ async fn only_computed_rewrites_since(view_ds: &Dataset, recorded: u64) -> Resul
|
||||
.all(|field| computed_fields.contains(&(*field as u32)))
|
||||
})
|
||||
}
|
||||
// The stamp `refresh_column` writes after its fill (see
|
||||
// `table::freshness`): field metadata on computed columns, no data.
|
||||
Operation::UpdateConfig {
|
||||
config_updates: None,
|
||||
table_metadata_updates: None,
|
||||
schema_metadata_updates: None,
|
||||
field_metadata_updates,
|
||||
} => {
|
||||
!field_metadata_updates.is_empty()
|
||||
&& field_metadata_updates
|
||||
.keys()
|
||||
.all(|field| computed_fields.contains(&(*field as u32)))
|
||||
}
|
||||
_ => false,
|
||||
};
|
||||
if !fill {
|
||||
@@ -3373,46 +3359,6 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
/// Field metadata on `field` only, the commit shape of the freshness
|
||||
/// stamp `refresh_column` leaves after its fill.
|
||||
async fn commit_field_metadata(view: &MaterializedView, field: &str, key: &str) {
|
||||
let native = view.table().as_native().unwrap();
|
||||
native.dataset.reload().await.unwrap();
|
||||
let mut dataset = native.dataset.get().await.unwrap().as_ref().clone();
|
||||
dataset
|
||||
.update_field_metadata()
|
||||
.update(field, [(key.to_string(), "{}".to_string())])
|
||||
.unwrap()
|
||||
.await
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
/// The stamp is metadata on the computed column and rewrites nothing
|
||||
/// refresh certifies, so it is not drift; the same commit shape on a
|
||||
/// projected column is, like any other write to it.
|
||||
#[tokio::test]
|
||||
async fn test_a_freshness_stamp_is_not_drift() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let view = refreshed_computed_view(&conn).await;
|
||||
|
||||
commit_field_metadata(
|
||||
&view,
|
||||
"emb",
|
||||
crate::table::computed_columns::SOURCE_SIGNATURE_META_KEY,
|
||||
)
|
||||
.await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::NoOp
|
||||
);
|
||||
|
||||
commit_field_metadata(&view, "id", "probe").await;
|
||||
assert_eq!(
|
||||
view.refresh().execute().await.unwrap().mode,
|
||||
RefreshMode::Rebuild
|
||||
);
|
||||
}
|
||||
|
||||
/// 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
|
||||
@@ -3578,9 +3524,9 @@ mod tests {
|
||||
}
|
||||
|
||||
/// A SQL declaration is filled by `refresh_column` on the view, which
|
||||
/// commits a data replacement and then its freshness stamp; the next
|
||||
/// refresh continues from its watermark and keeps what the fill wrote,
|
||||
/// and only rows the view added since come back unfilled.
|
||||
/// 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};
|
||||
|
||||
@@ -1299,7 +1299,7 @@ impl VectorQuery {
|
||||
/// This can be useful when there is a narrow filter to allow these queries to
|
||||
/// spend more time searching and avoid potential false negatives.
|
||||
///
|
||||
/// Set to None to search all partitions, if needed, to satisfy the limit
|
||||
/// Set to None to search all partitions, if needed, to satsify the limit
|
||||
pub fn maximum_nprobes(mut self, maximum_nprobes: Option<usize>) -> Result<Self> {
|
||||
if let Some(maximum_nprobes) = maximum_nprobes {
|
||||
if maximum_nprobes == 0 {
|
||||
|
||||
@@ -75,7 +75,6 @@ mod create_index;
|
||||
pub mod datafusion;
|
||||
pub(crate) mod dataset;
|
||||
pub mod delete;
|
||||
pub mod freshness;
|
||||
pub mod lsm_stats;
|
||||
pub mod merge;
|
||||
pub mod optimize;
|
||||
@@ -241,7 +240,7 @@ enum BadVectorHandling {
|
||||
/// An error is returned
|
||||
#[default]
|
||||
Error,
|
||||
/// The offending row is dropped
|
||||
/// The offending row is droppped
|
||||
Drop,
|
||||
/// The invalid/missing items are replaced by fill_value
|
||||
Fill(f32),
|
||||
@@ -779,8 +778,7 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
|
||||
message: "Function columns are supported only on LanceDB Cloud and Enterprise".into(),
|
||||
})
|
||||
}
|
||||
/// Fill a computed column's unfilled rows and recompute those whose
|
||||
/// inputs changed.
|
||||
/// Fill a computed column's unfilled rows.
|
||||
///
|
||||
/// The default returns `NotSupported`; Lance-backed tables override it.
|
||||
async fn refresh_column(&self, _column: &str) -> Result<RefreshColumnResult> {
|
||||
@@ -788,8 +786,8 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
|
||||
message: "computed columns are supported only on local tables".into(),
|
||||
})
|
||||
}
|
||||
/// Fill a computed column's unfilled rows and recompute those whose
|
||||
/// inputs changed, returning a [`Job`] tracking the operation.
|
||||
/// Fill a computed column's unfilled rows, returning a [`Job`] tracking
|
||||
/// the operation.
|
||||
async fn refresh_column_async(
|
||||
&self,
|
||||
_column: &str,
|
||||
@@ -1328,7 +1326,7 @@ impl Table {
|
||||
/// Note: if your condition is something like "some_id_column == 7" and
|
||||
/// you are updating many rows (with different ids) then you will get
|
||||
/// better performance with a single [`merge_insert`] call instead of
|
||||
/// repeatedly calling this method.
|
||||
/// repeatedly calilng this method.
|
||||
pub fn update(&self) -> UpdateBuilder {
|
||||
UpdateBuilder::new(self.inner.clone())
|
||||
}
|
||||
@@ -1751,10 +1749,9 @@ impl Table {
|
||||
/// Declared with
|
||||
/// [`AddColumnsBuilder::computed`](add_columns::AddColumnsBuilder::computed),
|
||||
/// a column starts empty and gets its values here. Fragments appended
|
||||
/// since the last refresh are filled by the next one, and fragments whose
|
||||
/// inputs changed since they were computed are recomputed (see
|
||||
/// [`freshness`](crate::table::freshness)); everything else is left as
|
||||
/// it is.
|
||||
/// since the last refresh are filled by the next one; fragments already
|
||||
/// filled are left as they are, so the call is idempotent and does not
|
||||
/// observe a mutated input.
|
||||
///
|
||||
/// Local tables only: a remote refresh runs as a server job, through
|
||||
/// [`Table::refresh_column_async`].
|
||||
|
||||
@@ -60,11 +60,10 @@ impl AddColumnsBuilder {
|
||||
/// every fragment that has none -- including fragments appended since the
|
||||
/// last refresh.
|
||||
///
|
||||
/// A refresh also recomputes the rows of a fragment whose inputs changed
|
||||
/// since it was computed (see [`freshness`](super::freshness)), so a
|
||||
/// mutated input is reflected by the next refresh. An input cannot be
|
||||
/// renamed, retyped or dropped while a declaration reads it, since the
|
||||
/// expression names it.
|
||||
/// Refresh does not revisit a fragment it has filled, so mutating an input
|
||||
/// leaves the value computed at fill time; recomputing means dropping the
|
||||
/// column and declaring it again. An input cannot be renamed, retyped or
|
||||
/// dropped while a declaration reads it, since the expression names it.
|
||||
///
|
||||
/// On LanceDB Cloud and Enterprise the expression is planned by the
|
||||
/// server, and the refresh runs as a server job -- see
|
||||
|
||||
@@ -71,22 +71,6 @@ pub const FUNCTION_BINDINGS_META_KEY: &str = "lancedb::function_bindings";
|
||||
/// Version of the schema-level Function binding envelope.
|
||||
pub const FUNCTION_BINDINGS_VERSION: u32 = 1;
|
||||
|
||||
/// Field metadata key holding `{fragment id -> input signature}` as JSON,
|
||||
/// recorded by the refresh that last computed each fragment. Outside the
|
||||
/// declaration namespace on purpose: a declaration is immutable through
|
||||
/// metadata edits, this is rewritten by every refresh. Seeded empty at
|
||||
/// declaration, so a column is tracked from birth; a column without it was
|
||||
/// declared before signatures existed.
|
||||
pub const SOURCE_SIGNATURE_META_KEY: &str = "computed_refresh.source_signature";
|
||||
|
||||
/// Field metadata key holding the definition digest a column was last
|
||||
/// computed under. A change to it makes every row stale.
|
||||
pub const DEFINITION_VERSION_META_KEY: &str = "computed_refresh.definition_version";
|
||||
|
||||
/// Field metadata key holding the table version the signature map describes:
|
||||
/// where a refresh starts following compactions to carry freshness forward.
|
||||
pub const RECORDED_AT_VERSION_META_KEY: &str = "computed_refresh.recorded_at_version";
|
||||
|
||||
/// Value of [`KIND_META_KEY`] for a column defined by a SQL expression.
|
||||
pub const SQL_KIND: &str = "sql";
|
||||
|
||||
@@ -155,7 +139,6 @@ fn computed_column_metadata(expression: &str, inputs: &[String]) -> HashMap<Stri
|
||||
INPUTS_META_KEY.to_string(),
|
||||
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
|
||||
),
|
||||
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
|
||||
])
|
||||
}
|
||||
|
||||
@@ -180,7 +163,6 @@ pub fn function_computed_column_metadata(
|
||||
INPUTS_META_KEY.to_string(),
|
||||
serde_json::to_string(inputs).unwrap_or_else(|_| "[]".to_string()),
|
||||
),
|
||||
(SOURCE_SIGNATURE_META_KEY.to_string(), "{}".to_string()),
|
||||
])
|
||||
}
|
||||
|
||||
@@ -1313,8 +1295,7 @@ pub(crate) fn ensure_not_an_input(schema: &SchemaRef, paths: &[&str]) -> Result<
|
||||
}
|
||||
|
||||
/// Reject a write that supplies values for a computed column directly:
|
||||
/// only refresh materializes one, and only refresh decides what it
|
||||
/// recomputes.
|
||||
/// only refresh materializes one, and refresh never revisits a filled row.
|
||||
pub(crate) fn ensure_not_written<'a>(
|
||||
schema: &ArrowSchema,
|
||||
written: impl IntoIterator<Item = &'a str>,
|
||||
@@ -1489,9 +1470,7 @@ fn ensure_no_foreign_declaration(field: &ArrowField) -> Result<()> {
|
||||
/// kind, the expression, the inputs -- would bypass that validation or move
|
||||
/// a binding out from under a refresh. Drop the column and declare it again.
|
||||
pub(crate) fn is_declaration_key(key: &str) -> bool {
|
||||
key == COMPUTED_COLUMN_META_KEY
|
||||
|| key.starts_with("computed_column.")
|
||||
|| key.starts_with("computed_refresh.")
|
||||
key == COMPUTED_COLUMN_META_KEY || key.starts_with("computed_column.")
|
||||
}
|
||||
|
||||
/// Reject retyping a computed column itself.
|
||||
|
||||
@@ -52,7 +52,7 @@ enum ConsistencyMode {
|
||||
/// refresh_window = min(3s, TTL/4)
|
||||
///
|
||||
/// | t < TTL - refresh_window | t < TTL | t >= TTL |
|
||||
/// | Return value | Background refresh & return value | synchronous refresh |
|
||||
/// | Return value | Background refresh & return value | syncronous refresh |
|
||||
Eventual(BackgroundCache<Arc<Dataset>, Error>),
|
||||
}
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -103,7 +103,7 @@ impl MergeInsertBuilder {
|
||||
/// but that behavior is subject to change.
|
||||
///
|
||||
/// An optional condition may be specified. If it is, then only
|
||||
/// matched rows that satisfy the condition will be updated. Any
|
||||
/// matched rows that satisfy the condtion will be updated. Any
|
||||
/// rows that do not satisfy the condition will be left as they
|
||||
/// are. Failing to satisfy the condition does not cause a
|
||||
/// "matched row" to become a "not matched" row.
|
||||
|
||||
@@ -904,7 +904,7 @@ fn unsharded_shard_id() -> Uuid {
|
||||
|
||||
/// Build a [`ShardWriterConfig`] from the persisted `writer_config_defaults`.
|
||||
///
|
||||
/// Unknown or unparsable keys are ignored; absent keys keep the
|
||||
/// Unknown or unparseable keys are ignored; absent keys keep the
|
||||
/// [`ShardWriterConfig`] default. The shard id is set by `mem_wal_writer`.
|
||||
fn shard_writer_config_from_defaults(defaults: &HashMap<String, String>) -> ShardWriterConfig {
|
||||
let mut config = ShardWriterConfig::default().with_shard_spec_id(SHARDING_SPEC_ID);
|
||||
|
||||
@@ -3,9 +3,9 @@
|
||||
|
||||
//! Filling computed columns.
|
||||
//!
|
||||
//! A row without a value gets one; a row that has one keeps it unless its
|
||||
//! fragment's inputs moved since it was computed, which `freshness` decides
|
||||
//! from the manifest and stamps after every fill.
|
||||
//! A row without a value gets one; a row that has one keeps it. Refresh is
|
||||
//! therefore idempotent and does not observe input mutation -- once a row is
|
||||
//! filled, changing what the expression reads leaves the stored result alone.
|
||||
//!
|
||||
//! A column's computed inputs are filled first -- the dependency graph is
|
||||
//! walked once, each reachable column filled once in dependency order, each
|
||||
@@ -31,7 +31,6 @@
|
||||
|
||||
use std::collections::HashSet;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
|
||||
use arrow_array::{
|
||||
Array, ArrayRef, BooleanArray, LargeBinaryArray, RecordBatch, RecordBatchOptions, StructArray,
|
||||
@@ -49,7 +48,6 @@ use lance_core::datatypes::{BlobHandling, Schema as LanceSchema};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use super::computed_columns::{BoundExpression, ComputedColumnKind, computed_column_from_field};
|
||||
use super::freshness::{self, SignatureMap, StalenessPlan};
|
||||
use super::{BaseTable, NativeTable};
|
||||
use crate::job::Job;
|
||||
use crate::{Error, Result};
|
||||
@@ -112,81 +110,28 @@ async fn execute_refresh_column_with_source(
|
||||
};
|
||||
let output_is_blob = field.is_blob_v2();
|
||||
|
||||
// Which fragments the null filter cannot speak for: their inputs moved
|
||||
// since they were computed, or the definition did. Decided once, from
|
||||
// the manifest the values are read from.
|
||||
let inputs = freshness::fields_for_paths(dataset.schema(), &bound.inputs)?;
|
||||
let definition = freshness::definition_version(&expression);
|
||||
let staleness = freshness::staleness_against(&dataset, column, &definition, &inputs).await?;
|
||||
|
||||
let mut rows_filled = 0u64;
|
||||
let mut replacements = Vec::new();
|
||||
// Fragments this refresh computed in full, signed at the version read.
|
||||
let mut computed = SignatureMap::new();
|
||||
for fragment in dataset.get_fragments() {
|
||||
let fragment_id = u32::try_from(fragment.id()).map_err(|_| Error::Runtime {
|
||||
message: format!("fragment id {} does not fit a signature map", fragment.id()),
|
||||
})?;
|
||||
// A recompute rewrites every live row, so it is staged without the
|
||||
// probe and counted as it fills; a null fill probes first, since a
|
||||
// fragment with nothing to gain is not worth a write.
|
||||
let recompute = staleness.is_dirty(fragment_id);
|
||||
let whole = recompute || {
|
||||
let (gained, unfilled) =
|
||||
count_fragment_gains(&dataset, &fragment, &bound, column).await?;
|
||||
if gained == 0 {
|
||||
continue;
|
||||
}
|
||||
rows_filled += gained;
|
||||
unfilled == u64::try_from(fragment.count_rows(None).await?).unwrap_or(u64::MAX)
|
||||
};
|
||||
if whole {
|
||||
computed.insert(
|
||||
fragment_id,
|
||||
freshness::fragment_input_signature(fragment.metadata(), &inputs)?,
|
||||
);
|
||||
let gained = count_fragment_gains(&dataset, &fragment, &bound, column).await?;
|
||||
if gained == 0 {
|
||||
continue;
|
||||
}
|
||||
let gained = Arc::new(AtomicU64::new(0));
|
||||
let values = fill_stream(
|
||||
&dataset,
|
||||
&fragment,
|
||||
bound.clone(),
|
||||
column,
|
||||
output_is_blob,
|
||||
recompute,
|
||||
gained.clone(),
|
||||
)
|
||||
.await?;
|
||||
rows_filled += gained;
|
||||
let values =
|
||||
fill_stream(&dataset, &fragment, bound.clone(), column, output_is_blob).await?;
|
||||
replacements.push(fragment.write_columns(values, &column_schema).await?);
|
||||
if recompute {
|
||||
rows_filled += gained.load(Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
|
||||
let source_version = dataset.version().version;
|
||||
if replacements.is_empty() {
|
||||
// Nothing to fill; the stamp may still have something to record -- a
|
||||
// column not yet enrolled, or fragments a compaction carried.
|
||||
let mut latest = (*dataset).clone();
|
||||
let stamped = record(
|
||||
&mut latest,
|
||||
(&dataset, &staleness),
|
||||
column,
|
||||
&definition,
|
||||
&inputs,
|
||||
computed,
|
||||
)
|
||||
.await;
|
||||
if stamped.is_some() {
|
||||
table.dataset.update(latest);
|
||||
}
|
||||
return Ok(RefreshExecution {
|
||||
result: RefreshColumnResult {
|
||||
rows_filled: 0,
|
||||
version: stamped.unwrap_or(source_version),
|
||||
version: source_version,
|
||||
},
|
||||
source_version,
|
||||
published_version: stamped,
|
||||
published_version: None,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -204,17 +149,7 @@ async fn execute_refresh_column_with_source(
|
||||
)
|
||||
.await?;
|
||||
|
||||
let mut new_dataset = new_dataset;
|
||||
let version = record(
|
||||
&mut new_dataset,
|
||||
(&dataset, &staleness),
|
||||
column,
|
||||
&definition,
|
||||
&inputs,
|
||||
computed,
|
||||
)
|
||||
.await
|
||||
.unwrap_or(new_dataset.version().version);
|
||||
let version = new_dataset.version().version;
|
||||
table.dataset.update(new_dataset);
|
||||
Ok(RefreshExecution {
|
||||
result: RefreshColumnResult {
|
||||
@@ -226,31 +161,6 @@ async fn execute_refresh_column_with_source(
|
||||
})
|
||||
}
|
||||
|
||||
/// Stamp the input state the refresh computed from (see
|
||||
/// [`freshness::record_freshness`]); the version the stamp landed at, which
|
||||
/// is the last one the refresh wrote. Never fails the refresh: the values
|
||||
/// are committed, and a missing stamp only costs a recompute next time.
|
||||
async fn record(
|
||||
latest: &mut Dataset,
|
||||
pinned: (&Dataset, &StalenessPlan),
|
||||
column: &str,
|
||||
definition: &str,
|
||||
inputs: &freshness::InputFields,
|
||||
computed: SignatureMap,
|
||||
) -> Option<u64> {
|
||||
match freshness::record_freshness(latest, Some(pinned), column, definition, inputs, computed)
|
||||
.await
|
||||
{
|
||||
Ok(record) => record.version,
|
||||
Err(error) => {
|
||||
log::warn!(
|
||||
"could not record the input state computed column '{column}' was refreshed from ({error}); its fragments will recompute on the next refresh"
|
||||
);
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Refuse while a computed input still has rows a refresh of it would fill:
|
||||
/// read now, its placeholder null would be evaluated as a value and kept.
|
||||
async fn ensure_inputs_filled(
|
||||
@@ -278,9 +188,7 @@ async fn ensure_inputs_filled(
|
||||
let input_bound = super::computed_columns::bind(schema.clone(), input, expression)?;
|
||||
let mut unfilled = 0u64;
|
||||
for fragment in dataset.get_fragments() {
|
||||
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input)
|
||||
.await?
|
||||
.0;
|
||||
unfilled += count_fragment_gains(dataset, &fragment, &input_bound, input).await?;
|
||||
}
|
||||
if unfilled > 0 {
|
||||
return Err(Error::InvalidInput {
|
||||
@@ -528,35 +436,32 @@ fn blob_array_from_binary(
|
||||
/// Scans only the unfilled live rows -- deleted rows never reach the
|
||||
/// expression here, the filter having already excluded them -- and counts the
|
||||
/// non-null results. Exact, so it is both the staging decision and the
|
||||
/// fragment's contribution to `rows_filled`. Returns the gains and the rows
|
||||
/// scanned.
|
||||
/// fragment's contribution to `rows_filled`.
|
||||
async fn count_fragment_gains(
|
||||
dataset: &Dataset,
|
||||
fragment: &FileFragment,
|
||||
bound: &BoundExpression,
|
||||
column: &str,
|
||||
) -> Result<(u64, u64)> {
|
||||
) -> Result<u64> {
|
||||
let mut scanner = dataset.scan();
|
||||
scanner
|
||||
.with_fragments(vec![fragment.metadata().clone()])
|
||||
.with_row_id()
|
||||
.project(&bound.roots)?
|
||||
.filter(&format!("{} IS NULL", quote_identifier(column)))?;
|
||||
.filter(&format!("{} IS NULL", quote_identifier(column)))?
|
||||
.project(&bound.roots)?;
|
||||
configure_blob_inputs(&mut scanner, dataset.schema(), bound, None)?;
|
||||
|
||||
let mut gained = 0u64;
|
||||
let mut considered = 0u64;
|
||||
let mut batches = scanner.try_into_stream().await?;
|
||||
while let Some(batch) = batches.try_next().await? {
|
||||
let evaluated = evaluate(bound, &evaluation_batch(&batch, bound, None)?)?;
|
||||
gained += (batch.num_rows() - evaluated.null_count()) as u64;
|
||||
considered += batch.num_rows() as u64;
|
||||
}
|
||||
Ok((gained, considered))
|
||||
Ok(gained)
|
||||
}
|
||||
|
||||
/// Stream one fragment's column in physical order, filling the unfilled live
|
||||
/// rows -- every live row, for a recompute -- and keeping every other value.
|
||||
/// rows and keeping every other value.
|
||||
///
|
||||
/// Deleted rows are carried through so the values line up positionally with
|
||||
/// the fragment's data files; they are never read back, but the column file
|
||||
@@ -567,8 +472,6 @@ async fn fill_stream(
|
||||
bound: Arc<BoundExpression>,
|
||||
column: &str,
|
||||
output_is_blob: bool,
|
||||
recompute: bool,
|
||||
gained: Arc<AtomicU64>,
|
||||
) -> Result<impl Stream<Item = lance_core::Result<RecordBatch>> + Send + use<>> {
|
||||
let mut projection: Vec<String> = bound.roots.clone();
|
||||
projection.push(column.to_string());
|
||||
@@ -618,20 +521,14 @@ async fn fill_stream(
|
||||
.column_by_name(ROW_ID)
|
||||
.ok_or_else(|| missing(ROW_ID))?;
|
||||
|
||||
// Only an unfilled live row gains a value, or every live row under a
|
||||
// recompute; a deleted row has a null row id and keeps its (null) slot.
|
||||
// Only an unfilled live row gains a value; a deleted row has a null
|
||||
// row id and keeps its (null) slot.
|
||||
let unfilled = arrow::compute::is_null(existing.as_ref())?;
|
||||
let live = arrow::compute::is_not_null(row_ids.as_ref())?;
|
||||
let fill = if recompute {
|
||||
live
|
||||
} else {
|
||||
let unfilled = arrow::compute::is_null(existing.as_ref())?;
|
||||
arrow::compute::and(&unfilled, &live)?
|
||||
};
|
||||
let fill = arrow::compute::and(&unfilled, &live)?;
|
||||
let keep = arrow::compute::not(&fill)?;
|
||||
|
||||
let computed = evaluate(&bound, &evaluation_batch(&batch, &bound, Some(&keep))?)?;
|
||||
let values = arrow::compute::and(&fill, &arrow::compute::is_not_null(&computed)?)?;
|
||||
gained.fetch_add(values.true_count() as u64, Ordering::Relaxed);
|
||||
let merged = arrow_select::zip::zip(&fill, &computed, existing)?;
|
||||
let merged = if output_is_blob {
|
||||
blob_array_from_binary(&merged, projected.field(0))?
|
||||
@@ -840,8 +737,7 @@ mod tests {
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(no_op.rows_assigned, 0);
|
||||
// The fill, then the stamp recording what it computed from.
|
||||
assert_eq!(no_op.source_version, 4);
|
||||
assert_eq!(no_op.source_version, 3);
|
||||
assert_eq!(no_op.published_version, None);
|
||||
}
|
||||
|
||||
@@ -881,8 +777,8 @@ mod tests {
|
||||
}
|
||||
|
||||
/// A row is filled only by gaining a value, so an expression yielding null
|
||||
/// settles at once instead of re-selecting the same rows forever: the
|
||||
/// second refresh finds the fragment signed and moves nothing.
|
||||
/// settles at once instead of re-selecting the same rows forever. Nothing
|
||||
/// is staged, so the version does not move either.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_converges_on_a_null_result() {
|
||||
let table = table_with("refresh_null_result", vec![1, 2, 3]).await;
|
||||
@@ -896,145 +792,28 @@ mod tests {
|
||||
|
||||
let first = table.refresh_column("maybe").await.unwrap();
|
||||
assert_eq!(first.rows_filled, 0);
|
||||
assert!(first.version > declared);
|
||||
assert_eq!(first.version, declared);
|
||||
assert_eq!(read(&table, "maybe").await, vec![None, None, None]);
|
||||
|
||||
let again = table.refresh_column("maybe").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 0);
|
||||
assert_eq!(again.version, first.version);
|
||||
assert_eq!(again.version, declared);
|
||||
}
|
||||
|
||||
/// A filled row whose input moved is recomputed: the update rewrites
|
||||
/// the row into a fragment the stamp never signed, and only that one.
|
||||
/// The contract's boundary: a filled fragment is not revisited, so
|
||||
/// mutating an input leaves the value computed at fill time.
|
||||
#[tokio::test]
|
||||
async fn test_refresh_recomputes_a_row_whose_input_moved() {
|
||||
let table = table_with("refresh_mutation", vec![1, 2]).await;
|
||||
async fn test_refresh_does_not_observe_input_mutation() {
|
||||
let table = table_with("refresh_mutation", vec![1]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
|
||||
|
||||
table
|
||||
.update()
|
||||
.column("x", "7")
|
||||
.only_if("x = 5")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 1);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(14)]
|
||||
);
|
||||
|
||||
let settled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(settled.rows_filled, 0);
|
||||
assert_eq!(settled.version, again.version);
|
||||
}
|
||||
|
||||
/// A deleted row is never computed and the rows that stay keep their
|
||||
/// values: a delete recomputes nothing and stamps nothing.
|
||||
#[tokio::test]
|
||||
async fn test_a_delete_recomputes_nothing() {
|
||||
let table = table_with("refresh_delete", vec![1, 2, 3]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
let filled = table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
table.delete("x = 2").await.unwrap();
|
||||
let deleted = table.version().await.unwrap();
|
||||
table.update().column("x", "3").execute().await.unwrap();
|
||||
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 0);
|
||||
assert_eq!(again.version, deleted);
|
||||
assert!(deleted > filled.version);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(6)]);
|
||||
}
|
||||
|
||||
/// Compaction copies inputs unchanged, so a fragment it builds from
|
||||
/// signed ones is fresh: the refresh recomputes nothing and only records
|
||||
/// the new fragment.
|
||||
#[tokio::test]
|
||||
async fn test_a_compaction_of_signed_fragments_recomputes_nothing() {
|
||||
let table = table_with("refresh_compact_signed", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
table
|
||||
.optimize(crate::table::OptimizeAction::Compact {
|
||||
options: crate::table::CompactionOptions::default(),
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
let compacted = table.version().await.unwrap();
|
||||
|
||||
let carried = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(carried.rows_filled, 0);
|
||||
assert_eq!(carried.version, compacted + 1);
|
||||
let settled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(settled.version, carried.version);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
}
|
||||
|
||||
/// A column declared before signatures existed has no map. Its first
|
||||
/// refresh keeps the null-fill contract and enrolls what it read from;
|
||||
/// from then on a moved input is recomputed like any other.
|
||||
#[tokio::test]
|
||||
async fn test_an_unsigned_column_is_enrolled_by_its_first_refresh() {
|
||||
let table = table_with("refresh_legacy", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
table
|
||||
.update()
|
||||
.column("x", "3")
|
||||
.only_if("x = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let native = table.as_native().unwrap();
|
||||
let mut dataset = (*native.dataset.get().await.unwrap()).clone();
|
||||
let declaration = dataset
|
||||
.schema()
|
||||
.field("doubled")
|
||||
.unwrap()
|
||||
.metadata
|
||||
.iter()
|
||||
.filter(|(key, _)| !key.starts_with("computed_refresh."))
|
||||
.map(|(key, value)| (key.clone(), value.clone()))
|
||||
.collect::<Vec<_>>();
|
||||
dataset
|
||||
.update_field_metadata()
|
||||
.replace("doubled", declaration)
|
||||
.unwrap()
|
||||
.await
|
||||
.unwrap();
|
||||
table.checkout_latest().await.unwrap();
|
||||
|
||||
// Null-fill only: the moved row keeps the value it was filled with.
|
||||
let enrolled = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(enrolled.rows_filled, 0);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2), Some(4)]);
|
||||
|
||||
table
|
||||
.update()
|
||||
.column("x", "5")
|
||||
.only_if("x = 3")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let again = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(again.rows_filled, 1);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(4), Some(10)]);
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(2)]);
|
||||
}
|
||||
|
||||
/// A row rewrite before the first refresh materializes the declared
|
||||
@@ -1052,11 +831,11 @@ mod tests {
|
||||
assert_eq!(read(&table, "doubled").await, vec![Some(6)]);
|
||||
}
|
||||
|
||||
/// A fragment compacted out of one the stamp never signed cannot vouch
|
||||
/// for any of its rows: every live row is recomputed, the moved one
|
||||
/// included.
|
||||
/// The contract holds row by row, not fragment by fragment: revisiting a
|
||||
/// fragment to fill one row must not recompute a filled row sitting beside
|
||||
/// it, even where the input behind it has since changed.
|
||||
#[tokio::test]
|
||||
async fn test_a_compaction_of_an_unsigned_fragment_recomputes_it() {
|
||||
async fn test_refresh_does_not_recompute_a_filled_row_beside_an_unfilled_one() {
|
||||
let table = table_with("refresh_mixed", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
@@ -1078,70 +857,18 @@ mod tests {
|
||||
.unwrap();
|
||||
|
||||
let result = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 3);
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(4), Some(10), Some(200)]
|
||||
);
|
||||
}
|
||||
|
||||
/// The gate's reproducer: a raw lance append may carry a value for the
|
||||
/// computed column. Compaction cannot certify it, so the product is
|
||||
/// recomputed and the supplied value replaced.
|
||||
#[tokio::test]
|
||||
async fn test_raw_append_values_are_not_trusted_after_compaction() {
|
||||
use arrow_array::RecordBatchIterator;
|
||||
use lance::Dataset;
|
||||
use lance::dataset::{WriteMode, WriteParams};
|
||||
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
let conn = connect(dir.path().to_str().unwrap())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let batch = record_batch!(("x", Int32, [1, 2])).unwrap();
|
||||
let table = conn
|
||||
.create_table("raw_append", batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
|
||||
let batch = record_batch!(("x", Int32, [5]), ("doubled", Int32, [Some(999_i32)])).unwrap();
|
||||
let schema = batch.schema();
|
||||
let uri = table.uri().await.unwrap();
|
||||
Dataset::write(
|
||||
RecordBatchIterator::new(vec![Ok(batch)], schema),
|
||||
&uri,
|
||||
Some(WriteParams {
|
||||
mode: WriteMode::Append,
|
||||
..Default::default()
|
||||
}),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
table.checkout_latest().await.unwrap();
|
||||
table
|
||||
.optimize(crate::table::OptimizeAction::Compact {
|
||||
options: crate::table::CompactionOptions::default(),
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let result = table.refresh_column("doubled").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 3);
|
||||
assert_eq!(result.rows_filled, 1);
|
||||
// 2 is the mutated row keeping the value it was filled with, not 200.
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
vec![Some(2), Some(4), Some(10)]
|
||||
);
|
||||
}
|
||||
|
||||
/// An appended fragment holds no values, so compacting it into a signed
|
||||
/// one leaves the product fresh: only the appended rows are filled.
|
||||
/// Filling a fragment must not disturb the values it already holds, which
|
||||
/// is what makes a compaction-mixed fragment safe to revisit.
|
||||
#[tokio::test]
|
||||
async fn test_a_compaction_with_an_appended_fragment_fills_only_its_rows() {
|
||||
async fn test_refresh_preserves_already_filled_rows() {
|
||||
let table = table_with("refresh_preserves", vec![1, 2]).await;
|
||||
declare_doubled(&table).await.unwrap();
|
||||
table.refresh_column("doubled").await.unwrap();
|
||||
@@ -1268,8 +995,7 @@ mod tests {
|
||||
assert_eq!(result.rows_failed, 0);
|
||||
assert_eq!(result.rows_remaining, 0);
|
||||
assert_eq!(result.source_version, 2);
|
||||
// The fill lands at 3; the stamp recording its inputs is published at 4.
|
||||
assert_eq!(result.published_version, Some(4));
|
||||
assert_eq!(result.published_version, Some(3));
|
||||
assert_eq!(job.status().await.unwrap(), "finished");
|
||||
assert_eq!(
|
||||
read(&table, "doubled").await,
|
||||
@@ -1333,32 +1059,31 @@ mod tests {
|
||||
assert_eq!(read(&table, "quotient").await, vec![Some(10)]);
|
||||
}
|
||||
|
||||
/// A filled row whose input moved is re-evaluated, and a row whose
|
||||
/// input did not move is not: the untouched fragment is never read, so
|
||||
/// its poison input is never reached.
|
||||
/// The gate's reproducer: an already-filled row's value must not be
|
||||
/// re-evaluated either -- its input may have mutated into one the
|
||||
/// expression chokes on.
|
||||
#[tokio::test]
|
||||
async fn test_only_a_moved_rows_value_is_re_evaluated() {
|
||||
let table = table_with("refresh_filled_poison", vec![1, 0]).await;
|
||||
async fn test_a_filled_rows_value_is_never_evaluated() {
|
||||
let table = table_with("refresh_filled_poison", vec![1, 2]).await;
|
||||
table
|
||||
.add_columns()
|
||||
.computed("quotient", "10 / coalesce(nullif(x, 0), 1)")
|
||||
.computed("quotient", "10 / x")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
table.refresh_column("quotient").await.unwrap();
|
||||
assert_eq!(read(&table, "quotient").await, vec![Some(10), Some(10)]);
|
||||
|
||||
append(&table, vec![5]).await;
|
||||
table
|
||||
.update()
|
||||
.column("x", "2")
|
||||
.column("x", "0")
|
||||
.only_if("x = 1")
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
append(&table, vec![5]).await;
|
||||
|
||||
let result = table.refresh_column("quotient").await.unwrap();
|
||||
assert_eq!(result.rows_filled, 2);
|
||||
assert_eq!(result.rows_filled, 1);
|
||||
assert_eq!(
|
||||
read(&table, "quotient").await,
|
||||
vec![Some(2), Some(5), Some(10)]
|
||||
|
||||
Reference in New Issue
Block a user