Compare commits

..

1 Commits

Author SHA1 Message Date
Lance Release 03e78a9d17 Bump version: 0.32.0-beta.2 → 0.32.0-beta.3 2026-07-24 22:03:48 +00:00
99 changed files with 812 additions and 4012 deletions
+1 -8
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.37.1-beta.0"
current_version = "0.32.0-beta.3"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
@@ -75,13 +75,6 @@ filename = "nodejs/Cargo.toml"
replace = "\nversion = \"{new_version}\""
search = "\nversion = \"{current_version}\""
# The Python package takes its version from here (pyproject.toml declares
# `dynamic = ["version"]`, so maturin reads it out of the crate manifest).
[[tool.bumpversion.files]]
filename = "python/Cargo.toml"
replace = "\nversion = \"{new_version}\""
search = "\nversion = \"{current_version}\""
# Java documentation
[[tool.bumpversion.files]]
filename = "docs/src/java/java.md"
@@ -27,31 +27,19 @@ runs:
# Extract failed job names
FAILED_JOBS=$(echo "$JOB_RESULTS" | jq -r 'to_entries | map(select(.value.result == "failure")) | map(.key) | join(", ")')
TITLE="$WORKFLOW_NAME Failed ($FAILED_JOBS)"
# This action now also runs on nightly schedules, so a breakage that
# persists for a few days would otherwise file one issue per night.
# Comment on the open report instead when one already exists.
EXISTING=$(gh issue list --state open --label ci --limit 100 --json number,title \
| jq -r --arg title "$TITLE" 'map(select(.title == $title)) | .[0].number // empty')
if [ -n "$EXISTING" ]; then
gh issue comment "$EXISTING" --body "Failed again: $RUN_URL"
echo "Commented on existing issue #$EXISTING"
else
gh issue create \
--title "$TITLE" \
--body "The workflow **$WORKFLOW_NAME** failed during execution.
# Create issue with workflow name, failed jobs, and run URL
gh issue create \
--title "$WORKFLOW_NAME Failed ($FAILED_JOBS)" \
--body "The workflow **$WORKFLOW_NAME** failed during execution.
**Failed jobs:** $FAILED_JOBS
**Run URL:** $RUN_URL
Please investigate the failed jobs and address any issues." \
--label "ci"
--label "ci"
echo "Issue created successfully"
fi
echo "Issue created successfully"
else
echo "No job failures detected, skipping issue creation"
fi
+1
View File
@@ -6,6 +6,7 @@ on:
# We don't publish pre-releases for Rust. Crates.io is just a source
# distribution, so we don't need to publish pre-releases.
- "v*-beta*"
- "*-v*" # for example, python-vX.Y.Z
env:
# This env var is used by Swatinem/rust-cache@v2 for the cache
-85
View File
@@ -1,85 +0,0 @@
name: GitHub Release
# All SDKs share one version, so a single `vX.Y.Z` tag produces a single GitHub
# release covering all of them. The per-package publish workflows (PyPI, NPM,
# Cargo, Maven) trigger off the same tag independently.
on:
push:
tags:
- "v*"
permissions:
contents: read
jobs:
gh-release:
runs-on: ubuntu-latest
permissions:
contents: write
steps:
- uses: actions/checkout@v6
with:
fetch-depth: 0
lfs: true
- name: Extract version
id: extract_version
env:
GITHUB_REF: ${{ github.ref }}
run: |
set -e
echo "Extracting tag and version from $GITHUB_REF"
if [[ $GITHUB_REF =~ refs/tags/v(.*) ]]; then
VERSION=${BASH_REMATCH[1]}
TAG=v$VERSION
echo "tag=$TAG" >> $GITHUB_OUTPUT
echo "version=$VERSION" >> $GITHUB_OUTPUT
else
echo "Failed to extract version from $GITHUB_REF"
exit 1
fi
echo "Extracted version $VERSION from $GITHUB_REF"
if [[ $VERSION =~ beta ]]; then
echo "This is a beta release"
echo "prerelease=true" >> $GITHUB_OUTPUT
# Get last release (that is not this one)
FROM_TAG=$(git tag --sort='version:refname' \
| grep ^v \
| grep -vF "$TAG" \
| python ci/semver_sort.py v \
| tail -n 1)
else
echo "This is a stable release"
echo "prerelease=false" >> $GITHUB_OUTPUT
# Get last stable tag (ignore betas)
FROM_TAG=$(git tag --sort='version:refname' \
| grep ^v \
| grep -vF "$TAG" \
| grep -v beta \
| python ci/semver_sort.py v \
| tail -n 1)
fi
echo "Found from tag $FROM_TAG"
echo "from_tag=$FROM_TAG" >> $GITHUB_OUTPUT
- name: Create Release Notes
id: release_notes
uses: mikepenz/release-changelog-builder-action@v4
with:
configuration: .github/release_notes.json
toTag: ${{ steps.extract_version.outputs.tag }}
fromTag: ${{ steps.extract_version.outputs.from_tag }}
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
- name: Create GH release
uses: softprops/action-gh-release@v2
with:
# Marking betas as pre-releases keeps them from taking the "Latest"
# badge on the releases page.
prerelease: ${{ steps.extract_version.outputs.prerelease }}
make_latest: ${{ steps.extract_version.outputs.prerelease == 'false' }}
tag_name: ${{ steps.extract_version.outputs.tag }}
token: ${{ secrets.GITHUB_TOKEN }}
generate_release_notes: false
name: LanceDB v${{ steps.extract_version.outputs.version }}
body: ${{ steps.release_notes.outputs.changelog }}
+29 -7
View File
@@ -1,14 +1,13 @@
name: Create release commit
# This workflow increments the version, tags it, and pushes it. All SDKs share
# a single version, so one tag releases all of them.
# This workflow increments versions, tags the version, and pushes it.
# When a tag is pushed, another workflow is triggered that creates a GH release
# and uploads the binaries. This workflow is only for creating the tag.
# This script will enforce that a minor version is incremented if there are any
# breaking changes since the last minor increment. A breaking change in any SDK
# bumps the minor version for all of them. If you wish to bypass this check, you
# can manually increment the version and push the tag.
# breaking changes since the last minor increment. However, it isn't able to
# differentiate between breaking changes in Node versus Python. If you wish to
# bypass this check, you can manually increment the version and push the tag.
on:
workflow_dispatch:
inputs:
@@ -25,6 +24,16 @@ on:
options:
- preview
- stable
python:
description: 'Make a Python release'
required: true
default: true
type: boolean
other:
description: 'Make a Node/Rust/Java release'
required: true
default: true
type: boolean
bump-minor:
description: 'Bump minor version'
required: true
@@ -56,12 +65,25 @@ jobs:
run: |
git config user.name 'Lance Release'
git config user.email 'lance-dev@lancedb.com'
- name: Bump version
- name: Bump Python version
if: ${{ inputs.python }}
working-directory: python
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
run: |
# Need to get the commit before bumping the version, so we can
# determine if there are breaking changes in the next step as well.
echo "COMMIT_BEFORE_BUMP=$(git rev-parse HEAD)" >> $GITHUB_ENV
pip install bump-my-version PyGithub packaging
bash ../ci/bump_version.sh ${{ inputs.type }} ${{ inputs.bump-minor }} python-v
- name: Bump Node/Rust version
if: ${{ inputs.other }}
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
run: |
pip install bump-my-version PyGithub packaging
bash ci/bump_version.sh ${{ inputs.type }} ${{ inputs.bump-minor }}
bash ci/bump_version.sh ${{ inputs.type }} ${{ inputs.bump-minor }} v $COMMIT_BEFORE_BUMP
bash ci/update_lockfiles.sh --amend
- name: Push new version tag
if: ${{ !inputs.dry_run }}
-15
View File
@@ -61,11 +61,6 @@ jobs:
sudo apt update
sudo apt install -y protobuf-compiler libssl-dev
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Format Rust
run: cargo fmt --all -- --check
- name: Lint Rust
@@ -108,11 +103,6 @@ jobs:
cache: 'pnpm'
cache-dependency-path: nodejs/pnpm-lock.yaml
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install dependencies
run: |
sudo apt update
@@ -192,11 +182,6 @@ jobs:
cache-dependency-path: nodejs/pnpm-lock.yaml
- uses: dtolnay/rust-toolchain@stable
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install dependencies
run: |
brew install protobuf
+89 -95
View File
@@ -10,16 +10,10 @@ permissions:
on:
push:
branches:
- main
tags:
- "v*"
# The cross-compiled targets (musl especially) break from toolchain and
# dependency changes that nothing else in CI catches, and discovering that
# mid-release is expensive. A nightly run keeps that signal while dropping
# the full 8-target release matrix from all ~90 pushes to main each month.
# `report-failure` files an issue when a nightly breaks.
schedule:
- cron: "0 8 * * *"
workflow_dispatch:
pull_request:
# This should trigger a dry run (we skip the final publish step)
paths:
@@ -32,6 +26,73 @@ concurrency:
cancel-in-progress: true
jobs:
gh-release:
if: startsWith(github.ref, 'refs/tags/v')
runs-on: ubuntu-latest
permissions:
contents: write
steps:
- uses: actions/checkout@v6
with:
fetch-depth: 0
lfs: true
- name: Extract version
id: extract_version
env:
GITHUB_REF: ${{ github.ref }}
run: |
set -e
echo "Extracting tag and version from $GITHUB_REF"
if [[ $GITHUB_REF =~ refs/tags/v(.*) ]]; then
VERSION=${BASH_REMATCH[1]}
TAG=v$VERSION
echo "tag=$TAG" >> $GITHUB_OUTPUT
echo "version=$VERSION" >> $GITHUB_OUTPUT
else
echo "Failed to extract version from $GITHUB_REF"
exit 1
fi
echo "Extracted version $VERSION from $GITHUB_REF"
if [[ $VERSION =~ beta ]]; then
echo "This is a beta release"
# Get last release (that is not this one)
FROM_TAG=$(git tag --sort='version:refname' \
| grep ^v \
| grep -vF "$TAG" \
| python ci/semver_sort.py v \
| tail -n 1)
else
echo "This is a stable release"
# Get last stable tag (ignore betas)
FROM_TAG=$(git tag --sort='version:refname' \
| grep ^v \
| grep -vF "$TAG" \
| grep -v beta \
| python ci/semver_sort.py v \
| tail -n 1)
fi
echo "Found from tag $FROM_TAG"
echo "from_tag=$FROM_TAG" >> $GITHUB_OUTPUT
- name: Create Release Notes
id: release_notes
uses: mikepenz/release-changelog-builder-action@v4
with:
configuration: .github/release_notes.json
toTag: ${{ steps.extract_version.outputs.tag }}
fromTag: ${{ steps.extract_version.outputs.from_tag }}
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
- name: Create GH release
uses: softprops/action-gh-release@v2
with:
prerelease: ${{ contains('beta', github.ref) }}
tag_name: ${{ steps.extract_version.outputs.tag }}
token: ${{ secrets.GITHUB_TOKEN }}
generate_release_notes: false
name: Node/Rust LanceDB v${{ steps.extract_version.outputs.version }}
body: ${{ steps.release_notes.outputs.changelog }}
build-lancedb:
strategy:
fail-fast: false
@@ -40,18 +101,9 @@ jobs:
- target: aarch64-apple-darwin
host: macos-latest
features: fp16kernels
pre_build: |-
brew install protobuf
# Fat LTO (the workspace default in .cargo/config.toml) is
# single-threaded and is the peak-memory step of the build. On
# this runner it accounted for ~111 of the job's ~113 minutes,
# making it the critical path of the entire publish pipeline.
# ThinLTO parallelizes it across the runner's cores, for a few
# percent of runtime performance.
export CARGO_PROFILE_RELEASE_LTO=thin
export CARGO_PROFILE_RELEASE_CODEGEN_UNITS=16
pre_build: brew install protobuf
- target: x86_64-pc-windows-msvc
host: windows-2025
host: windows-2025-8x-x64
features: ","
pre_build: |-
choco install --no-progress protoc ninja nasm
@@ -59,19 +111,19 @@ jobs:
# There is an issue where choco doesn't add nasm to the path
export PATH="$PATH:/c/Program Files/NASM"
nasm -v
# See the ThinLTO note on aarch64-apple-darwin above. Keeping
# peak memory down is also what lets this run on the standard
# 4-core runner: the 8-core larger runner was only needed to
# stop fat LTO from OOMing rustc-LLVM.
# Fat LTO of the cdylib is single-threaded and the peak-memory
# step of the build, and had started hitting rustc-LLVM OOM on the
# Windows runners. ThinLTO parallelizes it across the runner's
# cores and keeps peak memory well under the limit.
export CARGO_PROFILE_RELEASE_LTO=thin
export CARGO_PROFILE_RELEASE_CODEGEN_UNITS=16
- target: aarch64-pc-windows-msvc
host: windows-2025
host: windows-2025-8x-x64
features: ","
pre_build: |-
choco install --no-progress protoc
rustup target add aarch64-pc-windows-msvc
# See the ThinLTO note on aarch64-apple-darwin above.
# See ThinLTO note on the x86_64-pc-windows-msvc target above.
export CARGO_PROFILE_RELEASE_LTO=thin
export CARGO_PROFILE_RELEASE_CODEGEN_UNITS=16
- target: x86_64-unknown-linux-gnu
@@ -146,49 +198,16 @@ jobs:
with:
toolchain: stable
targets: ${{ matrix.settings.target }}
# These builds were entirely uncached: the old key was static, so
# `actions/cache` (which only writes on a miss) could never refresh it,
# and the multi-GB whole-`target/` copy it tried to store never fit the
# repo's cache budget, so no entry was ever saved. rust-cache prunes
# `target/` to dependency artifacts and keys on Cargo.lock plus the rustc
# version, which both fixes the key and keeps entries a sane size.
#
# This caches dependency *compilation* only. The LTO link of the cdylib
# re-runs regardless, since the local crate changes every time, so the
# win is larger on the non-LTO jobs than here.
- name: Cache cargo (native builds)
uses: Swatinem/rust-cache@v2
if: ${{ !matrix.settings.docker }}
- name: Cache cargo
uses: actions/cache@v5
with:
# The release profile and per-target dirs differ from what the test
# workflows cache, so these need to be separate entries.
key: release-${{ matrix.settings.target }}
# Only the nightly run on main writes, so tag and PR runs restore a
# warm entry without every dependabot PR writing its own (which would
# be unreadable elsewhere anyway, since GitHub scopes caches to the
# creating ref). The nightly cadence also keeps entries inside
# GitHub's 7-day eviction window, which a tag-only trigger would not.
save-if: ${{ github.ref == 'refs/heads/main' }}
# Docker builds can use rust-cache too. `target/` already lives on the
# host because the whole workspace is bind-mounted into the container, and
# rust-cache's prune and save run host-side, so they can manage it -- which
# is what keeps the entry to dependency artifacts rather than a multi-GB
# copy of everything.
#
# Two differences from the native builds. The container's CARGO_HOME is
# bind-mounted from `.cargo-cache` rather than the host's ~/.cargo, so that
# has to be cached explicitly. And the key is derived from the *host* rustc
# version, which is not the compiler that produced these artifacts; that is
# safe because cargo fingerprints the real compiler and rebuilds on a
# mismatch, it just means a base-image toolchain bump costs one cold build
# instead of invalidating the key.
- name: Cache cargo (docker builds)
uses: Swatinem/rust-cache@v2
if: ${{ matrix.settings.docker }}
with:
key: docker-${{ matrix.settings.target }}
cache-directories: .cargo-cache
save-if: ${{ github.ref == 'refs/heads/main' }}
path: |
~/.cargo/registry/index/
~/.cargo/registry/cache/
~/.cargo/git/db/
.cargo-cache
target/
key: nodejs-${{ matrix.settings.target }}-cargo-${{ matrix.settings.host }}
- name: Install dependencies
run: pnpm install --frozen-lockfile
- name: Install Zig
@@ -206,13 +225,9 @@ jobs:
if: ${{ matrix.settings.docker }}
with:
image: ${{ matrix.settings.docker }}
# All three mounts must live under `.cargo-cache`, which is what the
# cache step above saves. Previously the registry mounts pointed at
# `.cargo/...`, a path nothing cached, so the container re-downloaded
# the whole crate registry on every run.
options: "--user 0:0 -v ${{ github.workspace }}/.cargo-cache/git/db:/usr/local/cargo/git/db \
-v ${{ github.workspace }}/.cargo-cache/registry/cache:/usr/local/cargo/registry/cache \
-v ${{ github.workspace }}/.cargo-cache/registry/index:/usr/local/cargo/registry/index \
-v ${{ github.workspace }}/.cargo/registry/cache:/usr/local/cargo/registry/cache \
-v ${{ github.workspace }}/.cargo/registry/index:/usr/local/cargo/registry/index \
-v ${{ github.workspace }}:/build -w /build/nodejs"
run: |
set -e
@@ -224,16 +239,6 @@ jobs:
--js ../lancedb/native.js \
--strip \
--output-dir dist/
# The container runs as root (`--user 0:0`), so everything it wrote to the
# mounted cache dirs is root-owned. rust-cache's post step runs as the
# runner user and has to both read these and delete from them while
# pruning, so hand them back before it runs.
- name: Take ownership of docker build output
if: ${{ matrix.settings.docker }}
run: |
sudo chown -R "$(id -u):$(id -g)" \
"${{ github.workspace }}/.cargo-cache" \
"${{ github.workspace }}/target"
- name: Build
run: |
${{ matrix.settings.pre_build }}
@@ -247,15 +252,6 @@ jobs:
--output-dir dist/
if: ${{ !matrix.settings.docker }}
shell: bash
# The standard Windows runners have ~14 GB free, and a release `target/`
# for this workspace is a large fraction of that. Report the remaining
# headroom so a build that only just fits is visible before a dependency
# bump turns it into a failed release. `always()` so the numbers are
# still there when the build is what ran out of space.
- name: Report disk headroom
if: always()
run: df -h
shell: bash
- name: Upload artifact
uses: actions/upload-artifact@v7
with:
@@ -406,9 +402,7 @@ jobs:
name: Report Workflow Failure
runs-on: ubuntu-latest
needs: [build-lancedb, test-lancedb, publish]
# Nightly runs are the only thing watching the cross-compiled targets now,
# so they have to report failures too or the signal is silently lost.
if: always() && failure() && (startsWith(github.ref, 'refs/tags/v') || github.event_name == 'schedule')
if: always() && failure() && startsWith(github.ref, 'refs/tags/v')
permissions:
contents: read
issues: write
+72 -19
View File
@@ -3,7 +3,7 @@ name: PyPI Publish
on:
push:
tags:
- 'v*'
- 'python-v*'
pull_request:
# This should trigger a dry run (we skip the final publish step)
paths:
@@ -20,12 +20,6 @@ env:
permissions:
contents: read
# Without this, a force-push to a PR leaves the previous run going -- including
# a ~74 minute Windows job and a billed arm64 wheel build.
concurrency:
group: ${{ github.workflow }}-${{ github.event.pull_request.number || github.ref }}
cancel-in-progress: true
jobs:
linux:
name: Python ${{ matrix.config.package_name }} ${{ matrix.config.platform }} manylinux${{ matrix.config.manylinux }}
@@ -78,7 +72,7 @@ jobs:
package-name: ${{ matrix.config.package_name }}
rustflags: ${{ matrix.config.rustflags }}
- uses: actions/upload-artifact@v7
if: startsWith(github.ref, 'refs/tags/v')
if: startsWith(github.ref, 'refs/tags/python-v')
with:
name: wheels-linux-${{ matrix.config.package_name }}-${{ matrix.config.platform }}-${{ matrix.config.manylinux }}
path: target/wheels/*.whl
@@ -107,7 +101,7 @@ jobs:
python-minor-version: 10
args: "--release --strip --target ${{ matrix.config.target }} --features fp16kernels"
- uses: actions/upload-artifact@v7
if: startsWith(github.ref, 'refs/tags/v')
if: startsWith(github.ref, 'refs/tags/python-v')
with:
name: wheels-mac-${{ matrix.config.target }}
path: target/wheels/lancedb-*.whl
@@ -128,26 +122,19 @@ jobs:
uses: actions/setup-python@v6
with:
python-version: "3.13"
# NOTE: caching cargo here would be a no-op. This workflow only runs on
# tags and PRs, and GitHub only lets a run restore caches from its own ref
# or the default branch -- so with no run on main there is nothing that
# can populate an entry the release build would be allowed to read. Fixing
# this needs a main/nightly trigger (which would also catch wheel-build
# breakage before a release); the ~74 minutes here is otherwise dominated
# by the fat-LTO link, which no cache avoids.
- uses: ./.github/workflows/build_windows_wheel
with:
python-minor-version: 10
args: "--release --strip"
- uses: actions/upload-artifact@v7
if: startsWith(github.ref, 'refs/tags/v')
if: startsWith(github.ref, 'refs/tags/python-v')
with:
name: wheels-windows
path: target/wheels/lancedb-*.whl
if-no-files-found: error
publish:
name: Publish wheels
if: startsWith(github.ref, 'refs/tags/v')
if: startsWith(github.ref, 'refs/tags/python-v')
needs: [linux, mac, windows]
runs-on: ubuntu-latest
permissions:
@@ -196,6 +183,72 @@ jobs:
uses: pypa/gh-action-pypi-publish@release/v1
with:
packages-dir: target/wheels/
gh-release:
if: startsWith(github.ref, 'refs/tags/python-v')
runs-on: ubuntu-latest
permissions:
contents: write
steps:
- uses: actions/checkout@v6
with:
fetch-depth: 0
lfs: true
- name: Extract version
id: extract_version
env:
GITHUB_REF: ${{ github.ref }}
run: |
set -e
echo "Extracting tag and version from $GITHUB_REF"
if [[ $GITHUB_REF =~ refs/tags/python-v(.*) ]]; then
VERSION=${BASH_REMATCH[1]}
TAG=python-v$VERSION
echo "tag=$TAG" >> $GITHUB_OUTPUT
echo "version=$VERSION" >> $GITHUB_OUTPUT
else
echo "Failed to extract version from $GITHUB_REF"
exit 1
fi
echo "Extracted version $VERSION from $GITHUB_REF"
if [[ $VERSION =~ beta ]]; then
echo "This is a beta release"
# Get last release (that is not this one)
FROM_TAG=$(git tag --sort='version:refname' \
| grep ^python-v \
| grep -vF "$TAG" \
| python ci/semver_sort.py python-v \
| tail -n 1)
else
echo "This is a stable release"
# Get last stable tag (ignore betas)
FROM_TAG=$(git tag --sort='version:refname' \
| grep ^python-v \
| grep -vF "$TAG" \
| grep -v beta \
| python ci/semver_sort.py python-v \
| tail -n 1)
fi
echo "Found from tag $FROM_TAG"
echo "from_tag=$FROM_TAG" >> $GITHUB_OUTPUT
- name: Create Python Release Notes
id: python_release_notes
uses: mikepenz/release-changelog-builder-action@v4
with:
configuration: .github/release_notes.json
toTag: ${{ steps.extract_version.outputs.tag }}
fromTag: ${{ steps.extract_version.outputs.from_tag }}
env:
GITHUB_TOKEN: ${{ secrets.GITHUB_TOKEN }}
- name: Create Python GH release
uses: softprops/action-gh-release@v2
with:
prerelease: ${{ contains('beta', github.ref) }}
tag_name: ${{ steps.extract_version.outputs.tag }}
token: ${{ secrets.GITHUB_TOKEN }}
generate_release_notes: false
name: Python LanceDB v${{ steps.extract_version.outputs.version }}
body: ${{ steps.python_release_notes.outputs.changelog }}
report-failure:
name: Report Workflow Failure
runs-on: ubuntu-latest
@@ -203,7 +256,7 @@ jobs:
permissions:
contents: read
issues: write
if: always() && failure() && startsWith(github.ref, 'refs/tags/v')
if: always() && failure() && startsWith(github.ref, 'refs/tags/python-v')
steps:
- uses: actions/checkout@v6
- uses: ./.github/actions/create-failure-issue
-33
View File
@@ -108,15 +108,6 @@ jobs:
run: |
sudo apt update
sudo apt install -y protobuf-compiler
# `pip install -e .` builds the extension with maturin, which is most of
# this job's ~33 minutes. It had no Rust cache, so every dependency was
# recompiled from scratch on every run.
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install
run: |
pip install --extra-index-url https://pypi.fury.io/lance-format/ --extra-index-url https://pypi.fury.io/lancedb/ -e .[tests,dev,embeddings]
@@ -177,14 +168,6 @@ jobs:
uses: actions/setup-python@v6
with:
python-version: "3.13"
# maturin runs cargo natively on macOS (docker is Linux-only), so the host
# target dir is cacheable. This job had no Rust cache.
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- uses: ./.github/workflows/build_mac_wheel
with:
args: --profile ci
@@ -214,14 +197,6 @@ jobs:
uses: actions/setup-python@v6
with:
python-version: "3.13"
# maturin runs cargo natively on Windows (docker is Linux-only), so the
# host target dir is cacheable. This job had no Rust cache at all and so
# rebuilt every dependency from scratch on every run.
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. The repo sits at
# GitHub's cache cap, so per-PR saves just evict main's entries.
save-if: ${{ github.ref == 'refs/heads/main' }}
- uses: ./.github/workflows/build_windows_wheel
with:
args: --profile ci
@@ -249,14 +224,6 @@ jobs:
uses: actions/setup-python@v6
with:
python-version: "3.10"
# As with Doctest, `pip install -e .` compiles the extension and this job
# had no Rust cache, which is most of its ~37 minutes.
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install lancedb
run: |
pip install "pydantic<2"
+7 -45
View File
@@ -48,11 +48,6 @@ jobs:
with:
components: rustfmt, clippy
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install dependencies
run: |
sudo apt update
@@ -94,11 +89,6 @@ jobs:
run: rm -f Cargo.lock
- uses: rui314/setup-mold@v1
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install dependencies
run: |
sudo apt update
@@ -128,11 +118,6 @@ jobs:
fetch-depth: 0
lfs: true
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install dependencies
run: |
sudo apt update
@@ -190,11 +175,6 @@ jobs:
- name: CPU features
run: sysctl -a | grep cpu
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install dependencies
run: brew install protobuf
- name: Run tests
@@ -207,19 +187,12 @@ jobs:
cargo test --profile ci --features $ALL_FEATURES --locked
windows:
runs-on: windows-2022
strategy:
fail-fast: false
matrix:
include:
- target: x86_64-pc-windows-msvc
runner: windows-2022
# windows-11-arm is a standard runner, so it is free on public repos.
# Running natively lets the aarch64 tests actually execute -- this
# job used to cross-compile them and then skip the test step, paying
# full codegen and link cost for a compile check.
- target: aarch64-pc-windows-msvc
runner: windows-11-arm
runs-on: ${{ matrix.runner }}
target:
- x86_64-pc-windows-msvc
- aarch64-pc-windows-msvc
defaults:
run:
working-directory: rust/lancedb
@@ -228,11 +201,6 @@ jobs:
- name: Set target
run: rustup target add ${{ matrix.target }}
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Install Protoc v21.12
run: choco install --no-progress protoc
- name: Build
@@ -240,12 +208,11 @@ jobs:
$env:VCPKG_ROOT = $env:VCPKG_INSTALLATION_ROOT
cargo build --profile ci --features aws,remote --tests --locked --target ${{ matrix.target }}
- name: Run tests
# Can only run tests when target matches host
if: ${{ matrix.target == 'x86_64-pc-windows-msvc' }}
run: |
$env:VCPKG_ROOT = $env:VCPKG_INSTALLATION_ROOT
# `--target` has to match the build step above. Without it cargo uses
# target/ci/ rather than target/<triple>/ci/ and rebuilds the entire
# dependency graph a second time.
cargo test --profile ci --features aws,remote --locked --target ${{ matrix.target }}
cargo test --profile ci --features aws,remote --locked
msrv:
# Check the minimum supported Rust version
@@ -271,11 +238,6 @@ jobs:
with:
toolchain: ${{ matrix.msrv }}
- uses: Swatinem/rust-cache@v2
with:
# Restore everywhere, but only save from main. Per-PR saves are
# unreadable outside their own branch anyway, since GitHub scopes
# caches to the creating ref.
save-if: ${{ github.ref == 'refs/heads/main' }}
- name: Downgrade dependencies
# These packages have newer requirements for MSRV
run: |
-29
View File
@@ -92,8 +92,6 @@ Python bindings changes:
* Should use `LOOP.run()` to call the corresponding `AsyncTable` method.
6. Add concrete sync method to `RemoteTable` class in `python/python/lancedb/remote/table.py`.
7. Add unit test in `python/tests/test_table.py`.
8. If you added a new public class or module-level function (not just a method on an
existing class), expose it in the API reference. See "Python API reference" below.
TypeScript bindings changes:
@@ -105,33 +103,6 @@ TypeScript bindings changes:
5. Add test in `nodejs/__test__/table.test.ts`.
6. Run `npm run docs` to generate TypeScript documentation.
## Python API reference
`docs/src/python/python.md` is the entire Python API reference. It is maintained by
hand, and anything not listed there is not rendered at all, so new public classes and
module-level functions have to be added explicitly. How depends on the module:
* `lancedb.index`, `lancedb.embeddings`, `lancedb.remote`, and `lancedb.rerankers` are
rendered by a single directive each, driven by the module's `__all__`. Add the new
name to `__all__` and it appears; forget, and it is silently omitted.
* Everything else (`lancedb`, `lancedb.table`, `lancedb.query`, `lancedb.db`, ...) is
listed symbol by symbol. Add a `::: lancedb.<module>.<Name>` line to the matching
section, and remember that the page separates synchronous and asynchronous APIs.
Deliberately undocumented: concrete implementations reached through an abstract base
(`LanceTable`, `LanceDBConnection`, `RemoteDBConnection`), query base classes already
covered by `inherited_members`, and internal helpers.
Cross-references in docstrings use mkdocstrings syntax, `[text][lancedb.table.Table]`.
Plain relative links such as `[Table](Table)` do not resolve. To check your work:
```shell
pip install -r docs/requirements.txt
cd docs && PYTHONPATH=. mkdocs build
```
The docs site only builds on pushes to `main`, so this is not covered by PR CI.
## Review Guidelines
Please consider the following when reviewing code contributions.
Generated
+75 -88
View File
@@ -217,9 +217,9 @@ checksum = "7c02d123df017efcdfbd739ef81735b36c5ba83ec3c59c80a9d7ecc718f92e50"
[[package]]
name = "arrow"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6cfdd0833e32a9874d2b55089333ad310c0be208aafa277385ce2461dec90be3"
checksum = "378530e55cd479eda3c14eb345310799717e6f76d0c332041e8487022166b471"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -239,9 +239,9 @@ dependencies = [
[[package]]
name = "arrow-arith"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0a41203398f0eaa6f7ec8e62c0da742a21abf282c148fc157f6c35c90e29981a"
checksum = "a0ab212d2c1886e802f51c5212d78ebbcbb0bec980fff9dadc1eb8d45cd0b738"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -253,9 +253,9 @@ dependencies = [
[[package]]
name = "arrow-array"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ae33dad492b7df00a217563a7b0ef2874df68a0deea1b1a3acf628152f7f7a69"
checksum = "cfd33d3e92f207444098c75b42de99d329562be0cf686b307b097cc52b4e999e"
dependencies = [
"ahash",
"arrow-buffer",
@@ -272,9 +272,9 @@ dependencies = [
[[package]]
name = "arrow-buffer"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b9552f96391c005e6ab449fa941420935e7e062489b12b8b1b08879b2163f5b5"
checksum = "0c6cd424c2693bcdbc150d843dc9d4d137dd2de4782ce6df491ad11a3a0416c0"
dependencies = [
"bytes",
"half",
@@ -284,9 +284,9 @@ dependencies = [
[[package]]
name = "arrow-cast"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3a8a327c9649f30d8406995f27642b68df354713cca3baaaf100f076f18d5f34"
checksum = "4c5aefb56a2c02e9e2b30746241058b85f8983f0fcff2ba0c6d09006e1cded7f"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -306,9 +306,9 @@ dependencies = [
[[package]]
name = "arrow-csv"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "af0dd6d90d1955e9f9a014c1e563ee8aeffc21909085d25623e1da44d96eca26"
checksum = "e94e8cf7e517657a52b91ea1263acf38c4ca62a84655d72458a3359b12ab97de"
dependencies = [
"arrow-array",
"arrow-cast",
@@ -321,9 +321,9 @@ dependencies = [
[[package]]
name = "arrow-data"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2b24852db04738907e06c04ea61e42fe7fda962a34513022dc0d0e754fb7976b"
checksum = "3c88210023a2bfee1896af366309a3028fc3bcbd6515fa29a7990ee1baa08ee0"
dependencies = [
"arrow-buffer",
"arrow-schema",
@@ -334,9 +334,9 @@ dependencies = [
[[package]]
name = "arrow-ipc"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "29a908a11fcfb3fb2f6730f4ac15e367bc644e419155e96238f68cf3adde572b"
checksum = "238438f0834483703d88896db6fe5a7138b2230debc31b34c0336c2996e3c64f"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -350,9 +350,9 @@ dependencies = [
[[package]]
name = "arrow-json"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b8a96aed3931c076adee39ec2a40d8219fc7f09e79bcdaca1df16272993e1e14"
checksum = "205ca2119e6d679d5c133c6f30e68f027738d95ed948cf77677ea69c7800036b"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -375,9 +375,9 @@ dependencies = [
[[package]]
name = "arrow-ord"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "63a083ec750f5c043f02946b4baf05fcdbb55f4560a3277055caca5cc99f3eb0"
checksum = "1bffd8fd2579286a5d63bac898159873e5094a79009940bcb42bbfce4f19f1d0"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -388,9 +388,9 @@ dependencies = [
[[package]]
name = "arrow-pyarrow"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3ffb9be5a873590f825aef50df20e0f8dff5fd42a77058e28bbe7bd44bb53dec"
checksum = "d29abdf672a81c1aeb57fd2661457f9918964d49aed0e9f18932535f2a9e49ce"
dependencies = [
"arrow-array",
"arrow-data",
@@ -400,9 +400,9 @@ dependencies = [
[[package]]
name = "arrow-row"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "514ba0ef0d4c5896202dae736251ce415abb43a950bed570fb7981b8716c0e4c"
checksum = "bab5994731204603c73ba69267616c50f80780774c6bb0476f1f830625115e0c"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -413,9 +413,9 @@ dependencies = [
[[package]]
name = "arrow-schema"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "21ca356ad6425cecb6eb7b28e4f659f1ee7880fbb1a16127de7dd62901efee9e"
checksum = "f633dbfdf39c039ada1bf9e34c694816eb71fbb7dc78f613993b7245e078a1ed"
dependencies = [
"bitflags 2.11.1",
"serde_core",
@@ -424,9 +424,9 @@ dependencies = [
[[package]]
name = "arrow-select"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c58da39eb3d8350ad4a549e5c2bc49284dac554016c69829310350f1731b0aad"
checksum = "8cd065c54172ac787cf3f2f8d4107e0d3fdc26edba76fdf4f4cc170258942222"
dependencies = [
"ahash",
"arrow-array",
@@ -438,9 +438,9 @@ dependencies = [
[[package]]
name = "arrow-string"
version = "58.4.0"
version = "58.3.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b6789b388467525e3271326b6b4915666ecfdf5142aef09779445c954b67543c"
checksum = "29dd7cda3ab9692f43a2e4acc444d760cc17b12bb6d8232ddf64e9bab7c06b42"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -3421,8 +3421,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
[[package]]
name = "fsst"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"rand 0.9.5",
@@ -4777,8 +4777,8 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
[[package]]
name = "lance"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arc-swap",
"arrow",
@@ -4852,8 +4852,8 @@ dependencies = [
[[package]]
name = "lance-arrow"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4875,7 +4875,7 @@ dependencies = [
[[package]]
name = "lance-arrow-scalar"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4889,7 +4889,7 @@ dependencies = [
[[package]]
name = "lance-arrow-stats"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -4898,8 +4898,8 @@ dependencies = [
[[package]]
name = "lance-bitpacking"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrayref",
"crunchy",
@@ -4909,8 +4909,8 @@ dependencies = [
[[package]]
name = "lance-core"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4933,7 +4933,6 @@ dependencies = [
"object_store",
"pin-project",
"prost",
"quick_cache",
"rand 0.9.5",
"roaring",
"serde_json",
@@ -4949,8 +4948,8 @@ dependencies = [
[[package]]
name = "lance-datafusion"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow",
"arrow-array",
@@ -4980,8 +4979,8 @@ dependencies = [
[[package]]
name = "lance-datagen"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow",
"arrow-array",
@@ -4998,8 +4997,8 @@ dependencies = [
[[package]]
name = "lance-derive"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"proc-macro2",
"quote",
@@ -5008,8 +5007,8 @@ dependencies = [
[[package]]
name = "lance-encoding"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5044,8 +5043,8 @@ dependencies = [
[[package]]
name = "lance-file"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5075,8 +5074,8 @@ dependencies = [
[[package]]
name = "lance-index"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arc-swap",
"arrow",
@@ -5143,8 +5142,8 @@ dependencies = [
[[package]]
name = "lance-index-core"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5166,8 +5165,8 @@ dependencies = [
[[package]]
name = "lance-io"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow",
"arrow-arith",
@@ -5210,8 +5209,8 @@ dependencies = [
[[package]]
name = "lance-linalg"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5227,8 +5226,8 @@ dependencies = [
[[package]]
name = "lance-namespace"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow",
"async-trait",
@@ -5240,8 +5239,8 @@ dependencies = [
[[package]]
name = "lance-namespace-impls"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow",
"arrow-ipc",
@@ -5295,8 +5294,8 @@ dependencies = [
[[package]]
name = "lance-select"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5311,8 +5310,8 @@ dependencies = [
[[package]]
name = "lance-table"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow",
"arrow-array",
@@ -5351,8 +5350,8 @@ dependencies = [
[[package]]
name = "lance-testing"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5365,8 +5364,8 @@ dependencies = [
[[package]]
name = "lance-tokenizer"
version = "10.0.0-beta.5"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.5#ddb8e28ca238f29628b8e1795ddccbb7bf75e5c8"
version = "10.0.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v10.0.0-beta.3#ed078e7c19c9560906e1eae54b0baf01dc509986"
dependencies = [
"icu_segmenter",
"jieba-rs",
@@ -5379,7 +5378,7 @@ dependencies = [
[[package]]
name = "lancedb"
version = "0.37.1-beta.0"
version = "0.32.0-beta.2"
dependencies = [
"ahash",
"anyhow",
@@ -5467,7 +5466,7 @@ dependencies = [
[[package]]
name = "lancedb-nodejs"
version = "0.37.1-beta.0"
version = "0.32.0-beta.2"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5492,7 +5491,7 @@ dependencies = [
[[package]]
name = "lancedb-python"
version = "0.37.1-beta.0"
version = "0.35.0-beta.2"
dependencies = [
"arrow",
"async-trait",
@@ -7803,18 +7802,6 @@ dependencies = [
"memchr",
]
[[package]]
name = "quick_cache"
version = "0.6.24"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b9c6658afe513a3b484e3abfdaa0d03ef3c0bbf017542c178dd55f94eb3051f9"
dependencies = [
"ahash",
"equivalent",
"hashbrown 0.16.1",
"parking_lot",
]
[[package]]
name = "quinn"
version = "0.11.9"
+14 -14
View File
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
lance = { "version" = "=10.0.0-beta.5", default-features = false, "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=10.0.0-beta.5", default-features = false, "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=10.0.0-beta.5", default-features = false, "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=10.0.0-beta.5", "tag" = "v10.0.0-beta.5", "git" = "https://github.com/lance-format/lance.git" }
lance = { "version" = "=10.0.0-beta.3", default-features = false, "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=10.0.0-beta.3", default-features = false, "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=10.0.0-beta.3", default-features = false, "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=10.0.0-beta.3", "tag" = "v10.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
ahash = "0.8"
# Note that this one does not include pyarrow
arrow = { version = "58.0.0", optional = false }
+3 -3
View File
@@ -2,9 +2,9 @@ set -e
RELEASE_TYPE=${1:-"stable"}
BUMP_MINOR=${2:-false}
HEAD_SHA=$(git rev-parse HEAD)
TAG_PREFIX=${3:-"v"} # Such as "python-v"
HEAD_SHA=${4:-$(git rev-parse HEAD)}
readonly TAG_PREFIX="v"
readonly SELF_DIR=$(cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )
PREV_TAG=$(git tag --sort='version:refname' | grep ^$TAG_PREFIX | python $SELF_DIR/semver_sort.py $TAG_PREFIX | tail -n 1)
@@ -12,7 +12,7 @@ echo "Found previous tag $PREV_TAG"
# Initially, we don't want to tag if we are doing stable, because we will bump
# again later. See comment at end for why.
if [[ "$RELEASE_TYPE" == 'stable' ]]; then
if [[ "$RELEASE_TYPE" == 'stable' ]]; then
BUMP_ARGS="--no-tag"
fi
-5
View File
@@ -51,11 +51,6 @@ plugins:
paths: [../python/python]
options:
docstring_style: numpy
docstring_options:
# Attributes documented in a `Parameters` section, and pydantic
# dataclasses whose `__init__` griffe cannot see statically, both
# trip this check. It reports nothing actionable here.
warn_unknown_params: false
heading_level: 3
show_signature_annotations: true
show_root_heading: true
+1 -11
View File
@@ -453,16 +453,6 @@ paths:
The metric type to use for the index. l2, Cosine, Dot are supported.
index_type:
type: string
custom_stop_words:
type: [array, "null"]
items:
type: string
description: |
The custom stop-word list for an FTS index. A non-null
array replaces the language's built-in stop-word list and is only
applied when remove_stop_words is enabled. Null uses the built-in
language list, while an empty array explicitly replaces it with no
stop words.
responses:
"200":
description: Index successfully created
@@ -520,4 +510,4 @@ paths:
"401":
$ref: "#/components/responses/unauthorized"
"404":
$ref: "#/components/responses/not_found"
$ref: "#/components/responses/not_found"
+1 -1
View File
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
<dependency>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-core</artifactId>
<version>0.37.1-beta.0</version>
<version>0.32.0-beta.3</version>
</dependency>
```
+1 -1
View File
@@ -1,7 +1,7 @@
# Contributing to LanceDB Typescript
This document outlines the process for contributing to LanceDB Typescript.
For general contribution guidelines, see [CONTRIBUTING.md](https://github.com/lancedb/lancedb/blob/main/CONTRIBUTING.md).
For general contribution guidelines, see [CONTRIBUTING.md](../CONTRIBUTING.md).
## Project layout
+9 -8
View File
@@ -76,23 +76,24 @@ the query optimizer chooses a suboptimal path.
***
### useLsm()
### useLsmWrite()
```ts
useLsm(enable): MergeInsertBuilder
useLsmWrite(useLsmWrite): MergeInsertBuilder
```
Control MemWAL routing for this merge.
Controls whether the merge uses the MemWAL LSM write path.
By default (unset), a `mergeInsert` on a table with an LSM write spec is
routed through Lance's MemWAL shard writer, and a table without one uses the
standard path.
routed through Lance's MemWAL shard writer, and a table without one uses
the standard path. Pass `false` to force the standard path even when a
spec is set. Pass `true` to require a spec — `mergeInsert` rejects if none
is installed.
#### Parameters
* **enable**: `boolean`
`true` forces MemWAL routing and errors if the table has no
LSM write spec. `false` forces the standard write path even when a spec is set.
* **useLsmWrite**: `boolean`
Whether to use the LSM write path.
#### Returns
-36
View File
@@ -497,42 +497,6 @@ ArrowTable.
***
### useLsm()
```ts
useLsm(enable): this
```
Control MemWAL read routing for this query.
By default (unset), when the table carries a MemWAL write spec (see
[Table#setLsmWriteSpec](Table.md#setlsmwritespec)), reads are routed through the LSM scanner so
they also return data written via the `mergeInsert` LSM path that has not yet
been compacted into the base table (the active/frozen in-memory memtables and
the flushed generations), deduplicated by primary key; a table without a spec
reads the base table.
#### Parameters
* **enable**: `boolean`
`true` forces the LSM scanner and errors if the table has no
MemWAL write spec. `false` bypasses the MemWAL and reads the base table only,
even when a spec is present.
Note: the LSM scanner does not support every query shape (e.g. reranking,
hybrid search, `orderBy`). On a MemWAL table those shapes error unless
`useLsm(false)` is set, because a base-only read would silently exclude
un-compacted MemWAL data.
#### Returns
`this`
#### Inherited from
`StandardQueryBase.useLsm`
***
### where()
```ts
-23
View File
@@ -273,29 +273,6 @@ ArrowTable.
***
### useLsm()
```ts
useLsm(enable): this
```
Control MemWAL read routing for this take query.
`false` bypasses the MemWAL and reads the base table only — the escape hatch,
since take-by-row-id/offset is not supported on the LSM scanner and, on a
MemWAL table, auto-routes to it and errors otherwise.
#### Parameters
* **enable**: `boolean`
`false` reads the base table only.
#### Returns
`this`
***
### withRowId()
```ts
-36
View File
@@ -746,42 +746,6 @@ ArrowTable.
***
### useLsm()
```ts
useLsm(enable): this
```
Control MemWAL read routing for this query.
By default (unset), when the table carries a MemWAL write spec (see
[Table#setLsmWriteSpec](Table.md#setlsmwritespec)), reads are routed through the LSM scanner so
they also return data written via the `mergeInsert` LSM path that has not yet
been compacted into the base table (the active/frozen in-memory memtables and
the flushed generations), deduplicated by primary key; a table without a spec
reads the base table.
#### Parameters
* **enable**: `boolean`
`true` forces the LSM scanner and errors if the table has no
MemWAL write spec. `false` bypasses the MemWAL and reads the base table only,
even when a spec is present.
Note: the LSM scanner does not support every query shape (e.g. reranking,
hybrid search, `orderBy`). On a MemWAL table those shapes error unless
`useLsm(false)` is set, because a base-only read would silently exclude
un-compacted MemWAL data.
#### Returns
`this`
#### Inherited from
`StandardQueryBase.useLsm`
***
### where()
```ts
-15
View File
@@ -56,21 +56,6 @@ the experimental FTS V3 format and may introduce breaking changes.
***
### customStopWords?
```ts
optional customStopWords: string[];
```
Custom stop words that replace the built-in list for `language`.
This option only affects tokenization when `removeStopWords` is true.
`undefined` keeps the built-in language list. An empty array explicitly
replaces it with no stop words.
***
### language?
```ts
-15
View File
@@ -30,21 +30,6 @@ The tokenizer to use. The default is "simple".
***
### customStopWords?
```ts
optional customStopWords: string[];
```
Custom stop words that replace the built-in list for `language`.
This option only affects tokenization when `removeStopWords` is true.
`undefined` keeps the built-in language list. An empty array explicitly
replaces it with no stop words.
***
### language?
```ts
+52 -141
View File
@@ -26,18 +26,6 @@ is also an [asynchronous API client](#connections-asynchronous).
::: lancedb.db.DBConnection
::: lancedb.Session
## Namespaces (Synchronous)
A namespace-backed connection resolves tables through a
[Lance namespace](https://lancedb.github.io/lance-namespace/) service instead of
listing a storage directory.
::: lancedb.connect_namespace
::: lancedb.namespace.LanceNamespaceDBConnection
## Tables (Synchronous)
::: lancedb.table.Table
@@ -46,12 +34,8 @@ listing a storage directory.
::: lancedb.table.FragmentSummaryStats
::: lancedb.table.TableStatistics
::: lancedb.table.Tags
::: lancedb.table.Branches
## Expressions
Type-safe expression builder for filters and projections. Use these instead
@@ -78,46 +62,29 @@ of raw SQL strings with [where][lancedb.query.LanceQueryBuilder.where] and
::: lancedb.query.LanceHybridQueryBuilder
::: lancedb.query.LanceEmptyQueryBuilder
::: lancedb.query.LanceTakeQueryBuilder
## Full text queries
Structured full text queries can be passed to
[Table.search][lancedb.table.Table.search] or
[AsyncTable.search][lancedb.table.AsyncTable.search] in place of a query string,
and combined with [BooleanQuery][lancedb.query.BooleanQuery].
::: lancedb.query.FullTextQuery
::: lancedb.query.MatchQuery
::: lancedb.query.PhraseQuery
::: lancedb.query.BoostQuery
::: lancedb.query.MultiMatchQuery
::: lancedb.query.BooleanQuery
::: lancedb.query.FullTextOperator
::: lancedb.query.Occur
## Embeddings
::: lancedb.embeddings
options:
show_root_heading: false
show_root_toc_entry: false
::: lancedb.embeddings.registry.EmbeddingFunctionRegistry
::: lancedb.embeddings.base.EmbeddingFunctionConfig
::: lancedb.embeddings.base.EmbeddingFunction
::: lancedb.embeddings.base.TextEmbeddingFunction
::: lancedb.embeddings.sentence_transformers.SentenceTransformerEmbeddings
::: lancedb.embeddings.openai.OpenAIEmbeddings
::: lancedb.embeddings.open_clip.OpenClipEmbeddings
## Remote configuration
::: lancedb.remote
options:
show_root_heading: false
show_root_toc_entry: false
::: lancedb.remote.ClientConfig
::: lancedb.remote.TimeoutConfig
::: lancedb.remote.RetryConfig
## Context
@@ -127,50 +94,11 @@ and combined with [BooleanQuery][lancedb.query.BooleanQuery].
## Full text search
Pass `custom_stop_words` to [lancedb.index.FTS][]:
Use [lancedb.table.Table.create_fts_index][] for the synchronous API or
[lancedb.table.AsyncTable.create_index][] with [lancedb.index.FTS][] for the
asynchronous API.
```python
from lancedb.index import FTS
table.create_index(
"text",
config=FTS(remove_stop_words=True, custom_stop_words=["acme", "internal"]),
)
```
The list replaces the built-in stop words and is used only when
`remove_stop_words=True`:
- `custom_stop_words=None` uses the built-in list for `language`.
- `custom_stop_words=[]` removes no words.
- Values are passed through without trimming, lowercasing, or other rewriting.
The same option is available on `lancedb.tokenize(...)` and the deprecated
[lancedb.table.Table.create_fts_index][] compatibility helper:
```python
import lancedb
tokens = list(lancedb.tokenize("acme makes searchable data",
custom_stop_words=["acme"]))
```
::: lancedb.tokenize
::: lancedb.FtsToken
## Blobs
Blob columns store large binary values out of line so they can be read lazily
instead of being materialized with the rest of the row.
::: lancedb.blob
::: lancedb.BlobType
::: lancedb._blob.BlobFile
options:
show_root_full_path: false
::: lancedb.index.FTS
## Utilities
@@ -178,14 +106,6 @@ instead of being materialized with the rest of the row.
::: lancedb.merge.LanceMergeInsertBuilder
::: lancedb.otel.instrument_lancedb_metrics
## Exceptions
::: lancedb.exceptions.MissingValueError
::: lancedb.exceptions.MissingColumnError
## Integrations
## Pydantic
@@ -194,30 +114,19 @@ instead of being materialized with the rest of the row.
::: lancedb.pydantic.vector
::: lancedb.pydantic.Vector
::: lancedb.pydantic.MultiVector
::: lancedb.pydantic.LanceModel
## PyTorch
::: lancedb.streaming.StreamingDataset
::: lancedb.permutation.permutation_builder
::: lancedb.permutation.PermutationBuilder
::: lancedb.permutation.Permutation
::: lancedb.permutation.Transforms
## Reranking
::: lancedb.rerankers
options:
show_root_heading: false
show_root_toc_entry: false
::: lancedb.rerankers.linear_combination.LinearCombinationReranker
::: lancedb.rerankers.cohere.CohereReranker
::: lancedb.rerankers.colbert.ColbertReranker
::: lancedb.rerankers.cross_encoder.CrossEncoderReranker
::: lancedb.rerankers.openai.OpenaiReranker
## Connections (Asynchronous)
@@ -228,12 +137,6 @@ can be used to create, list, or open tables.
::: lancedb.db.AsyncConnection
## Namespaces (Asynchronous)
::: lancedb.connect_namespace_async
::: lancedb.namespace.AsyncLanceNamespaceDBConnection
## Tables (Asynchronous)
Table hold your actual data as a collection of records / rows.
@@ -242,20 +145,32 @@ Table hold your actual data as a collection of records / rows.
::: lancedb.table.AsyncTags
::: lancedb.table.AsyncBranches
## Indices (Asynchronous)
Indices can be created on a table to speed up queries. This section
lists the indices that LanceDb supports.
::: lancedb.index
options:
show_root_heading: false
show_root_toc_entry: false
# `lang_mapping` is defined in the module rather than imported, so it is
# picked up despite not being in `__all__`. It is an internal lookup table.
filters: ["!^_", "!^lang_mapping$"]
::: lancedb.index.BTree
::: lancedb.index.Bitmap
::: lancedb.index.LabelList
::: lancedb.index.FTS
::: lancedb.index.IvfPq
::: lancedb.index.HnswPq
::: lancedb.index.HnswSq
::: lancedb.index.IvfFlat
::: lancedb.index.IvfSq
::: lancedb.index.IvfRq
::: lancedb.index.HnswFlat
::: lancedb.table.IndexStatistics
@@ -283,7 +198,3 @@ rows nearest to a query vector and can be created with the
::: lancedb.query.AsyncHybridQuery
options:
inherited_members: true
::: lancedb.query.AsyncTakeQuery
options:
inherited_members: true
+1 -1
View File
@@ -8,7 +8,7 @@
<parent>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.37.1-beta.0</version>
<version>0.32.0-beta.3</version>
<relativePath>../pom.xml</relativePath>
</parent>
+2 -2
View File
@@ -6,7 +6,7 @@
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.37.1-beta.0</version>
<version>0.32.0-beta.3</version>
<packaging>pom</packaging>
<name>${project.artifactId}</name>
<description>LanceDB Java SDK Parent POM</description>
@@ -28,7 +28,7 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<arrow.version>15.0.0</arrow.version>
<lance-core.version>10.0.0-beta.5</lance-core.version>
<lance-core.version>10.0.0-beta.3</lance-core.version>
<spotless.skip>false</spotless.skip>
<spotless.version>2.30.0</spotless.version>
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
+1 -1
View File
@@ -1,7 +1,7 @@
# Contributing to LanceDB Typescript
This document outlines the process for contributing to LanceDB Typescript.
For general contribution guidelines, see [CONTRIBUTING.md](https://github.com/lancedb/lancedb/blob/main/CONTRIBUTING.md).
For general contribution guidelines, see [CONTRIBUTING.md](../CONTRIBUTING.md).
## Project layout
+1 -1
View File
@@ -1,7 +1,7 @@
[package]
name = "lancedb-nodejs"
edition.workspace = true
version = "0.37.1-beta.0"
version = "0.32.0-beta.3"
publish = false
license.workspace = true
description.workspace = true
-31
View File
@@ -991,37 +991,6 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
expectValidMapField(roundTripped.schema.fields[0]);
});
it("preserves string schema metadata", function () {
const metadata = new Map([["source", "fixture"]]);
const schema = new Schema(
[new Field("value", new Int32(), true)],
metadata,
);
expect(makeEmptyTable(schema).schema.metadata.get("source")).toBe(
"fixture",
);
});
it.each([
["non-string keys", new Map<unknown, unknown>([[42, "fixture"]])],
["non-string values", new Map<unknown, unknown>([["source", 42]])],
[
"non-string keys and values",
new Map<unknown, unknown>([[42, false]]),
],
])("rejects schema metadata with %s", function (_, metadataLike) {
const metadata = metadataLike as unknown as Map<string, string>;
const schema = new Schema(
[new Field("value", new Int32(), true)],
metadata,
);
expect(() => makeEmptyTable(schema)).toThrow(
"Expected metadata, if present, to be a Map<string, string> but it had non-string keys or values",
);
});
});
describe("when using two versions of arrow", function () {
+2 -7
View File
@@ -226,7 +226,7 @@ describe("remote connection", () => {
);
});
it("sends FTS options to remote tables", async () => {
it("sends the FTS posting block size to remote tables", async () => {
let createIndexBody: Record<string, unknown> | undefined;
await withMockDatabase(
@@ -264,11 +264,7 @@ describe("remote connection", () => {
async (db) => {
const table = await db.openTable("t");
await table.createIndex("text", {
config: Index.fts({
blockSize: 256,
removeStopWords: true,
customStopWords: ["the"],
}),
config: Index.fts({ blockSize: 256 }),
});
},
);
@@ -276,7 +272,6 @@ describe("remote connection", () => {
expect(createIndexBody?.["column"]).toBe("text");
expect(createIndexBody?.["index_type"]).toBe("FTS");
expect(createIndexBody?.["block_size"]).toBe(256);
expect(createIndexBody?.["custom_stop_words"]).toEqual(["the"]);
});
it("diffs and merges remote branches", async () => {
+2 -51
View File
@@ -527,14 +527,6 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
);
});
it("should expose useLsm on takeRowIds as the base-only escape hatch", async () => {
await table.add([{ id: 1 }, { id: 2 }, { id: 3 }]);
// useLsm(false) is reachable on TakeQuery (the escape hatch for MemWAL tables,
// where take-by-row-id auto-routes to the LSM scanner and is rejected).
const res = await table.takeRowIds([0, 2]).useLsm(false).toArray();
expect(res.map((r) => r.id)).toEqual([1, 3]);
});
it("should throw for negative number in takeRowIds", () => {
expect(() => table.takeRowIds([-1])).toThrow("Row id cannot be negative");
expect(() => table.takeRowIds([0, -5, 2])).toThrow(
@@ -2769,15 +2761,6 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
},
);
test("tokenize supports custom stop words", async () => {
const tokens = await tokenize("the lance data", {
stem: false,
removeStopWords: true,
customStopWords: ["lance"],
});
expect(tokens.map((token) => token.text)).toEqual(["the", "data"]);
});
describe("when calling explainPlan", () => {
let tmpDir: tmp.DirResult;
let table: Table;
@@ -3216,14 +3199,14 @@ describe("LSM merge insert", () => {
await table.closeLsmWriters();
});
it("falls back to the standard path with useLsm(false)", async () => {
it("falls back to the standard path with useLsmWrite(false)", async () => {
const conn = await connect(tmpDir.name);
const table = await bucketTable(conn);
const res = await table
.mergeInsert("id")
.whenNotMatchedInsertAll()
.useLsm(false)
.useLsmWrite(false)
.execute([
{ id: "b", value: 9 },
{ id: "e", value: 5 },
@@ -3257,36 +3240,4 @@ describe("LSM merge insert", () => {
.execute([{ id: "g", value: 7 }]),
).rejects.toThrow();
});
it("auto-routes reads through the MemWAL scanner", async () => {
const conn = await connect(tmpDir.name);
const table = await bucketTable(conn); // base ids "a", "b"
await table
.mergeInsert("id")
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute([{ id: "c", value: 3 }]);
// Default read auto-routes and includes the active memtable row.
const lsm = await table.query().toArray();
expect(lsm.map((r) => r.id).sort()).toEqual(["a", "b", "c"]);
// useLsm(false) bypasses the MemWAL and reads the base table only.
const baseOnly = await table.query().useLsm(false).toArray();
expect(baseOnly.map((r) => r.id).sort()).toEqual(["a", "b"]);
});
it("reads the base table when no LSM spec is installed", async () => {
const conn = await connect(tmpDir.name);
const table = await conn.createEmptyTable(
"plain",
new arrow.Schema([new arrow.Field("id", new arrow.Utf8(), false)]),
);
// No spec: default read and useLsm(false) both succeed against the base table.
await expect(table.query().toArray()).resolves.toBeDefined();
await expect(table.query().useLsm(false).toArray()).resolves.toBeDefined();
// useLsm(true) demands MemWAL routing; without a spec it errors.
await expect(table.query().useLsm(true).toArray()).rejects.toThrow();
});
});
+1 -7
View File
@@ -29,14 +29,8 @@ test("full text search", async () => {
const tbl = await db.createTable("myVectors", data, { mode: "overwrite" });
await tbl.createIndex("doc", {
config: lancedb.Index.fts({
stem: false,
removeStopWords: true,
customStopWords: ["banana"],
}),
config: lancedb.Index.fts(),
});
const tokens = await tbl.tokenize("apple banana", { column: "doc" });
expect(tokens.map((token) => token.text)).toEqual(["apple"]);
// --8<-- [start:full_text_search]
const result = await tbl
-11
View File
@@ -194,16 +194,6 @@ export interface TokenizeOptions {
/** Whether to remove stop words. */
removeStopWords?: boolean;
/**
* Custom stop words that replace the built-in list for `language`.
*
* This option only affects tokenization when `removeStopWords` is true.
*
* `undefined` keeps the built-in language list. An empty array explicitly
* replaces it with no stop words.
*/
customStopWords?: string[];
/** Whether to fold ASCII characters. */
asciiFolding?: boolean;
@@ -235,7 +225,6 @@ export async function tokenize(
options?.lowercase,
options?.stem,
options?.removeStopWords,
options?.customStopWords,
options?.asciiFolding,
options?.ngramMinLength,
options?.ngramMaxLength,
-11
View File
@@ -553,16 +553,6 @@ export interface FtsOptions {
*/
removeStopWords?: boolean;
/**
* Custom stop words that replace the built-in list for `language`.
*
* This option only affects tokenization when `removeStopWords` is true.
*
* `undefined` keeps the built-in language list. An empty array explicitly
* replaces it with no stop words.
*/
customStopWords?: string[];
/**
* whether to remove punctuation
*/
@@ -765,7 +755,6 @@ export class Index {
options?.lowercase,
options?.stem,
options?.removeStopWords,
options?.customStopWords,
options?.asciiFolding,
options?.ngramMinLength,
options?.ngramMaxLength,
+11 -7
View File
@@ -88,17 +88,21 @@ export class MergeInsertBuilder {
);
}
/**
* Control MemWAL routing for this merge.
* Controls whether the merge uses the MemWAL LSM write path.
*
* By default (unset), a `mergeInsert` on a table with an LSM write spec is
* routed through Lance's MemWAL shard writer, and a table without one uses the
* standard path.
* routed through Lance's MemWAL shard writer, and a table without one uses
* the standard path. Pass `false` to force the standard path even when a
* spec is set. Pass `true` to require a spec — `mergeInsert` rejects if none
* is installed.
*
* @param enable - `true` forces MemWAL routing and errors if the table has no
* LSM write spec. `false` forces the standard write path even when a spec is set.
* @param useLsmWrite - Whether to use the LSM write path.
*/
useLsm(enable: boolean): MergeInsertBuilder {
return new MergeInsertBuilder(this.#native.useLsm(enable), this.#schema);
useLsmWrite(useLsmWrite: boolean): MergeInsertBuilder {
return new MergeInsertBuilder(
this.#native.useLsmWrite(useLsmWrite),
this.#schema,
);
}
/**
* Controls how an LSM merge checks that its input targets a single shard.
-38
View File
@@ -460,30 +460,6 @@ export class StandardQueryBase<
this.doCall((inner: NativeQueryType) => inner.fastSearch());
return this;
}
/**
* Control MemWAL read routing for this query.
*
* By default (unset), when the table carries a MemWAL write spec (see
* {@link Table#setLsmWriteSpec}), reads are routed through the LSM scanner so
* they also return data written via the `mergeInsert` LSM path that has not yet
* been compacted into the base table (the active/frozen in-memory memtables and
* the flushed generations), deduplicated by primary key; a table without a spec
* reads the base table.
*
* @param enable - `true` forces the LSM scanner and errors if the table has no
* MemWAL write spec. `false` bypasses the MemWAL and reads the base table only,
* even when a spec is present.
*
* Note: the LSM scanner does not support every query shape (e.g. reranking,
* hybrid search, `orderBy`). On a MemWAL table those shapes error unless
* `useLsm(false)` is set, because a base-only read would silently exclude
* un-compacted MemWAL data.
*/
useLsm(enable: boolean): this {
this.doCall((inner: NativeQueryType) => inner.useLsm(enable));
return this;
}
}
/**
@@ -772,20 +748,6 @@ export class TakeQuery extends QueryBase<NativeTakeQuery> {
constructor(inner: NativeTakeQuery) {
super(inner);
}
/**
* Control MemWAL read routing for this take query.
*
* `false` bypasses the MemWAL and reads the base table only — the escape hatch,
* since take-by-row-id/offset is not supported on the LSM scanner and, on a
* MemWAL table, auto-routes to it and errors otherwise.
*
* @param enable - `false` reads the base table only.
*/
useLsm(enable: boolean): this {
this.doCall((inner: NativeTakeQuery) => inner.useLsm(enable));
return this;
}
}
/** A builder for LanceDB queries.
+1 -1
View File
@@ -84,7 +84,7 @@ export function sanitizeMetadata(
throw Error("Expected metadata, if present, to be a Map<string, string>");
}
for (const item of metadataLike) {
if (typeof item[0] !== "string" || typeof item[1] !== "string") {
if (!(typeof item[0] === "string" || !(typeof item[1] === "string"))) {
throw Error(
"Expected metadata, if present, to be a Map<string, string> but it had non-string keys or values",
);
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-darwin-arm64",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"os": ["darwin"],
"cpu": ["arm64"],
"main": "lancedb.darwin-arm64.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-gnu",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-musl",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-gnu",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-musl",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-arm64-msvc",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"os": [
"win32"
],
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-x64-msvc",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"os": ["win32"],
"cpu": ["x64"],
"main": "lancedb.win32-x64-msvc.node",
+2 -2
View File
@@ -1,12 +1,12 @@
{
"name": "@lancedb/lancedb",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.2",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@lancedb/lancedb",
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.2",
"cpu": [
"x64",
"arm64"
+1 -1
View File
@@ -11,7 +11,7 @@
"ann"
],
"private": false,
"version": "0.37.1-beta.0",
"version": "0.32.0-beta.3",
"main": "dist/index.js",
"exports": {
".": "./dist/index.js",
-4
View File
@@ -43,7 +43,6 @@ pub fn tokenize(
lower_case: Option<bool>,
stem: Option<bool>,
remove_stop_words: Option<bool>,
custom_stop_words: Option<Vec<String>>,
ascii_folding: Option<bool>,
ngram_min_length: Option<u32>,
ngram_max_length: Option<u32>,
@@ -73,7 +72,6 @@ pub fn tokenize(
if let Some(remove_stop_words) = remove_stop_words {
opts = opts.remove_stop_words(remove_stop_words);
}
opts = opts.custom_stop_words(custom_stop_words);
if let Some(ascii_folding) = ascii_folding {
opts = opts.ascii_folding(ascii_folding);
}
@@ -224,7 +222,6 @@ impl Index {
lower_case: Option<bool>,
stem: Option<bool>,
remove_stop_words: Option<bool>,
custom_stop_words: Option<Vec<String>>,
ascii_folding: Option<bool>,
ngram_min_length: Option<u32>,
ngram_max_length: Option<u32>,
@@ -253,7 +250,6 @@ impl Index {
if let Some(remove_stop_words) = remove_stop_words {
opts = opts.remove_stop_words(remove_stop_words);
}
opts = opts.custom_stop_words(custom_stop_words);
if let Some(ascii_folding) = ascii_folding {
opts = opts.ascii_folding(ascii_folding);
}
+2 -2
View File
@@ -51,9 +51,9 @@ impl NativeMergeInsertBuilder {
}
#[napi]
pub fn use_lsm(&self, enable: bool) -> Self {
pub fn use_lsm_write(&self, use_lsm_write: bool) -> Self {
let mut this = self.clone();
this.inner.use_lsm(enable);
this.inner.use_lsm_write(use_lsm_write);
this
}
-15
View File
@@ -168,11 +168,6 @@ impl Query {
self.inner = self.inner.clone().with_row_id();
}
#[napi]
pub fn use_lsm(&mut self, enable: bool) {
self.inner = self.inner.clone().use_lsm(enable);
}
#[napi]
pub fn order_by(&mut self, ordering: Option<Vec<ColumnOrdering>>) -> napi::Result<()> {
let ordering = ordering.map(|ordering| {
@@ -379,11 +374,6 @@ impl VectorQuery {
self.inner = self.inner.clone().with_row_id();
}
#[napi]
pub fn use_lsm(&mut self, enable: bool) {
self.inner = self.inner.clone().use_lsm(enable);
}
#[napi]
pub fn rerank(
&mut self,
@@ -489,11 +479,6 @@ impl TakeQuery {
self.inner = self.inner.clone().with_row_id();
}
#[napi]
pub fn use_lsm(&mut self, enable: bool) {
self.inner = self.inner.clone().use_lsm(enable);
}
#[napi(catch_unwind)]
pub async fn output_schema(&self) -> napi::Result<Buffer> {
let schema = self.inner.output_schema().await.default_error()?;
+49
View File
@@ -0,0 +1,49 @@
[tool.bumpversion]
current_version = "0.35.0-beta.3"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
(?P<patch>0|[1-9]\\d*)
(?:-(?P<pre_l>[a-zA-Z-]+)\\.(?P<pre_n>0|[1-9]\\d*))?
"""
serialize = [
"{major}.{minor}.{patch}-{pre_l}.{pre_n}",
"{major}.{minor}.{patch}",
]
search = "{current_version}"
replace = "{new_version}"
regex = false
ignore_missing_version = false
ignore_missing_files = false
tag = true
sign_tags = false
tag_name = "python-v{new_version}"
tag_message = "Bump version: {current_version} → {new_version}"
allow_dirty = true
commit = true
message = "Bump version: {current_version} → {new_version}"
commit_args = ""
# bump-my-version >=1.4.0 rejects pre_commit_hooks containing shell syntax unless opted in.
allow_shell_hooks = true
# Update Cargo.lock after version bump
pre_commit_hooks = [
"""
cd python && cargo update -p lancedb-python
if git diff --quiet ../Cargo.lock; then
echo "Cargo.lock unchanged"
else
git add ../Cargo.lock
echo "Updated and staged Cargo.lock"
fi
""",
]
[tool.bumpversion.parts.pre_l]
values = ["beta", "final"]
optional_value = "final"
[[tool.bumpversion.files]]
filename = "Cargo.toml"
search = "\nversion = \"{current_version}\""
replace = "\nversion = \"{new_version}\""
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb-python"
version = "0.37.1-beta.0"
version = "0.35.0-beta.3"
publish = false
edition.workspace = true
description = "Python bindings for LanceDB"
+2 -5
View File
@@ -258,7 +258,6 @@ def tokenize(
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -266,10 +265,9 @@ def tokenize(
) -> Iterable[FtsToken]:
"""Tokenize a full-text search query using an explicit tokenizer.
This does not require an FTS index. The tokenizer options match
:class:`lancedb.index.FTS`. ``custom_stop_words`` accepts a list of strings.
This does not require a table or FTS index. The tokenizer options match
:class:`lancedb.index.FTS`.
"""
return _tokenize(
query,
base_tokenizer=base_tokenizer,
@@ -278,7 +276,6 @@ def tokenize(
lower_case=lower_case,
stem=stem,
remove_stop_words=remove_stop_words,
custom_stop_words=custom_stop_words,
ascii_folding=ascii_folding,
ngram_min_length=ngram_min_length,
ngram_max_length=ngram_max_length,
-12
View File
@@ -59,7 +59,6 @@ def tokenize(
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -298,11 +297,6 @@ class Table:
async def fetch_blobs(
self, column: str, row_ids: list[int]
) -> pa.LargeBinaryArray: ...
async def fetch_blob_ranges(
self,
column: str,
requests: List[Tuple[int, int, int]],
) -> pa.LargeBinaryArray: ...
async def fetch_blob_files(
self, column: str, row_ids: list[int]
) -> list[Optional[BlobFile]]: ...
@@ -397,7 +391,6 @@ class Query:
def fast_search(self): ...
def with_row_id(self): ...
def postfilter(self): ...
def use_lsm(self, enable: bool): ...
def nearest_to(self, query_vec: pa.Array) -> VectorQuery: ...
def nearest_to_text(self, query: dict) -> FTSQuery: ...
def order_by(self, ordering: Optional[List[ColumnOrdering]]): ...
@@ -414,7 +407,6 @@ class Query:
class TakeQuery:
def select(self, columns: List[str]): ...
def with_row_id(self): ...
def use_lsm(self, enable: bool): ...
async def output_schema(self) -> pa.Schema: ...
async def execute(self) -> RecordBatchStream: ...
async def explain_plan(self, verbose: Optional[bool]) -> str: ...
@@ -433,7 +425,6 @@ class FTSQuery:
def fast_search(self): ...
def with_row_id(self): ...
def postfilter(self): ...
def use_lsm(self, enable: bool): ...
def get_query(self) -> str: ...
def add_query_vector(self, query_vec: pa.Array) -> None: ...
def nearest_to(self, query_vec: pa.Array) -> HybridQuery: ...
@@ -461,7 +452,6 @@ class VectorQuery:
def column(self, column: str): ...
def distance_type(self, distance_type: str): ...
def postfilter(self): ...
def use_lsm(self, enable: bool): ...
def refine_factor(self, refine_factor: int): ...
def nprobes(self, nprobes: int): ...
def minimum_nprobes(self, minimum_nprobes: int): ...
@@ -485,7 +475,6 @@ class HybridQuery:
def fast_search(self): ...
def with_row_id(self): ...
def postfilter(self): ...
def use_lsm(self, enable: bool): ...
def distance_type(self, distance_type: str): ...
def refine_factor(self, refine_factor: int): ...
def nprobes(self, nprobes: int): ...
@@ -510,7 +499,6 @@ class PyQueryRequest:
select: Optional[Union[str, List[str]]]
fast_search: Optional[bool]
with_row_id: Optional[bool]
use_lsm: Optional[bool]
column: Optional[str]
query_vector: Optional[List[pa.Array]]
minimum_nprobes: Optional[int]
+2 -2
View File
@@ -359,7 +359,7 @@ class DBConnection(EnforceOverrides):
Data is converted to Arrow before being written to disk. For maximum
control over how data is saved, either provide the PyArrow schema to
convert to or else provide a [PyArrow Table][pyarrow.Table] directly.
convert to or else provide a [PyArrow Table](pyarrow.Table) directly.
>>> import pyarrow as pa
>>> custom_schema = pa.schema([
@@ -1529,7 +1529,7 @@ class AsyncConnection(object):
Data is converted to Arrow before being written to disk. For maximum
control over how data is saved, either provide the PyArrow schema to
convert to or else provide a [PyArrow Table][pyarrow.Table] directly.
convert to or else provide a [PyArrow Table](pyarrow.Table) directly.
>>> import pyarrow as pa
>>> custom_schema = pa.schema([
@@ -21,32 +21,3 @@ from .watsonx import WatsonxEmbeddings
from .voyageai import VoyageAIEmbeddingFunction
from .colpali import ColPaliEmbeddings
from .siglip import SigLipEmbeddings
# The API reference renders this package with a single mkdocstrings directive,
# which only picks up names listed here. New embedding functions must be added
# to both the imports above and this list, or they will silently go undocumented.
__all__ = [
"EmbeddingFunction",
"EmbeddingFunctionConfig",
"TextEmbeddingFunction",
"EmbeddingFunctionRegistry",
"get_registry",
"register",
"SentenceTransformerEmbeddings",
"OpenAIEmbeddings",
"OpenClipEmbeddings",
"BedRockText",
"CohereEmbeddingFunction",
"GeminiText",
"GteEmbeddings",
"InstructorEmbeddingFunction",
"JinaEmbeddings",
"OllamaEmbeddings",
"TransformersEmbeddingFunction",
"ColbertEmbeddings",
"VoyageAIEmbeddingFunction",
"WatsonxEmbeddings",
"ColPaliEmbeddings",
"ImageBindEmbeddings",
"SigLipEmbeddings",
]
+5 -5
View File
@@ -21,20 +21,20 @@ class BedRockText(TextEmbeddingFunction):
"""
Parameters
----------
name : str, default "amazon.titan-embed-text-v1"
name: str, default "amazon.titan-embed-text-v1"
The model ID of the bedrock model to use. Supported models for are:
- amazon.titan-embed-text-v1
- cohere.embed-english-v3
- cohere.embed-multilingual-v3
region : str, default "us-east-1"
region: str, default "us-east-1"
Optional name of the AWS Region in which the service should be called.
profile_name : str, default None
profile_name: str, default None
Optional name of the AWS profile to use for calling the Bedrock service.
If not specified, the default profile will be used.
assumed_role : str, default None
assumed_role: str, default None
Optional ARN of an AWS IAM role to assume for calling the Bedrock service.
If not specified, the current active credentials will be used.
role_session_name : str, default "lancedb-embeddings"
role_session_name: str, default "lancedb-embeddings"
Optional name of the AWS IAM role session to use for calling the Bedrock
service. If not specified, "lancedb-embeddings" name will be used.
+3 -5
View File
@@ -22,7 +22,7 @@ class CohereEmbeddingFunction(TextEmbeddingFunction):
Parameters
----------
name : str, default "embed-multilingual-v2.0"
name: str, default "embed-multilingual-v2.0"
The name of the model to use. List of acceptable models:
* embed-english-v3.0
@@ -33,14 +33,12 @@ class CohereEmbeddingFunction(TextEmbeddingFunction):
* embed-english-light-v2.0
* embed-multilingual-v2.0
source_input_type : str, default "search_document"
source_input_type: str, default "search_document"
The input type for the source column in the database
query_input_type : str, default "search_query"
query_input_type: str, default "search_query"
The input type for the query column in the database
Notes
-----
Cohere supports following input types:
| Input Type | Description |
+2 -2
View File
@@ -44,7 +44,7 @@ class ColPaliEmbeddings(EmbeddingFunction):
The token pooling strategy to use, by default "hierarchical".
- "hierarchical": Progressively pools tokens to reduce sequence length.
- "lambda": A simpler pooling that uses a custom `pooling_func`.
pooling_func : typing.Callable, optional
pooling_func: typing.Callable, optional
A function to use for pooling when `pooling_strategy` is "lambda".
pool_factor : int
Factor to reduce sequence length if token pooling is enabled (default 2).
@@ -52,7 +52,7 @@ class ColPaliEmbeddings(EmbeddingFunction):
Quantization configuration for the model. (default None, bitsandbytes needed)
batch_size : int
Batch size for processing inputs (default 2).
offload_folder : str, optional
offload_folder: str, optional
Folder to offload model weights if using CPU offloading (default None). This is
useful for large models that do not fit in memory.
"""
@@ -48,16 +48,16 @@ class GeminiText(TextEmbeddingFunction):
Parameters
----------
name : str, default "gemini-embedding-001"
name: str, default "gemini-embedding-001"
The name of the model to use. Supported models include:
- "gemini-embedding-001" (768 dimensions)
Note: The legacy "models/embedding-001" format is also supported but
"gemini-embedding-001" is recommended.
query_task_type : str, default "retrieval_query"
query_task_type: str, default "retrieval_query"
Sets the task type for the queries.
source_task_type : str, default "retrieval_document"
source_task_type: str, default "retrieval_document"
Sets the task type for ingestion.
Examples
+4 -4
View File
@@ -26,13 +26,13 @@ class GteEmbeddings(TextEmbeddingFunction):
Parameters
----------
name : str, default "thenlper/gte-large"
name: str, default "thenlper/gte-large"
The name of the model to use.
device : str, default "cpu"
device: str, default "cpu"
Sets the device type for the model.
normalize : str, default "True"
normalize: str, default "True"
Controls normalize param in encode function for the transformer.
mlx : bool, default False
mlx: bool, default False
Controls which model to use. False for gte-large,True for the mlx version.
Examples
@@ -35,23 +35,23 @@ class InstructorEmbeddingFunction(TextEmbeddingFunction):
Parameters
----------
name : str
name: str
The name of the model to use. Available models are listed at
https://github.com/xlang-ai/instructor-embedding#model-list;
The default model is hkunlp/instructor-base
batch_size : int, default 32
batch_size: int, default 32
The batch size to use when generating embeddings
device : str, default "cpu"
device: str, default "cpu"
The device to use when generating embeddings
show_progress_bar : bool, default True
show_progress_bar: bool, default True
Whether to show a progress bar when generating embeddings
normalize_embeddings : bool, default True
normalize_embeddings: bool, default True
Whether to normalize the embeddings
quantize : bool, default False
quantize: bool, default False
Whether to quantize the model
source_instruction : str, default "represent the document for retrieval"
source_instruction: str, default "represent the document for retrieval"
The instruction for the source column
query_instruction : str, default "represent the document for retrieving the most
query_instruction: str, default "represent the document for retrieving the most
similar documents"
The instruction for the query
+2 -2
View File
@@ -40,10 +40,10 @@ class JinaEmbeddings(EmbeddingFunction):
Parameters
----------
name : str, default "jina-clip-v1". Note that some models support both image
name: str, default "jina-clip-v1". Note that some models support both image
and text embeddings and some just text embedding
api_key : str, default None
api_key: str, default None
The api key to access Jina API. If you pass None, you can set JINA_API_KEY
environment variable
@@ -21,13 +21,13 @@ class SentenceTransformerEmbeddings(TextEmbeddingFunction):
Parameters
----------
name : str, default "all-MiniLM-L6-v2"
name: str, default "all-MiniLM-L6-v2"
The name of the model to use.
device : str, default "cpu"
device: str, default "cpu"
The device to use for the model
normalize : bool, default True
normalize: bool, default True
Whether to normalize the embeddings
trust_remote_code : bool, default True
trust_remote_code: bool, default True
Whether to trust the remote code
"""
+2 -2
View File
@@ -167,7 +167,7 @@ class VoyageAIEmbeddingFunction(EmbeddingFunction):
Parameters
----------
name : str
name: str
The name of the model to use. List of acceptable models:
* voyage-4 (1024 dims, general-purpose and multilingual retrieval)
@@ -185,7 +185,7 @@ class VoyageAIEmbeddingFunction(EmbeddingFunction):
* voyage-law-2
* voyage-code-2
output_dimension : int, optional
output_dimension: int, optional
The output dimension for models that support flexible dimensions.
Currently only voyage-multimodal-3.5 supports this feature.
Valid options: 256, 512, 1024 (default), 2048.
+24 -33
View File
@@ -2,7 +2,7 @@
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
from dataclasses import dataclass
from typing import List, Literal, Optional
from typing import Literal, Optional
from ._lancedb import (
IndexConfig,
@@ -151,11 +151,6 @@ class FTS:
remove_stop_words : bool, default True
Whether to remove stop words. Stop words are common words that are often
removed from text before indexing. For example, in English "the" and "and".
custom_stop_words : list of str, optional
Custom words replace the built-in language stop words
and only take effect when ``remove_stop_words`` is True. ``None`` uses
the built-in language list, while an empty list explicitly uses no
stop words.
ascii_folding : bool, default True
Whether to fold ASCII characters. This converts accented characters to
their ASCII equivalent. For example, "café" would be converted to "cafe".
@@ -184,7 +179,6 @@ class FTS:
ngram_max_length: int = 3
prefix_only: bool = False
block_size: int = 128
custom_stop_words: Optional[List[str]] = None
@dataclass
@@ -219,7 +213,7 @@ class HnswPq:
distance has a range of (-, ). If the vectors are normalized (i.e. their
l2 norm is 1), then dot distance is equivalent to the cosine distance.
num_partitions: int, default sqrt(num_rows)
num_partitions, default sqrt(num_rows)
The number of IVF partitions to create.
@@ -228,7 +222,7 @@ class HnswPq:
will require too much memory. Each partition becomes its own HNSW graph, so
setting this value higher reduces the peak memory use of training.
num_sub_vectors: int, default is vector dimension / 16
num_sub_vectors, default is vector dimension / 16
Number of sub-vectors of PQ.
@@ -244,13 +238,13 @@ class HnswPq:
If the dimension is not visible by 8 then we use 1 subvector. This is not
ideal and will likely result in poor performance.
num_bits: int, default 8
num_bits: int, default 8
Number of bits to encode each sub-vector.
This value controls how much the sub-vectors are compressed. The more bits
the more accurate the index but the slower search. Only 4 and 8 are supported.
max_iterations: int, default 50
max_iterations, default 50
Max iterations to train kmeans.
@@ -263,7 +257,7 @@ class HnswPq:
those cases it is unlikely that setting this larger will lead to the index
converging anyways.
sample_rate: int, default 256
sample_rate, default 256
The rate used to calculate the number of training vectors for kmeans.
@@ -279,14 +273,14 @@ class HnswPq:
Increasing this value might improve the quality of the index but in
most cases the default should be sufficient.
m: int, default 20
m, default 20
The number of neighbors to select for each vector in the HNSW graph.
This value controls the tradeoff between search speed and accuracy.
The higher the value the more accurate the search but the slower it will be.
ef_construction: int, default 300
ef_construction, default 300
The number of candidates to evaluate during the construction of the HNSW graph.
@@ -297,7 +291,7 @@ class HnswPq:
This value should be set to a value that is not less than `ef` in the
search phase.
target_partition_size: int, default is 1,048,576
target_partition_size, default is 1,048,576
The target size of each partition.
@@ -351,7 +345,7 @@ class HnswSq:
distance has a range of (-, ). If the vectors are normalized (i.e. their
l2 norm is 1), then dot distance is equivalent to the cosine distance.
num_partitions: int, default sqrt(num_rows)
num_partitions, default sqrt(num_rows)
The number of IVF partitions to create.
@@ -360,7 +354,7 @@ class HnswSq:
will require too much memory. Each partition becomes its own HNSW graph, so
setting this value higher reduces the peak memory use of training.
max_iterations: int, default 50
max_iterations, default 50
Max iterations to train kmeans.
@@ -373,7 +367,7 @@ class HnswSq:
In those cases it is unlikely that setting this larger will lead to
the index converging anyways.
sample_rate: int, default 256
sample_rate, default 256
The rate used to calculate the number of training vectors for kmeans.
@@ -389,14 +383,14 @@ class HnswSq:
Increasing this value might improve the quality of the index but in
most cases the default should be sufficient.
m: int, default 20
m, default 20
The number of neighbors to select for each vector in the HNSW graph.
This value controls the tradeoff between search speed and accuracy.
The higher the value the more accurate the search but the slower it will be.
ef_construction: int, default 300
ef_construction, default 300
The number of candidates to evaluate during the construction of the HNSW graph.
@@ -407,7 +401,7 @@ class HnswSq:
This value should be set to a value that is not less than `ef` in the search
phase.
target_partition_size: int, default is 1,048,576
target_partition_size, default is 1,048,576
The target size of each partition.
@@ -460,7 +454,7 @@ class HnswFlat:
distance has a range of (-, ). If the vectors are normalized (i.e. their
l2 norm is 1), then dot distance is equivalent to the cosine distance.
num_partitions: int, default sqrt(num_rows)
num_partitions, default sqrt(num_rows)
The number of IVF partitions to create.
@@ -470,18 +464,18 @@ class HnswFlat:
graph, so setting this value higher reduces the peak memory use of
training.
max_iterations: int, default 50
max_iterations, default 50
Max iterations to train kmeans.
When training an IVF index we use kmeans to calculate the partitions.
This parameter controls how many iterations of kmeans to run.
sample_rate: int, default 256
sample_rate, default 256
The rate used to calculate the number of training vectors for kmeans.
m: int, default 20
m, default 20
The number of neighbors to select for each vector in the HNSW graph.
@@ -489,7 +483,7 @@ class HnswFlat:
The higher the value the more accurate the search but the slower it
will be.
ef_construction: int, default 300
ef_construction, default 300
The number of candidates to evaluate during the construction of the HNSW
graph.
@@ -501,7 +495,7 @@ class HnswFlat:
than 500. This value should be set to a value that is not less than `ef`
in the search phase.
target_partition_size: int, default is 1,048,576
target_partition_size, default is 1,048,576
The target size of each partition.
"""
@@ -605,7 +599,7 @@ class IvfFlat:
The default value is 256.
target_partition_size: int, default is 8192
target_partition_size, default is 8192
The target size of each partition.
@@ -769,7 +763,7 @@ class IvfPq:
The default value is 256.
target_partition_size: int, default is 8192
target_partition_size, default is 8192
The target size of each partition.
@@ -830,7 +824,7 @@ class IvfRq:
sample_rate: int, default 256
Controls the number of training vectors: sample_rate * num_partitions.
target_partition_size: int, default is 8192
target_partition_size, default is 8192
Target size of each partition.
"""
@@ -845,9 +839,6 @@ class IvfRq:
accelerator: Optional[str] = None
# The API reference renders this module with a single mkdocstrings directive,
# which only picks up names listed here. New public names must be added to this
# list, or they will silently go undocumented.
__all__ = [
"BTree",
"IvfPq",
+11 -11
View File
@@ -37,7 +37,7 @@ class LanceMergeInsertBuilder(object):
self._when_not_matched_by_source_condition_expr = None
self._timeout = None
self._use_index = True
self._use_lsm = None
self._use_lsm_write = None
self._validate_single_shard = None
def when_matched_update_all(
@@ -113,22 +113,22 @@ class LanceMergeInsertBuilder(object):
self._use_index = use_index
return self
def use_lsm(self, enable: bool) -> LanceMergeInsertBuilder:
def use_lsm_write(self, use_lsm_write: bool) -> LanceMergeInsertBuilder:
"""
Control MemWAL routing for this merge.
Controls whether the merge uses the MemWAL LSM write path.
By default (unset), a `merge_insert` on a table with an LSM write spec is
routed through Lance's MemWAL shard writer, and a table without one uses
the standard path.
By default (unset), a `merge_insert` on a table with an LSM write spec
is routed through Lance's MemWAL shard writer, and a table without one
uses the standard path. Pass `False` to force the standard path even
when a spec is set. Pass `True` to require a spec `merge_insert`
raises an error if none is installed.
Parameters
----------
enable: bool
``True`` forces MemWAL routing and errors if the table has no LSM
write spec. ``False`` forces the standard write path even when a spec
is set.
use_lsm_write: bool
Whether to use the LSM write path.
"""
self._use_lsm = enable
self._use_lsm_write = use_lsm_write
return self
def validate_single_shard(
+6 -11
View File
@@ -438,8 +438,7 @@ class Permutation:
_reader: Optional[PermutationReader] = None,
):
"""
Internal constructor. Use
[from_tables][lancedb.permutation.Permutation.from_tables] instead.
Internal constructor. Use [from_tables](#from_tables) instead.
"""
assert base_table is not None, "base_table is required"
assert selection is not None, "selection is required"
@@ -986,9 +985,8 @@ class Permutation:
types. Conversion of strings, lists, and structs will require creating python
objects and this is not zero-copy.
For custom formatting, use
[with_transform][lancedb.permutation.Permutation.with_transform] which
overrides this method.
For custom formatting, use [with_transform](#with_transform) which overrides
this method.
"""
assert format is not None, "format is required"
if format == "python":
@@ -1063,8 +1061,7 @@ class Permutation:
Note: this method returns a new permutation and does not modify `self`
It is provided for compatibility with the huggingface Dataset API.
Use [with_skip][lancedb.permutation.Permutation.with_skip] instead to
avoid confusion.
Use [with_skip](#with_skip) instead to avoid confusion.
"""
return self.with_skip(skip)
@@ -1087,8 +1084,7 @@ class Permutation:
Note: this method returns a new permutation and does not modify `self`
It is provided for compatibility with the huggingface Dataset API.
Use [with_take][lancedb.permutation.Permutation.with_take] instead to
avoid confusion.
Use [with_take](#with_take) instead to avoid confusion.
"""
return self.with_take(limit)
@@ -1111,8 +1107,7 @@ class Permutation:
Note: this method returns a new permutation and does not modify `self`
It is provided for compatibility with the huggingface Dataset API.
Use [with_repeat][lancedb.permutation.Permutation.with_repeat] instead
to avoid confusion.
Use [with_repeat](#with_repeat) instead to avoid confusion.
"""
return self.with_repeat(times)
+14 -98
View File
@@ -651,8 +651,7 @@ class Query(pydantic.BaseModel):
distance_type : Optional[str]
the distance type to use for vector search
This can be l2 (default), cosine and dot. See
[metric definitions](https://lancedb.com/docs/search/vector-search/) for
This can be l2 (default), cosine and dot. See [metric definitions][search] for
more details.
If this is not a vector search this will be None.
@@ -665,9 +664,8 @@ class Query(pydantic.BaseModel):
- A higher number makes search more accurate but also slower.
- See discussion in
[Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
- See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
Will be None if this is not a vector search.
refine_factor : Optional[int]
@@ -675,9 +673,8 @@ class Query(pydantic.BaseModel):
- A higher number makes search more accurate but also slower.
- See discussion in
[Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
- See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
Will be None if this is not a vector search.
lower_bound : Optional[float]
@@ -781,11 +778,6 @@ class Query(pydantic.BaseModel):
# if true, will only search the indexed data
fast_search: Optional[bool] = None
# MemWAL LSM read routing: None auto-routes when the table carries a write
# spec, True forces the LSM scanner (errors without a spec), False reads the
# base table only
use_lsm: Optional[bool] = None
# size of the nearest neighbor list maintained during HNSW search
ef: Optional[int] = None
@@ -803,9 +795,6 @@ class Query(pydantic.BaseModel):
query.full_text_query = req.full_text_search
query.columns = req.select
query.with_row_id = req.with_row_id
# use_lsm is a genuine tri-state (None / True / False); preserve it as-is
# so a round-tripped query keeps an explicit False.
query.use_lsm = req.use_lsm
query.vector_column = req.column
query.vector = req.query_vector
query.distance_type = req.distance_type
@@ -978,7 +967,6 @@ class LanceQueryBuilder(ABC):
self._with_row_address = None
self._fragments = None
self._fragment_ids = None
self._use_lsm = None
self._vector = None
self._text = None
self._ef = None
@@ -1338,30 +1326,6 @@ class LanceQueryBuilder(ABC):
self._fragment_ids = fragment_ids
return self
def use_lsm(self, enable: bool) -> Self:
"""Control MemWAL LSM read routing for this query.
By default (unset), a query against a table with an LSM write spec is
routed through the LSM scanner so it also returns data written via the
``merge_insert`` LSM path that has not yet been compacted into the base
table (active/frozen memtables + flushed generations); a table without a
spec reads the base table.
Parameters
----------
enable : bool
``True`` forces the LSM scanner and errors if the table has no LSM
write spec. ``False`` bypasses the MemWAL and reads the base table
only, even when a spec is present.
Returns
-------
LanceQueryBuilder
The LanceQueryBuilder object.
"""
self._use_lsm = enable
return self
def explain_plan(self, verbose: Optional[bool] = False) -> str:
"""Return the execution plan for this query.
@@ -1654,8 +1618,8 @@ class LanceVectorQueryBuilder(LanceQueryBuilder):
Higher values will yield better recall (more likely to find vectors if
they exist) at the expense of latency.
See discussion in [Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
This method sets both the minimum and maximum number of probes to the same
value. See `minimum_nprobes` and `maximum_nprobes` for more fine-grained
@@ -1755,8 +1719,8 @@ class LanceVectorQueryBuilder(LanceQueryBuilder):
As an example, a refine factor of 2 will sample 2x as many vectors as
requested, re-ranks them, and returns the top half most relevant results.
See discussion in [Querying an ANN Index](https://lancedb.com/docs/indexing/)
for tuning advice.
See discussion in [Querying an ANN Index][querying-an-ann-index] for
tuning advice.
Parameters
----------
@@ -1824,7 +1788,6 @@ class LanceVectorQueryBuilder(LanceQueryBuilder):
with_row_address=self._with_row_address,
fragments=self._fragments,
fragment_ids=self._fragment_ids,
use_lsm=self._use_lsm,
offset=self._offset,
fast_search=self._fast_search,
ef=self._ef,
@@ -2049,7 +2012,6 @@ class LanceFtsQueryBuilder(LanceQueryBuilder):
with_row_address=self._with_row_address,
fragments=self._fragments,
fragment_ids=self._fragment_ids,
use_lsm=self._use_lsm,
full_text_query=FullTextSearchQuery(
query=self._query_with_phrase_semantics(), columns=self._fts_columns
),
@@ -2116,7 +2078,6 @@ class LanceEmptyQueryBuilder(LanceQueryBuilder):
with_row_address=self._with_row_address,
fragments=self._fragments,
fragment_ids=self._fragment_ids,
use_lsm=self._use_lsm,
offset=self._offset,
order_by=self._order_by,
)
@@ -2694,9 +2655,6 @@ class LanceHybridQueryBuilder(LanceQueryBuilder):
if self._with_row_id:
self._vector_query.with_row_id(True)
self._fts_query.with_row_id(True)
if self._use_lsm is not None:
self._vector_query.use_lsm(self._use_lsm)
self._fts_query.use_lsm(self._use_lsm)
if self._phrase_query:
self._fts_query.phrase_query(True)
if self._distance_type:
@@ -3044,7 +3002,7 @@ class AsyncQueryBase(object):
if blob_mode == "bytes"
else {}
)
dataset = await self._table.to_lance()
dataset = await self._table._to_lance()
scanner = dataset.scanner(
**_scanner_kwargs_for_query(
query,
@@ -3273,27 +3231,6 @@ class AsyncStandardQuery(AsyncQueryBase):
self._inner.fast_search()
return self
def use_lsm(self, enable: bool) -> Self:
"""
Control MemWAL LSM read routing for this query.
By default (unset), a query against a table with an LSM write spec (see
[AsyncTable.set_lsm_write_spec][lancedb.table.AsyncTable.set_lsm_write_spec])
is routed through the LSM scanner so it also returns data written via the
``merge_insert`` LSM path that has not yet been compacted into the base
table (the active/frozen in-memory memtables and the flushed generations),
deduplicated by primary key; a table without a spec reads the base table.
Parameters
----------
enable : bool
``True`` forces the LSM scanner and errors if the table has no LSM
write spec. ``False`` bypasses the MemWAL and reads the base table
only, even when a spec is present.
"""
self._inner.use_lsm(enable)
return self
def postfilter(self) -> Self:
"""
If this is called then filtering will happen after the search instead of
@@ -3382,9 +3319,8 @@ class AsyncQuery(AsyncStandardQuery):
are various ANN search parameters that will let you fine tune your recall
accuracy vs search latency.
Vector searches always have a
[limit][lancedb.query.AsyncVectorQuery.limit]. If `limit` has not been
called then a default `limit` of 10 will be used.
Vector searches always have a [limit][]. If `limit` has not been called then
a default `limit` of 10 will be used.
Typically, a single vector is passed in as the query. However, you can also
pass in multiple vectors. When multiple vectors are passed in, if the vector
@@ -3515,9 +3451,8 @@ class AsyncFTSQuery(AsyncStandardQuery):
are various ANN search parameters that will let you fine tune your recall
accuracy vs search latency.
Hybrid searches always have a
[limit][lancedb.query.AsyncHybridQuery.limit]. If `limit` has not been
called then a default `limit` of 10 will be used.
Hybrid searches always have a [limit][]. If `limit` has not been called then
a default `limit` of 10 will be used.
Typically, a single vector is passed in as the query. However, you can also
pass in multiple vectors. This can be useful if you want to find the nearest
@@ -4009,15 +3944,6 @@ class AsyncTakeQuery(AsyncQueryBase):
def __init__(self, inner: LanceTakeQuery, table: Optional["AsyncTable"] = None):
super().__init__(inner, table)
def use_lsm(self, enable: bool) -> "AsyncTakeQuery":
"""Control MemWAL LSM read routing for this take query.
``False`` bypasses the MemWAL and reads the base table only the escape
hatch, since take-by-row-id/offset is not supported on the LSM scanner.
"""
self._inner.use_lsm(enable)
return self
async def _plain_scan_to_pandas(
self,
blob_mode: BlobMode,
@@ -4076,16 +4002,6 @@ class BaseQueryBuilder(object):
self._inner.with_row_id()
return self
def use_lsm(self, enable: bool) -> Self:
"""
Control MemWAL LSM read routing for this query.
``False`` bypasses the MemWAL and reads the base table only, the escape
hatch for shapes the LSM scanner cannot honor (e.g. take-by-row-id).
"""
self._inner.use_lsm(enable)
return self
def with_row_address(self, with_row_address: bool = True) -> Self:
"""
Include the _rowaddr column in scanner-backed plain query results.
-3
View File
@@ -11,9 +11,6 @@ from lancedb import __version__
from .header import HeaderProvider
from .oauth import OAuthConfig, OAuthFlowType
# The API reference renders this module with a single mkdocstrings directive,
# which only picks up names listed here. New public names must be added to this
# list, or they will silently go undocumented.
__all__ = [
"TimeoutConfig",
"RetryConfig",
+2 -2
View File
@@ -53,9 +53,9 @@ class RetryError(LanceDBClientError):
"""An error that occurs when the client has exceeded the maximum number of retries.
The retry strategy can be adjusted by setting the
[retry_config][lancedb.remote.ClientConfig.retry_config] in the client
[retry_config](lancedb.remote.ClientConfig.retry_config) in the client
configuration. This is passed in the `client_config` argument of
[connect][lancedb.connect] and [connect_async][lancedb.connect_async].
[connect](lancedb.connect) and [connect_async](lancedb.connect_async).
The __cause__ attribute of this exception will be the last exception that
caused the retry to fail. It will be an
+3 -12
View File
@@ -340,7 +340,6 @@ class RemoteTable(Table):
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -362,7 +361,6 @@ class RemoteTable(Table):
lower_case=lower_case,
stem=stem,
remove_stop_words=remove_stop_words,
custom_stop_words=custom_stop_words,
ascii_folding=ascii_folding,
ngram_min_length=ngram_min_length,
ngram_max_length=ngram_max_length,
@@ -580,9 +578,8 @@ class RemoteTable(Table):
progress: Optional[Union[bool, Callable, Any]] = None,
write_parallelism: Optional[int] = None,
) -> AddResult:
"""Add more data to the [Table][lancedb.table.Table].
It has the same API signature as the OSS version.
"""Add more data to the [Table](Table). It has the same API signature as
the OSS version.
Parameters
----------
@@ -642,8 +639,7 @@ class RemoteTable(Table):
fast_search: bool = False,
) -> LanceVectorQueryBuilder:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support
[vector search](https://lancedb.com/docs/search/vector-search/)
of the given query vector. We currently support [vector search][search]
All query options are defined in
[LanceVectorQueryBuilder][lancedb.query.LanceVectorQueryBuilder].
@@ -1046,11 +1042,6 @@ class RemoteTable(Table):
def fetch_blobs(self, column: str, row_ids) -> pa.LargeBinaryArray:
raise NotImplementedError("fetch_blobs() is not supported on LanceDB Cloud")
def fetch_blob_ranges(self, column: str, requests) -> pa.LargeBinaryArray:
raise NotImplementedError(
"fetch_blob_ranges() is not supported on LanceDB Cloud"
)
def fetch_blob_files(self, column: str, row_ids):
raise NotImplementedError(
"fetch_blob_files() is not supported on LanceDB Cloud"
@@ -14,9 +14,6 @@ from .answerdotai import AnswerdotaiRerankers
from .voyageai import VoyageAIReranker
from .watsonx import WatsonxReranker
# The API reference renders this module with a single mkdocstrings directive,
# which only picks up names listed here. New public names must be added to this
# list, or they will silently go undocumented.
__all__ = [
"Reranker",
"CrossEncoderReranker",
+8 -21
View File
@@ -59,10 +59,9 @@ class StreamingDataset(IterableDataset):
- **Stage 1 (I/O)**: one thread pool with ``num_splits * prefetch_batches``
workers fetches raw ``RecordBatch`` objects from LanceDB in parallel
across all splits and places them in a per-split raw-batch queue.
- **Stage 2 (transform)**: a second thread pool with
``transform_parallelism`` workers picks up raw batches, applies the
transform, and places the results in a per-split cooked-row queue. By
default, the number of workers is determined by ``os.cpu_count()``.
- **Stage 2 (transform)**: a second thread pool with ``os.cpu_count()``
workers picks up raw batches, applies the transform, and places the
results in a per-split cooked-row queue.
The main thread round-robins over the cooked queues, yielding one row per
split per cycle.
@@ -123,10 +122,6 @@ class StreamingDataset(IterableDataset):
are yielded. Receives one batch at a time and must return an iterable
whose length equals the number of rows in the batch. When ``None``
(the default) rows are returned as plain Python dicts.
transform_parallelism:
Maximum number of transforms to run concurrently. Must be greater
than zero. When ``None`` (the default), uses ``os.cpu_count()`` or 1
when the CPU count is unavailable.
worker_info_override:
If set, used in place of ``torch.utils.data.get_worker_info()`` to
determine the DataLoader worker assignment. Intended for unit tests
@@ -151,7 +146,6 @@ class StreamingDataset(IterableDataset):
shuffle_clump_size: Optional[int] = None,
filter: Optional[str] = None,
transform: Optional[Callable] = None,
transform_parallelism: Optional[int] = None,
connection_factory: Optional[Callable[[str], Any]] = None,
worker_info_override=None,
):
@@ -165,8 +159,6 @@ class StreamingDataset(IterableDataset):
f"num_splits ({num_splits}) must be divisible by "
f"world_size ({world_size})"
)
if transform_parallelism is not None and transform_parallelism <= 0:
raise ValueError("transform_parallelism must be greater than 0")
self._table = table
self._num_splits = num_splits
@@ -181,7 +173,6 @@ class StreamingDataset(IterableDataset):
self._shuffle_clump_size = shuffle_clump_size
self._filter = filter
self._transform = transform
self._transform_parallelism = transform_parallelism
self._connection_factory = connection_factory
self._worker_info_override = worker_info_override
@@ -293,11 +284,7 @@ class StreamingDataset(IterableDataset):
batch_size = self._read_batch_size
max_prefetch = self._prefetch_batches
transform_workers = (
self._transform_parallelism
if self._transform_parallelism is not None
else (os.cpu_count() or 1)
)
cpu_workers = os.cpu_count() or 1
final_transform = (
self._transform if self._transform is not None else Transforms.arrow2python
)
@@ -309,8 +296,8 @@ class StreamingDataset(IterableDataset):
tx_pending = [deque() for _ in range(n)] # Future[list[Any]]
cooked = [deque() for _ in range(n)] # rows ready to yield
# Limit simultaneous transforms to transform_workers across all splits.
tx_semaphore = threading.Semaphore(transform_workers)
# Limit simultaneous transforms to cpu_workers across all splits.
tx_semaphore = threading.Semaphore(cpu_workers)
# ── Stage 1 helpers ───────────────────────────────────────────────────
@@ -382,7 +369,7 @@ class StreamingDataset(IterableDataset):
_advance(i)
elif raw_batches[i]:
# Acquire a transform slot (may block briefly if all
# transform_workers are busy with other splits).
# cpu_workers are busy with other splits).
tx_semaphore.acquire()
batch = raw_batches[i].popleft()
tx_pending[i].append(tx_pool.submit(_tx_call_guarded, batch))
@@ -396,7 +383,7 @@ class StreamingDataset(IterableDataset):
# ── Main loop ─────────────────────────────────────────────────────────
with ThreadPoolExecutor(max_workers=n * max_prefetch) as io_pool:
with ThreadPoolExecutor(max_workers=transform_workers) as tx_pool:
with ThreadPoolExecutor(max_workers=cpu_workers) as tx_pool:
self._raw_batches_ref = raw_batches
self._cooked_ref = cooked
self._fetch_head_ref = fetch_head
+20 -88
View File
@@ -20,7 +20,6 @@ from typing import (
List,
Literal,
Optional,
Sequence,
Tuple,
Union,
overload,
@@ -1103,7 +1102,6 @@ class Table(ABC):
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -1171,9 +1169,6 @@ class Table(ABC):
remove_stop_words : bool, default True
Whether to remove stop words. Stop words are common words that are often
removed from text before indexing. For example, in English "the" and "and".
custom_stop_words : list of str, optional
Custom words that replace the built-in language stop words. ``None``
uses the built-in list; an empty list explicitly uses no stop words.
ascii_folding : bool, default True
Whether to fold ASCII characters. This converts accented characters to
their ASCII equivalent. For example, "café" would be converted to "cafe".
@@ -1211,7 +1206,7 @@ class Table(ABC):
progress: Optional[Union[bool, Callable, Any]] = None,
write_parallelism: Optional[int] = None,
) -> AddResult:
"""Add more data to the [Table][lancedb.table.Table].
"""Add more data to the [Table](Table).
Parameters
----------
@@ -1343,8 +1338,8 @@ class Table(ABC):
fts_columns: Optional[Union[str, List[str]]] = None,
) -> LanceQueryBuilder:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support [vector search](https://lancedb.com/docs/search/vector-search/)
and [full-text search](https://lancedb.com/docs/search/full-text-search/).
of the given query vector. We currently support [vector search][search]
and [full-text search][experimental-full-text-search].
All query options are defined in
[LanceQueryBuilder][lancedb.query.LanceQueryBuilder].
@@ -1543,30 +1538,10 @@ class Table(ABC):
) -> pa.LargeBinaryArray:
"""Materialize full blob bytes for ``column`` at the given rows.
The result has the same length and order as ``row_ids``. Null blobs
produce null slots; valid empty blobs produce ``b""``.
Convenience for small payloads. For large values use
:meth:`fetch_blob_files`.
"""
@abstractmethod
def fetch_blob_ranges(
self,
column: str,
requests: Sequence[Tuple[int, int, int]],
) -> pa.LargeBinaryArray:
"""Materialize row-specific byte ranges from a blob v2 column.
Each request is a ``(row_id, offset, length)`` tuple. Requests may be
repeated or reordered, including multiple ranges for the same blob.
The result has the same length and order as ``requests``; null blobs
produce null slots and empty ranges on non-null blobs produce ``b""``.
Row IDs can be obtained from a query with ``with_row_id(True)``. This
API is currently supported only by local tables.
"""
@abstractmethod
def fetch_blob_files(
self, column: str, row_ids: Union[list[int], pa.Table]
@@ -1778,7 +1753,7 @@ class Table(ABC):
for faster reads.
Arguments are passed onto Lance's
`lance.dataset.DatasetOptimizer.compact_files`.
[compact_files][lance.dataset.DatasetOptimizer.compact_files].
For most cases, the default should be fine.
See Also
@@ -1832,8 +1807,6 @@ class Table(ABC):
retrain: bool, default False
This parameter is no longer used and is deprecated.
Notes
-----
The frequency an application should call optimize is based on the frequency of
data modifications. If data is frequently added, deleted, or updated then
optimize should be run frequently. A good rule of thumb is to run optimize if
@@ -1988,14 +1961,15 @@ class Table(ABC):
change permanent you can use the `[Self::restore]` method.
Any operation that modifies the table will fail while the table is in a checked
out state. To return the table to a normal state use
`[Self::checkout_latest]`.
out state.
Parameters
----------
version: int | str,
The version to check out. A version number (`int`) or a tag
(`str`) can be provided.
To return the table to a normal state use `[Self::checkout_latest]`
"""
@abstractmethod
@@ -2291,13 +2265,6 @@ class LanceTable(Table):
) -> pa.LargeBinaryArray:
return LOOP.run(self._table.fetch_blobs(column, row_ids))
def fetch_blob_ranges(
self,
column: str,
requests: Sequence[Tuple[int, int, int]],
) -> pa.LargeBinaryArray:
return LOOP.run(self._table.fetch_blob_ranges(column, list(requests)))
def fetch_blob_files(
self, column: str, row_ids: Union[list[int], pa.Table]
) -> "list[Optional[BlobFile]]":
@@ -3060,7 +3027,6 @@ class LanceTable(Table):
lower_case: bool = True,
stem: bool = True,
remove_stop_words: bool = True,
custom_stop_words: Optional[List[str]] = None,
ascii_folding: bool = True,
ngram_min_length: int = 3,
ngram_max_length: int = 3,
@@ -3107,7 +3073,6 @@ class LanceTable(Table):
"lower_case": lower_case,
"stem": stem,
"remove_stop_words": remove_stop_words,
"custom_stop_words": custom_stop_words,
"ascii_folding": ascii_folding,
"ngram_min_length": ngram_min_length,
"ngram_max_length": ngram_max_length,
@@ -3115,7 +3080,6 @@ class LanceTable(Table):
}
else:
tokenizer_configs = self.infer_tokenizer_configs(tokenizer_name)
tokenizer_configs["custom_stop_words"] = custom_stop_words
config = FTS(block_size=block_size, **tokenizer_configs)
@@ -3388,8 +3352,8 @@ class LanceTable(Table):
fts_columns: Optional[Union[str, List[str]]] = None,
) -> LanceQueryBuilder:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support [vector search](https://lancedb.com/docs/search/vector-search/)
and [full-text search](https://lancedb.com/docs/search/full-text-search/).
of the given query vector. We currently support [vector search][search]
and [full-text search][search].
Examples
--------
@@ -3419,9 +3383,8 @@ class LanceTable(Table):
- *default None*.
Acceptable types are: list, np.ndarray, PIL.Image.Image
- If None then the
select/[where][lancedb.query.LanceQueryBuilder.where]/limit clauses
are applied to filter the table
- If None then the select/[where][sql]/limit clauses are applied
to filter the table
vector_column_name: str, optional
The name of the vector column to search.
@@ -3815,8 +3778,6 @@ class LanceTable(Table):
retrain: bool, default False
This parameter is no longer used and is deprecated.
Notes
-----
The frequency an application should call optimize is based on the frequency of
data modifications. If data is frequently added, deleted, or updated then
optimize should be run frequently. A good rule of thumb is to run optimize if
@@ -4689,24 +4650,7 @@ class AsyncTable:
"""
return AsyncQuery(self._inner.query(), self)
async def to_lance(self, **kwargs) -> lance.LanceDataset:
"""Return the Lance dataset backing this table.
Parameters
----------
**kwargs
Forwarded to `lance.dataset`.
Returns
-------
lance.LanceDataset
The Lance dataset at this table handle's version and branch.
Examples
--------
>>> async def get_lance_dataset(table):
... return await table.to_lance()
"""
async def _to_lance(self, **kwargs) -> lance.LanceDataset:
try:
import lance
except ImportError:
@@ -4756,7 +4700,7 @@ class AsyncTable:
return (await self.to_arrow()).to_pandas(**kwargs)
if blob_mode == "bytes" and blob_v2_column_paths(schema):
return await self.query().to_pandas(blob_mode=blob_mode, **kwargs)
return (await self.to_lance()).to_pandas(blob_mode=blob_mode, **kwargs)
return (await self._to_lance()).to_pandas(blob_mode=blob_mode, **kwargs)
async def to_arrow(self) -> pa.Table:
"""Return the table as a pyarrow Table.
@@ -5014,7 +4958,7 @@ class AsyncTable:
progress: Optional[Union[bool, Callable, Any]] = None,
write_parallelism: Optional[int] = None,
) -> AddResult:
"""Add more data to the [AsyncTable][lancedb.table.AsyncTable].
"""Add more data to the [Table](Table).
Parameters
----------
@@ -5216,8 +5160,8 @@ class AsyncTable:
fts_columns: Optional[Union[str, List[str]]] = None,
) -> Union[AsyncHybridQuery, AsyncFTSQuery, AsyncVectorQuery]:
"""Create a search query to find the nearest neighbors
of the given query vector. We currently support [vector search](https://lancedb.com/docs/search/vector-search/)
and [full-text search](https://lancedb.com/docs/search/full-text-search/).
of the given query vector. We currently support [vector search][search]
and [full-text search][experimental-full-text-search].
All query options are defined in [AsyncQuery][lancedb.query.AsyncQuery].
@@ -5419,8 +5363,6 @@ class AsyncTable:
async_query = async_query.where(query.filter)
if query.fast_search:
async_query = async_query.fast_search()
if query.use_lsm is not None:
async_query = async_query.use_lsm(query.use_lsm)
if query.with_row_id:
async_query = async_query.with_row_id()
if query.order_by:
@@ -5541,7 +5483,7 @@ class AsyncTable:
when_not_matched_by_source_condition_expr=merge._when_not_matched_by_source_condition_expr,
timeout=merge._timeout,
use_index=merge._use_index,
use_lsm=merge._use_lsm,
use_lsm_write=merge._use_lsm_write,
validate_single_shard=merge._validate_single_shard,
),
)
@@ -5778,14 +5720,15 @@ class AsyncTable:
change permanent you can use the `[Self::restore]` method.
Any operation that modifies the table will fail while the table is in a checked
out state. To return the table to a normal state use
`[Self::checkout_latest]`.
out state.
Parameters
----------
version: int | str,
The version to check out. A version number (`int`) or a tag
(`str`) can be provided.
To return the table to a normal state use `[Self::checkout_latest]`
"""
try:
await self._inner.checkout(version)
@@ -5884,13 +5827,6 @@ class AsyncTable:
column, _normalize_blob_row_ids(row_ids, column)
)
async def fetch_blob_ranges(
self,
column: str,
requests: Sequence[Tuple[int, int, int]],
) -> pa.LargeBinaryArray:
return await self._inner.fetch_blob_ranges(column, list(requests))
async def fetch_blob_files(
self, column: str, row_ids: Union[list[int], pa.Table]
) -> "list[Optional[BlobFile]]":
@@ -5969,8 +5905,6 @@ class AsyncTable:
retrain: bool, default False
This parameter is no longer used and is deprecated.
Notes
-----
The frequency an application should call optimize is based on the frequency of
data modifications. If data is frequently added, deleted, or updated then
optimize should be run frequently. A good rule of thumb is to run optimize if
@@ -6351,8 +6285,6 @@ class Branches:
dry_run: bool, default False
When True, only preview. When False, attempt the merge.
Notes
-----
A rejected merge returns ``status="rejected"`` instead of raising.
"""
return LOOP.run(self._table.branches.merge(from_branch, dry_run))
+4 -61
View File
@@ -184,75 +184,18 @@ def test_fetch_blobs_accepts_query_result():
assert {blobs[i].as_py() for i in range(len(blobs))} == {b"gamma"}
def test_fetch_blobs_preserves_null_and_empty_values():
def test_fetch_blobs_null_alignment():
table = _blob_table(
"nulls",
[
{"id": 1, "image": b"present"},
{"id": 2, "image": None},
{"id": 3, "image": b""},
],
[{"id": 1, "image": b"present"}, {"id": 2, "image": None}],
)
by_id = _row_ids_by_id(table)
request = [by_id[1], by_id[2], by_id[3], by_id[1]]
request = [by_id[1], by_id[2], by_id[1]]
blobs = table.fetch_blobs("image", request)
assert len(blobs) == len(request)
assert blobs[0].as_py() == b"present"
assert blobs[1].as_py() is None
assert blobs[2].as_py() == b""
assert blobs[3].as_py() == b"present"
def test_fetch_blob_ranges_aligns_repeated_ranges_and_nulls():
table = _blob_table(
"range_alignment",
[{"id": 1, "image": b"abcdefghij"}, {"id": 2, "image": None}],
)
by_id = _row_ids_by_id(table)
requests = [
(by_id[1], 2, 3),
(by_id[2], 0, 0),
(by_id[1], 0, 2),
(by_id[1], 2, 3),
(by_id[1], 10, 0),
]
ranges = table.fetch_blob_ranges("image", requests)
assert ranges.to_pylist() == [b"cde", None, b"ab", b"cde", b""]
def test_fetch_blob_ranges_validates_requests():
table = _blob_table("range_validation", [{"id": 1, "image": b"abc"}])
row_id = _row_ids_by_id(table)[1]
with pytest.raises(RuntimeError, match="exceeds blob size"):
table.fetch_blob_ranges("image", [(row_id, 2, 2)])
with pytest.raises(RuntimeError, match="offset \\+ length overflowed"):
table.fetch_blob_ranges("image", [(row_id, 2**64 - 1, 1)])
with pytest.raises(ValueError, match="row ids"):
table.fetch_blob_ranges("image", [(2**64 - 1, 0, 1)])
def test_fetch_blob_ranges_empty_requests_returns_empty_array():
table = _blob_table("range_empty", [{"id": 1, "image": b"x"}])
assert table.fetch_blob_ranges("image", []).to_pylist() == []
@pytest.mark.asyncio
async def test_async_fetch_blob_ranges():
db = await lancedb.connect_async("memory:///")
schema = pa.schema([pa.field("id", pa.int64()), lancedb.blob("image")])
table = await db.create_table("range_async", schema=schema)
await table.add([{"id": 1, "image": b"abcdefghij"}])
hits = await table.query().with_row_id().to_arrow()
row_id = hits["_rowid"][0].as_py()
ranges = await table.fetch_blob_ranges("image", [(row_id, 1, 3), (row_id, 6, 2)])
assert ranges.to_pylist() == [b"bcd", b"gh"]
assert blobs[2].as_py() == b"present"
def test_fetch_blobs_nested_path():
@@ -1333,42 +1333,6 @@ def test_transform_none_yields_dicts(lance_table):
assert all("id" in item for item in items)
@pytest.mark.parametrize(
("configured", "detected", "expected"),
[(2, 8, 2), (None, 3, 3), (None, None, 1)],
)
def test_transform_parallelism_configures_executor(
lance_table, monkeypatch, configured, detected, expected
):
"""Explicit transform parallelism overrides the detected CPU count."""
real_executor = streaming.ThreadPoolExecutor
monkeypatch.setattr(streaming.os, "cpu_count", lambda: detected)
with patch.object(streaming, "ThreadPoolExecutor", wraps=real_executor) as executor:
list(
StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle_seed=SHUFFLE_SEED,
transform_parallelism=configured,
)
)
assert executor.call_args_list[-1].kwargs["max_workers"] == expected
@pytest.mark.parametrize("transform_parallelism", [0, -1])
def test_transform_parallelism_must_be_positive(lance_table, transform_parallelism):
with pytest.raises(
ValueError, match="transform_parallelism must be greater than 0"
):
StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
transform_parallelism=transform_parallelism,
)
def test_filter_limits_rows(tmp_path):
"""A filter expression is applied to the permutation so only matching rows
are yielded. IDs 0..59 pass ``id < 60``; the other 60 are excluded."""
-20
View File
@@ -219,13 +219,11 @@ def test_create_inverted_index(table, with_position):
table.create_fts_index(
"text",
with_position=with_position,
custom_stop_words=["puppy"],
name="custom_fts_index",
)
indices = table.list_indices()
fts_indices = [i for i in indices if i.index_type == "FTS"]
assert any(i.name == "custom_fts_index" for i in fts_indices)
assert fts_indices[0].index_details["custom_stop_words"] == ["puppy"]
@pytest.mark.parametrize("block_size", [128, 256])
@@ -245,24 +243,6 @@ def test_create_inverted_index_rejects_invalid_block_size(table):
table.create_index("text", config=FTS(block_size=129))
def test_custom_stop_words_list(table):
table.create_index(
"text",
config=FTS(stem=False, custom_stop_words=["lance"]),
)
assert table.list_indices()[0].index_details["custom_stop_words"] == ["lance"]
tokens = table.tokenize("the lance data", column="text")
assert [token.text for token in tokens] == ["the", "data"]
empty_tokens = ldb.tokenize("the lance data", stem=False, custom_stop_words=[])
assert [token.text for token in empty_tokens] == ["the", "lance", "data"]
with pytest.raises(TypeError, match=r"custom_stop_words.*int"):
ldb.tokenize(
"the lance data",
custom_stop_words=["lance", 42],
)
def test_search_fts(table):
table.create_fts_index("text")
results = table.search("puppy").select(["id", "text"]).limit(5).to_list()
+12 -472
View File
@@ -9,7 +9,6 @@ import lancedb
import pyarrow as pa
import pytest
from lancedb._lancedb import LsmWriteSpec
from lancedb.index import FTS, IvfPq
SCHEMA = pa.schema(
[
@@ -103,35 +102,19 @@ def test_lsm_merge_insert_identity(tmp_path):
assert result.num_rows == 2
def test_lsm_merge_insert_use_lsm_false(tmp_path):
def test_lsm_merge_insert_use_lsm_write_false(tmp_path):
table = _bucket_table(tmp_path) # rows id = 1, 2, 3
# use_lsm(False) opts out: the standard path runs and commits even with a spec.
# use_lsm_write(False) opts out: the standard path runs and commits.
result = (
table.merge_insert("id")
.when_not_matched_insert_all()
.use_lsm(False)
.use_lsm_write(False)
.execute(_reader([3, 4, 5]))
)
assert result.num_inserted_rows == 2
assert table.count_rows() == 5
def test_lsm_merge_insert_use_lsm_true_without_spec_errors(tmp_path):
# A table with a primary key but no LSM write spec installed.
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3]))
table.set_unenforced_primary_key("id")
# use_lsm(True) demands MemWAL routing; without a spec it errors.
with pytest.raises(Exception, match="use_lsm"):
(
table.merge_insert("id")
.when_matched_update_all()
.when_not_matched_insert_all()
.use_lsm(True)
.execute(_reader([3, 4, 5]))
)
def test_lsm_merge_insert_validate_single_shard_off(tmp_path):
table = _bucket_table(tmp_path)
result = (
@@ -144,20 +127,19 @@ def test_lsm_merge_insert_validate_single_shard_off(tmp_path):
assert result.num_rows == 3
def test_lsm_merge_insert_no_spec_uses_standard_path(tmp_path):
def test_lsm_merge_insert_use_lsm_write_true_requires_spec(tmp_path):
# A table with a primary key but no LSM write spec installed.
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3]))
table.set_unenforced_primary_key("id")
# With no spec, a default merge_insert uses the standard path and commits.
result = (
table.merge_insert("id")
.when_matched_update_all()
.when_not_matched_insert_all()
.execute(_reader([3, 4, 5]))
)
assert result.num_inserted_rows == 2
assert table.count_rows() == 5
with pytest.raises(Exception, match="use_lsm_write"):
(
table.merge_insert("id")
.when_matched_update_all()
.when_not_matched_insert_all()
.use_lsm_write(True)
.execute(_reader([4]))
)
def test_lsm_merge_insert_rejects_on_not_primary_key(tmp_path):
@@ -212,445 +194,3 @@ async def test_async_lsm_merge_insert(tmp_path):
result = await builder.execute(_reader([3, 4, 5]))
assert result.num_rows == 3
await table.close_lsm_writers()
def _lsm_upsert(table, ids):
"""Upsert ``ids`` (value = 0..n) through the LSM merge_insert path."""
(
table.merge_insert([])
.when_matched_update_all()
.when_not_matched_insert_all()
.execute(_reader(ids))
)
def test_lsm_read_sees_active_memtable(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3])) # base ids 1,2,3
table.set_unenforced_primary_key("id")
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
_lsm_upsert(table, [4, 5]) # active memtable only, not committed to base
# Default read auto-routes through the LSM scanner: base active memtable.
lsm = table.search().to_arrow()
assert sorted(lsm["id"].to_pylist()) == [1, 2, 3, 4, 5]
# use_lsm(False) bypasses the MemWAL and reads the base table only.
base_only = table.search().use_lsm(False).to_arrow()
assert sorted(base_only["id"].to_pylist()) == [1, 2, 3]
def test_lsm_read_dedup_newest_wins(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3])) # id 2 -> value 1
table.set_unenforced_primary_key("id")
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
_lsm_upsert(table, [2, 3, 4]) # ids 2,3,4 -> values 0,1,2
lsm = table.search().to_arrow().sort_by("id")
assert lsm["id"].to_pylist() == [1, 2, 3, 4]
# id 1 from base (value 0); 2,3,4 from memtable (values 0,1,2).
assert lsm["value"].to_pylist() == [0, 0, 1, 2]
def test_lsm_read_without_spec_reads_base(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3]))
table.set_unenforced_primary_key("id") # no LSM write spec
# No spec: default read and use_lsm(False) both read the base table, no error.
assert sorted(table.search().to_arrow()["id"].to_pylist()) == [1, 2, 3]
assert sorted(table.search().use_lsm(False).to_arrow()["id"].to_pylist()) == [
1,
2,
3,
]
def test_lsm_read_unsupported_shape_errors_without_use_lsm_false(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3]))
table.set_unenforced_primary_key("id")
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
_lsm_upsert(table, [4])
# with_row_id is unsupported by the LSM scanner; on a MemWAL table the default
# (auto-routed) read hard-errors instead of silently reading a stale base.
with pytest.raises(Exception):
table.search().with_row_id(True).to_arrow()
# use_lsm(False) is the escape hatch: it reads the base table only.
base = table.search().with_row_id(True).use_lsm(False).to_arrow()
assert sorted(base["id"].to_pylist()) == [1, 2, 3]
@pytest.mark.asyncio
async def test_async_lsm_read(tmp_path):
db = await lancedb.connect_async(
tmp_path, read_consistency_interval=timedelta(seconds=0)
)
table = await db.create_table("t", _reader([1, 2, 3]))
await table.set_unenforced_primary_key("id")
await table.set_lsm_write_spec(LsmWriteSpec.unsharded())
builder = (
table.merge_insert([]).when_matched_update_all().when_not_matched_insert_all()
)
await builder.execute(_reader([4, 5]))
arrow = await table.query().to_arrow()
assert sorted(arrow["id"].to_pylist()) == [1, 2, 3, 4, 5]
VECTOR_DIM = 8
VECTOR_SCHEMA = pa.schema(
[
pa.field("id", pa.int64(), nullable=False),
pa.field("category", pa.utf8(), nullable=False),
pa.field("vector", pa.list_(pa.float32(), VECTOR_DIM), nullable=False),
]
)
def _vector_reader(rows):
"""Rows are ``(id, category, [f32; VECTOR_DIM])`` tuples."""
batch = pa.RecordBatch.from_arrays(
[
pa.array([row[0] for row in rows], type=pa.int64()),
pa.array([row[1] for row in rows], type=pa.utf8()),
pa.array([row[2] for row in rows], type=pa.list_(pa.float32(), VECTOR_DIM)),
],
schema=VECTOR_SCHEMA,
)
return pa.RecordBatchReader.from_batches(VECTOR_SCHEMA, [batch])
def _vector_table(tmp_path):
"""Base table whose vector column is indexed so its rows are visible to the LSM
vector scanner (the base arm uses ``fast_search`` indexed data only), plus an
unsharded LSM spec that maintains that index for the memtable.
Rows 1,2 are category ``a``, row 3 is ``b``, and 4..60 are filler ``c`` that
give the tiny IVF index enough data to train.
"""
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
rows = [
(
i,
"a" if i in (1, 2) else "b" if i == 3 else "c",
[float((i * 7 + j) % 13) for j in range(VECTOR_DIM)],
)
for i in range(1, 61)
]
table = db.create_table("t", _vector_reader(rows))
table.set_unenforced_primary_key("id")
# num_partitions=1 makes the search exhaustive within the single partition
# (deterministic); num_bits=4 keeps PQ training viable on a tiny dataset.
table.create_index(
"vector", config=IvfPq(num_partitions=1, num_sub_vectors=2, num_bits=4)
)
index_name = table.list_indices()[0].name
table.set_lsm_write_spec(
LsmWriteSpec.unsharded().with_maintained_indexes([index_name])
)
return table
def _vector_upsert(table, rows):
(
table.merge_insert([])
.when_matched_update_all()
.when_not_matched_insert_all()
.execute(_vector_reader(rows))
)
def test_lsm_read_vector_sees_memtable(tmp_path):
table = _vector_table(tmp_path)
# id 1000 lands in the active memtable, not committed to the base table.
_vector_upsert(table, [(1000, "a", [1.0] * VECTOR_DIM)])
query = [1.0] * VECTOR_DIM
# Vector search auto-routes through the LSM scanner: indexed base memtable.
ids = set(table.search(query).limit(100).to_arrow()["id"].to_pylist())
assert {1, 2, 3} <= ids # indexed base rows
assert 1000 in ids # in-flight memtable row
# use_lsm(False) bypasses the MemWAL, so the in-flight row is not visible.
base_ids = set(
table.search(query).use_lsm(False).limit(100).to_arrow()["id"].to_pylist()
)
assert {1, 2, 3} <= base_ids
assert 1000 not in base_ids
def test_lsm_read_vector_prefilter(tmp_path):
table = _vector_table(tmp_path)
# in-flight rows in both categories.
_vector_upsert(
table, [(1000, "a", [1.0] * VECTOR_DIM), (1001, "b", [1.0] * VECTOR_DIM)]
)
query = [1.0] * VECTOR_DIM
# The `where` predicate must apply as a prefilter across base memtable —
# regression test for the vector arm silently dropping the filter.
rows = table.search(query).where("category = 'a'").limit(100).to_arrow()
assert set(rows["id"].to_pylist()) == {1, 2, 1000}
assert set(rows["category"].to_pylist()) == {"a"}
# Sanity: without the filter, other categories are returned too.
unfiltered = set(table.search(query).limit(100).to_arrow()["category"].to_pylist())
assert unfiltered != {"a"}
def test_lsm_read_plain_prefilter(tmp_path):
table = _vector_table(tmp_path)
_vector_upsert(
table, [(1000, "a", [1.0] * VECTOR_DIM), (1001, "b", [1.0] * VECTOR_DIM)]
)
# Plain scan + filter over base memtable: base 'a' rows 1,2 and memtable 1000.
rows = table.search().where("category = 'a'").to_arrow()
assert set(rows["id"].to_pylist()) == {1, 2, 1000}
FTS_SCHEMA = pa.schema(
[
pa.field("id", pa.int64(), nullable=False),
pa.field("text", pa.utf8(), nullable=False),
]
)
def _fts_reader(rows):
"""Rows are ``(id, text)`` tuples."""
batch = pa.RecordBatch.from_arrays(
[
pa.array([row[0] for row in rows], type=pa.int64()),
pa.array([row[1] for row in rows], type=pa.utf8()),
],
schema=FTS_SCHEMA,
)
return pa.RecordBatchReader.from_batches(FTS_SCHEMA, [batch])
def test_lsm_read_fts_sees_memtable(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table(
"t",
_fts_reader(
[
(1, "the quick brown fox"),
(2, "lazy dog sleeps"),
(3, "quick red fox"),
]
),
)
table.set_unenforced_primary_key("id")
# Native FTS index (tantivy is not compatible with the LSM memtable index).
table.create_index("text", config=FTS())
index_name = table.list_indices()[0].name
table.set_lsm_write_spec(
LsmWriteSpec.unsharded().with_maintained_indexes([index_name])
)
# in-flight doc 4 lands in the memtable's maintained FTS index.
(
table.merge_insert([])
.when_matched_update_all()
.when_not_matched_insert_all()
.execute(_fts_reader([(4, "brown fox jumps")]))
)
# Full-text search auto-routes through the LSM scanner: base memtable.
ids = set(
table.search("fox", query_type="fts", fts_columns="text")
.limit(10)
.to_arrow()["id"]
.to_pylist()
)
assert ids == {1, 3, 4}
# Prefilter restricts the FTS results across both tiers.
filtered = set(
table.search("fox", query_type="fts", fts_columns="text")
.where("id > 1")
.limit(10)
.to_arrow()["id"]
.to_pylist()
)
assert filtered == {3, 4}
def test_lsm_read_vector_unsupported_knobs_error(tmp_path):
table = _vector_table(tmp_path)
_vector_upsert(table, [(1000, "a", [1.0] * VECTOR_DIM)])
query = [1.0] * VECTOR_DIM
# distance_range and use_index(False) change the vector result set/mode, which
# the LSM scanner can't honor, so it hard-errors instead of silently returning
# wrong results (matching the prefilter / unsupported-shape contract).
with pytest.raises(Exception, match="distance_range"):
table.search(query).distance_range(0.0, 0.5).to_arrow()
with pytest.raises(Exception, match="use_index"):
table.search(query).bypass_vector_index().to_arrow()
# use_lsm(False) is the escape hatch: the base-only standard path honors them.
base = table.search(query).distance_range(0.0, 100.0).use_lsm(False).to_arrow()
assert 1000 not in set(base["id"].to_pylist())
def test_lsm_read_vector_limit_offset(tmp_path):
table = _vector_table(tmp_path)
_vector_upsert(table, [(1000, "a", [1.0] * VECTOR_DIM)])
query = [1.0] * VECTOR_DIM
# Lance's plan_vector over-fetches k + offset internally, so paging is correct:
# the second page is a full page (not truncated) and disjoint from the first.
page1 = table.search(query).limit(3).offset(0).to_arrow()["id"].to_pylist()
page2 = table.search(query).limit(3).offset(3).to_arrow()["id"].to_pylist()
assert len(page1) == 3
# If k ignored offset, page2 would be empty (limit - offset = 0); a full second
# page that differs from the first proves offset widens the candidate pool.
assert len(page2) == 3
assert set(page1) != set(page2)
def test_lsm_read_vector_postfilter_errors(tmp_path):
table = _vector_table(tmp_path)
_vector_upsert(table, [(1000, "a", [1.0] * VECTOR_DIM)])
query = [1.0] * VECTOR_DIM
# The LSM scanner always prefilters; a requested postfilter changes results, so
# it hard-errors rather than silently prefiltering.
with pytest.raises(Exception, match="postfilter"):
table.search(query).where("category = 'a'").postfilter().to_arrow()
def test_lsm_read_projection_excludes_pk(tmp_path):
table = _vector_table(tmp_path)
_vector_upsert(table, [(1000, "a", [1.0] * VECTOR_DIM)])
# Selecting only 'category' must not leak the 'id' primary key Lance appends
# internally for dedup.
rows = table.search().select(["category"]).where("category = 'a'").to_arrow()
assert rows.column_names == ["category"]
def test_lsm_read_fts_unmaintained_index_errors(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _fts_reader([(1, "quick fox"), (2, "lazy dog")]))
table.set_unenforced_primary_key("id")
table.create_index("text", config=FTS())
# No maintained indexes: the active memtable FTS arm cannot serve un-compacted
# docs, so the search would silently omit them — reject instead.
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
with pytest.raises(Exception, match="maintained"):
table.search("fox", query_type="fts", fts_columns="text").to_arrow()
def test_lsm_read_time_travel_errors(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3]))
table.set_unenforced_primary_key("id")
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
pinned = table.version
table.add(_reader([4, 5])) # standard add commits a newer version
table.checkout(pinned) # detached head at the historical version
# The WAL/manifest expose current live state, so an LSM read at a pinned
# historical version is rejected.
with pytest.raises(Exception, match="time-travel"):
table.search().to_arrow()
# use_lsm(False) reads the base table at the pinned version.
base = table.search().use_lsm(False).to_arrow()
assert sorted(base["id"].to_pylist()) == [1, 2, 3]
def test_lsm_read_take_row_ids_errors(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _reader([1, 2, 3]))
table.set_unenforced_primary_key("id")
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
_lsm_upsert(table, [4])
# take-by-row-id auto-routes through the LSM scanner, which has no stable _rowid,
# so it hard-errors instead of failing with an opaque column-not-found error.
with pytest.raises(Exception, match="row id"):
table.take_row_ids([0, 1]).to_arrow()
# use_lsm(False) is the escape hatch: it reads the base table.
base = table.take_row_ids([0, 1]).use_lsm(False).to_arrow()
assert base.num_rows == 2
def test_lsm_read_fts_postfilter_errors(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _fts_reader([(1, "quick fox"), (2, "lazy dog")]))
table.set_unenforced_primary_key("id")
table.create_index("text", config=FTS())
index_name = table.list_indices()[0].name
table.set_lsm_write_spec(
LsmWriteSpec.unsharded().with_maintained_indexes([index_name])
)
# The LSM scanner always prefilters; postfilter on FTS changes result semantics,
# so it hard-errors (previously only the vector arm rejected it).
with pytest.raises(Exception, match="postfilter"):
(
table.search("fox", query_type="fts", fts_columns="text")
.where("id > 0")
.postfilter()
.to_arrow()
)
def test_lsm_read_fts_multiple_same_type_indexes_errors(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _fts_reader([(1, "quick fox"), (2, "lazy dog")]))
table.set_unenforced_primary_key("id")
table.create_index("text", config=FTS(), name="fts_a")
table.create_index("text", config=FTS(), name="fts_b", replace=False)
table.set_lsm_write_spec(
LsmWriteSpec.unsharded().with_maintained_indexes(["fts_a"])
)
# Two FTS indexes on the column: the base planner's chosen index is ambiguous, so
# the scanner can't pick a catch-up watermark and rejects rather than risk
# dropping rows the actually-used index has not caught up to.
with pytest.raises(Exception, match="multiple"):
table.search("fox", query_type="fts", fts_columns="text").to_arrow()
def test_lsm_read_vector_unmaintained_index_errors(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
rows = [
(i, "a", [float((i * 7 + j) % 13) for j in range(VECTOR_DIM)])
for i in range(1, 61)
]
table = db.create_table("t", _vector_reader(rows))
table.set_unenforced_primary_key("id")
table.create_index(
"vector", config=IvfPq(num_partitions=1, num_sub_vectors=2, num_bits=4)
)
# Spec with NO maintained indexes: the base vector index's catch-up is untracked,
# so the scanner rejects rather than risk dropping compacted-but-unindexed rows.
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
with pytest.raises(Exception, match="maintained"):
table.search([1.0] * VECTOR_DIM).to_arrow()
def test_lsm_read_fts_optimized_index_not_rejected(tmp_path):
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(seconds=0))
table = db.create_table("t", _fts_reader([(i, "quick fox") for i in range(1, 6)]))
table.set_unenforced_primary_key("id")
table.create_index("text", config=FTS())
table.add(_fts_reader([(i, "lazy fox") for i in range(6, 11)]))
table.optimize() # may split the FTS index into multiple physical segments
name = table.list_indices()[0].name
table.set_lsm_write_spec(LsmWriteSpec.unsharded().with_maintained_indexes([name]))
# Multiple physical segments of one logical index must not be miscounted as
# multiple indexes and rejected.
ids = set(
table.search("fox", query_type="fts", fts_columns="text")
.limit(20)
.to_arrow()["id"]
.to_pylist()
)
assert ids == set(range(1, 11))
-2
View File
@@ -771,7 +771,6 @@ def test_table_create_indices():
"text",
wait_timeout=timedelta(seconds=2),
block_size=256,
custom_stop_words=["cloud"],
name="custom_fts_idx",
)
@@ -796,7 +795,6 @@ def test_table_create_indices():
assert "name" in fts_req
assert fts_req["name"] == "custom_fts_idx"
assert fts_req["block_size"] == 256
assert fts_req["custom_stop_words"] == ["cloud"]
# Check vector index request has custom name
vector_req = received_requests[2]
-47
View File
@@ -1257,53 +1257,6 @@ def test_branch_to_lance_targets_branch(tmp_path):
assert table.to_lance().count_rows() == 1
@pytest.mark.asyncio
async def test_async_to_lance(tmp_path):
pytest.importorskip("lance")
db = await lancedb.connect_async(tmp_path)
table = await db.create_table("t", [{"i": 1}])
dataset = await table.to_lance()
assert dataset.count_rows() == 1
@pytest.mark.asyncio
async def test_async_branch_to_lance_targets_branch(tmp_path):
pytest.importorskip("lance")
db = await lancedb.connect_async(tmp_path)
table = await db.create_table("t", [{"i": 1}])
branch = await table.branches.create("exp")
await branch.add([{"i": 2}])
assert (await branch.to_lance()).count_rows() == 2
assert (await table.to_lance()).count_rows() == 1
@pytest.mark.asyncio
async def test_async_to_lance_targets_checked_out_version(tmp_path):
pytest.importorskip("lance")
db = await lancedb.connect_async(tmp_path)
table = await db.create_table("t", [{"i": 1}])
version = await table.version()
await table.add([{"i": 2}])
checked_out = await db.open_table("t", version=version)
assert (await checked_out.to_lance()).count_rows() == 1
assert (await table.to_lance()).count_rows() == 2
@pytest.mark.asyncio
async def test_async_to_lance_forwards_dataset_options(tmp_path):
pytest.importorskip("lance")
db = await lancedb.connect_async(tmp_path)
table = await db.create_table("t", [{"i": 1}])
dataset = await table.to_lance(default_scan_options={"with_row_id": True})
assert "_rowid" in dataset.schema.names
@pytest.mark.asyncio
async def test_async_branches(tmp_path):
db = await lancedb.connect_async(tmp_path)
+1 -3
View File
@@ -59,8 +59,7 @@ pub fn extract_index_params(source: &Option<Bound<'_, PyAny>>) -> PyResult<Lance
.ascii_folding(params.ascii_folding)
.ngram_min_length(params.ngram_min_length)
.ngram_max_length(params.ngram_max_length)
.ngram_prefix_only(params.prefix_only)
.custom_stop_words(params.custom_stop_words);
.ngram_prefix_only(params.prefix_only);
let inner_opts = inner_opts
.block_size(params.block_size)
.map_err(|err| PyValueError::new_err(err.to_string()))?;
@@ -207,7 +206,6 @@ struct FtsParams {
lower_case: bool,
stem: bool,
remove_stop_words: bool,
custom_stop_words: Option<Vec<String>>,
ascii_folding: bool,
ngram_min_length: u32,
ngram_max_length: u32,
-24
View File
@@ -294,7 +294,6 @@ pub struct PyQueryRequest {
pub select: PySelect,
pub fast_search: Option<bool>,
pub with_row_id: Option<bool>,
pub use_lsm: Option<bool>,
pub column: Option<String>,
pub query_vector: Option<PyQueryVectors>,
pub minimum_nprobes: Option<usize>,
@@ -325,7 +324,6 @@ impl From<AnyQuery> for PyQueryRequest {
select: PySelect(query_request.select),
fast_search: Some(query_request.fast_search),
with_row_id: Some(query_request.with_row_id),
use_lsm: query_request.use_lsm,
column: None,
query_vector: None,
minimum_nprobes: None,
@@ -350,7 +348,6 @@ impl From<AnyQuery> for PyQueryRequest {
select: PySelect(vector_query.base.select),
fast_search: Some(vector_query.base.fast_search),
with_row_id: Some(vector_query.base.with_row_id),
use_lsm: vector_query.base.use_lsm,
column: vector_query.column,
query_vector: Some(PyQueryVectors(vector_query.query_vector)),
minimum_nprobes: Some(vector_query.minimum_nprobes),
@@ -477,10 +474,6 @@ impl Query {
self.inner = self.inner.clone().fast_search();
}
pub fn use_lsm(&mut self, enable: bool) {
self.inner = self.inner.clone().use_lsm(enable);
}
pub fn with_row_id(&mut self) {
self.inner = self.inner.clone().with_row_id();
}
@@ -643,10 +636,6 @@ impl TakeQuery {
self.inner = self.inner.clone().with_row_id();
}
pub fn use_lsm(&mut self, enable: bool) {
self.inner = self.inner.clone().use_lsm(enable);
}
#[pyo3(signature = ())]
pub fn output_schema(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner.clone();
@@ -756,10 +745,6 @@ impl FTSQuery {
self.inner = self.inner.clone().fast_search();
}
pub fn use_lsm(&mut self, enable: bool) {
self.inner = self.inner.clone().use_lsm(enable);
}
pub fn with_row_id(&mut self) {
self.inner = self.inner.clone().with_row_id();
}
@@ -907,10 +892,6 @@ impl VectorQuery {
self.inner = self.inner.clone().fast_search();
}
pub fn use_lsm(&mut self, enable: bool) {
self.inner = self.inner.clone().use_lsm(enable);
}
pub fn with_row_id(&mut self) {
self.inner = self.inner.clone().with_row_id();
}
@@ -1105,11 +1086,6 @@ impl HybridQuery {
self.inner_fts.postfilter();
}
pub fn use_lsm(&mut self, enable: bool) {
self.inner_vec.use_lsm(enable);
self.inner_fts.use_lsm(enable);
}
pub fn add_query_vector(&mut self, vector: Bound<'_, PyAny>) -> PyResult<()> {
self.inner_vec.add_query_vector(vector)
}
+5 -29
View File
@@ -17,7 +17,7 @@ use arrow::{
ffi_stream::ArrowArrayStreamReader,
pyarrow::{FromPyArrow, PyArrowType, ToPyArrow},
};
use lancedb::blob::{BlobFile, BlobRangeRequest};
use lancedb::blob::BlobFile;
use lancedb::index::scalar::FtsIndexBuilder;
use lancedb::table::{
AddDataMode, ColumnAlteration, Duration, FieldMetadataUpdate, FtsToken as LanceDbFtsToken,
@@ -520,7 +520,6 @@ impl From<LanceDbFtsToken> for FtsToken {
lower_case = true,
stem = true,
remove_stop_words = true,
custom_stop_words = None,
ascii_folding = true,
ngram_min_length = 3,
ngram_max_length = 3,
@@ -535,7 +534,6 @@ pub fn tokenize(
lower_case: bool,
stem: bool,
remove_stop_words: bool,
custom_stop_words: Option<Vec<String>>,
ascii_folding: bool,
ngram_min_length: u32,
ngram_max_length: u32,
@@ -557,8 +555,7 @@ pub fn tokenize(
.ascii_folding(ascii_folding)
.ngram_min_length(ngram_min_length)
.ngram_max_length(ngram_max_length)
.ngram_prefix_only(prefix_only)
.custom_stop_words(custom_stop_words);
.ngram_prefix_only(prefix_only);
let tokens = lancedb_tokenize(&query, &params).infer_error()?;
Ok(tokens.into_iter().map(FtsToken::from).collect())
}
@@ -1104,27 +1101,6 @@ impl Table {
})
}
/// Read row-specific blob-local byte ranges in one planned operation.
#[pyo3(signature = (column, requests))]
pub fn fetch_blob_ranges(
self_: PyRef<'_, Self>,
column: String,
requests: Vec<(u64, u64, u64)>,
) -> PyResult<Bound<'_, PyAny>> {
let inner = self_.inner_ref()?.clone();
future_into_py(self_.py(), async move {
let requests = requests
.into_iter()
.map(|(row_id, offset, length)| BlobRangeRequest::new(row_id, offset, length))
.collect::<Vec<_>>();
let blobs: LargeBinaryArray = inner
.fetch_blob_ranges(column, requests)
.await
.infer_error()?;
Python::attach(|py| blobs.to_data().to_pyarrow(py).map(|obj| obj.unbind()))
})
}
/// Open lazy blob handles for `row_ids` from blob v2 column `column`.
#[pyo3(signature = (column, row_ids))]
pub fn fetch_blob_files(
@@ -1239,8 +1215,8 @@ impl Table {
if let Some(use_index) = parameters.use_index {
builder.use_index(use_index);
}
if let Some(use_lsm) = parameters.use_lsm {
builder.use_lsm(use_lsm);
if let Some(use_lsm_write) = parameters.use_lsm_write {
builder.use_lsm_write(use_lsm_write);
}
if let Some(validate_single_shard) = parameters.validate_single_shard {
builder.validate_single_shard(validate_single_shard);
@@ -1478,7 +1454,7 @@ pub struct MergeInsertParams {
when_not_matched_by_source_condition_expr: Option<PyExpr>,
timeout: Option<std::time::Duration>,
use_index: Option<bool>,
use_lsm: Option<bool>,
use_lsm_write: Option<bool>,
validate_single_shard: Option<bool>,
}
+20 -26
View File
@@ -1,17 +1,12 @@
# Release process
We release four `lancedb` packages: Python, Rust, Java, and Node.js.
There are five total packages we release. Four are the `lancedb` packages
for Python, Rust, Java, and Node.js. The other one is the legacy `vectordb`
package node.js.
All four share a single version number, defined by `current_version` in
`.bumpversion.toml`. One `vX.Y.Z` tag releases all of them, so a breaking change
in any SDK bumps the minor version for every SDK.
> [!NOTE]
> Python used to be versioned separately, under `python-vX.Y.Z` tags. It ran
> three minor versions ahead of the other SDKs, which made the two numbers hard
> to reason about. Both tracks were merged at `v0.37.0`: the Python line went
> `0.36``0.37` as usual, while Rust, Java, and Node.js jumped `0.33``0.37`
> to catch up. Tags before `v0.37.0` follow the old split scheme.
The Python package is versioned and released separately from the Rust, Java, and Node.js
ones. For Node.js the release process is shared between `lancedb` and
`vectordb` for now.
## Preview releases
@@ -32,21 +27,20 @@ The release process uses a handful of GitHub actions to automate the process.
┌─────────────────────┐
│Create Release Commit│
└─┬───────────────────┘
│ ┌──────────────
──►(tag) vX.Y.Z ──►│GitHub Release├───►GH Release
└──────────────
│ ┌────────────┐
├─►│PyPI Publish├─────►Python Wheels
│ └────────────┘
───────────
├─►│NPM Publish├──────►NPM Packages
───────────
─────────────┐
├─►│Cargo Publish├────►Cargo Release
─────────────
┌─────────────┐
└─►│Maven Publish├────►Java Maven Repo Release
└─────────────┘
┌────────────┐ ┌──►Python GH Release
──►(tag) python-vX.Y.Z ──►│PyPI Publish├─┤
└────────────┘ └──►Python Wheels
│ ┌───────────┐
└──►(tag) vX.Y.Z ───┬──────►│NPM Publish├──┬──►Rust/Node GH Release
───────────┘ │
│ └──►NPM Packages
┌─────────────
──────►│Cargo Publish├───►Cargo Release
│ └─────────────┘
─────────────
──────►│Maven Publish├───►Java Maven Repo Release
└─────────────┘
```
To start a release, trigger a `Create Release Commit` action from
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb"
version = "0.37.1-beta.0"
version = "0.32.0-beta.3"
edition.workspace = true
description = "LanceDB: A serverless, low-latency vector database for AI applications"
license.workspace = true
+1 -6
View File
@@ -76,12 +76,7 @@ async fn create_table(db: &Connection) -> Result<Table> {
async fn create_index(table: &Table) -> Result<()> {
table
.create_index(
&["doc"],
Index::FTS(
FtsIndexBuilder::default().custom_stop_words(Some(vec!["example".to_owned()])),
),
)
.create_index(&["doc"], Index::FTS(FtsIndexBuilder::default()))
.execute()
.await?;
Ok(())
+142 -84
View File
@@ -11,42 +11,18 @@
use std::sync::Arc;
use arrow_array::LargeBinaryArray;
use arrow_array::builder::LargeBinaryBuilder;
use arrow_array::{Array, LargeBinaryArray, RecordBatch, StructArray, UInt8Array, UInt64Array};
use arrow_schema::{DataType, Field, Schema};
use lance::dataset::{BlobRangeRequest as LanceBlobRangeRequest, Dataset, WriteParams};
use lance::dataset::{Dataset, WriteParams};
use lance_arrow::FieldExt;
use lance_core::datatypes::parse_field_path;
use lance_encoding::version::LanceFileVersion;
use crate::error::{Error, Result};
pub use lance::dataset::BlobFile;
/// One row-specific blob range read request.
///
/// `row_id` is obtained from a query with row ids enabled.
/// `offset` and `length` are relative to the beginning of the logical blob.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct BlobRangeRequest {
/// Row id of the blob value to read.
pub row_id: u64,
/// Byte offset from the beginning of the blob value.
pub offset: u64,
/// Number of bytes to read.
pub length: u64,
}
impl BlobRangeRequest {
/// Create a row-specific blob range request.
pub const fn new(row_id: u64, offset: u64, length: u64) -> Self {
Self {
row_id,
offset,
length,
}
}
}
/// Creates an Arrow field for a Lance blob v2 column.
///
/// `Struct<data, uri>` with the `lance.blob.v2` marker. Same layout Lance
@@ -169,57 +145,91 @@ pub(crate) fn ensure_blob_v2_column(
}
}
fn ensure_all_row_ids_resolved(column: &str, requested: usize, resolved: usize) -> Result<()> {
if requested == resolved {
return Ok(());
}
if resolved < requested {
Err(Error::InvalidInput {
message: format!(
"blob read for column '{column}' requested {requested} row ids but only {resolved} \
exist in the table; pass row ids collected from this table"
),
})
} else {
Err(Error::Runtime {
message: format!(
"blob read for column '{column}' returned {resolved} results for {requested} row ids"
),
})
/// Returns the leaf descriptor `StructArray` for `column` in a descriptor batch.
fn leaf_descriptor_struct<'a>(batch: &'a RecordBatch, column: &str) -> Result<&'a StructArray> {
let path = parse_field_path(column).map_err(|e| Error::InvalidInput {
message: format!("invalid blob column path '{column}': {e}"),
})?;
let not_struct = || Error::Runtime {
message: format!("blob column '{column}' did not read back as a descriptor struct"),
};
let mut current = batch
.column_by_name(&path[0])
.and_then(|c| c.as_any().downcast_ref::<StructArray>())
.ok_or_else(not_struct)?;
for segment in &path[1..] {
current = current
.column_by_name(segment)
.and_then(|c| c.as_any().downcast_ref::<StructArray>())
.ok_or_else(not_struct)?;
}
Ok(current)
}
/// Materialize blob-local ranges (same length and order as `requests`, nulls preserved).
pub(crate) async fn take_blob_ranges_aligned(
/// Null rows in `row_ids`, from a descriptor take.
///
/// Lance `read_blobs` / `take_blobs` skip null rows (`kind == 0 && position == 0 && size == 0`).
/// TODO(lance): aligned read API would drop this pass.
async fn blob_null_mask(
dataset: &Arc<Dataset>,
column: &str,
requests: &[BlobRangeRequest],
) -> Result<LargeBinaryArray> {
ensure_blob_v2_column(dataset.schema(), column)?;
if requests.is_empty() {
return Ok(LargeBinaryBuilder::new().finish());
row_ids: &[u64],
) -> Result<Vec<bool>> {
let projection = dataset.schema().project(&[column])?;
let descriptors = dataset.take_builder(row_ids, projection)?.execute().await?;
if descriptors.num_rows() != row_ids.len() {
return Err(Error::InvalidInput {
message: format!(
"blob take for column '{column}' requested {} row ids but only {} exist in the \
table; pass row ids collected from this table",
row_ids.len(),
descriptors.num_rows()
),
});
}
let descriptor_struct = leaf_descriptor_struct(&descriptors, column)?;
let child = |name: &str| {
descriptor_struct
.column_by_name(name)
.ok_or_else(|| Error::Runtime {
message: format!("blob descriptor for '{column}' is missing the '{name}' field"),
})
};
let kinds = child("kind")?
.as_any()
.downcast_ref::<UInt8Array>()
.ok_or_else(|| Error::Runtime {
message: format!("blob descriptor 'kind' for '{column}' is not a UInt8 array"),
})?;
let positions = child("position")?
.as_any()
.downcast_ref::<UInt64Array>()
.ok_or_else(|| Error::Runtime {
message: format!("blob descriptor 'position' for '{column}' is not a UInt64 array"),
})?;
let sizes = child("size")?
.as_any()
.downcast_ref::<UInt64Array>()
.ok_or_else(|| Error::Runtime {
message: format!("blob descriptor 'size' for '{column}' is not a UInt64 array"),
})?;
let lance_requests = requests
// Match Lance `collect_blob_entries_v2` skip condition (`BlobKind::Inline` == 0).
Ok((0..descriptor_struct.len())
.map(|i| {
descriptor_struct.is_null(i)
|| kinds.is_null(i)
|| (kinds.value(i) == 0 && positions.value(i) == 0 && sizes.value(i) == 0)
})
.collect())
}
fn non_null_row_ids(row_ids: &[u64], null_mask: &[bool]) -> Vec<u64> {
row_ids
.iter()
.map(|request| LanceBlobRangeRequest::new(request.row_id, request.offset, request.length))
.collect::<Vec<_>>();
let payloads = dataset
.read_blob_ranges(column)?
.with_row_ids(lance_requests)
.preserve_order(true)
.execute()
.await?;
ensure_all_row_ids_resolved(column, requests.len(), payloads.len())?;
let mut builder = LargeBinaryBuilder::new();
for payload in payloads {
match payload.data {
Some(data) => builder.append_value(data),
None => builder.append_null(),
}
}
Ok(builder.finish())
.zip(null_mask)
.filter_map(|(row_id, is_null)| (!is_null).then_some(*row_id))
.collect()
}
/// Materialize blob bytes for `row_ids` (same length and order, nulls preserved).
@@ -233,19 +243,42 @@ pub(crate) async fn take_blobs_aligned(
return Ok(LargeBinaryBuilder::new().finish());
}
let payloads = dataset
.read_blobs(column)?
.with_row_ids(row_ids.to_vec())
.preserve_order(true)
.execute()
.await?;
ensure_all_row_ids_resolved(column, row_ids.len(), payloads.len())?;
let null_mask = blob_null_mask(dataset, column, row_ids).await?;
let non_null_row_ids = non_null_row_ids(row_ids, &null_mask);
let non_null_count = non_null_row_ids.len();
let payloads = if non_null_count == 0 {
Vec::new()
} else {
dataset
.read_blobs(column)?
.with_row_ids(non_null_row_ids)
.preserve_order(true)
.execute()
.await?
};
if payloads.len() != non_null_count {
return Err(Error::Runtime {
message: format!(
"blob read for column '{column}' returned {} payloads for {} non-null rows",
payloads.len(),
non_null_count
),
});
}
let mut builder = LargeBinaryBuilder::new();
for payload in payloads {
match payload.data {
Some(data) => builder.append_value(data),
None => builder.append_null(),
let mut payload_idx = 0;
for is_null in &null_mask {
if *is_null {
builder.append_null();
} else {
if let Some(data) = &payloads[payload_idx].data {
builder.append_value(data);
} else {
builder.append_null();
}
payload_idx += 1;
}
}
Ok(builder.finish())
@@ -262,9 +295,34 @@ pub(crate) async fn take_blob_files_aligned(
return Ok(Vec::new());
}
let handles = dataset.take_blobs(row_ids, column).await?;
ensure_all_row_ids_resolved(column, row_ids.len(), handles.len())?;
Ok(handles)
let null_mask = blob_null_mask(dataset, column, row_ids).await?;
let non_null_row_ids = non_null_row_ids(row_ids, &null_mask);
let handles = if non_null_row_ids.is_empty() {
Vec::new()
} else {
dataset.take_blobs(&non_null_row_ids, column).await?
};
if handles.len() != non_null_row_ids.len() {
return Err(Error::Runtime {
message: format!(
"blob take for column '{column}' returned {} handles for {} non-null rows",
handles.len(),
non_null_row_ids.len()
),
});
}
let mut handles = handles.into_iter();
Ok(null_mask
.iter()
.map(|is_null| {
if *is_null {
None
} else {
handles.next().flatten()
}
})
.collect())
}
#[cfg(test)]
+1 -6
View File
@@ -167,11 +167,6 @@
//! # }
//! ```
// The MemWAL LSM read path (`table::query::lsm`) deepens the `create_plan` future's
// type graph enough to overflow the default trait-recursion limit while evaluating
// auto-traits (`Send`) through the Linux io_uring build's moka cache. Raise it.
#![recursion_limit = "256"]
pub mod arrow;
pub mod blob;
pub mod connection;
@@ -201,7 +196,7 @@ use std::{fmt::Display, str::FromStr};
use serde::{Deserialize, Serialize};
pub use blob::{BlobRangeRequest, blob, is_blob};
pub use blob::{blob, is_blob};
pub use connection::{ConnectNamespaceBuilder, Connection};
pub use error::{Error, Result};
use lance_index::vector::ApproxMode as LanceApproxMode;
-40
View File
@@ -523,26 +523,6 @@ pub trait QueryBase {
///
/// This allows ordering query results by one or more columns in either ascending or descending order.
fn order_by(self, ordering: Option<Vec<ColumnOrdering>>) -> Self;
/// Control MemWAL read routing for this query.
///
/// By default (unset), when the table carries a MemWAL write spec (see
/// [`crate::Table::set_lsm_write_spec`]), reads are routed through the LSM
/// scanner so they also return data written via the `merge_insert` LSM path
/// that has not yet been compacted into the base table (active/frozen
/// memtables and flushed generations); a table without a spec reads the base
/// table.
///
/// - `use_lsm(true)` forces LSM routing and errors if the table has no
/// MemWAL write spec.
/// - `use_lsm(false)` bypasses the MemWAL and reads the base table only,
/// even when a spec is present.
///
/// Note: the LSM scanner does not support every query shape (e.g. reranking,
/// hybrid search, `order_by`). On a MemWAL table those shapes error unless
/// `use_lsm(false)` is set, because a base-only read would silently
/// exclude un-compacted MemWAL data.
fn use_lsm(self, enable: bool) -> Self;
}
pub trait HasQuery {
@@ -613,11 +593,6 @@ impl<T: HasQuery> QueryBase for T {
self.mut_query().order_by = ordering;
self
}
fn use_lsm(mut self, enable: bool) -> Self {
self.mut_query().use_lsm = Some(enable);
self
}
}
/// Options for controlling the execution of a query
@@ -869,20 +844,6 @@ pub struct QueryRequest {
///
/// This allows ordering query results by one or more columns in either ascending or descending order.
pub order_by: Option<Vec<ColumnOrdering>>,
/// Controls MemWAL read routing. When unset (the default), a query against a
/// table that carries a MemWAL write spec (see
/// [`crate::Table::set_lsm_write_spec`]) is routed through the LSM scanner so
/// it also sees data written via the `merge_insert` LSM path that has not yet
/// been compacted into the base table — the active and frozen in-memory
/// memtables and the flushed (L0) generations, deduplicated by primary key
/// against the base table (newest generation wins); a table without a spec
/// reads the base table.
///
/// - `Some(true)` forces LSM routing and errors if the table has no MemWAL
/// write spec.
/// - `Some(false)` reads only the base table, bypassing the MemWAL.
pub use_lsm: Option<bool>,
}
impl Default for QueryRequest {
@@ -901,7 +862,6 @@ impl Default for QueryRequest {
norm: None,
disable_scoring_autoprojection: false,
order_by: None,
use_lsm: None,
}
}
}
+8 -116
View File
@@ -617,11 +617,6 @@ impl<S: HttpSend> RemoteTable<S> {
) -> Result<()> {
params.check_filter()?;
body["prefilter"] = params.prefilter.into();
// Only forward use_lsm when explicitly set; a server that predates it
// ignores the field and routes as it would by default.
if let Some(use_lsm) = params.use_lsm {
body["use_lsm"] = serde_json::Value::Bool(use_lsm);
}
if let Some(offset) = params.offset {
body["offset"] = serde_json::Value::Number(serde_json::Number::from(offset));
}
@@ -1631,21 +1626,9 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
let (request_id, response) = self.send(request, true).await?;
let response = self.check_table_response(&request_id, response).await?;
// Servers report the creation time either as an RFC 3339 `timestamp`
// (direct-table path) or as `timestamp_millis` in milliseconds since
// epoch (namespace-backed path), and may omit `metadata`.
#[derive(Deserialize)]
struct VersionEntry {
version: u64,
timestamp: Option<DateTime<Utc>>,
timestamp_millis: Option<i64>,
#[serde(default)]
metadata: std::collections::BTreeMap<String, String>,
}
#[derive(Deserialize)]
struct ListVersionsResponse {
versions: Vec<VersionEntry>,
versions: Vec<Version>,
}
let body = response.text().await.err_to_http(request_id.clone())?;
@@ -1656,37 +1639,11 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
err, body
)
.into(),
request_id: request_id.clone(),
request_id,
status_code: None,
})?;
body.versions
.into_iter()
.map(|entry| {
let timestamp = entry
.timestamp
.or_else(|| {
entry
.timestamp_millis
.and_then(DateTime::<Utc>::from_timestamp_millis)
})
.ok_or_else(|| Error::Http {
source: format!(
"list_versions response for version {} has neither a valid \
`timestamp` nor `timestamp_millis` field",
entry.version
)
.into(),
request_id: request_id.clone(),
status_code: None,
})?;
Ok(Version {
version: entry.version,
timestamp,
metadata: entry.metadata,
})
})
.collect()
Ok(body.versions)
}
async fn schema(&self) -> Result<SchemaRef> {
@@ -2886,10 +2843,6 @@ struct MergeInsertRequest {
// (the default is true)
#[serde(skip_serializing_if = "is_true")]
use_index: bool,
// Only serialize use_lsm when explicitly set (Some); a server that predates
// it ignores the field and routes as it would by default.
#[serde(skip_serializing_if = "Option::is_none")]
use_lsm: Option<bool>,
}
fn is_true(b: &bool) -> bool {
@@ -2941,7 +2894,6 @@ impl TryFrom<MergeInsertBuilder> for MergeInsertRequest {
when_not_matched_by_source_delete_filt,
// Only serialize use_index when it's false for backwards compatibility
use_index: value.use_index,
use_lsm: value.use_lsm,
})
}
}
@@ -4553,19 +4505,6 @@ mod tests {
},
Index::FTS(InvertedIndexParams::default().block_size(256).unwrap()),
),
(
"FTS",
{
let mut body = serde_json::to_value(InvertedIndexParams::default()).unwrap();
body["custom_stop_words"] = json!(["cat", " cat ", "CAT"]);
body
},
Index::FTS(InvertedIndexParams::default().custom_stop_words(Some(vec![
"cat".to_string(),
" cat ".to_string(),
"CAT".to_string(),
]))),
),
];
for (index_type, expected_body, index) in cases {
@@ -5097,9 +5036,8 @@ mod tests {
"max_token_length": 40,
"lower_case": true,
"stem": false,
"remove_stop_words": true,
"remove_stop_words": false,
"ascii_folding": true,
"custom_stop_words": ["hello"],
})
.to_string();
let table = Table::new_with_handler("my_table", move |request| {
@@ -5137,6 +5075,10 @@ mod tests {
assert_eq!(
tokens,
vec![
FtsToken {
text: "hello".to_string(),
position: 0,
},
FtsToken {
text: "こんにちは".to_string(),
position: 1,
@@ -5294,56 +5236,6 @@ mod tests {
// assert_eq!(versions, expected);
}
/// Namespace-backed servers report `timestamp_millis` instead of
/// `timestamp`, and may omit `metadata` entirely.
#[tokio::test]
async fn test_list_versions_timestamp_millis() {
let table = Table::new_with_handler("my_table", |request| {
assert_eq!(request.method(), "POST");
assert_eq!(request.url().path(), "/v1/table/my_table/version/list/");
let response_body = serde_json::json!({
"versions": [
{
"version": 1,
"manifest_path": "path/to/_versions/1.manifest",
"timestamp_millis": 1704067200000i64,
},
{
"version": 2,
"manifest_path": "path/to/_versions/2.manifest",
"timestamp_millis": 1706745600000i64,
"metadata": {"key": "value"},
},
]
});
let response_body = serde_json::to_string(&response_body).unwrap();
http::Response::builder()
.status(200)
.body(response_body)
.unwrap()
});
let versions = table.list_versions().await.unwrap();
assert_eq!(versions.len(), 2);
assert_eq!(versions[0].version, 1);
assert_eq!(
versions[0].timestamp,
"2024-01-01T00:00:00Z".parse::<DateTime<Utc>>().unwrap()
);
assert!(versions[0].metadata.is_empty());
assert_eq!(versions[1].version, 2);
assert_eq!(
versions[1].timestamp,
"2024-02-01T00:00:00Z".parse::<DateTime<Utc>>().unwrap()
);
assert_eq!(
versions[1].metadata.get("key").map(String::as_str),
Some("value")
);
}
#[tokio::test]
async fn test_index_stats() {
let table = Table::new_with_handler("my_table", |request| {
+9 -102
View File
@@ -21,7 +21,6 @@ use lance::dataset::WriteMode;
use lance::dataset::builder::DatasetBuilder;
use lance::dataset::{InsertBuilder, WriteParams};
use lance::index::DatasetIndexExt;
use lance::index::scalar::load_segment_params;
use lance::io::{ObjectStoreParams, WrappingObjectStore};
use lance_datafusion::utils::StreamingWriteSource;
use lance_index::IndexCriteria;
@@ -47,7 +46,6 @@ use std::sync::Arc;
use crate::connection::NamespaceClientPushdownOperation;
use crate::DistanceType;
use crate::blob::BlobRangeRequest;
use crate::data::scannable::{PeekedScannable, Scannable, estimate_write_partitions};
use crate::database::Database;
use crate::database::read_freshness::TableFreshness;
@@ -648,16 +646,6 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
message: "fetch_blobs is not supported on this table type".into(),
})
}
/// Materialize blob-local ranges. See [`Table::fetch_blob_ranges`].
async fn fetch_blob_ranges(
&self,
_column: &str,
_requests: &[BlobRangeRequest],
) -> Result<LargeBinaryArray> {
Err(Error::NotSupported {
message: "fetch_blob_ranges is not supported on this table type".into(),
})
}
/// Open lazy blob handles for the given row ids. See [`Table::fetch_blob_files`].
async fn fetch_blob_files(
&self,
@@ -1031,9 +1019,8 @@ impl Table {
/// Materialize blob bytes for the given row ids.
///
/// Output matches `row_ids` in length and order. Null blobs are null;
/// valid empty blobs contain empty byte strings. Prefer
/// [`Self::fetch_blob_files`] for large selections.
/// Output matches `row_ids` in length and order. Null and zero-length rows
/// are null. Prefer [`Self::fetch_blob_files`] for large selections.
///
/// ```
/// use arrow_array::UInt64Array;
@@ -1068,47 +1055,6 @@ impl Table {
self.inner.fetch_blobs(column.as_ref(), row_ids).await
}
/// Materialize row-specific ranges from a blob v2 column.
///
/// Each request contains a row id and a blob-local offset and length.
/// Requests may be duplicated or reordered, including multiple
/// ranges for the same blob. The output has the same length and order as
/// the requests. Null blobs produce null output slots; empty ranges on
/// non-null blobs produce empty byte strings.
///
/// ```
/// use lancedb::blob::BlobRangeRequest;
///
/// # use lancedb::Table;
/// # async fn read_ranges(table: &Table, row_id: u64) -> Result<(), Box<dyn std::error::Error>> {
/// let ranges = table
/// .fetch_blob_ranges(
/// "image",
/// [
/// BlobRangeRequest::new(row_id, 0, 1024),
/// BlobRangeRequest::new(row_id, 4096, 1024),
/// ],
/// )
/// .await?;
/// # let _ = ranges;
/// # Ok(())
/// # }
/// ```
///
/// Returns an error when a range is invalid, a requested row id does not
/// exist, or the column is not a blob v2 column. Returns
/// [`Error::NotSupported`] on table types without blob support.
pub async fn fetch_blob_ranges(
&self,
column: impl AsRef<str>,
requests: impl IntoIterator<Item = BlobRangeRequest>,
) -> Result<LargeBinaryArray> {
let requests = requests.into_iter().collect::<Vec<_>>();
self.inner
.fetch_blob_ranges(column.as_ref(), &requests)
.await
}
/// Open lazy [`BlobFile`] handles for the given row ids.
///
/// Same length and order as `row_ids`. Null rows are `None`. Bytes are not
@@ -3124,15 +3070,6 @@ impl BaseTable for NativeTable {
crate::blob::take_blobs_aligned(&dataset, column, row_ids).await
}
async fn fetch_blob_ranges(
&self,
column: &str,
requests: &[BlobRangeRequest],
) -> Result<LargeBinaryArray> {
let dataset = self.dataset.get().await?;
crate::blob::take_blob_ranges_aligned(&dataset, column, requests).await
}
async fn fetch_blob_files(
&self,
column: &str,
@@ -3194,9 +3131,10 @@ impl BaseTable for NativeTable {
async fn list_indices(&self) -> Result<Vec<IndexConfig>> {
let dataset = self.dataset.get().await?;
let total_rows = dataset.count_rows(None).await? as u64;
let descriptions = dataset.describe_indices(None).await?;
let mut indices: Vec<IndexConfig> = descriptions
.iter()
let indices = dataset
.describe_indices(None)
.await?
.into_iter()
.filter_map(|idx_desc| {
let index_type: crate::index::IndexType = idx_desc
.index_type()
@@ -3254,31 +3192,6 @@ impl BaseTable for NativeTable {
})
})
.collect();
for index in indices
.iter_mut()
.filter(|index| index.index_type == crate::index::IndexType::FTS)
{
let Some(description) = descriptions
.iter()
.find(|description| description.name() == index.name)
else {
continue;
};
let segments = description.segments();
let Some(segment) = segments.first() else {
continue;
};
let params = load_segment_params(&dataset, segment).await?;
let details = serde_json::to_string(&params).map_err(|source| Error::Other {
message: format!(
"Failed to serialize full text search configuration for index '{}'",
index.name
),
source: Some(Box::new(source)),
})?;
index.index_details = Some(details);
}
Ok(indices)
}
@@ -4136,10 +4049,10 @@ mod tests {
Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema))
}
// Windows does not support precise sleep durations due to timer resolution limitations.
#[cfg(not(target_os = "windows"))]
#[tokio::test]
async fn test_read_consistency_interval() {
use crate::utils::background_cache::clock;
let intervals = vec![
None,
Some(0),
@@ -4166,12 +4079,6 @@ mod tests {
let conn2 = conn2.execute().await.unwrap();
let table2 = conn2.open_table("my_table").execute().await.unwrap();
// Freeze the consistency clock now that `table2` has seeded its cache, so the
// interval only elapses when this test advances it. Otherwise the write and
// count_rows calls below race the real 100ms interval, which a loaded CI
// runner loses. Must come after open_table: creating the cache clears the mock.
clock::pin();
assert_eq!(table1.count_rows(None).await.unwrap(), 0);
assert_eq!(table2.count_rows(None).await.unwrap(), 0);
@@ -4189,7 +4096,7 @@ mod tests {
}
Some(100) => {
assert_eq!(table2.count_rows(None).await.unwrap(), 0);
clock::advance_by(Duration::from_millis(100));
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(table2.count_rows(None).await.unwrap(), 1);
}
_ => unreachable!(),
+1 -35
View File
@@ -1366,19 +1366,11 @@ mod tests {
table
.create_index(
&["text"],
Index::FTS(
FtsIndexBuilder::default()
.stem(false)
.custom_stop_words(Some(vec!["cat".to_string()]))
.block_size(256)
.unwrap(),
),
Index::FTS(FtsIndexBuilder::default().block_size(256).unwrap()),
)
.execute()
.await
.unwrap();
drop(table);
let table = conn.open_table("test_bitmap").execute().await.unwrap();
let index_configs = table.list_indices().await.unwrap();
assert_eq!(index_configs.len(), 1);
let index = index_configs.into_iter().next().unwrap();
@@ -1389,32 +1381,6 @@ mod tests {
let index_params: FtsIndexBuilder =
serde_json::from_str(index.index_details.as_deref().unwrap()).unwrap();
assert_eq!(index_params.posting_block_size(), 256);
assert_eq!(
serde_json::to_value(&index_params).unwrap()["custom_stop_words"],
serde_json::json!(["cat"])
);
assert_eq!(
table
.tokenize("cat dog", "text_idx")
.await
.unwrap()
.into_iter()
.map(|token| token.text)
.collect::<Vec<_>>(),
vec!["dog"]
);
let batches = table
.query()
.full_text_search(FullTextSearchQuery::new("cat dog".to_string()))
.limit(120)
.execute()
.await
.unwrap()
.try_collect::<Vec<_>>()
.await
.unwrap();
assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 40);
let num_rows = 120;
let stats = table.index_stats("text_idx").await.unwrap().unwrap();
+2 -61
View File
@@ -12,7 +12,7 @@ pub mod udtf;
use std::{collections::HashMap, sync::Arc};
use arrow_array::{RecordBatch, RecordBatchOptions};
use arrow_array::RecordBatch;
use arrow_schema::Schema as ArrowSchema;
use async_trait::async_trait;
use datafusion_catalog::{Session, TableProvider};
@@ -126,12 +126,7 @@ impl ExecutionPlan for MetadataEraserExec {
let stream = self.input.execute(partition, context)?;
let schema = self.schema.clone();
let stream = stream.map_ok(move |batch| {
RecordBatch::try_new_with_options(
schema.clone(),
batch.columns().to_vec(),
&RecordBatchOptions::new().with_row_count(Some(batch.num_rows())),
)
.unwrap()
RecordBatch::try_new(schema.clone(), batch.columns().to_vec()).unwrap()
});
Ok(
Box::pin(RecordBatchStreamAdapter::new(self.schema.clone(), stream))
@@ -549,60 +544,6 @@ pub mod tests {
}
}
/// A scan with an EMPTY projection (the shape a filtered `COUNT(*)` feeds in) yields
/// zero-column batches. `MetadataEraserExec::execute` rebuilt each batch with
/// `RecordBatch::try_new(schema, cols).unwrap()`; for a column-less batch that errors with
/// "must either specify a row count or at least one column" and the `.unwrap()` panics.
#[tokio::test]
async fn test_metadata_eraser_empty_projection_preserves_row_count() {
let fixture = TestFixture::new().await;
// Empty projection => zero output columns over N rows. Table "foo" has 10 rows.
let plan =
LogicalPlanBuilder::scan("foo", provider_as_source(fixture.adapter), Some(vec![]))
.unwrap()
.build()
.unwrap();
let mut stream = TestFixture::plan_to_stream(plan).await;
let mut rows = 0usize;
while let Some(batch) = stream.try_next().await.unwrap() {
assert_eq!(
batch.num_columns(),
0,
"empty projection must yield zero columns"
);
rows += batch.num_rows();
}
assert_eq!(
rows, 10,
"row count must survive MetadataEraserExec on a zero-column batch"
);
// End-to-end SQL regression for the previous panic.
let fixture = TestFixture::new().await;
let ctx = SessionContext::new();
ctx.register_table("foo", fixture.adapter.clone()).unwrap();
let batches = ctx
.sql("SELECT COUNT(*) FROM foo WHERE i < 5")
.await
.unwrap()
.collect()
.await
.unwrap();
let count = batches[0]
.column(0)
.as_any()
.downcast_ref::<Int64Array>()
.unwrap()
.value(0);
assert_eq!(count, 5, "COUNT(*) WHERE i < 5 over 0..10 must be 5");
}
#[tokio::test]
async fn test_filter_pushdown() {
let fixture = TestFixture::new().await;
+18 -427
View File
@@ -73,7 +73,7 @@ pub struct MergeInsertBuilder {
pub(crate) when_not_matched_by_source_delete_filt: Option<MergeFilter>,
pub(crate) timeout: Option<Duration>,
pub(crate) use_index: bool,
pub(crate) use_lsm: Option<bool>,
pub(crate) use_lsm_write: Option<bool>,
pub(crate) validate_single_shard: bool,
}
@@ -89,7 +89,7 @@ impl MergeInsertBuilder {
when_not_matched_by_source_delete_filt: None,
timeout: None,
use_index: true,
use_lsm: None,
use_lsm_write: None,
validate_single_shard: true,
}
}
@@ -187,17 +187,16 @@ impl MergeInsertBuilder {
self
}
/// Control MemWAL routing for this `merge_insert`.
/// Controls whether `merge_insert` uses the MemWAL LSM write path.
///
/// By default (unset), a `merge_insert` on a table with an
/// [`LsmWriteSpec`](super::LsmWriteSpec) installed is routed through Lance's
/// MemWAL shard writer; a table without one uses the standard path.
///
/// - `use_lsm(true)` forces MemWAL routing and errors if the table has no
/// LSM write spec.
/// - `use_lsm(false)` forces the standard write path even when a spec is set.
pub fn use_lsm(&mut self, enable: bool) -> &mut Self {
self.use_lsm = Some(enable);
/// [`LsmWriteSpec`](super::LsmWriteSpec) installed is routed through
/// Lance's MemWAL shard writer, and a table without one uses the standard
/// path. Calling this with `false` forces the standard path even when a
/// spec is set. Calling it with `true` requires a spec — `merge_insert`
/// errors if none is installed.
pub fn use_lsm_write(&mut self, use_lsm_write: bool) -> &mut Self {
self.use_lsm_write = Some(use_lsm_write);
self
}
@@ -627,7 +626,7 @@ mod lsm_tests {
}
#[tokio::test]
async fn lsm_merge_insert_use_lsm_false_falls_back() {
async fn lsm_merge_insert_use_lsm_write_false_falls_back() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await;
table
@@ -635,10 +634,9 @@ mod lsm_tests {
.await
.unwrap();
// use_lsm(false) opts out: the standard path runs and commits even though
// a spec is installed.
// use_lsm_write(false) opts out: the standard path runs and commits.
let mut builder = table.merge_insert(&["id"]);
builder.when_not_matched_insert_all().use_lsm(false);
builder.when_not_matched_insert_all().use_lsm_write(false);
let result = builder
.execute(id_value_reader(vec![3, 4, 5]))
.await
@@ -648,25 +646,6 @@ mod lsm_tests {
assert_eq!(table.count_rows(None).await.unwrap(), 5);
}
#[tokio::test]
async fn lsm_merge_insert_use_lsm_true_without_spec_errors() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await;
// use_lsm(true) demands MemWAL routing; without a write spec it errors
// rather than silently falling back to the standard path.
let mut builder = table.merge_insert(&["id"]);
builder
.when_matched_update_all(None)
.when_not_matched_insert_all()
.use_lsm(true);
let err = builder
.execute(id_value_reader(vec![3, 4, 5]))
.await
.unwrap_err();
assert!(matches!(err, Error::InvalidInput { .. }), "got {err:?}");
}
#[tokio::test]
async fn lsm_merge_insert_rejects_on_not_primary_key() {
let dir = tempdir().unwrap();
@@ -775,23 +754,18 @@ mod lsm_tests {
}
#[tokio::test]
async fn lsm_merge_insert_no_spec_uses_standard_path() {
async fn lsm_merge_insert_use_lsm_write_true_requires_spec() {
let dir = tempdir().unwrap();
// id_value_table sets a primary key but no LSM write spec.
let table = id_value_table(&dir).await;
// Without a spec, a default merge_insert (use_lsm unset) simply uses
// the standard path and commits — no opt-out required, no error.
let mut builder = table.merge_insert(&["id"]);
builder
.when_matched_update_all(None)
.when_not_matched_insert_all();
let result = builder
.execute(id_value_reader(vec![3, 4, 5]))
.await
.unwrap();
assert_eq!(result.num_inserted_rows, 2);
assert_eq!(table.count_rows(None).await.unwrap(), 5);
.when_not_matched_insert_all()
.use_lsm_write(true);
let err = builder.execute(id_value_reader(vec![4])).await.unwrap_err();
assert!(matches!(err, Error::InvalidInput { .. }), "got {err:?}");
}
#[tokio::test]
@@ -844,387 +818,4 @@ mod lsm_tests {
.await
.unwrap();
}
// ---------------------------------------------------------------------
// LSM read path
// ---------------------------------------------------------------------
use crate::arrow::SendableRecordBatchStream;
use crate::query::{ExecutableQuery, QueryBase};
use arrow::array::AsArray;
use arrow::datatypes::Int64Type;
use futures::TryStreamExt;
/// Collect `(id, value)` pairs from a result stream, sorted by id.
async fn collect_id_value(stream: SendableRecordBatchStream) -> Vec<(i64, i64)> {
let batches: Vec<_> = stream.try_collect().await.unwrap();
let mut rows = Vec::new();
for batch in &batches {
let ids = batch
.column_by_name("id")
.unwrap()
.as_primitive::<Int64Type>();
let values = batch
.column_by_name("value")
.unwrap()
.as_primitive::<Int64Type>();
for i in 0..batch.num_rows() {
rows.push((ids.value(i), values.value(i)));
}
}
rows.sort();
rows
}
/// Upsert `ids` (value = 0..n) through the LSM `merge_insert` path.
async fn lsm_upsert(table: &Table, ids: Vec<i64>) {
let mut builder = table.merge_insert(&[]);
builder
.when_matched_update_all(None)
.when_not_matched_insert_all();
builder.execute(id_value_reader(ids)).await.unwrap();
}
#[tokio::test]
async fn lsm_read_sees_active_memtable() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await; // base: ids 1,2,3 (value 0,1,2)
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
// Insert ids 4,5 into the active memtable (not committed to base).
lsm_upsert(&table, vec![4, 5]).await;
// Default read auto-routes through the LSM scanner: base active memtable.
let lsm = table.query().execute().await.unwrap();
let rows = collect_id_value(lsm).await;
assert_eq!(
rows.iter().map(|(id, _)| *id).collect::<Vec<_>>(),
vec![1, 2, 3, 4, 5]
);
// use_lsm(false) bypasses the MemWAL and reads the base table only.
let base_only = table.query().use_lsm(false).execute().await.unwrap();
let rows = collect_id_value(base_only).await;
assert_eq!(
rows.iter().map(|(id, _)| *id).collect::<Vec<_>>(),
vec![1, 2, 3]
);
}
#[tokio::test]
async fn lsm_read_dedup_newest_wins() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await; // base: id 2 -> value 1
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
// Upsert ids 2,3,4 with values 0,1,2. id 2 and 3 shadow the base rows.
lsm_upsert(&table, vec![2, 3, 4]).await;
let lsm = table.query().execute().await.unwrap();
let rows = collect_id_value(lsm).await;
// id 1 from base (value 0); ids 2,3,4 from memtable (values 0,1,2).
assert_eq!(rows, vec![(1, 0), (2, 0), (3, 1), (4, 2)]);
}
#[tokio::test]
async fn lsm_read_point_lookup_filter() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await;
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
lsm_upsert(&table, vec![2, 3, 4]).await; // id 2 -> value 0 (shadows base)
let lsm = table.query().only_if("id = 2").execute().await.unwrap();
let rows = collect_id_value(lsm).await;
assert_eq!(rows, vec![(2, 0)]);
}
#[tokio::test]
async fn lsm_read_multi_shard() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await;
table
.set_lsm_write_spec(LsmWriteSpec::bucket("id", 8))
.await
.unwrap();
// Two single-row upserts that route to (likely) different buckets; each
// closes the writer so the next opens a fresh shard.
lsm_upsert(&table, vec![10]).await;
table.close_lsm_writers().await.unwrap();
lsm_upsert(&table, vec![11]).await;
let lsm = table.query().execute().await.unwrap();
let rows = collect_id_value(lsm).await;
let ids: Vec<i64> = rows.iter().map(|(id, _)| *id).collect();
// Base 1,2,3 + flushed/active shards for 10 and 11.
assert_eq!(ids, vec![1, 2, 3, 10, 11]);
}
#[tokio::test]
async fn lsm_read_after_close_sees_flushed() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await;
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
lsm_upsert(&table, vec![4, 5]).await;
// close flushes the active memtable to an on-disk generation and drops
// the cached writer; the read must still see those rows via the shard
// manifest snapshot.
table.close_lsm_writers().await.unwrap();
let lsm = table.query().execute().await.unwrap();
let ids: Vec<i64> = collect_id_value(lsm)
.await
.iter()
.map(|(id, _)| *id)
.collect();
assert_eq!(ids, vec![1, 2, 3, 4, 5]);
}
#[tokio::test]
async fn lsm_read_without_spec_reads_base() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await; // no LSM write spec
// With no spec installed there is nothing to route: the default read and
// an explicit use_lsm(false) both read the base table without error.
for query in [table.query(), table.query().use_lsm(false)] {
let rows = collect_id_value(query.execute().await.unwrap()).await;
assert_eq!(
rows.iter().map(|(id, _)| *id).collect::<Vec<_>>(),
vec![1, 2, 3]
);
}
}
#[tokio::test]
async fn lsm_read_unsupported_shape_errors_without_use_lsm_false() {
let dir = tempdir().unwrap();
let table = id_value_table(&dir).await;
table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap();
lsm_upsert(&table, vec![4]).await;
// `with_row_id` is a shape the LSM scanner cannot honor. On a MemWAL
// table the default (auto-routed) read hard-errors rather than silently
// reading a stale base-only result that would exclude un-compacted row 4.
let err = table
.query()
.with_row_id()
.execute()
.await
.err()
.expect("unsupported shape on a MemWAL table must error");
assert!(matches!(err, Error::NotSupported { .. }), "got {err:?}");
// use_lsm(false) is the escape hatch: it reads the base table only.
let rows = collect_id_value(
table
.query()
.with_row_id()
.use_lsm(false)
.execute()
.await
.unwrap(),
)
.await;
assert_eq!(
rows.iter().map(|(id, _)| *id).collect::<Vec<_>>(),
vec![1, 2, 3]
);
}
/// A reader of `[id: Int64, text: Utf8]` rows.
fn id_text_reader(rows: Vec<(i64, &str)>) -> Box<dyn RecordBatchReader + Send> {
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new("text", DataType::Utf8, false),
]));
let ids: Vec<i64> = rows.iter().map(|(id, _)| *id).collect();
let texts: Vec<&str> = rows.iter().map(|(_, t)| *t).collect();
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(Int64Array::from(ids)),
Arc::new(StringArray::from(texts)),
],
)
.unwrap();
Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema))
}
#[tokio::test]
async fn lsm_read_full_text_search() {
use crate::index::Index;
use lance_index::scalar::FullTextSearchQuery;
let dir = tempdir().unwrap();
let conn = connect(dir.path().to_str().unwrap())
.execute()
.await
.unwrap();
let table = conn
.create_table(
"t",
id_text_reader(vec![(1, "alpha"), (2, "beta"), (3, "gamma")]),
)
.execute()
.await
.unwrap();
table.set_unenforced_primary_key(["id"]).await.unwrap();
table
.create_index(&["text"], Index::FTS(Default::default()))
.execute()
.await
.unwrap();
let fts_index = table.list_indices().await.unwrap()[0].name.clone();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes([fts_index]))
.await
.unwrap();
// Insert a row whose term ("zebra") exists in no base row.
let mut builder = table.merge_insert(&[]);
builder
.when_matched_update_all(None)
.when_not_matched_insert_all();
builder
.execute(id_text_reader(vec![(99, "zebra")]))
.await
.unwrap();
let search = |term: &str| {
let q = FullTextSearchQuery::new(term.to_string())
.with_column("text".to_string())
.unwrap();
table.query().full_text_search(q)
};
// "zebra" lives only in the active memtable; LSM read finds it.
let stream = search("zebra").execute().await.unwrap();
let batches: Vec<_> = stream.try_collect().await.unwrap();
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 1, "LSM FTS must surface the memtable row");
// A base-only term still matches the base table through the LSM scan.
let stream = search("alpha").execute().await.unwrap();
let batches: Vec<_> = stream.try_collect().await.unwrap();
let rows: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(rows, 1, "LSM FTS must still see base rows");
}
#[tokio::test]
async fn lsm_read_vector_search() {
use crate::index::Index;
use crate::index::vector::IvfPqIndexBuilder;
use arrow::array::{FixedSizeListBuilder, Float32Builder};
use arrow::datatypes::Int64Type;
const DIM: i32 = 8;
const N: i64 = 256;
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new(
"vec",
DataType::FixedSizeList(Arc::new(Field::new("item", DataType::Float32, true)), DIM),
false,
),
]));
let make_batch = |rows: Vec<(i64, f32)>| -> RecordBatch {
let ids: Vec<i64> = rows.iter().map(|(id, _)| *id).collect();
let mut vb = FixedSizeListBuilder::new(Float32Builder::new(), DIM);
for (_, fill) in &rows {
for _ in 0..DIM {
vb.values().append_value(*fill);
}
vb.append(true);
}
RecordBatch::try_new(
schema.clone(),
vec![Arc::new(Int64Array::from(ids)), Arc::new(vb.finish())],
)
.unwrap()
};
let dir = tempdir().unwrap();
let conn = connect(dir.path().to_str().unwrap())
.execute()
.await
.unwrap();
// Base rows fill each vector with its own id (0..256); all far from 1000.
let base = make_batch((0..N).map(|i| (i, i as f32)).collect());
let base_reader: Box<dyn RecordBatchReader + Send> =
Box::new(RecordBatchIterator::new(vec![Ok(base)], schema.clone()));
let table = conn.create_table("t", base_reader).execute().await.unwrap();
table.set_unenforced_primary_key(["id"]).await.unwrap();
table
.create_index(
&["vec"],
Index::IvfPq(
IvfPqIndexBuilder::default()
.num_partitions(1)
.num_sub_vectors(2),
),
)
.execute()
.await
.unwrap();
let vec_index = table.list_indices().await.unwrap()[0].name.clone();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes([vec_index]))
.await
.unwrap();
// Insert a vector (filled with 1000) that is nearest to the query.
let mut builder = table.merge_insert(&[]);
builder
.when_matched_update_all(None)
.when_not_matched_insert_all();
let insert_reader: Box<dyn RecordBatchReader + Send> = Box::new(RecordBatchIterator::new(
vec![Ok(make_batch(vec![(9999, 1000.0)]))],
schema.clone(),
));
builder.execute(insert_reader).await.unwrap();
// KNN near [1000; DIM]: the default (auto-routed) read surfaces the
// memtable row.
let stream = table
.query()
.nearest_to(&[1000.0_f32; 8])
.unwrap()
.limit(1)
.execute()
.await
.unwrap();
let batches: Vec<_> = stream.try_collect().await.unwrap();
let ids: Vec<i64> = batches
.iter()
.flat_map(|b| {
b.column_by_name("id")
.unwrap()
.as_primitive::<Int64Type>()
.values()
.to_vec()
})
.collect();
assert_eq!(
ids,
vec![9999],
"LSM vector search must rank the memtable row first"
);
}
}
+8 -72
View File
@@ -306,39 +306,6 @@ impl ShardWriterEntry {
}
Ok(())
}
/// The cached writer's latest in-memory manifest (current generation +
/// flushed generations). `Ok(None)` if the writer was already closed.
/// Used by the LSM read path to snapshot this shard authoritatively
/// without re-reading the on-disk manifest.
async fn manifest(&self) -> Result<Option<lance_index::mem_wal::ShardManifest>> {
let guard = self.inner.read().await;
let Some(writer) = guard.as_ref() else {
return Ok(None);
};
writer.manifest().await.map_err(|e| Error::Runtime {
message: format!("read: shard writer manifest read failed: {}", e),
})
}
/// Atomically capture the cached writer's active + frozen-awaiting-flush
/// memtables for unified LSM scanning. `Ok(None)` if the writer was
/// already closed.
async fn in_memory_memtable_refs(
&self,
) -> Result<Option<lance::dataset::mem_wal::scanner::InMemoryMemTables>> {
let guard = self.inner.read().await;
let Some(writer) = guard.as_ref() else {
return Ok(None);
};
writer
.in_memory_memtable_refs()
.await
.map(Some)
.map_err(|e| Error::Runtime {
message: format!("read: shard writer memtable capture failed: {}", e),
})
}
}
impl ShardWriterCache {
@@ -378,36 +345,6 @@ impl ShardWriterCache {
Ok(entry)
}
/// Snapshot the cached writer's shard for the LSM read path: its shard id,
/// authoritative in-memory manifest, and active + frozen memtable refs.
/// Returns `None` when no writer is currently cached (e.g. nothing has been
/// written this session, or the writer was closed).
#[allow(clippy::redundant_pub_crate)]
pub(crate) async fn read_snapshot(
&self,
) -> Result<
Option<(
Uuid,
Option<lance_index::mem_wal::ShardManifest>,
Option<lance::dataset::mem_wal::scanner::InMemoryMemTables>,
)>,
> {
let cached = {
let guard = self.slot.read().await;
guard.as_ref().map(|(id, entry)| (*id, entry.clone()))
};
let Some((shard_id, entry)) = cached else {
return Ok(None);
};
// Capture memtables before the manifest. If a flush interleaves, dedup
// tolerates the same rows appearing in both a memtable and a freshly
// flushed generation, but would drop rows present in neither. Manifest
// last guarantees any generation flushed mid-capture is still covered.
let memtables = entry.in_memory_memtable_refs().await?;
let manifest = entry.manifest().await?;
Ok(Some((shard_id, manifest, memtables)))
}
/// Close the cached writer, if any, and clear the slot.
#[allow(clippy::redundant_pub_crate)]
pub(crate) async fn drain_and_close(&self) -> Result<()> {
@@ -471,20 +408,19 @@ pub(crate) async fn lsm_dispatch_decision(
table: &NativeTable,
params: &MergeInsertBuilder,
) -> Result<LsmDispatch> {
// Explicit opt-out: use the standard path regardless of any installed spec.
if params.use_lsm == Some(false) {
// `Some(false)` is an explicit opt-out: use the standard path.
if params.use_lsm_write == Some(false) {
return Ok(LsmDispatch::Standard);
}
let dataset = table.dataset.get().await?;
let Some(details) = dataset.mem_wal_index_details().await? else {
// No write spec installed. `use_lsm(true)` demanded MemWAL routing, so
// that is an error; otherwise fall back to the standard path.
if params.use_lsm == Some(true) {
// No LSM write spec installed. `Some(true)` explicitly asked for the
// LSM path, which is meaningless without a spec; `None` (the default)
// just falls back to the standard path.
if params.use_lsm_write == Some(true) {
return Err(Error::InvalidInput {
message: "use_lsm(true) was set but the table has no MemWAL write spec; \
install one with set_lsm_write_spec or leave use_lsm unset"
.to_string(),
message: "merge_insert: use_lsm_write(true) requires an LSM write spec on the table; call set_lsm_write_spec first".to_string(),
});
}
return Ok(LsmDispatch::Standard);
@@ -513,7 +449,7 @@ pub(crate) async fn lsm_dispatch_decision(
if !is_upsert_only(params) {
return Err(Error::InvalidInput {
message: "merge_insert: when an LSM write spec is set, only the upsert form (when_matched_update_all without a filter + when_not_matched_insert_all, no by-source delete) is supported; call use_lsm(false) to use the standard merge_insert path".to_string(),
message: "merge_insert: when an LSM write spec is set, only the upsert form (when_matched_update_all without a filter + when_not_matched_insert_all, no by-source delete) is supported; call use_lsm_write(false) to use the standard merge_insert path".to_string(),
});
}
+6 -108
View File
@@ -3,8 +3,6 @@
use std::sync::Arc;
mod lsm;
use super::NativeTable;
use crate::connection::NamespaceClientPushdownOperation;
use crate::error::{Error, Result};
@@ -22,7 +20,6 @@ use datafusion_physical_plan::projection::ProjectionExec;
use datafusion_physical_plan::repartition::RepartitionExec;
use datafusion_physical_plan::union::UnionExec;
use futures::future::try_join_all;
use lance::dataset::mem_wal::DatasetMemWalExt;
use lance::dataset::scanner::DatasetRecordBatchStream;
use lance::dataset::scanner::Scanner;
use lance_datafusion::exec::{analyze_plan as lance_analyze_plan, execute_plan};
@@ -56,7 +53,7 @@ pub async fn execute_query(
// QueryTable pushdown runs the query server-side, but only on the main
// branch: the namespace request carries no branch yet, so a branch handle
// must fall through to local execution.
if can_execute_namespace_query(table, query).await?
if can_execute_namespace_query(table, query)
&& let Some(ref namespace_client) = table.namespace_client
{
return execute_namespace_query(table, namespace_client.clone(), query, options).await;
@@ -64,35 +61,18 @@ pub async fn execute_query(
execute_generic_query(table, query, options).await
}
async fn can_execute_namespace_query(table: &NativeTable, query: &AnyQuery) -> Result<bool> {
if !(table
fn can_execute_namespace_query(table: &NativeTable, query: &AnyQuery) -> bool {
table
.pushdown_operations
.contains(&NamespaceClientPushdownOperation::QueryTable)
&& table.namespace_client.is_some()
&& table.dataset.current_branch().is_none()
&& !requires_local_namespace_execution(query))
{
return Ok(false);
}
// A MemWAL write spec means reads auto-route through the LSM scanner in
// `create_plan` even when `use_lsm` is unset. The namespace request has no
// use_lsm field, so pushing the default query down would silently omit
// un-compacted rows — force local execution whenever a spec is installed.
let dataset = table.dataset.get().await?;
if dataset.mem_wal_index_details().await?.is_some() {
return Ok(false);
}
Ok(true)
&& !requires_local_namespace_execution(query)
}
fn requires_local_namespace_execution(query: &AnyQuery) -> bool {
// The namespace QueryTable request has no approx_mode or use_lsm field yet, so
// pushing these down would silently ignore the user's setting. For use_lsm that
// is worse than a tuning miss: MemWAL read routing lives only in `create_plan`,
// so a pushed-down query would return stale base-only data with no error.
if query.base().use_lsm.is_some() {
return true;
}
// The namespace QueryTable request has no approx_mode field yet, so
// pushing this query down would silently ignore the user's setting.
matches!(
query,
AnyQuery::VectorQuery(VectorQueryRequest {
@@ -140,29 +120,6 @@ pub async fn create_plan(
query.base.check_filter()?;
let ds_ref = table.dataset.get().await?;
// MemWAL read routing driven by `use_lsm`:
// * unset — route through the LSM scanner iff the table carries a write spec
// * Some(true) — force LSM routing; error if the table has no write spec
// * Some(false) — read the base table only, bypassing the MemWAL
// The LSM scanner surfaces in-flight `merge_insert` data (active/frozen
// memtables + flushed generations); validation and dispatch live in `lsm`.
let has_spec = ds_ref.mem_wal_index_details().await?.is_some();
let use_lsm = match query.base.use_lsm {
Some(true) if !has_spec => {
return Err(Error::InvalidInput {
message: "use_lsm(true) was set but the table has no MemWAL write spec; \
install one with set_lsm_write_spec or leave use_lsm unset"
.to_string(),
});
}
Some(enable) => enable,
None => has_spec,
};
if use_lsm {
return lsm::create_lsm_plan(table, ds_ref, query).await;
}
let schema = ds_ref.schema();
let mut column = query.column.clone();
@@ -947,65 +904,6 @@ mod tests {
assert_eq!(namespace_client.query_table_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn test_execute_query_use_lsm_with_namespace_pushdown_runs_locally() {
use crate::connect;
use crate::table::query::execute_query;
use arrow_array::{Int32Array, RecordBatch};
use arrow_schema::{DataType, Field, Schema};
let conn = connect("memory://").execute().await.unwrap();
let vectors = Arc::new(fixed_size_list_array(
vec![0.0, 0.0, 10.0, 10.0, 20.0, 20.0],
2,
));
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("vector", vectors.data_type().clone(), false),
]));
let batch = RecordBatch::try_new(
schema,
vec![Arc::new(Int32Array::from(vec![1, 2, 3])), vectors],
)
.unwrap();
let table = conn
.create_table("test_use_lsm_namespace_fallback", batch)
.execute()
.await
.unwrap();
let namespace_client = Arc::new(CountingNamespaceClient::default());
let mut native_table = table.as_native().unwrap().clone();
native_table.namespace_client = Some(namespace_client.clone());
native_table
.pushdown_operations
.insert(NamespaceClientPushdownOperation::QueryTable);
// `use_lsm` set (even to false) must force local execution — the namespace
// request has no use_lsm field, so a pushdown would silently ignore it.
let query_vector = Arc::new(Float32Array::from(vec![0.0, 0.0]));
let query = AnyQuery::VectorQuery(VectorQueryRequest {
base: QueryRequest {
limit: Some(1),
use_lsm: Some(false),
..Default::default()
},
column: Some("vector".to_string()),
query_vector: vec![query_vector as ArrayRef],
..Default::default()
});
let stream = execute_query(&native_table, &query, QueryExecutionOptions::default())
.await
.unwrap();
let batches = stream.try_collect::<Vec<_>>().await.unwrap();
let count: usize = batches.iter().map(|b| b.num_rows()).sum();
assert_eq!(count, 1);
assert_eq!(namespace_client.query_table_calls.load(Ordering::SeqCst), 0);
}
#[tokio::test]
async fn test_create_plan_multivector_structure() {
use arrow_array::{Float32Array, RecordBatch};
-786
View File
@@ -1,786 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! MemWAL LSM read path.
//!
//! When a table has an LSM write spec installed (see [`set_lsm_write_spec`]),
//! reads are routed through Lance's [`LsmScanner`] instead of the plain
//! base-table scan unless the query sets
//! [`use_lsm(false)`](crate::query::QueryBase::use_lsm). This makes data
//! written via the LSM `merge_insert` path — which lives in the active/frozen
//! in-memory memtables and the flushed SSTable generations until an external
//! compaction merges it into the base table — visible to queries, deduplicated by
//! primary key (newest generation wins).
//!
//! Three query shapes are supported, mirroring the standard scan: a plain scan
//! (filter / projection / limit), full-text search, and vector (ANN) search. All
//! three run through a single [`LsmScanner`], so a `where` filter is honored as a
//! prefilter uniformly — including for vector search, where `LsmScanner` threads
//! it into the vector planner's prefilter. Shapes the LSM path cannot honor are
//! rejected with [`Error::NotSupported`]; the caller must set `use_lsm(false)` to
//! run those against the base table.
//!
//! [`set_lsm_write_spec`]: crate::Table::set_lsm_write_spec
use std::collections::HashMap;
use std::sync::Arc;
use arrow_array::Array;
use arrow_schema::{DataType, Schema as ArrowSchema};
use datafusion_physical_plan::expressions::Column;
use datafusion_physical_plan::projection::ProjectionExec;
use datafusion_physical_plan::{ExecutionPlan, PhysicalExpr};
use lance::Dataset;
use lance::dataset::mem_wal::scanner::InMemoryMemTables;
use lance::dataset::mem_wal::{
DatasetMemWalExt, LsmScanner, ShardManifestStore, ShardSnapshot, ShardWriterConfig,
};
use lance_index::mem_wal::{MemWalIndexDetails, ShardManifest};
use uuid::Uuid;
use super::NativeTable;
use crate::DistanceType;
use crate::error::{Error, Result};
use crate::query::{DEFAULT_TOP_K, QueryFilter, Select, VectorQueryRequest};
use crate::utils::default_vector_column;
/// Over-fetch factor for the LSM vector/FTS arms. With the default of `1.0` a
/// source blocked by cross-generation PK dedup fetches exactly `k` and can return
/// fewer than `k` live rows; upserts routinely create such blocked candidates, so
/// widen the per-source fetch to keep result pages filled.
const LSM_OVERFETCH_FACTOR: f64 = 2.0;
/// Build the LSM read plan for a MemWAL-routed query.
///
/// The caller guarantees `ds_ref` carries a MemWAL write spec (routing is decided
/// in [`create_plan`](super::create_plan)). Errors with [`Error::NotSupported`]
/// for query shapes the LSM scanner cannot honor — the caller must set
/// `use_lsm(false)` to run those against the base table.
pub(super) async fn create_lsm_plan(
table: &NativeTable,
ds_ref: Arc<Dataset>,
query: VectorQueryRequest,
) -> Result<Arc<dyn ExecutionPlan>> {
reject_unsupported(&query)?;
// A time-traveled (checked-out) handle pins an older dataset version, but the
// WAL manifests and cached writer expose current live state — mixing them would
// surface WAL rows written after the requested version. `use_lsm(false)` reads
// the base table at the pinned version.
if table.dataset.time_travel_version().is_some() {
return Err(Error::NotSupported {
message: "the MemWAL LSM scanner cannot read from a time-traveled dataset version; set use_lsm(false) to read the base table at this version".to_string(),
});
}
// Routing guarantees a write spec is installed (see `create_plan`).
let details = ds_ref
.mem_wal_index_details()
.await?
.ok_or_else(|| Error::Runtime {
message: "the MemWAL LSM write spec disappeared during read planning".to_string(),
})?;
let pk_columns = pk_columns(&ds_ref)?;
// The base index an indexed arm relies on may lag compaction; resolve it so the
// snapshot retains SSTables the index has not yet caught up to.
let arm_index = arm_maintained_index_name(&ds_ref, &query, &details).await?;
let (snapshots, in_memory) =
build_read_context(table, &ds_ref, &details, arm_index.as_deref()).await?;
let limit = query.base.limit;
let offset = query.base.offset;
let plan = if !query.query_vector.is_empty() {
vector_plan(
&ds_ref,
&query,
&details,
pk_columns.clone(),
snapshots,
in_memory,
limit,
offset,
)
.await?
} else if let Some(fts) = &query.base.full_text_search {
fts_plan(
&ds_ref,
fts.clone(),
&query,
&details,
pk_columns.clone(),
snapshots,
in_memory,
limit,
offset,
)
.await?
} else {
plain_plan(
&ds_ref,
&query,
pk_columns.clone(),
snapshots,
in_memory,
limit,
offset,
)
.await?
};
// Lance appends the primary-key columns internally for dedup and keeps them in
// the output; drop the ones the user did not request so the projection matches.
restore_projection(plan, &query, &pk_columns)
}
/// Reject query shapes the LSM read path does not implement. On a MemWAL table
/// reads route through the LSM scanner by default, so an unsupported shape is a
/// hard error rather than a silent fallback to the base-only scan — which would
/// exclude un-compacted MemWAL data. The caller must set `use_lsm(false)` to
/// run these against the base table, accepting that the results omit un-compacted
/// MemWAL data.
///
/// A `where` filter is intentionally *not* rejected: every arm routes through
/// [`LsmScanner`], which applies it as a prefilter (see [`base_scanner`]).
fn reject_unsupported(query: &VectorQueryRequest) -> Result<()> {
let unsupported = |what: &str| {
Err(Error::NotSupported {
message: format!(
"the MemWAL LSM scanner does not support {what}; set use_lsm(false) to read the base table only (results will exclude un-compacted MemWAL data)"
),
})
};
if query.query_vector.len() > 1 {
return unsupported("multiple query vectors");
}
if !query.query_vector.is_empty() && query.base.full_text_search.is_some() {
return unsupported("hybrid (vector + full-text) search");
}
if query.base.with_row_id {
return unsupported("with_row_id (the LSM scanner exposes _rowaddr, not a stable _rowid)");
}
if query.base.reranker.is_some() {
return unsupported("reranking / hybrid search");
}
if query.base.order_by.is_some() {
return unsupported("order_by");
}
// Vector-only knobs the LSM scanner cannot honor. Both change results rather
// than just recall, so error instead of silently ignoring them: distance_range
// would return rows outside the bound, and use_index(false) asks for a
// brute-force search the index-only base arm can't do. (ef / approx_mode /
// maximum_nprobes are recall/speed knobs and are left to no-op — and
// maximum_nprobes defaults to Some, so it cannot be rejected on presence.)
if !query.query_vector.is_empty() {
if query.lower_bound.is_some() || query.upper_bound.is_some() {
return unsupported("distance_range on vector search");
}
if !query.use_index {
return unsupported(
"use_index(false) / brute-force vector search (the LSM base arm is index-only)",
);
}
}
// Postfilter changes result semantics for both vector and full-text search, and
// the LSM scanner always prefilters — reject a requested postfilter for either.
if (!query.query_vector.is_empty() || query.base.full_text_search.is_some())
&& !query.base.prefilter
{
return unsupported(
"postfilter on vector or full-text search (the LSM scanner always prefilters)",
);
}
match &query.base.select {
Select::All | Select::Columns(_) => {}
Select::Dynamic(_) | Select::Expr(_) => return unsupported("dynamic column projection"),
}
if let Some(QueryFilter::Substrait(_)) = &query.base.filter {
return unsupported("Substrait filters");
}
// Take-by-row-id / row-offset queries carry a `_rowid` / `_rowoffset` filter,
// columns the LSM scanner never exposes (only `_rowaddr`); reject with guidance
// rather than failing deep in datafusion with a column-not-found error.
if let Some(QueryFilter::Datafusion(expr)) = &query.base.filter
&& expr
.column_refs()
.iter()
.any(|c| c.name == "_rowid" || c.name == "_rowoffset")
{
return unsupported(
"take by row id or row offset (the LSM scanner has no stable _rowid / _rowoffset)",
);
}
Ok(())
}
/// Primary-key column names from the dataset's unenforced primary key.
fn pk_columns(dataset: &Dataset) -> Result<Vec<String>> {
let pk: Vec<String> = dataset
.schema()
.unenforced_primary_key()
.iter()
.map(|f| f.name.clone())
.collect();
if pk.is_empty() {
return Err(Error::InvalidInput {
message:
"the MemWAL LSM scanner requires an unenforced primary key, but the table has none"
.to_string(),
});
}
Ok(pk)
}
/// Per-shard SSTable exclusion watermark: the generation at or below which SSTables
/// are safe to drop for this arm. A generation is droppable only once it is
/// compacted into the base table AND covered by `index_name`'s catch-up (for an
/// indexed arm); a plain scan (`index_name == None`) uses the compaction watermark
/// alone. Capping at the index catch-up keeps rows the base index has not yet
/// indexed visible through their SSTable. First occurrence per shard mirrors Lance's
/// `compacted_generation_for_shard`.
fn exclusion_watermarks(
details: &MemWalIndexDetails,
index_name: Option<&str>,
) -> HashMap<Uuid, u64> {
let mut exclude: HashMap<Uuid, u64> = HashMap::new();
for entry in &details.compacted_sstables {
let mut watermark = entry.generation;
if let Some(name) = index_name
&& let Some(caught_up) = details
.index_catchup
.iter()
.find(|icp| icp.index_name == name)
.and_then(|icp| icp.caught_up_generation_for_shard(&entry.shard_id))
{
watermark = watermark.min(caught_up);
}
exclude.entry(entry.shard_id).or_insert(watermark);
}
exclude
}
/// Assemble the per-shard snapshots (flushed SSTable generations) and the
/// in-memory memtables (active + frozen) for the table.
///
/// Snapshots for all shards come from their on-disk manifests; for the shard
/// with a live cached `ShardWriter` (this session's in-flight writes) the
/// writer's authoritative in-memory manifest and memtables override the
/// on-disk view so a read sees data not yet flushed.
async fn build_read_context(
table: &NativeTable,
dataset: &Dataset,
details: &MemWalIndexDetails,
index_name: Option<&str>,
) -> Result<(Vec<ShardSnapshot>, HashMap<Uuid, InMemoryMemTables>)> {
let exclude = exclusion_watermarks(details, index_name);
let shard_ids = dataset.list_mem_wal_latest_shard_ids().await?;
// Use the dataset's own object store (not `ObjectStore::from_uri`, which
// builds a fresh registry and would miss `memory://` and custom-registered
// stores). The base path matches `list_mem_wal_latest_shard_ids`.
let store = dataset.object_store(None).await?;
let base_path = dataset.branch_location().path;
let scan_batch_size = ShardWriterConfig::default().manifest_scan_batch_size;
let mut snapshots: Vec<ShardSnapshot> = Vec::new();
for shard_id in shard_ids {
let manifest_store =
ShardManifestStore::new(store.clone(), &base_path, shard_id, scan_batch_size);
if let Some(manifest) = manifest_store.read_latest().await? {
snapshots.push(snapshot_from_manifest(shard_id, &manifest, &exclude));
}
}
// WAL-only writers (enable_memtable=false) keep no in-memory memtable, and
// `in_memory_memtable_refs` errors in that mode; the on-disk manifests above
// already cover their flushed SSTables, so skip the live-writer snapshot. (Lance
// forbids maintained indexes in WAL-only mode, so only plain scans reach here.)
let wal_only = details
.writer_config_defaults
.get("enable_memtable")
.map(|v| v == "false")
.unwrap_or(false);
// Override the active shard with the cached writer's in-memory view.
let mut in_memory: HashMap<Uuid, InMemoryMemTables> = HashMap::new();
if !wal_only
&& let Some((shard_id, manifest, memtables)) =
table.dataset.shard_writer().read_snapshot().await?
{
if let Some(manifest) = manifest {
let snapshot = snapshot_from_manifest(shard_id, &manifest, &exclude);
match snapshots.iter_mut().find(|s| s.shard_id == shard_id) {
Some(existing) => *existing = snapshot,
None => snapshots.push(snapshot),
}
}
if let Some(memtables) = memtables {
in_memory.insert(shard_id, memtables);
}
}
Ok((snapshots, in_memory))
}
/// Convert a shard manifest into a read snapshot (current + not-yet-compacted
/// flushed SSTables). SSTable generations at or below the shard's compaction
/// watermark are already in the base table and are skipped.
fn snapshot_from_manifest(
shard_id: Uuid,
manifest: &ShardManifest,
compacted: &HashMap<Uuid, u64>,
) -> ShardSnapshot {
let mut snapshot = ShardSnapshot::new(shard_id)
.with_spec_id(manifest.shard_spec_id)
.with_current_generation(manifest.current_generation);
let watermark = compacted.get(&shard_id).copied();
for sstable in &manifest.sstables {
if watermark.is_some_and(|w| sstable.generation <= w) {
continue;
}
snapshot = snapshot.with_sstable(sstable.generation, sstable.path.clone());
}
snapshot
}
/// Columns selected by the query, if an explicit projection was requested.
fn selected_columns(query: &VectorQueryRequest) -> Option<Vec<String>> {
match &query.base.select {
Select::Columns(columns) => Some(columns.clone()),
_ => None,
}
}
/// Non-negative `Option<usize>` limit/offset as the `Option<i64>` the scanner
/// expects.
fn as_i64(value: Option<usize>) -> Option<i64> {
value.map(|v| v as i64)
}
/// Build a base `LsmScanner` configured with sources, filter, and projection.
///
/// The filter set here is applied as a prefilter across every arm — plain scan,
/// full-text search, and vector search — since all three terminate on this
/// scanner's `create_plan`.
fn base_scanner(
dataset: &Dataset,
query: &VectorQueryRequest,
pk_columns: Vec<String>,
snapshots: Vec<ShardSnapshot>,
in_memory: HashMap<Uuid, InMemoryMemTables>,
) -> Result<LsmScanner> {
let mut scanner = LsmScanner::new(Arc::new(dataset.clone()), snapshots, pk_columns);
for (shard_id, memtables) in in_memory {
scanner = scanner.with_in_memory_memtables(shard_id, memtables);
}
if let Some(columns) = selected_columns(query) {
let refs: Vec<&str> = columns.iter().map(String::as_str).collect();
scanner = scanner.project(&refs)?;
}
if let Some(filter) = &query.base.filter {
scanner = match filter {
QueryFilter::Sql(sql) => scanner.filter(sql)?,
QueryFilter::Datafusion(expr) => scanner.filter_expr(expr.clone()),
QueryFilter::Substrait(_) => {
return Err(Error::NotSupported {
message: "the MemWAL LSM scanner does not support Substrait filters; set use_lsm(false) to read the base table only".to_string(),
});
}
};
}
Ok(scanner)
}
/// Plain scan: filter / projection / limit over base SSTables in-memory.
/// The plain scan applies limit and offset inside the planner.
async fn plain_plan(
dataset: &Dataset,
query: &VectorQueryRequest,
pk_columns: Vec<String>,
snapshots: Vec<ShardSnapshot>,
in_memory: HashMap<Uuid, InMemoryMemTables>,
limit: Option<usize>,
offset: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
let scanner = base_scanner(dataset, query, pk_columns, snapshots, in_memory)?
.limit(as_i64(limit), as_i64(offset))?;
Ok(scanner.create_plan().await?)
}
/// Full-text search over base SSTables in-memory, merged by local BM25 score.
/// The scanner threads the query filter in as a prefilter and pages via limit/offset.
#[allow(clippy::too_many_arguments)]
async fn fts_plan(
dataset: &Dataset,
fts: lance_index::scalar::FullTextSearchQuery,
query: &VectorQueryRequest,
details: &MemWalIndexDetails,
pk_columns: Vec<String>,
snapshots: Vec<ShardSnapshot>,
in_memory: HashMap<Uuid, InMemoryMemTables>,
limit: Option<usize>,
offset: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
// Pre-check for a lancedb-flavored error; `LsmScanner` also validates the
// single-column requirement, but without the `use_lsm(false)` guidance.
let columns: Vec<String> = fts.columns().into_iter().collect();
if columns.len() > 1 {
return Err(Error::NotSupported {
message: "the MemWAL LSM scanner full-text search supports a single column; set use_lsm(false) to read the base table only".to_string(),
});
}
let column = columns.first().ok_or_else(|| Error::NotSupported {
message: "the MemWAL LSM scanner full-text search requires an explicit FTS column"
.to_string(),
})?;
// Without a maintained in-memory FTS index for this column, the active memtable
// arm produces an empty plan (`active_source_can_execute_fts` returns false), so
// the search silently omits un-compacted documents. Reject rather than mislead.
if !index_maintained(
dataset,
column,
&details.maintained_indexes,
"InvertedIndexDetails",
)
.await?
{
return Err(Error::NotSupported {
message: format!(
"the MemWAL LSM scanner full-text search requires the FTS index on '{column}' to be maintained by the write spec (LsmWriteSpec::with_maintained_indexes); otherwise un-compacted documents are omitted. set use_lsm(false) to read the base table only"
),
});
}
let scanner = base_scanner(dataset, query, pk_columns, snapshots, in_memory)?
.with_overfetch_factor(LSM_OVERFETCH_FACTOR)
.full_text_search(fts)?
.limit(as_i64(limit), as_i64(offset))?;
Ok(scanner.create_plan().await?)
}
/// Whether an index of `type_url_suffix` covering `column` is in the MemWAL spec's
/// maintained set. Only a maintained index has its catch-up tracked (so exclusion is
/// gated correctly) and its in-memory arm kept current; an unmaintained base index
/// falls back to the compaction watermark and can drop rows it has not re-indexed.
/// The type must match specifically — a maintained BTree on the same column is not
/// the FTS/vector index the arm relies on.
async fn index_maintained(
dataset: &Dataset,
column: &str,
maintained: &[String],
type_url_suffix: &str,
) -> Result<bool> {
use lance::index::DatasetIndexExt;
let Some(field) = dataset.schema().field(column) else {
return Ok(false);
};
let indices = dataset.load_indices().await?;
Ok(indices.iter().any(|idx| {
idx.fields.contains(&field.id)
&& maintained.iter().any(|m| m == &idx.name)
&& idx
.index_details
.as_ref()
.is_some_and(|d| d.type_url.ends_with(type_url_suffix))
}))
}
/// The maintained base index the query's arm relies on (vector index for ANN, FTS
/// index for full-text), used to gate SSTable compaction exclusion by index catch-up.
/// `None` for a plain scan or when no maintained index covers the searched column.
async fn arm_maintained_index_name(
dataset: &Dataset,
query: &VectorQueryRequest,
details: &MemWalIndexDetails,
) -> Result<Option<String>> {
use lance::index::DatasetIndexExt;
// Resolve the arm's searched column, the index-detail type it relies on, and a
// label for diagnostics — catch-up is taken from the vector/FTS index
// specifically, not a BTree on the same column.
let (column, type_url_suffix, arm) = if !query.query_vector.is_empty() {
let arrow_schema = ArrowSchema::from(dataset.schema());
let column = match &query.column {
Some(column) => column.clone(),
None => {
let dim = query.query_vector.first().map(|v| v.len() as i32);
default_vector_column(&arrow_schema, dim)?
}
};
(column, "VectorIndexDetails", "vector")
} else if let Some(fts) = &query.base.full_text_search {
match fts.columns().into_iter().next() {
Some(column) => (column, "InvertedIndexDetails", "full-text"),
None => return Ok(None),
}
} else {
return Ok(None);
};
let Some(field) = dataset.schema().field(&column) else {
return Ok(None);
};
let indices = dataset.load_indices().await?;
let segment_names: Vec<String> = indices
.iter()
.filter(|idx| {
idx.fields.contains(&field.id)
&& idx
.index_details
.as_ref()
.is_some_and(|d| d.type_url.ends_with(type_url_suffix))
})
.map(|idx| idx.name.clone())
.collect();
resolve_single_index(segment_names, &details.maintained_indexes, arm, &column)
}
/// Resolve the single logical index from the names of its matching physical
/// segments. `load_indices` returns one entry per segment, so one logical index can
/// appear multiple times (same name); dedupe by name before counting. Errors when
/// more than one *distinct* index covers the field — the base planner's choice is
/// ambiguous and their catch-up watermarks can diverge, so gating exclusion on the
/// wrong one could drop SSTables the used index has not caught up to. Otherwise
/// returns the name only when it is maintained (else the caller falls back to the
/// compaction watermark).
fn resolve_single_index(
mut names: Vec<String>,
maintained: &[String],
arm: &str,
column: &str,
) -> Result<Option<String>> {
names.sort();
names.dedup();
if names.len() > 1 {
return Err(Error::NotSupported {
message: format!(
"the MemWAL LSM scanner cannot resolve the {arm} index catch-up watermark for '{column}': it has multiple {arm} indexes; set use_lsm(false) to read the base table only"
),
});
}
Ok(names
.into_iter()
.next()
.filter(|name| maintained.contains(name)))
}
/// Drop the primary-key columns Lance appends internally for dedup when the user's
/// explicit projection did not request them, restoring the requested output schema.
/// `Select::All` legitimately includes the pk columns and is left untouched.
fn restore_projection(
plan: Arc<dyn ExecutionPlan>,
query: &VectorQueryRequest,
pk_columns: &[String],
) -> Result<Arc<dyn ExecutionPlan>> {
let Select::Columns(selected) = &query.base.select else {
return Ok(plan);
};
let schema = plan.schema();
// Keep a column unless it is a pk column the user did not select (this preserves
// user columns and score columns like `_distance`, dropping only leaked pk).
let keep: Vec<(Arc<dyn PhysicalExpr>, String)> = schema
.fields()
.iter()
.enumerate()
.filter(|(_, f)| {
selected.iter().any(|c| c == f.name()) || !pk_columns.iter().any(|pk| pk == f.name())
})
.map(|(i, f)| {
(
Arc::new(Column::new(f.name(), i)) as Arc<dyn PhysicalExpr>,
f.name().clone(),
)
})
.collect();
if keep.len() == schema.fields().len() {
return Ok(plan);
}
Ok(Arc::new(ProjectionExec::try_new(keep, plan)?))
}
/// Vector (ANN) search over base SSTables in-memory, routed through the same
/// [`LsmScanner`] as the other arms so the query filter is applied as a prefilter.
///
/// Note: the base and SSTable arms use `fast_search` (indexed data only), so a
/// base-table row not covered by a vector index is invisible here — it surfaces
/// only via the memtable or `use_lsm(false)`.
#[allow(clippy::too_many_arguments)]
async fn vector_plan(
dataset: &Dataset,
query: &VectorQueryRequest,
details: &MemWalIndexDetails,
pk_columns: Vec<String>,
snapshots: Vec<ShardSnapshot>,
in_memory: HashMap<Uuid, InMemoryMemTables>,
limit: Option<usize>,
offset: Option<usize>,
) -> Result<Arc<dyn ExecutionPlan>> {
let query_vector = query
.query_vector
.first()
.cloned()
.ok_or_else(|| Error::InvalidInput {
message: "vector search requires a query vector".to_string(),
})?;
let arrow_schema = ArrowSchema::from(dataset.schema());
let column = match &query.column {
Some(column) => column.clone(),
None => default_vector_column(&arrow_schema, Some(query_vector.len() as i32))?,
};
// The base arm relies on the column's vector index (`fast_search`). Unless it is
// maintained, its catch-up is untracked and exclusion falls back to the
// compaction watermark — dropping compacted SSTables the (lagging) base index has
// not re-indexed. Reject rather than silently omit rows, mirroring the FTS arm.
if !index_maintained(
dataset,
&column,
&details.maintained_indexes,
"VectorIndexDetails",
)
.await?
{
return Err(Error::NotSupported {
message: format!(
"the MemWAL LSM scanner requires the vector index on '{column}' to be maintained by the write spec (LsmWriteSpec::with_maintained_indexes); otherwise compacted rows not yet re-indexed are omitted. set use_lsm(false) to read the base table only"
),
});
}
// The LSM vector planner is Float32-only; reject binary (uint8) vectors with a
// clear error rather than failing deep in the planner.
if is_binary_vector_column(&arrow_schema, &column) {
return Err(Error::NotSupported {
message: "the MemWAL LSM scanner does not support binary (uint8) vector search; set use_lsm(false) to read the base table only".to_string(),
});
}
let distance_type = resolve_distance_type(dataset, query, &column).await?;
// `nearest` takes a flat query vector and builds the fixed-size list itself.
// Guard `k` to at least 1 so a degenerate `limit(0)` is trimmed by `limit`
// below rather than rejected by `nearest`.
let k = limit.unwrap_or(DEFAULT_TOP_K).max(1);
let mut scanner = base_scanner(dataset, query, pk_columns, snapshots, in_memory)?
.with_overfetch_factor(LSM_OVERFETCH_FACTOR)
.nearest(&column, query_vector.as_ref(), k)?
.nprobes(query.minimum_nprobes)
.distance_metric(distance_type.into());
if let Some(refine_factor) = query.refine_factor {
scanner = scanner.refine(refine_factor);
}
scanner = scanner.limit(as_i64(limit), as_i64(offset))?;
Ok(scanner.create_plan().await?)
}
/// Whether `column` stores binary (uint8) vectors, which the LSM vector planner
/// does not support.
fn is_binary_vector_column(schema: &ArrowSchema, column: &str) -> bool {
matches!(
schema.field_with_name(column).map(|f| f.data_type()),
Ok(DataType::FixedSizeList(field, _)) if matches!(field.data_type(), DataType::UInt8)
)
}
/// Resolve the distance metric for the vector arm: the explicit query metric if
/// set, else the metric of the column's vector index, else L2.
async fn resolve_distance_type(
dataset: &Dataset,
query: &VectorQueryRequest,
column: &str,
) -> Result<DistanceType> {
if let Some(dt) = query.distance_type {
return Ok(dt);
}
// Inherit the column's vector-index metric so cross-source distances match
// the metric the maintained memtable index was built with.
use lance::index::{DatasetIndexExt, DatasetIndexInternalExt};
use lance_index::metrics::NoOpMetricsCollector;
let field = dataset.schema().field(column);
if let Some(field) = field {
let indices = dataset.load_indices().await?;
for index in indices.iter() {
if index.fields.contains(&field.id)
&& let Ok(vector_index) = dataset
.open_vector_index(column, &index.uuid, &NoOpMetricsCollector)
.await
{
return Ok(vector_index.metric_type().into());
}
}
}
Ok(DistanceType::L2)
}
#[cfg(test)]
mod tests {
use super::*;
use lance_index::mem_wal::{CompactedSsTable, IndexCatchupProgress};
#[test]
fn exclusion_watermark_gates_on_lagging_index_catchup() {
let shard = Uuid::from_u128(1);
let details = MemWalIndexDetails {
// Compaction has drained generations through 5 into the base table...
compacted_sstables: vec![CompactedSsTable::new(shard, 5)],
// ...but the FTS index has only caught up through generation 2.
index_catchup: vec![IndexCatchupProgress::new(
"fts_idx".to_string(),
vec![CompactedSsTable::new(shard, 2)],
)],
maintained_indexes: vec!["fts_idx".to_string()],
..Default::default()
};
// Plain scan: drop every compacted generation (through 5).
assert_eq!(exclusion_watermarks(&details, None).get(&shard), Some(&5));
// FTS arm with a lagging index: exclusion is capped at the index catch-up
// (2), so SSTable generations 3..=5 are retained until the index covers
// them — otherwise those documents would silently vanish from FTS results.
assert_eq!(
exclusion_watermarks(&details, Some("fts_idx")).get(&shard),
Some(&2)
);
// A caught-up index — or one untracked in index_catchup — falls back to the
// compaction watermark.
assert_eq!(
exclusion_watermarks(&details, Some("caught_up_idx")).get(&shard),
Some(&5)
);
}
#[test]
fn resolve_single_index_dedupes_segments() {
let maintained = vec!["fts_idx".to_string()];
// Two physical segments of ONE logical index must not count as "multiple".
assert_eq!(
resolve_single_index(
vec!["fts_idx".to_string(), "fts_idx".to_string()],
&maintained,
"full-text",
"text"
)
.unwrap(),
Some("fts_idx".to_string())
);
// Two distinct indexes on the field are ambiguous → error.
assert!(
resolve_single_index(
vec!["fts_a".to_string(), "fts_b".to_string()],
&maintained,
"full-text",
"text"
)
.is_err()
);
// A single unmaintained index resolves to None (compaction-watermark fallback).
assert_eq!(
resolve_single_index(vec!["other".to_string()], &maintained, "full-text", "text")
.unwrap(),
None
);
}
}
+3 -131
View File
@@ -12,7 +12,7 @@ use futures::TryStreamExt;
use lance_encoding::version::LanceFileVersion;
use lancedb::{
Connection, Error, Result, Table,
blob::{BlobRangeRequest, blob},
blob::blob,
connect, connect_namespace,
database::listing::OPT_NEW_TABLE_ENABLE_STABLE_ROW_IDS,
query::{ExecutableQuery, QueryBase},
@@ -595,73 +595,6 @@ async fn fetch_blobs_aligns_with_reordered_and_duplicate_ids() -> Result<()> {
Ok(())
}
#[tokio::test]
async fn fetch_blob_ranges_aligns_repeated_ranges_and_nulls() -> Result<()> {
let tmp = tempdir().unwrap();
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
let table =
create_inline_blob_table(&db, "t", &[1, 2], &[Some(b"abcdefghij".as_slice()), None])
.await?;
let pairs = collect_id_rowid(&table).await?;
let by_id = |want: i64| pairs.iter().find(|(id, _)| *id == want).unwrap().1;
let requests = [
BlobRangeRequest::new(by_id(1), 2, 3),
BlobRangeRequest::new(by_id(2), 0, 0),
BlobRangeRequest::new(by_id(1), 0, 2),
BlobRangeRequest::new(by_id(1), 2, 3),
BlobRangeRequest::new(by_id(1), 10, 0),
];
let bytes = table.fetch_blob_ranges("image", requests).await?;
assert_eq!(bytes.len(), requests.len());
assert_eq!(bytes.value(0), b"cde");
assert!(bytes.is_null(1));
assert_eq!(bytes.value(2), b"ab");
assert_eq!(bytes.value(3), b"cde");
assert_eq!(bytes.value(4), b"");
Ok(())
}
#[tokio::test]
async fn fetch_blob_ranges_validates_requests() -> Result<()> {
let tmp = tempdir().unwrap();
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
let table = create_inline_blob_table(&db, "t", &[1], &[Some(b"abc".as_slice())]).await?;
let row_id = collect_row_ids(&table).await?[0];
let err = table
.fetch_blob_ranges("image", [BlobRangeRequest::new(row_id, 2, 2)])
.await
.unwrap_err();
assert!(err.to_string().contains("exceeds blob size"));
let err = table
.fetch_blob_ranges("image", [BlobRangeRequest::new(row_id, u64::MAX, 1)])
.await
.unwrap_err();
assert!(err.to_string().contains("offset + length overflowed"));
let err = table
.fetch_blob_ranges("image", [BlobRangeRequest::new(u64::MAX, 0, 1)])
.await
.unwrap_err();
assert!(matches!(&err, Error::InvalidInput { .. }), "got {err:?}");
assert!(err.to_string().contains("row ids"));
Ok(())
}
#[tokio::test]
async fn fetch_blob_ranges_empty_requests_returns_empty_array() -> Result<()> {
let tmp = tempdir().unwrap();
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
let table = create_inline_blob_table(&db, "t", &[1], &[Some(b"x".as_slice())]).await?;
let bytes = table.fetch_blob_ranges("image", std::iter::empty()).await?;
assert!(bytes.is_empty());
Ok(())
}
#[tokio::test]
async fn fetch_blobs_empty_ids_returns_empty() -> Result<()> {
let tmp = tempdir().unwrap();
@@ -684,32 +617,6 @@ async fn fetch_blobs_out_of_range_id_errors_without_panic() -> Result<()> {
Ok(())
}
#[tokio::test]
async fn fetch_blob_apis_reject_mixed_valid_and_missing_row_ids() -> Result<()> {
let tmp = tempdir().unwrap();
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
let table = create_inline_blob_table(&db, "t", &[1], &[Some(b"x".as_slice())]).await?;
let row_id = collect_row_ids(&table).await?[0];
let row_ids = [u64::MAX, row_id];
let err = table.fetch_blobs("image", &row_ids).await.unwrap_err();
assert!(matches!(&err, Error::InvalidInput { .. }), "got {err:?}");
assert!(err.to_string().contains("row ids"));
let err = table.fetch_blob_files("image", &row_ids).await.unwrap_err();
assert!(matches!(&err, Error::InvalidInput { .. }), "got {err:?}");
assert!(err.to_string().contains("row ids"));
let requests = row_ids.map(|row_id| BlobRangeRequest::new(row_id, 0, 1));
let err = table
.fetch_blob_ranges("image", requests)
.await
.unwrap_err();
assert!(matches!(&err, Error::InvalidInput { .. }), "got {err:?}");
assert!(err.to_string().contains("row ids"));
Ok(())
}
#[tokio::test]
async fn fetch_blobs_rejects_non_blob_column() -> Result<()> {
let tmp = tempdir().unwrap();
@@ -936,24 +843,11 @@ async fn fetch_blobs_with_precompaction_row_ids_survives_compaction() -> Result<
_ => unreachable!(),
}
}
let ranges = ids_before
.iter()
.map(|row_id| BlobRangeRequest::new(*row_id, 5, 3));
let ranges_after = table.fetch_blob_ranges("image", ranges).await?;
assert_eq!(ranges_after.len(), 2);
for (i, (id, _)) in pairs_before.iter().enumerate() {
match id {
1 => assert_eq!(ranges_after.value(i), b"one"),
2 => assert_eq!(ranges_after.value(i), b"two"),
_ => unreachable!(),
}
}
Ok(())
}
#[tokio::test]
async fn empty_blob_reads_back_as_empty_bytes() -> Result<()> {
async fn zero_length_blob_reads_back_as_null() -> Result<()> {
let tmp = tempdir().unwrap();
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
let table = create_inline_blob_table(&db, "t", &[1], &[Some(b"".as_slice())]).await?;
@@ -961,8 +855,7 @@ async fn empty_blob_reads_back_as_empty_bytes() -> Result<()> {
let ids = collect_row_ids(&table).await?;
let bytes = table.fetch_blobs("image", &ids).await?;
assert_eq!(bytes.len(), 1);
assert!(!bytes.is_null(0));
assert!(bytes.value(0).is_empty());
assert!(bytes.is_null(0));
Ok(())
}
@@ -1034,27 +927,6 @@ async fn fetch_blobs_aligns_across_fragments_with_nulls_and_dups() -> Result<()>
Ok(())
}
#[tokio::test]
async fn fetch_blob_ranges_aligns_across_fragments_with_nulls_and_dups() -> Result<()> {
let tmp = tempdir().unwrap();
let db = connect(tmp.path().to_str().unwrap()).execute().await?;
let table = multi_fragment_dedicated_blob_table(&db).await?;
let row_ids = row_ids_for_logical(&table, &SCRAMBLED_LOGICAL_IDS).await?;
let requests = row_ids
.iter()
.map(|row_id| BlobRangeRequest::new(*row_id, 123, 8));
let bytes = table.fetch_blob_ranges("image", requests).await?;
assert_eq!(bytes.len(), SCRAMBLED_LOGICAL_IDS.len());
for (slot, logical_id) in SCRAMBLED_LOGICAL_IDS.iter().enumerate() {
match logical_id {
3 | 5 => assert!(bytes.is_null(slot)),
id => assert_eq!(bytes.value(slot), [*id as u8; 8]),
}
}
Ok(())
}
#[tokio::test]
async fn fetch_blob_files_aligns_across_fragments_with_nulls_and_dups() -> Result<()> {
let tmp = tempdir().unwrap();