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
+45
-45
@@ -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",
|
||||
@@ -5567,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",
|
||||
@@ -5592,7 +5592,7 @@ dependencies = [
|
||||
|
||||
[[package]]
|
||||
name = "lancedb-python"
|
||||
version = "0.39.0-beta.6"
|
||||
version = "0.39.0-beta.5"
|
||||
dependencies = [
|
||||
"arrow",
|
||||
"async-trait",
|
||||
|
||||
+14
-14
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
lance = { "version" = "=12.0.0-beta.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.
|
||||
|
||||
@@ -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>
|
||||
|
||||
+1
-1
@@ -44,6 +44,6 @@ aws-lc-rs = "=1.16.3"
|
||||
napi-build = "2.3.1"
|
||||
|
||||
[features]
|
||||
default = ["remote", "lancedb/sql", "lancedb/aws", "lancedb/gcs", "lancedb/azure", "lancedb/dynamodb", "lancedb/oss", "lancedb/huggingface", "lancedb/goosefs", "lancedb/metrics-otel"]
|
||||
default = ["remote", "lancedb/aws", "lancedb/gcs", "lancedb/azure", "lancedb/dynamodb", "lancedb/oss", "lancedb/huggingface", "lancedb/goosefs", "lancedb/metrics-otel"]
|
||||
fp16kernels = ["lancedb/fp16kernels"]
|
||||
remote = ["lancedb/remote"]
|
||||
|
||||
@@ -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
|
||||
|
||||
+1
-1
@@ -47,6 +47,6 @@ libc = "0.2"
|
||||
pyo3-build-config = { version = "0.28", features = ["abi3-py310"] }
|
||||
|
||||
[features]
|
||||
default = ["remote", "lancedb/sql", "lancedb/aws", "lancedb/gcs", "lancedb/azure", "lancedb/dynamodb", "lancedb/oss", "lancedb/huggingface", "lancedb/cos", "lancedb/goosefs", "lancedb/metrics-otel"]
|
||||
default = ["remote", "lancedb/aws", "lancedb/gcs", "lancedb/azure", "lancedb/dynamodb", "lancedb/oss", "lancedb/huggingface", "lancedb/cos", "lancedb/goosefs", "lancedb/metrics-otel"]
|
||||
fp16kernels = ["lancedb/fp16kernels"]
|
||||
remote = ["lancedb/remote"]
|
||||
|
||||
@@ -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
|
||||
@@ -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
|
||||
@@ -5638,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
|
||||
@@ -5814,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
|
||||
|
||||
@@ -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
|
||||
|
||||
+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>,
|
||||
|
||||
@@ -120,13 +120,7 @@ pprof = { version = "0.14", features = ["flamegraph"] }
|
||||
|
||||
|
||||
[features]
|
||||
default = ["sql"]
|
||||
# The SQL statement extension seam (`lancedb::sql`): the registry a host adds
|
||||
# dialect statements through, the access vocabulary those statements declare,
|
||||
# and the write-commit observer. It pulls in no dependency that is not already
|
||||
# required, so it is on by default; the flag exists so an embedder that does
|
||||
# not want the surface can opt out of it.
|
||||
sql = []
|
||||
default = []
|
||||
aws = [
|
||||
"lance/aws",
|
||||
"lance-io/aws",
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -1,35 +1,7 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! SQL: handles to queries running on a remote database, and the seam an
|
||||
//! embedder extends the dialect through.
|
||||
//!
|
||||
//! The extension seam is behind the default-on `sql` feature. It lets a host
|
||||
//! add statements to the dialect from outside this crate: a statement brings
|
||||
//! its own grammar ([`CustomSqlHandler`]), and declares its audit label and
|
||||
//! the access it needs ([`SqlStatement`]), so the host's authorization and
|
||||
//! auditing do not have to know each statement by name.
|
||||
|
||||
#[cfg(feature = "sql")]
|
||||
mod dml;
|
||||
#[cfg(feature = "sql")]
|
||||
mod observer;
|
||||
#[cfg(feature = "sql")]
|
||||
mod parser;
|
||||
#[cfg(feature = "sql")]
|
||||
mod statement;
|
||||
|
||||
#[cfg(feature = "sql")]
|
||||
pub use dml::{DmlOperation, DmlResult, dml_result_schema};
|
||||
#[cfg(feature = "sql")]
|
||||
pub use observer::{CommittedWrite, DmlEventKind, WriteObserver, observe_write};
|
||||
#[cfg(feature = "sql")]
|
||||
pub use parser::route_custom_sql;
|
||||
#[cfg(feature = "sql")]
|
||||
pub use statement::{
|
||||
AccessRequirement, CreateKind, CustomSqlHandler, DatabaseScope, RelationKind,
|
||||
RequirementContext, SqlStatement, StatementRegistry, SystemScope, WriteMode,
|
||||
};
|
||||
//! Handles to SQL queries running on a remote database.
|
||||
|
||||
use std::{fmt, sync::Arc};
|
||||
|
||||
@@ -1,150 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! The result of a DML statement, as a one-row record batch.
|
||||
//!
|
||||
//! A DML statement has to answer over the same channel a query does, so its
|
||||
//! result is carried as an ordinary [`RecordBatch`] with a fixed schema. The
|
||||
//! round trip is lossless, which is what lets a caller recover the typed form
|
||||
//! after the batch has crossed a transport such as Arrow Flight.
|
||||
|
||||
use std::fmt;
|
||||
use std::sync::Arc;
|
||||
|
||||
use arrow_array::{Int64Array, RecordBatch, StringArray};
|
||||
use arrow_schema::{DataType, Field, Schema};
|
||||
|
||||
/// Which DML statement produced a result.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum DmlOperation {
|
||||
Insert,
|
||||
Update,
|
||||
Delete,
|
||||
}
|
||||
|
||||
impl DmlOperation {
|
||||
pub fn as_str(self) -> &'static str {
|
||||
match self {
|
||||
Self::Insert => "INSERT",
|
||||
Self::Update => "UPDATE",
|
||||
Self::Delete => "DELETE",
|
||||
}
|
||||
}
|
||||
|
||||
fn parse(s: &str) -> Option<Self> {
|
||||
match s {
|
||||
"INSERT" => Some(Self::Insert),
|
||||
"UPDATE" => Some(Self::Update),
|
||||
"DELETE" => Some(Self::Delete),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl fmt::Display for DmlOperation {
|
||||
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
|
||||
f.write_str(self.as_str())
|
||||
}
|
||||
}
|
||||
|
||||
/// What a DML statement did.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub struct DmlResult {
|
||||
pub table: String,
|
||||
pub operation: DmlOperation,
|
||||
pub rows_affected: i64,
|
||||
pub version: i64,
|
||||
}
|
||||
|
||||
/// The schema every [`DmlResult`] batch carries.
|
||||
pub fn dml_result_schema() -> Schema {
|
||||
Schema::new(vec![
|
||||
Field::new("table", DataType::Utf8, false),
|
||||
Field::new("operation", DataType::Utf8, false),
|
||||
Field::new("rows_affected", DataType::Int64, false),
|
||||
Field::new("version", DataType::Int64, false),
|
||||
])
|
||||
}
|
||||
|
||||
impl DmlResult {
|
||||
pub fn new(
|
||||
table: impl Into<String>,
|
||||
operation: DmlOperation,
|
||||
rows_affected: i64,
|
||||
version: i64,
|
||||
) -> Self {
|
||||
Self {
|
||||
table: table.into(),
|
||||
operation,
|
||||
rows_affected,
|
||||
version,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn to_record_batch(&self) -> RecordBatch {
|
||||
RecordBatch::try_new(
|
||||
Arc::new(dml_result_schema()),
|
||||
vec![
|
||||
Arc::new(StringArray::from(vec![self.table.as_str()])),
|
||||
Arc::new(StringArray::from(vec![self.operation.as_str()])),
|
||||
Arc::new(Int64Array::from(vec![self.rows_affected])),
|
||||
Arc::new(Int64Array::from(vec![self.version])),
|
||||
],
|
||||
)
|
||||
.expect("static schema")
|
||||
}
|
||||
|
||||
/// Recover a result from a batch, or `None` if the batch is not one.
|
||||
///
|
||||
/// A query result can arrive on the same channel, so this has to be able
|
||||
/// to say "not a DML result" rather than fail.
|
||||
pub fn try_from_batch(batch: &RecordBatch) -> Option<Self> {
|
||||
if *batch.schema().as_ref() != dml_result_schema() || batch.num_rows() != 1 {
|
||||
return None;
|
||||
}
|
||||
Some(Self {
|
||||
table: batch
|
||||
.column(0)
|
||||
.as_any()
|
||||
.downcast_ref::<StringArray>()?
|
||||
.value(0)
|
||||
.to_string(),
|
||||
operation: DmlOperation::parse(
|
||||
batch
|
||||
.column(1)
|
||||
.as_any()
|
||||
.downcast_ref::<StringArray>()?
|
||||
.value(0),
|
||||
)?,
|
||||
rows_affected: batch
|
||||
.column(2)
|
||||
.as_any()
|
||||
.downcast_ref::<Int64Array>()?
|
||||
.value(0),
|
||||
version: batch
|
||||
.column(3)
|
||||
.as_any()
|
||||
.downcast_ref::<Int64Array>()?
|
||||
.value(0),
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn roundtrip() {
|
||||
let original = DmlResult::new("foo", DmlOperation::Insert, 1, 3);
|
||||
let batch = original.to_record_batch();
|
||||
assert_eq!(DmlResult::try_from_batch(&batch), Some(original));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn non_dml_returns_none() {
|
||||
let s = Arc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
|
||||
let batch = RecordBatch::try_new(s, vec![Arc::new(Int64Array::from(vec![1]))]).unwrap();
|
||||
assert_eq!(DmlResult::try_from_batch(&batch), None);
|
||||
}
|
||||
}
|
||||
@@ -1,70 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! Notification of committed writes.
|
||||
//!
|
||||
//! A statement that writes rows often has to tell its host that it did, so
|
||||
//! that follow-up work can be scheduled. What it should *not* have to know is
|
||||
//! how the host represents that notification. [`WriteObserver`] is the seam:
|
||||
//! the statement reports what it wrote, and the host decides what that means
|
||||
//! -- an event on a bus, a metric, or nothing at all.
|
||||
|
||||
use std::sync::Arc;
|
||||
|
||||
use async_trait::async_trait;
|
||||
|
||||
/// Which DML operation committed.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
pub enum DmlEventKind {
|
||||
Insert,
|
||||
Update,
|
||||
Delete,
|
||||
}
|
||||
|
||||
/// A write that has already been made durable.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct CommittedWrite {
|
||||
/// The database holding the table.
|
||||
pub database: String,
|
||||
/// The schema the statement named the table through.
|
||||
pub schema: String,
|
||||
/// The table written.
|
||||
pub table: String,
|
||||
/// The table's storage location, when the statement resolved one.
|
||||
pub table_uri: Option<String>,
|
||||
/// Which operation committed.
|
||||
pub kind: DmlEventKind,
|
||||
}
|
||||
|
||||
/// Notified after a statement's write commits.
|
||||
///
|
||||
/// Implementations are best-effort by contract: the write is already durable
|
||||
/// when this is called, so an observer that fails must not fail the statement.
|
||||
/// That is why the method cannot report an error.
|
||||
#[async_trait]
|
||||
pub trait WriteObserver: Send + Sync {
|
||||
async fn write_committed(&self, write: CommittedWrite);
|
||||
}
|
||||
|
||||
/// Report a committed write, if anything is observing.
|
||||
pub async fn observe_write(
|
||||
observer: Option<&Arc<dyn WriteObserver>>,
|
||||
database: &str,
|
||||
schema: &str,
|
||||
table: &str,
|
||||
table_uri: Option<String>,
|
||||
kind: DmlEventKind,
|
||||
) {
|
||||
let Some(observer) = observer else {
|
||||
return;
|
||||
};
|
||||
observer
|
||||
.write_committed(CommittedWrite {
|
||||
database: database.to_string(),
|
||||
schema: schema.to_string(),
|
||||
table: table.to_string(),
|
||||
table_uri,
|
||||
kind,
|
||||
})
|
||||
.await;
|
||||
}
|
||||
@@ -1,174 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! Routing a statement to the grammar that owns it.
|
||||
|
||||
use datafusion::error::{DataFusionError, Result};
|
||||
use datafusion::logical_expr::LogicalPlan;
|
||||
use datafusion::sql::sqlparser::{
|
||||
dialect::GenericDialect,
|
||||
parser::{Parser, ParserError},
|
||||
tokenizer::{Token, Tokenizer, TokenizerError},
|
||||
};
|
||||
|
||||
use super::statement::StatementRegistry;
|
||||
|
||||
/// Route a statement through the registry's grammars.
|
||||
///
|
||||
/// Returns `Ok(None)` when no grammar claims the statement, which is the
|
||||
/// caller's cue to hand it to DataFusion's own planner.
|
||||
///
|
||||
/// The first grammar whose `matches` accepts the tokens is the only one given
|
||||
/// the statement: a grammar that matches and then returns `Ok(None)` declines
|
||||
/// the form rather than falling through to the next grammar. Registration
|
||||
/// order therefore decides reachability, which is why [`StatementRegistry`]
|
||||
/// fixes it explicitly.
|
||||
pub fn route_custom_sql(registry: &StatementRegistry, sql: &str) -> Result<Option<LogicalPlan>> {
|
||||
let dialect = GenericDialect {};
|
||||
let mut tokenizer = Tokenizer::new(&dialect, sql);
|
||||
let tokens = tokenizer.tokenize().map_err(|e: TokenizerError| {
|
||||
DataFusionError::SQL(Box::new(ParserError::TokenizerError(e.to_string())), None)
|
||||
})?;
|
||||
|
||||
// Handlers match on keywords, so layout must not change the decision.
|
||||
let word_tokens: Vec<&Token> = tokens
|
||||
.iter()
|
||||
.filter(|t| !matches!(t, Token::Whitespace(_)))
|
||||
.collect();
|
||||
|
||||
for handler in registry.parsers() {
|
||||
if handler.matches(&word_tokens) {
|
||||
// `Parser` takes ownership of the tokens, so it is built only once
|
||||
// a handler has claimed the statement.
|
||||
let mut parser = Parser::new(&dialect).with_tokens(tokens.clone());
|
||||
if let Some(plan) = handler.parse(&mut parser)? {
|
||||
return Ok(Some(plan));
|
||||
}
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
Ok(None)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::sync::{
|
||||
Arc,
|
||||
atomic::{AtomicUsize, Ordering},
|
||||
};
|
||||
|
||||
use datafusion::common::DFSchema;
|
||||
use datafusion::error::Result as DfResult;
|
||||
use datafusion::logical_expr::{EmptyRelation, LogicalPlan};
|
||||
use datafusion::sql::sqlparser::keywords::Keyword;
|
||||
|
||||
use super::*;
|
||||
use crate::sql::statement::CustomSqlHandler;
|
||||
|
||||
fn empty_plan() -> LogicalPlan {
|
||||
LogicalPlan::EmptyRelation(EmptyRelation {
|
||||
produce_one_row: false,
|
||||
schema: Arc::new(DFSchema::empty()),
|
||||
})
|
||||
}
|
||||
|
||||
/// Matches on a leading keyword, and reports whether it was asked to parse.
|
||||
struct Handler {
|
||||
keyword: Keyword,
|
||||
outcome: Outcome,
|
||||
parsed: Arc<AtomicUsize>,
|
||||
}
|
||||
|
||||
enum Outcome {
|
||||
Plans,
|
||||
Declines,
|
||||
}
|
||||
|
||||
impl Handler {
|
||||
fn new(keyword: Keyword, outcome: Outcome) -> (Arc<Self>, Arc<AtomicUsize>) {
|
||||
let parsed = Arc::new(AtomicUsize::new(0));
|
||||
let handler = Arc::new(Self {
|
||||
keyword,
|
||||
outcome,
|
||||
parsed: parsed.clone(),
|
||||
});
|
||||
(handler, parsed)
|
||||
}
|
||||
}
|
||||
|
||||
impl CustomSqlHandler for Handler {
|
||||
fn matches(&self, tokens: &[&Token]) -> bool {
|
||||
matches!(tokens.first(), Some(Token::Word(w)) if w.keyword == self.keyword)
|
||||
}
|
||||
|
||||
fn parse(&self, _parser: &mut Parser) -> DfResult<Option<LogicalPlan>> {
|
||||
self.parsed.fetch_add(1, Ordering::SeqCst);
|
||||
Ok(match self.outcome {
|
||||
Outcome::Plans => Some(empty_plan()),
|
||||
Outcome::Declines => None,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_unclaimed_statement_is_left_for_datafusion() {
|
||||
let registry = StatementRegistry::new();
|
||||
assert!(route_custom_sql(®istry, "SELECT 1").unwrap().is_none());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn whitespace_does_not_change_which_handler_matches() {
|
||||
let (handler, parsed) = Handler::new(Keyword::EXPLAIN, Outcome::Plans);
|
||||
let mut registry = StatementRegistry::new();
|
||||
registry.register_parser(handler);
|
||||
|
||||
for sql in ["EXPLAIN t", " EXPLAIN\n\t t "] {
|
||||
assert!(route_custom_sql(®istry, sql).unwrap().is_some());
|
||||
}
|
||||
assert_eq!(parsed.load(Ordering::SeqCst), 2);
|
||||
}
|
||||
|
||||
/// Front-insertion is what lets an extension get ahead of a catch-all that
|
||||
/// would otherwise swallow the same keyword.
|
||||
#[test]
|
||||
fn the_last_registered_handler_is_consulted_first() {
|
||||
let (first, first_parsed) = Handler::new(Keyword::EXPLAIN, Outcome::Plans);
|
||||
let (second, second_parsed) = Handler::new(Keyword::EXPLAIN, Outcome::Plans);
|
||||
|
||||
let mut registry = StatementRegistry::new();
|
||||
registry.register_parser(first).register_parser(second);
|
||||
|
||||
assert!(route_custom_sql(®istry, "EXPLAIN t").unwrap().is_some());
|
||||
assert_eq!(second_parsed.load(Ordering::SeqCst), 1);
|
||||
assert_eq!(first_parsed.load(Ordering::SeqCst), 0);
|
||||
}
|
||||
|
||||
/// A handler that matches and declines vetoes the statement rather than
|
||||
/// letting a later handler see it. Shadowing is silent, which is why
|
||||
/// registration order is part of the contract.
|
||||
#[test]
|
||||
fn a_handler_that_declines_shadows_the_handlers_behind_it() {
|
||||
let (shadowed, shadowed_parsed) = Handler::new(Keyword::EXPLAIN, Outcome::Plans);
|
||||
let (decliner, decliner_parsed) = Handler::new(Keyword::EXPLAIN, Outcome::Declines);
|
||||
|
||||
let mut registry = StatementRegistry::new();
|
||||
registry.register_parser(shadowed).register_parser(decliner);
|
||||
|
||||
assert!(route_custom_sql(®istry, "EXPLAIN t").unwrap().is_none());
|
||||
assert_eq!(decliner_parsed.load(Ordering::SeqCst), 1);
|
||||
assert_eq!(shadowed_parsed.load(Ordering::SeqCst), 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_parts_keeps_the_order_it_was_given() {
|
||||
let (first, first_parsed) = Handler::new(Keyword::EXPLAIN, Outcome::Plans);
|
||||
let (second, second_parsed) = Handler::new(Keyword::EXPLAIN, Outcome::Plans);
|
||||
|
||||
let registry = StatementRegistry::from_parts(vec![first, second], vec![]);
|
||||
|
||||
assert!(route_custom_sql(®istry, "EXPLAIN t").unwrap().is_some());
|
||||
assert_eq!(first_parsed.load(Ordering::SeqCst), 1);
|
||||
assert_eq!(second_parsed.load(Ordering::SeqCst), 0);
|
||||
}
|
||||
}
|
||||
@@ -1,365 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! The statement registry: the seam between the SQL dialect and the behaviour
|
||||
//! an embedder adds to it.
|
||||
//!
|
||||
//! A statement owns three things that are otherwise easy to spread across
|
||||
//! parallel `downcast_ref` chains: the grammar that produces its node, the
|
||||
//! audit label it reports, and the access it requires. Keeping them in one
|
||||
//! place is what makes a statement addable from outside this crate.
|
||||
//!
|
||||
//! Two registries, because the axes differ. Grammar is matched against tokens
|
||||
//! before a node exists, and several statements can share one handler -- an
|
||||
//! `ALTER TABLE` handler may yield a different node per subcommand. A planned
|
||||
//! node, by contrast, is claimed by exactly one statement.
|
||||
|
||||
use std::any::Any;
|
||||
use std::sync::Arc;
|
||||
|
||||
use datafusion::common::{ResolvedTableReference, TableReference};
|
||||
use datafusion::error::Result as DfResult;
|
||||
use datafusion::logical_expr::LogicalPlan;
|
||||
use datafusion::sql::sqlparser::{parser::Parser, tokenizer::Token};
|
||||
|
||||
/// A pluggable handler for custom SQL statements.
|
||||
pub trait CustomSqlHandler: Send + Sync {
|
||||
/// Whether this handler wants to handle these tokens.
|
||||
///
|
||||
/// The tokens have had whitespace removed, so a handler can match on
|
||||
/// leading keywords without accounting for layout.
|
||||
fn matches(&self, tokens: &[&Token]) -> bool;
|
||||
|
||||
/// Parse the statement into a logical plan.
|
||||
///
|
||||
/// Returning `Ok(None)` declines a form this handler matched on; the
|
||||
/// statement then goes to DataFusion's own planner. See
|
||||
/// [`StatementRegistry`] for why that stops routing rather than falling
|
||||
/// through to the next handler.
|
||||
fn parse(&self, parser: &mut Parser) -> DfResult<Option<LogicalPlan>>;
|
||||
}
|
||||
|
||||
/// What kind of relation a requirement is about.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum RelationKind {
|
||||
Table,
|
||||
View,
|
||||
}
|
||||
|
||||
/// What kind of object a DDL statement brings into existence.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum CreateKind {
|
||||
Table,
|
||||
View,
|
||||
MaterializedView,
|
||||
}
|
||||
|
||||
/// What a statement needs authorized before it runs.
|
||||
///
|
||||
/// The vocabulary is deliberately generic: it names *what is being reached
|
||||
/// for*, not the privilege that grants it. An embedder maps these onto its own
|
||||
/// privilege model and audit labels, so no access-control concept has to live
|
||||
/// in the dialect.
|
||||
///
|
||||
/// The variants are finer-grained than a bare read/write split because the
|
||||
/// distinctions are load-bearing for that mapping -- appending to a table and
|
||||
/// redefining it are different grants, and collapsing them would silently
|
||||
/// widen what a statement is allowed to do.
|
||||
#[derive(Clone, Debug, PartialEq, Eq)]
|
||||
pub enum AccessRequirement {
|
||||
/// Read the contents of a relation.
|
||||
Read {
|
||||
relation: ResolvedTableReference,
|
||||
kind: RelationKind,
|
||||
},
|
||||
/// Change the rows of a relation.
|
||||
Write {
|
||||
relation: ResolvedTableReference,
|
||||
kind: RelationKind,
|
||||
mode: WriteMode,
|
||||
},
|
||||
/// Change a relation's definition, or anything about it other than its
|
||||
/// rows. Index and column changes land here.
|
||||
Own {
|
||||
relation: ResolvedTableReference,
|
||||
kind: RelationKind,
|
||||
},
|
||||
/// Bring a new relation into existence.
|
||||
CreateIn {
|
||||
relation: ResolvedTableReference,
|
||||
kind: CreateKind,
|
||||
},
|
||||
/// Reach the connected database itself rather than a relation in it.
|
||||
Database { name: String, scope: DatabaseScope },
|
||||
/// Reach a namespace's metadata.
|
||||
Namespace { database: String, namespace: String },
|
||||
/// Reach the deployment rather than any one database.
|
||||
System { scope: SystemScope },
|
||||
}
|
||||
|
||||
/// How an [`AccessRequirement::Write`] changes a relation's rows.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum WriteMode {
|
||||
/// Add rows.
|
||||
Append,
|
||||
/// Change existing rows.
|
||||
Modify,
|
||||
/// Take rows away.
|
||||
Remove,
|
||||
}
|
||||
|
||||
/// How far into a database an [`AccessRequirement::Database`] reaches.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum DatabaseScope {
|
||||
/// See that the database exists and list what is in it.
|
||||
Usage,
|
||||
/// Change what the database contains.
|
||||
Ownership,
|
||||
}
|
||||
|
||||
/// How far into the deployment an [`AccessRequirement::System`] reaches.
|
||||
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
|
||||
pub enum SystemScope {
|
||||
/// Observe deployment-wide state.
|
||||
Usage,
|
||||
/// Act on deployment-wide state.
|
||||
Operate,
|
||||
}
|
||||
|
||||
impl AccessRequirement {
|
||||
/// The relation this requirement is about, for the relation-shaped
|
||||
/// variants.
|
||||
///
|
||||
/// An embedder's privilege mapping is written against the variants
|
||||
/// directly; this is the shortcut for the common case of needing the
|
||||
/// relation without caring which shape asked for it.
|
||||
pub fn relation(&self) -> Option<(&ResolvedTableReference, RelationKind)> {
|
||||
match self {
|
||||
Self::Read { relation, kind }
|
||||
| Self::Write { relation, kind, .. }
|
||||
| Self::Own { relation, kind } => Some((relation, *kind)),
|
||||
Self::CreateIn { .. }
|
||||
| Self::Database { .. }
|
||||
| Self::Namespace { .. }
|
||||
| Self::System { .. } => None,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Everything a statement needs in order to state its requirements, without
|
||||
/// reaching for the engine running it.
|
||||
pub struct RequirementContext<'a> {
|
||||
/// The database a bare relation name resolves against.
|
||||
pub default_database: &'a str,
|
||||
/// The schema a bare relation name resolves against.
|
||||
pub default_schema: &'a str,
|
||||
}
|
||||
|
||||
impl RequirementContext<'_> {
|
||||
/// Resolve a possibly-bare reference against the request's defaults.
|
||||
pub fn resolve(&self, relation: TableReference) -> ResolvedTableReference {
|
||||
relation.resolve(self.default_database, self.default_schema)
|
||||
}
|
||||
|
||||
/// Resolve a bare relation name against the request's defaults.
|
||||
pub fn resolve_bare(&self, name: impl Into<String>) -> ResolvedTableReference {
|
||||
self.resolve(TableReference::bare(name.into()))
|
||||
}
|
||||
}
|
||||
|
||||
/// One statement in the dialect: the node it plans to, what it is called in an
|
||||
/// audit log, and what it needs authorized.
|
||||
pub trait SqlStatement: Send + Sync {
|
||||
/// Whether this statement owns the given planned node.
|
||||
fn claims(&self, node: &dyn Any) -> bool;
|
||||
|
||||
/// The audit label for this statement.
|
||||
///
|
||||
/// This is an open string rather than an enum so that an embedder can add
|
||||
/// a statement -- and a label for it -- without changing this crate.
|
||||
fn audit_operation(&self) -> &'static str;
|
||||
|
||||
/// What must be authorized before the node runs.
|
||||
///
|
||||
/// Returning an empty set means the statement needs nothing beyond
|
||||
/// whatever the engine already collects from the plan's scans.
|
||||
fn access_requirements(
|
||||
&self,
|
||||
node: &dyn Any,
|
||||
context: &RequirementContext<'_>,
|
||||
) -> DfResult<Vec<AccessRequirement>>;
|
||||
}
|
||||
|
||||
/// The set of statements and grammars an engine knows about.
|
||||
///
|
||||
/// Ordering is load-bearing on the parse side and stays explicit. A handler
|
||||
/// may be a catch-all over its leading keyword -- erroring on any form of that
|
||||
/// keyword it does not recognize, or matching on the first token alone -- so a
|
||||
/// handler registered *after* such a one can never be reached for that
|
||||
/// keyword. Extensions are therefore consulted before whatever is already
|
||||
/// registered.
|
||||
///
|
||||
/// A handler that matches and then returns `Ok(None)` stops routing entirely
|
||||
/// rather than falling through to the next handler; the statement then goes to
|
||||
/// DataFusion's own planner. That veto is intentional -- it is how a handler
|
||||
/// declines a form it matched on -- but it means an overlapping handler
|
||||
/// registered later is shadowed rather than reported, which is the other
|
||||
/// reason ordering is explicit here.
|
||||
#[derive(Default)]
|
||||
pub struct StatementRegistry {
|
||||
parsers: Vec<Arc<dyn CustomSqlHandler>>,
|
||||
statements: Vec<Arc<dyn SqlStatement>>,
|
||||
}
|
||||
|
||||
impl StatementRegistry {
|
||||
/// An empty registry.
|
||||
pub fn new() -> Self {
|
||||
Self::default()
|
||||
}
|
||||
|
||||
/// Build a registry from an explicit, already-ordered set.
|
||||
///
|
||||
/// The ordering is used as given -- unlike [`Self::register_parser`], this
|
||||
/// does not reverse anything. It is how an embedder that owns the whole
|
||||
/// dialect states the order once.
|
||||
pub fn from_parts(
|
||||
parsers: Vec<Arc<dyn CustomSqlHandler>>,
|
||||
statements: Vec<Arc<dyn SqlStatement>>,
|
||||
) -> Self {
|
||||
Self {
|
||||
parsers,
|
||||
statements,
|
||||
}
|
||||
}
|
||||
|
||||
/// Add a grammar, consulted before every grammar already registered.
|
||||
///
|
||||
/// Registration is front-insertion because an existing catch-all handler
|
||||
/// would otherwise shadow anything added later; see the type docs.
|
||||
pub fn register_parser(&mut self, parser: Arc<dyn CustomSqlHandler>) -> &mut Self {
|
||||
self.parsers.insert(0, parser);
|
||||
self
|
||||
}
|
||||
|
||||
/// Add a statement, consulted before every statement already registered.
|
||||
pub fn register_statement(&mut self, statement: Arc<dyn SqlStatement>) -> &mut Self {
|
||||
self.statements.insert(0, statement);
|
||||
self
|
||||
}
|
||||
|
||||
/// The grammars, in the order they are consulted.
|
||||
pub fn parsers(&self) -> &[Arc<dyn CustomSqlHandler>] {
|
||||
&self.parsers
|
||||
}
|
||||
|
||||
/// The statement owning this planned node, if any.
|
||||
pub fn claim(&self, node: &dyn Any) -> Option<&Arc<dyn SqlStatement>> {
|
||||
self.statements.iter().find(|s| s.claims(node))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
struct Claimant {
|
||||
label: &'static str,
|
||||
claims_everything: bool,
|
||||
}
|
||||
|
||||
impl SqlStatement for Claimant {
|
||||
fn claims(&self, node: &dyn Any) -> bool {
|
||||
self.claims_everything && node.is::<u8>()
|
||||
}
|
||||
|
||||
fn audit_operation(&self) -> &'static str {
|
||||
self.label
|
||||
}
|
||||
|
||||
fn access_requirements(
|
||||
&self,
|
||||
_node: &dyn Any,
|
||||
_context: &RequirementContext<'_>,
|
||||
) -> DfResult<Vec<AccessRequirement>> {
|
||||
Ok(vec![])
|
||||
}
|
||||
}
|
||||
|
||||
fn claimant(label: &'static str, claims_everything: bool) -> Arc<dyn SqlStatement> {
|
||||
Arc::new(Claimant {
|
||||
label,
|
||||
claims_everything,
|
||||
})
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn an_unclaimed_node_has_no_statement() {
|
||||
let mut registry = StatementRegistry::new();
|
||||
registry.register_statement(claimant("never", false));
|
||||
assert!(registry.claim(&0u8).is_none());
|
||||
}
|
||||
|
||||
/// Front-insertion on the claim side too: an extension must be able to
|
||||
/// take over a node shape that something already registered also claims.
|
||||
#[test]
|
||||
fn the_last_registered_statement_claims_first() {
|
||||
let mut registry = StatementRegistry::new();
|
||||
registry
|
||||
.register_statement(claimant("first", true))
|
||||
.register_statement(claimant("second", true));
|
||||
|
||||
assert_eq!(registry.claim(&0u8).unwrap().audit_operation(), "second");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn from_parts_keeps_the_claim_order_it_was_given() {
|
||||
let registry =
|
||||
StatementRegistry::from_parts(vec![], vec![claimant("a", true), claimant("b", true)]);
|
||||
assert_eq!(registry.claim(&0u8).unwrap().audit_operation(), "a");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_bare_name_resolves_against_the_request_defaults() {
|
||||
let context = RequirementContext {
|
||||
default_database: "db",
|
||||
default_schema: "public",
|
||||
};
|
||||
let resolved = context.resolve_bare("t");
|
||||
assert_eq!(&*resolved.catalog, "db");
|
||||
assert_eq!(&*resolved.schema, "public");
|
||||
assert_eq!(&*resolved.table, "t");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_qualified_name_keeps_its_own_parts() {
|
||||
let context = RequirementContext {
|
||||
default_database: "db",
|
||||
default_schema: "public",
|
||||
};
|
||||
let resolved = context.resolve(TableReference::partial("other", "t"));
|
||||
assert_eq!(&*resolved.catalog, "db");
|
||||
assert_eq!(&*resolved.schema, "other");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn only_the_relation_shaped_requirements_name_a_relation() {
|
||||
let relation = TableReference::bare("t").resolve("db", "public");
|
||||
|
||||
let read = AccessRequirement::Read {
|
||||
relation: relation.clone(),
|
||||
kind: RelationKind::Table,
|
||||
};
|
||||
assert_eq!(read.relation().unwrap().1, RelationKind::Table);
|
||||
|
||||
let create = AccessRequirement::CreateIn {
|
||||
relation,
|
||||
kind: CreateKind::MaterializedView,
|
||||
};
|
||||
assert!(create.relation().is_none());
|
||||
|
||||
let system = AccessRequirement::System {
|
||||
scope: SystemScope::Operate,
|
||||
};
|
||||
assert!(system.relation().is_none());
|
||||
}
|
||||
}
|
||||
@@ -240,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),
|
||||
@@ -1326,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())
|
||||
}
|
||||
|
||||
@@ -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>),
|
||||
}
|
||||
|
||||
|
||||
@@ -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);
|
||||
|
||||
Reference in New Issue
Block a user