Compare commits

..
Author SHA1 Message Date
lancedb automation 2478993501 chore: update lance dependency to v4.1.0-beta.3 2026-03-23 20:55:53 +00:00
Esteban GutierrezandEsteban Gutierrez a0228036ae ci: fix unused PreprocessingOutput (#3180)
Simple fix to for CI due unused import of PreprocessingOutput in
table.rs

Co-authored-by: Esteban Gutierrez <esteban@lancedb.com>
2026-03-23 13:45:44 -07:00
Esteban GutierrezandEsteban Gutierrez d8fc071a7d fix(ci): bump AWS SDK MSRV pins to March 2025 release (#3179)
Lance v4.1.0-beta requires the default-https-client feature on
aws-sdk-dynamodb and aws-sdk-s3, which was introduced in the March
2025 AWS SDK release. Update all AWS SDK pins to versions from the
same AWS SDK release to maintain internal dependency compatibility.

Co-authored-by: Esteban Gutierrez <esteban@lancedb.com>
2026-03-23 15:30:33 -05:00
Will JonesandClaude Opus 4.6 e6fd8d071e feat(rust): parallel inserts for remote tables via multipart write (#3071)
Similar to https://github.com/lancedb/lancedb/pull/3062, we can write in
parallel to remote tables if the input data source is large enough.

We take advantage of new endpoints coming in server version 0.4.0, which
allow writing data in multiple requests, and the committing at the end
in a single request.

To make testing easier, I also introduce a `write_parallelism`
parameter. In the future, we can expose that in Python and NodeJS so
users can manually specify the parallelism they get.

Closes #2861

---------

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-20 13:19:07 -07:00
LanceDB RobotandEsteban Gutierrez 670dcca551 feat: update lance dependency to v3.0.1 (#3168)
## Summary
- Updated Lance Rust workspace dependencies to `3.0.1` using
`ci/set_lance_version.py`.
- Updated Java `lance-core` dependency property in `java/pom.xml` to
`3.0.1`.
- Refreshed `Cargo.lock` entries for Lance crates at `3.0.1`.

## Verification
- `cargo clippy --workspace --tests --all-features -- -D warnings`
- `cargo fmt --all`

## Trigger
- Tag:
[`refs/tags/v3.0.1`](https://github.com/lancedb/lance/tree/v3.0.1)

Co-authored-by: Esteban Gutierrez <estebangtz@gmail.com>
2026-03-20 09:53:20 -07:00
Prashanth Rao ed7e01a58b docs: fix rendering issues with missing index types in API docs (#3143)
## Problem

The generated Python API docs for
`lancedb.table.IndexStatistics.index_type` were misleading because
mkdocstrings renders that field’s type annotation directly, and the
existing `Literal[...]` listed only a subset of the actual canonical SDK
index type strings.

Current (missing index types):
<img width="823" height="83" alt="image"
src="https://github.com/user-attachments/assets/f6f29fe3-4c16-4d00-a4e9-28a7cd6e19ec"
/>


## Fix

- Update the `IndexStatistics.index_type` annotation in
`python/python/lancedb/table.py` to include the full supported set of
canonical values, so the generated docs show all valid index_type
strings inline.
- Add a small regression test in `python/python/tests/test_index.py` to
ensure the docs-facing annotation does not drift silently again in case
we add a new index/quantization type in the future.
- Bumps mkdocs and material theme versions to mkdocs 1.6 to allow access
to more features like hooks

After fix (all index types are included and tested for in the
annotations):
<img width="1017" height="93" alt="image"
src="https://github.com/user-attachments/assets/66c74d5c-34b3-4b44-8173-3ee23e3648ac"
/>
2026-03-20 09:34:42 -07:00
Lance Release 3450ccaf7f Bump version: 0.27.1-beta.0 → 0.27.1 2026-03-20 00:35:36 +00:00
Lance Release 9b229f1e7c Bump version: 0.27.0 → 0.27.1-beta.0 2026-03-20 00:35:19 +00:00
Lance Release f5b21c0aa4 Bump version: 0.30.1-beta.0 → 0.30.1 2026-03-20 00:35:03 +00:00
Lance Release e927924d26 Bump version: 0.30.0 → 0.30.1-beta.0 2026-03-20 00:35:02 +00:00
Weston PaceandClaude Opus 4.6 11a4966bfc feat: upgrade lance dependency to v3.0.1 (#3157)
## Summary
- Upgrade all lance-* dependencies from v3.0.0 to v3.0.1 (stable, from
crates.io)

## Test plan
- [x] `cargo check --features remote --tests --examples` passes
- [x] `cargo clippy --features remote --tests --examples` passes
- [x] `cargo fmt --all --check` passes
- [ ] CI tests pass

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-19 17:30:46 -07:00
Weston Pace dd5aaa72dc ci: modify check_lance_release.py to prefer stable releases over betas (#3146)
When Lance 3.0.0 released the check_lance_release.py script did not make
a PR for it because it was a pre-release. This change may not be perfect
but it always ranks stable releases above non-stable releases.
2026-03-17 09:21:30 -07:00
marca116 3a200d77ef fix: pre-filtering on hybrid search (#3096)
When using hybrid search with a where filter, the prefilter argument is
silently inverted. Passing prefilter=True actually performs
post-filtering, and prefilter=False actually performs pre-filtering.
2026-03-16 21:48:42 -07:00
Lance Release bd09c53938 Bump version: 0.27.0-beta.6 → 0.27.0 2026-03-16 22:47:06 +00:00
Lance Release 0b18e33180 Bump version: 0.27.0-beta.5 → 0.27.0-beta.6 2026-03-16 22:46:48 +00:00
Lance Release c89240b16c Bump version: 0.30.0-beta.6 → 0.30.0 2026-03-16 22:46:19 +00:00
Lance Release 099ff355a4 Bump version: 0.30.0-beta.5 → 0.30.0-beta.6 2026-03-16 22:46:17 +00:00
Weston PaceandClaude Opus 4.6 c5995fda67 feat: update lance dependency to 3.0.0 release (#3137)
## Summary
- Update all 14 lance crates from `3.0.0-rc.3` (git source) to `3.0.0`
(crates.io release)
- Remove git/tag source references since 3.0.0 is published on crates.io

## Test plan
- [x] `cargo check --features remote --tests --examples` passes
- [x] `cargo clippy --features remote --tests --examples` passes
- [ ] CI passes

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-16 15:29:18 -07:00
Weston Pace 25eb1fbfa4 fix: restore storage options on copy in localstack tests (#3148) 2026-03-16 14:02:19 -07:00
Weston PaceandClaude Opus 4.6 4ac41c5c3f fix(ci): upgrade LocalStack to 4.0 for S3 integration tests (#3147)
## Summary
- Upgrade LocalStack from 3.3 to 4.0 in `docker-compose.yml` to fix S3
integration test failures in CI
- Version 3.3 has compatibility issues with newer Python 3.13 and
updated boto3 dependencies
- Matches the LocalStack version used successfully in the lance
repository

## Test plan
- [ ] Verify `docker compose up --detach --wait` completes successfully
in CI
- [ ] All tests in `test_s3.py` pass (5 tests)
- [ ] All `@pytest.mark.s3_test` tests in
`test_namespace_integration.py` pass (7 tests)
- [ ] No regressions in non-integration test jobs (Mac, Windows)

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-16 09:02:11 -07:00
Will JonesandClaude Opus 4.6 9a5b0398ec chore: fix ci (#3139)
* Move away from buildjet, which is shutting down runners for GHA [^1]
* Add `Cargo.lock` to build jobs, so when we upgrade locked dependencies
we check the builds actually pass. CI started failing because
dependencies were changed in #3116 without running all build jobs.
* Add fixes for aws-lc-rs build in NodeJS.

[^1]: https://buildjet.com/for-github-actions/blog/we-are-shutting-down

---------

Co-authored-by: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
2026-03-16 06:25:40 -07:00
Pratik DeyandClaude d1d720d08a feat(nodejs): support field/data type input in add_columns() method (#3114)
Add support for passing field/data type information into add_columns()
method, bringing parity with Python bindings. The method now accepts:

- AddColumnsSql[] - SQL expressions (existing functionality)
- Field - single Arrow field with explicit data type
- Field[] - array of Arrow fields with explicit data types
- Schema - Arrow schema with explicit data types

New columns added via Field/Schema are initialized with null values. All
field-based columns must be nullable due to null initialization.

Resolves #3107

---------

Signed-off-by: Pratik <pratikrocks.dey11@gmail.com>
Co-authored-by: Claude <noreply@anthropic.com>
2026-03-13 12:57:14 -07:00
Mesut-Doner c2e543f1b7 feat(rust): support Expr in projection query (#3069)
Referred and followed [`Select::Dynamic`] implementation. 

Closes #3039
2026-03-13 12:54:26 -07:00
Weston PaceandClaude Opus 4.6 216c1b5f77 docs: remove experimental label from optimize and warn about delete_unverified (#3128)
## Summary
- Removes the "Experimental API" section from `optimize` method
documentation across Rust, Python, and TypeScript
- Adds a warning to `delete_unverified` documentation in all bindings:
this should only be set to true if you can guarantee no other process is
working on the dataset, otherwise it could be corrupted
- Fixes a typo ("shoudl" → "should")

Closes #3125


🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-12 14:37:42 +08:00
47 changed files with 1524 additions and 295 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.27.0-beta.5"
current_version = "0.27.1"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
+1
View File
@@ -7,6 +7,7 @@ on:
pull_request:
paths:
- Cargo.toml
- Cargo.lock
- nodejs/**
- rust/**
- docs/src/js/**
+12 -5
View File
@@ -19,6 +19,7 @@ on:
paths:
- .github/workflows/npm-publish.yml
- Cargo.toml # Change in dependency frequently breaks builds
- Cargo.lock
concurrency:
group: ${{ github.workflow }}-${{ github.ref }}
@@ -124,7 +125,12 @@ jobs:
pre_build: |-
set -e &&
apt-get update &&
apt-get install -y protobuf-compiler pkg-config
apt-get install -y protobuf-compiler pkg-config &&
# The base image (manylinux2014-cross) sets TARGET_CC to the old
# GCC 4.8 cross-compiler. aws-lc-sys checks TARGET_CC before CC,
# so it picks up GCC even though the napi-rs image sets CC=clang.
# Override to use the image's clang-18 which supports -fuse-ld=lld.
export TARGET_CC=clang TARGET_CXX=clang++
- target: x86_64-unknown-linux-musl
# This one seems to need some extra memory
host: ubuntu-2404-8x-x64
@@ -144,9 +150,10 @@ jobs:
set -e &&
apt-get update &&
apt-get install -y protobuf-compiler pkg-config &&
# https://github.com/aws/aws-lc-rs/issues/737#issuecomment-2725918627
ln -s /usr/aarch64-unknown-linux-gnu/lib/gcc/aarch64-unknown-linux-gnu/4.8.5/crtbeginS.o /usr/aarch64-unknown-linux-gnu/aarch64-unknown-linux-gnu/sysroot/usr/lib/crtbeginS.o &&
ln -s /usr/aarch64-unknown-linux-gnu/lib/gcc /usr/aarch64-unknown-linux-gnu/aarch64-unknown-linux-gnu/sysroot/usr/lib/gcc &&
export TARGET_CC=clang TARGET_CXX=clang++ &&
# The manylinux2014 sysroot has glibc 2.17 headers which lack
# AT_HWCAP2 (added in Linux 3.17). Define it for aws-lc-sys.
export CFLAGS="$CFLAGS -DAT_HWCAP2=26" &&
rustup target add aarch64-unknown-linux-gnu
- target: aarch64-unknown-linux-musl
host: ubuntu-2404-8x-x64
@@ -266,7 +273,7 @@ jobs:
- target: x86_64-unknown-linux-gnu
host: ubuntu-latest
- target: aarch64-unknown-linux-gnu
host: buildjet-16vcpu-ubuntu-2204-arm
host: ubuntu-2404-8x-arm64
node:
- '20'
runs-on: ${{ matrix.settings.host }}
+1
View File
@@ -9,6 +9,7 @@ on:
paths:
- .github/workflows/pypi-publish.yml
- Cargo.toml # Change in dependency frequently breaks builds
- Cargo.lock
env:
PIP_EXTRA_INDEX_URL: "https://pypi.fury.io/lance-format/ https://pypi.fury.io/lancedb/"
+1
View File
@@ -7,6 +7,7 @@ on:
pull_request:
paths:
- Cargo.toml
- Cargo.lock
- python/**
- rust/**
- .github/workflows/python.yml
+9 -8
View File
@@ -7,6 +7,7 @@ on:
pull_request:
paths:
- Cargo.toml
- Cargo.lock
- rust/**
- .github/workflows/rust.yml
@@ -206,14 +207,14 @@ jobs:
- name: Downgrade dependencies
# These packages have newer requirements for MSRV
run: |
cargo update -p aws-sdk-bedrockruntime --precise 1.64.0
cargo update -p aws-sdk-dynamodb --precise 1.55.0
cargo update -p aws-config --precise 1.5.10
cargo update -p aws-sdk-kms --precise 1.51.0
cargo update -p aws-sdk-s3 --precise 1.65.0
cargo update -p aws-sdk-sso --precise 1.50.0
cargo update -p aws-sdk-ssooidc --precise 1.51.0
cargo update -p aws-sdk-sts --precise 1.51.0
cargo update -p aws-sdk-bedrockruntime --precise 1.77.0
cargo update -p aws-sdk-dynamodb --precise 1.68.0
cargo update -p aws-config --precise 1.6.0
cargo update -p aws-sdk-kms --precise 1.63.0
cargo update -p aws-sdk-s3 --precise 1.79.0
cargo update -p aws-sdk-sso --precise 1.62.0
cargo update -p aws-sdk-ssooidc --precise 1.63.0
cargo update -p aws-sdk-sts --precise 1.63.0
cargo update -p home --precise 0.5.9
- name: cargo +${{ matrix.msrv }} check
env:
Generated
+45 -43
View File
@@ -3070,8 +3070,8 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
[[package]]
name = "fsst"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow-array",
"rand 0.9.2",
@@ -3852,7 +3852,7 @@ dependencies = [
"libc",
"percent-encoding",
"pin-project-lite",
"socket2 0.5.10",
"socket2 0.6.0",
"system-configuration",
"tokio",
"tower-service",
@@ -4241,8 +4241,8 @@ dependencies = [
[[package]]
name = "lance"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"arrow-arith",
@@ -4308,8 +4308,8 @@ dependencies = [
[[package]]
name = "lance-arrow"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4329,8 +4329,8 @@ dependencies = [
[[package]]
name = "lance-bitpacking"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrayref",
"paste",
@@ -4339,8 +4339,8 @@ dependencies = [
[[package]]
name = "lance-core"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4377,8 +4377,8 @@ dependencies = [
[[package]]
name = "lance-datafusion"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"arrow-array",
@@ -4408,8 +4408,8 @@ dependencies = [
[[package]]
name = "lance-datagen"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"arrow-array",
@@ -4427,8 +4427,8 @@ dependencies = [
[[package]]
name = "lance-encoding"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -4465,8 +4465,8 @@ dependencies = [
[[package]]
name = "lance-file"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -4498,8 +4498,8 @@ dependencies = [
[[package]]
name = "lance-index"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"arrow-arith",
@@ -4513,6 +4513,7 @@ dependencies = [
"bitpacking",
"bitvec",
"bytes",
"chrono",
"crossbeam-queue",
"datafusion",
"datafusion-common",
@@ -4562,8 +4563,8 @@ dependencies = [
[[package]]
name = "lance-io"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"arrow-arith",
@@ -4604,8 +4605,8 @@ dependencies = [
[[package]]
name = "lance-linalg"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4621,21 +4622,22 @@ dependencies = [
[[package]]
name = "lance-namespace"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"async-trait",
"bytes",
"lance-core",
"lance-namespace-reqwest-client",
"serde",
"snafu 0.9.0",
]
[[package]]
name = "lance-namespace-impls"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"arrow-ipc",
@@ -4666,9 +4668,9 @@ dependencies = [
[[package]]
name = "lance-namespace-reqwest-client"
version = "0.5.2"
version = "0.6.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3ad4c947349acd6e37e984eba0254588bd894e6128434338b9e6904e56fb4633"
checksum = "ee2e48de899e2931afb67fcddd0a08e439bf5d8b6ea2a2ed9cb8f4df669bd5cc"
dependencies = [
"reqwest",
"serde",
@@ -4679,8 +4681,8 @@ dependencies = [
[[package]]
name = "lance-table"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow",
"arrow-array",
@@ -4719,8 +4721,8 @@ dependencies = [
[[package]]
name = "lance-testing"
version = "4.0.0-beta.9"
source = "git+https://github.com/lance-format/lance.git?tag=v4.0.0-beta.9#e133d8210b82e7dd43234bc6b38e5ad1863a4665"
version = "4.1.0-beta.3"
source = "git+https://github.com/lance-format/lance.git?tag=v4.1.0-beta.3#244c721504c6ef0b4c2f9700a342509976898d6e"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -4731,7 +4733,7 @@ dependencies = [
[[package]]
name = "lancedb"
version = "0.27.0-beta.5"
version = "0.27.1"
dependencies = [
"ahash",
"anyhow",
@@ -4813,7 +4815,7 @@ dependencies = [
[[package]]
name = "lancedb-nodejs"
version = "0.27.0-beta.5"
version = "0.27.1"
dependencies = [
"arrow-array",
"arrow-ipc",
@@ -4833,7 +4835,7 @@ dependencies = [
[[package]]
name = "lancedb-python"
version = "0.30.0-beta.5"
version = "0.30.1"
dependencies = [
"arrow",
"async-trait",
@@ -6443,7 +6445,7 @@ version = "0.14.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "343d3bd7056eda839b03204e68deff7d1b13aba7af2b2fd16890697274262ee7"
dependencies = [
"heck 0.4.1",
"heck 0.5.0",
"itertools 0.14.0",
"log",
"multimap",
@@ -6632,7 +6634,7 @@ dependencies = [
"quinn-udp",
"rustc-hash",
"rustls 0.23.31",
"socket2 0.5.10",
"socket2 0.6.0",
"thiserror 2.0.17",
"tokio",
"tracing",
@@ -6669,7 +6671,7 @@ dependencies = [
"cfg_aliases",
"libc",
"once_cell",
"socket2 0.5.10",
"socket2 0.6.0",
"tracing",
"windows-sys 0.60.2",
]
@@ -7735,7 +7737,7 @@ version = "0.8.9"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c1c97747dbf44bb1ca44a561ece23508e99cb592e862f22222dcf42f51d1e451"
dependencies = [
"heck 0.4.1",
"heck 0.5.0",
"proc-macro2",
"quote",
"syn 2.0.114",
@@ -7747,7 +7749,7 @@ version = "0.9.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "54254b8531cafa275c5e096f62d48c81435d1015405a91198ddb11e967301d40"
dependencies = [
"heck 0.4.1",
"heck 0.5.0",
"proc-macro2",
"quote",
"syn 2.0.114",
+14 -14
View File
@@ -15,20 +15,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
lance = { "version" = "=4.0.0-beta.9", default-features = false, "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=4.0.0-beta.9", default-features = false, "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=4.0.0-beta.9", default-features = false, "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=4.0.0-beta.9", "tag" = "v4.0.0-beta.9", "git" = "https://github.com/lance-format/lance.git" }
lance = { "version" = "=4.1.0-beta.3", default-features = false, "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=4.1.0-beta.3", default-features = false, "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=4.1.0-beta.3", default-features = false, "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=4.1.0-beta.3", "tag" = "v4.1.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
ahash = "0.8"
# Note that this one does not include pyarrow
arrow = { version = "57.2", optional = false }
+5 -1
View File
@@ -3,6 +3,7 @@
from __future__ import annotations
import argparse
import functools
import json
import os
import re
@@ -26,6 +27,7 @@ SEMVER_RE = re.compile(
)
@functools.total_ordering
@dataclass(frozen=True)
class SemVer:
major: int
@@ -156,7 +158,9 @@ def read_current_version(repo_root: Path) -> str:
def determine_latest_tag(tags: Iterable[TagInfo]) -> TagInfo:
return max(tags, key=lambda tag: tag.semver)
# Stable releases (no prerelease) are always preferred over pre-releases.
# Within each group, standard semver ordering applies.
return max(tags, key=lambda tag: (not tag.semver.prerelease, tag.semver))
def write_outputs(args: argparse.Namespace, payload: dict) -> None:
+1 -1
View File
@@ -1,7 +1,7 @@
version: "3.9"
services:
localstack:
image: localstack/localstack:3.3
image: localstack/localstack:4.0
ports:
- 4566:4566
environment:
+3 -3
View File
@@ -1,8 +1,8 @@
mkdocs==1.5.3
mkdocs==1.6.1
mkdocs-jupyter==0.24.1
mkdocs-material==9.5.3
mkdocs-material==9.6.23
mkdocs-autorefs>=0.5,<=1.0
mkdocstrings[python]==0.25.2
mkdocstrings[python]>=0.24,<1.0
griffe>=0.40,<1.0
mkdocs-render-swagger-plugin>=0.1.0
pydantic>=2.0,<3.0
+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.27.0-beta.5</version>
<version>0.27.1</version>
</dependency>
```
+7 -18
View File
@@ -71,11 +71,12 @@ Add new columns with defined values.
#### Parameters
* **newColumnTransforms**: [`AddColumnsSql`](../interfaces/AddColumnsSql.md)[]
pairs of column names and
the SQL expression to use to calculate the value of the new column. These
expressions will be evaluated for each row in the table, and can
reference existing columns in the table.
* **newColumnTransforms**: `Field`&lt;`any`&gt; \| `Field`&lt;`any`&gt;[] \| `Schema`&lt;`any`&gt; \| [`AddColumnsSql`](../interfaces/AddColumnsSql.md)[]
Either:
- An array of objects with column names and SQL expressions to calculate values
- A single Arrow Field defining one column with its data type (column will be initialized with null values)
- An array of Arrow Fields defining columns with their data types (columns will be initialized with null values)
- An Arrow Schema defining columns with their data types (columns will be initialized with null values)
#### Returns
@@ -484,19 +485,7 @@ Modeled after ``VACUUM`` in PostgreSQL.
- Prune: Removes old versions of the dataset
- Index: Optimizes the indices, adding new data to existing indices
Experimental API
----------------
The optimization process is undergoing active development and may change.
Our goal with these changes is to improve the performance of optimization and
reduce the complexity.
That being said, it is essential today to run optimize if you want the best
performance. It should be stable and safe to use in production, but it our
hope that the API may be simplified (or not even need to be called) in the
future.
The frequency an application shoudl call optimize is based on the frequency of
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
you have added or modified 100,000 or more records or run more than 20 data
@@ -37,3 +37,12 @@ tbl.optimize({cleanupOlderThan: new Date()});
```ts
deleteUnverified: boolean;
```
Because they may be part of an in-progress transaction, files newer than
7 days old are not deleted by default. If you are sure that there are no
in-progress transactions, then you can set this to true to delete all
files older than `cleanupOlderThan`.
**WARNING**: This should only be set to true if you can guarantee that
no other process is currently working on this dataset. Otherwise the
dataset could be put into a corrupted state.
+1 -1
View File
@@ -8,7 +8,7 @@
<parent>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.27.0-beta.5</version>
<version>0.27.1-final.0</version>
<relativePath>../pom.xml</relativePath>
</parent>
+2 -2
View File
@@ -6,7 +6,7 @@
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.27.0-beta.5</version>
<version>0.27.1-final.0</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>4.0.0-beta.9</lance-core.version>
<lance-core.version>4.1.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 @@
[package]
name = "lancedb-nodejs"
edition.workspace = true
version = "0.27.0-beta.5"
version = "0.27.1"
license.workspace = true
description.workspace = true
repository.workspace = true
+92
View File
@@ -1259,6 +1259,98 @@ describe("schema evolution", function () {
expect(await table.schema()).toEqual(expectedSchema);
});
it("can add columns with schema for explicit data types", async function () {
const con = await connect(tmpDir.name);
const table = await con.createTable("vectors", [
{ id: 1n, vector: [0.1, 0.2] },
]);
// Define schema for new columns with explicit data types
// Note: All columns must be nullable when using addColumns with Schema
// because they are initially populated with null values
const newColumnsSchema = new Schema([
new Field("price", new Float64(), true),
new Field("category", new Utf8(), true),
new Field("rating", new Int32(), true),
]);
const result = await table.addColumns(newColumnsSchema);
expect(result).toHaveProperty("version");
expect(result.version).toBe(2);
const expectedSchema = new Schema([
new Field("id", new Int64(), true),
new Field(
"vector",
new FixedSizeList(2, new Field("item", new Float32(), true)),
true,
),
new Field("price", new Float64(), true),
new Field("category", new Utf8(), true),
new Field("rating", new Int32(), true),
]);
expect(await table.schema()).toEqual(expectedSchema);
// Verify that new columns are populated with null values
const results = await table.query().toArray();
expect(results).toHaveLength(1);
expect(results[0].price).toBeNull();
expect(results[0].category).toBeNull();
expect(results[0].rating).toBeNull();
});
it("can add a single column using Field", async function () {
const con = await connect(tmpDir.name);
const table = await con.createTable("vectors", [
{ id: 1n, vector: [0.1, 0.2] },
]);
// Add a single field
const priceField = new Field("price", new Float64(), true);
const result = await table.addColumns(priceField);
expect(result).toHaveProperty("version");
expect(result.version).toBe(2);
const expectedSchema = new Schema([
new Field("id", new Int64(), true),
new Field(
"vector",
new FixedSizeList(2, new Field("item", new Float32(), true)),
true,
),
new Field("price", new Float64(), true),
]);
expect(await table.schema()).toEqual(expectedSchema);
});
it("can add multiple columns using array of Fields", async function () {
const con = await connect(tmpDir.name);
const table = await con.createTable("vectors", [
{ id: 1n, vector: [0.1, 0.2] },
]);
// Add multiple fields as array
const fields = [
new Field("price", new Float64(), true),
new Field("category", new Utf8(), true),
];
const result = await table.addColumns(fields);
expect(result).toHaveProperty("version");
expect(result.version).toBe(2);
const expectedSchema = new Schema([
new Field("id", new Int64(), true),
new Field(
"vector",
new FixedSizeList(2, new Field("item", new Float32(), true)),
true,
),
new Field("price", new Float64(), true),
new Field("category", new Utf8(), true),
]);
expect(await table.schema()).toEqual(expectedSchema);
});
it("can alter the columns in the schema", async function () {
const con = await connect(tmpDir.name);
const schema = new Schema([
+53 -20
View File
@@ -5,12 +5,15 @@ import {
Table as ArrowTable,
Data,
DataType,
Field,
IntoVector,
MultiVector,
Schema,
dataTypeToJson,
fromDataToBuffer,
fromTableToBuffer,
isMultiVector,
makeEmptyTable,
tableFromIPC,
} from "./arrow";
@@ -84,6 +87,16 @@ export interface OptimizeOptions {
* tbl.optimize({cleanupOlderThan: new Date()});
*/
cleanupOlderThan: Date;
/**
* Because they may be part of an in-progress transaction, files newer than
* 7 days old are not deleted by default. If you are sure that there are no
* in-progress transactions, then you can set this to true to delete all
* files older than `cleanupOlderThan`.
*
* **WARNING**: This should only be set to true if you can guarantee that
* no other process is currently working on this dataset. Otherwise the
* dataset could be put into a corrupted state.
*/
deleteUnverified: boolean;
}
@@ -381,15 +394,16 @@ export abstract class Table {
abstract vectorSearch(vector: IntoVector | MultiVector): VectorQuery;
/**
* Add new columns with defined values.
* @param {AddColumnsSql[]} newColumnTransforms pairs of column names and
* the SQL expression to use to calculate the value of the new column. These
* expressions will be evaluated for each row in the table, and can
* reference existing columns in the table.
* @param {AddColumnsSql[] | Field | Field[] | Schema} newColumnTransforms Either:
* - An array of objects with column names and SQL expressions to calculate values
* - A single Arrow Field defining one column with its data type (column will be initialized with null values)
* - An array of Arrow Fields defining columns with their data types (columns will be initialized with null values)
* - An Arrow Schema defining columns with their data types (columns will be initialized with null values)
* @returns {Promise<AddColumnsResult>} A promise that resolves to an object
* containing the new version number of the table after adding the columns.
*/
abstract addColumns(
newColumnTransforms: AddColumnsSql[],
newColumnTransforms: AddColumnsSql[] | Field | Field[] | Schema,
): Promise<AddColumnsResult>;
/**
@@ -501,19 +515,7 @@ export abstract class Table {
* - Index: Optimizes the indices, adding new data to existing indices
*
*
* Experimental API
* ----------------
*
* The optimization process is undergoing active development and may change.
* Our goal with these changes is to improve the performance of optimization and
* reduce the complexity.
*
* That being said, it is essential today to run optimize if you want the best
* performance. It should be stable and safe to use in production, but it our
* hope that the API may be simplified (or not even need to be called) in the
* future.
*
* The frequency an application shoudl call optimize is based on the frequency of
* 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
* you have added or modified 100,000 or more records or run more than 20 data
@@ -806,9 +808,40 @@ export class LocalTable extends Table {
// TODO: Support BatchUDF
async addColumns(
newColumnTransforms: AddColumnsSql[],
newColumnTransforms: AddColumnsSql[] | Field | Field[] | Schema,
): Promise<AddColumnsResult> {
return await this.inner.addColumns(newColumnTransforms);
// Handle single Field -> convert to array of Fields
if (newColumnTransforms instanceof Field) {
newColumnTransforms = [newColumnTransforms];
}
// Handle array of Fields -> convert to Schema
if (
Array.isArray(newColumnTransforms) &&
newColumnTransforms.length > 0 &&
newColumnTransforms[0] instanceof Field
) {
const fields = newColumnTransforms as Field[];
newColumnTransforms = new Schema(fields);
}
// Handle Schema -> use schema-based approach
if (newColumnTransforms instanceof Schema) {
const schema = newColumnTransforms;
// Convert schema to buffer using Arrow IPC format
const emptyTable = makeEmptyTable(schema);
const schemaBuf = await fromTableToBuffer(emptyTable);
return await this.inner.addColumnsWithSchema(schemaBuf);
}
// Handle SQL expressions (existing functionality)
if (Array.isArray(newColumnTransforms)) {
return await this.inner.addColumns(
newColumnTransforms as AddColumnsSql[],
);
}
throw new Error("Invalid input type for addColumns");
}
async alterColumns(
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-darwin-arm64",
"version": "0.27.0-beta.5",
"version": "0.27.1",
"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.27.0-beta.5",
"version": "0.27.1",
"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.27.0-beta.5",
"version": "0.27.1",
"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.27.0-beta.5",
"version": "0.27.1",
"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.27.0-beta.5",
"version": "0.27.1",
"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.27.0-beta.5",
"version": "0.27.1",
"os": [
"win32"
],
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-x64-msvc",
"version": "0.27.0-beta.5",
"version": "0.27.1",
"os": ["win32"],
"cpu": ["x64"],
"main": "lancedb.win32-x64-msvc.node",
+2 -2
View File
@@ -1,12 +1,12 @@
{
"name": "@lancedb/lancedb",
"version": "0.27.0-beta.5",
"version": "0.27.1",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@lancedb/lancedb",
"version": "0.27.0-beta.5",
"version": "0.27.1",
"cpu": [
"x64",
"arm64"
+1 -1
View File
@@ -11,7 +11,7 @@
"ann"
],
"private": false,
"version": "0.27.0-beta.5",
"version": "0.27.1",
"main": "dist/index.js",
"exports": {
".": "./dist/index.js",
+18 -1
View File
@@ -3,7 +3,7 @@
use std::collections::HashMap;
use lancedb::ipc::ipc_file_to_batches;
use lancedb::ipc::{ipc_file_to_batches, ipc_file_to_schema};
use lancedb::table::{
AddDataMode, ColumnAlteration as LanceColumnAlteration, Duration, NewColumnTransform,
OptimizeAction, OptimizeOptions, Table as LanceDbTable,
@@ -279,6 +279,23 @@ impl Table {
Ok(res.into())
}
#[napi(catch_unwind)]
pub async fn add_columns_with_schema(
&self,
schema_buf: Buffer,
) -> napi::Result<AddColumnsResult> {
let schema = ipc_file_to_schema(schema_buf.to_vec())
.map_err(|e| napi::Error::from_reason(format!("Failed to read IPC schema: {}", e)))?;
let transforms = NewColumnTransform::AllNulls(schema);
let res = self
.inner_ref()?
.add_columns(transforms, None)
.await
.default_error()?;
Ok(res.into())
}
#[napi(catch_unwind)]
pub async fn alter_columns(
&self,
+1 -1
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.30.0-beta.5"
current_version = "0.30.1"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb-python"
version = "0.30.0-beta.5"
version = "0.30.1"
edition.workspace = true
description = "Python bindings for LanceDB"
license.workspace = true
+2 -2
View File
@@ -2205,8 +2205,8 @@ class LanceHybridQueryBuilder(LanceQueryBuilder):
self._vector_query.select(self._columns)
self._fts_query.select(self._columns)
if self._where:
self._vector_query.where(self._where, self._postfilter)
self._fts_query.where(self._where, self._postfilter)
self._vector_query.where(self._where, not self._postfilter)
self._fts_query.where(self._where, not self._postfilter)
if self._with_row_id:
self._vector_query.with_row_id(True)
self._fts_query.with_row_id(True)
+34 -40
View File
@@ -1506,22 +1506,17 @@ class Table(ABC):
in-progress operation (e.g. appending new data) and these files will not
be deleted unless they are at least 7 days old. If delete_unverified is True
then these files will be deleted regardless of their age.
.. warning::
This should only be set to True if you can guarantee that no other
process is currently working on this dataset. Otherwise the dataset
could be put into a corrupted state.
retrain: bool, default False
This parameter is no longer used and is deprecated.
Experimental API
----------------
The optimization process is undergoing active development and may change.
Our goal with these changes is to improve the performance of optimization and
reduce the complexity.
That being said, it is essential today to run optimize if you want the best
performance. It should be stable and safe to use in production, but it our
hope that the API may be simplified (or not even need to be called) in the
future.
The frequency an application shoudl call optimize is based on the frequency of
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
you have added or modified 100,000 or more records or run more than 20 data
@@ -3047,22 +3042,17 @@ class LanceTable(Table):
in-progress operation (e.g. appending new data) and these files will not
be deleted unless they are at least 7 days old. If delete_unverified is True
then these files will be deleted regardless of their age.
.. warning::
This should only be set to True if you can guarantee that no other
process is currently working on this dataset. Otherwise the dataset
could be put into a corrupted state.
retrain: bool, default False
This parameter is no longer used and is deprecated.
Experimental API
----------------
The optimization process is undergoing active development and may change.
Our goal with these changes is to improve the performance of optimization and
reduce the complexity.
That being said, it is essential today to run optimize if you want the best
performance. It should be stable and safe to use in production, but it our
hope that the API may be simplified (or not even need to be called) in the
future.
The frequency an application shoudl call optimize is based on the frequency of
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
you have added or modified 100,000 or more records or run more than 20 data
@@ -4630,22 +4620,17 @@ class AsyncTable:
in-progress operation (e.g. appending new data) and these files will not
be deleted unless they are at least 7 days old. If delete_unverified is True
then these files will be deleted regardless of their age.
.. warning::
This should only be set to True if you can guarantee that no other
process is currently working on this dataset. Otherwise the dataset
could be put into a corrupted state.
retrain: bool, default False
This parameter is no longer used and is deprecated.
Experimental API
----------------
The optimization process is undergoing active development and may change.
Our goal with these changes is to improve the performance of optimization and
reduce the complexity.
That being said, it is essential today to run optimize if you want the best
performance. It should be stable and safe to use in production, but it our
hope that the API may be simplified (or not even need to be called) in the
future.
The frequency an application shoudl call optimize is based on the frequency of
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
you have added or modified 100,000 or more records or run more than 20 data
@@ -4766,7 +4751,16 @@ class IndexStatistics:
num_indexed_rows: int
num_unindexed_rows: int
index_type: Literal[
"IVF_PQ", "IVF_HNSW_PQ", "IVF_HNSW_SQ", "FTS", "BTREE", "BITMAP", "LABEL_LIST"
"IVF_FLAT",
"IVF_SQ",
"IVF_PQ",
"IVF_RQ",
"IVF_HNSW_SQ",
"IVF_HNSW_PQ",
"FTS",
"BTREE",
"BITMAP",
"LABEL_LIST",
]
distance_type: Optional[Literal["l2", "cosine", "dot"]] = None
num_indices: Optional[int] = None
+54
View File
@@ -177,6 +177,60 @@ async def test_analyze_plan(table: AsyncTable):
assert "metrics=" in res
@pytest.fixture
def table_with_id(tmpdir_factory) -> Table:
tmp_path = str(tmpdir_factory.mktemp("data"))
db = lancedb.connect(tmp_path)
data = pa.table(
{
"id": pa.array([1, 2, 3, 4], type=pa.int64()),
"text": pa.array(["a", "b", "cat", "dog"]),
"vector": pa.array(
[[0.1, 0.1], [2, 2], [-0.1, -0.1], [0.5, -0.5]],
type=pa.list_(pa.float32(), list_size=2),
),
}
)
table = db.create_table("test_with_id", data)
table.create_fts_index("text", with_position=False, use_tantivy=False)
return table
def test_hybrid_prefilter_explain_plan(table_with_id: Table):
"""
Verify that the prefilter logic is not inverted in LanceHybridQueryBuilder.
"""
plan_prefilter = (
table_with_id.search(query_type="hybrid")
.vector([0.0, 0.0])
.text("dog")
.where("id = 1", prefilter=True)
.limit(2)
.explain_plan(verbose=True)
)
plan_postfilter = (
table_with_id.search(query_type="hybrid")
.vector([0.0, 0.0])
.text("dog")
.where("id = 1", prefilter=False)
.limit(2)
.explain_plan(verbose=True)
)
# prefilter=True: filter is pushed into the LanceRead scan.
# The FTS sub-plan exposes this as "full_filter=id = Int64(1)" inside LanceRead.
assert "full_filter=id = Int64(1)" in plan_prefilter, (
f"Should push the filter into the scan.\nPlan:\n{plan_prefilter}"
)
# prefilter=False: filter is applied as a separate FilterExec after the search.
# The filter must NOT be embedded in the scan.
assert "full_filter=id = Int64(1)" not in plan_postfilter, (
f"Should NOT push the filter into the scan.\nPlan:\n{plan_postfilter}"
)
def test_normalize_scores():
cases = [
(pa.array([0.1, 0.4]), pa.array([0.0, 1.0])),
+22
View File
@@ -3,6 +3,7 @@
from datetime import timedelta
import random
from typing import get_args, get_type_hints
import pyarrow as pa
import pytest
@@ -22,6 +23,7 @@ from lancedb.index import (
HnswSq,
FTS,
)
from lancedb.table import IndexStatistics
@pytest_asyncio.fixture
@@ -283,3 +285,23 @@ async def test_create_index_with_binary_vectors(binary_table: AsyncTable):
for v in range(256):
res = await binary_table.query().nearest_to([v] * 128).to_arrow()
assert res["id"][0].as_py() == v
def test_index_statistics_index_type_lists_all_supported_values():
expected_index_types = {
"IVF_FLAT",
"IVF_SQ",
"IVF_PQ",
"IVF_RQ",
"IVF_HNSW_SQ",
"IVF_HNSW_PQ",
"FTS",
"BTREE",
"BITMAP",
"LABEL_LIST",
}
assert (
set(get_args(get_type_hints(IndexStatistics)["index_type"]))
== expected_index_types
)
@@ -147,7 +147,12 @@ class TrackingNamespace(LanceNamespace):
This simulates a credential rotation system where each call returns
new credentials that expire after credential_expires_in_seconds.
"""
modified = copy.deepcopy(storage_options) if storage_options else {}
# Start from base storage options (endpoint, region, allow_http, etc.)
# because DirectoryNamespace returns None for storage_options from
# describe_table/declare_table when no credential vendor is configured.
modified = copy.deepcopy(self.base_storage_options)
if storage_options:
modified.update(storage_options)
# Increment credentials to simulate rotation
modified["aws_access_key_id"] = f"AKID_{count}"
+13
View File
@@ -316,6 +316,19 @@ impl<'py> IntoPyObject<'py> for PySelect {
Select::All => Ok(py.None().into_bound(py).into_any()),
Select::Columns(columns) => Ok(columns.into_pyobject(py)?.into_any()),
Select::Dynamic(columns) => Ok(columns.into_pyobject(py)?.into_any()),
Select::Expr(pairs) => {
// Serialize DataFusion Expr -> SQL string so Python sees the same
// format as Select::Dynamic: a list of (name, sql_string) tuples.
let sql_pairs: PyResult<Vec<(String, String)>> = pairs
.into_iter()
.map(|(name, expr)| {
lancedb::expr::expr_to_sql_string(&expr)
.map(|sql| (name, sql))
.map_err(|e| PyRuntimeError::new_err(e.to_string()))
})
.collect();
Ok(sql_pairs?.into_pyobject(py)?.into_any())
}
}
}
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb"
version = "0.27.0-beta.5"
version = "0.27.1"
edition.workspace = true
description = "LanceDB: A serverless, low-latency vector database for AI applications"
license.workspace = true
@@ -339,6 +339,12 @@ impl PermutationReader {
}
Ok(false)
}
Select::Expr(columns) => {
// For Expr projections, we check if any alias is _rowid.
// We can't validate the expression itself (it may differ from _rowid)
// but we allow it through; the column will be included.
Ok(columns.iter().any(|(alias, _)| alias == ROW_ID))
}
}
}
+94
View File
@@ -47,6 +47,25 @@ pub enum Select {
///
/// See [`Query::select`] for more details and examples
Dynamic(Vec<(String, String)>),
/// Advanced selection using type-safe DataFusion expressions
///
/// Similar to [`Select::Dynamic`] but uses [`datafusion_expr::Expr`] instead of
/// raw SQL strings. Use [`crate::expr`] helpers to build expressions:
///
/// ```
/// use lancedb::expr::{col, lit};
/// use lancedb::query::Select;
///
/// // SELECT id, id * 2 AS id2 FROM ...
/// let selection = Select::expr_projection(&[
/// ("id", col("id")),
/// ("id2", col("id") * lit(2)),
/// ]);
/// ```
///
/// Note: For remote/server-side queries the expressions are serialized to SQL strings
/// automatically (same as [`Select::Dynamic`]).
Expr(Vec<(String, datafusion_expr::Expr)>),
}
impl Select {
@@ -69,6 +88,29 @@ impl Select {
.collect(),
)
}
/// Create a typed-expression projection.
///
/// This is a convenience method for creating a [`Select::Expr`] variant from
/// a slice of `(name, expr)` pairs where each `expr` is a [`datafusion_expr::Expr`].
///
/// # Example
/// ```
/// use lancedb::expr::{col, lit};
/// use lancedb::query::Select;
///
/// let selection = Select::expr_projection(&[
/// ("id", col("id")),
/// ("id2", col("id") * lit(2)),
/// ]);
/// ```
pub fn expr_projection(columns: &[(impl AsRef<str>, datafusion_expr::Expr)]) -> Self {
Self::Expr(
columns
.iter()
.map(|(name, expr)| (name.as_ref().to_string(), expr.clone()))
.collect(),
)
}
}
/// A trait for converting a type to a query vector
@@ -1591,6 +1633,58 @@ mod tests {
});
}
#[tokio::test]
async fn test_select_with_expr_projection() {
// Mirrors test_select_with_transform but uses Select::Expr instead of Select::Dynamic
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test_expr.lance");
let uri = dataset_path.to_str().unwrap();
let batches = make_non_empty_batches();
let conn = connect(uri).execute().await.unwrap();
let table = conn
.create_table("my_table", batches)
.execute()
.await
.unwrap();
use crate::expr::{col, lit};
let query = table.query().limit(10).select(Select::expr_projection(&[
("id2", col("id") * lit(2i32)),
("id", col("id")),
]));
let schema = query.output_schema().await.unwrap();
assert_eq!(
schema,
Arc::new(ArrowSchema::new(vec![
ArrowField::new("id2", DataType::Int32, true),
ArrowField::new("id", DataType::Int32, true),
]))
);
let result = query.execute().await;
let mut batches = result
.expect("should have result")
.try_collect::<Vec<_>>()
.await
.unwrap();
assert_eq!(batches.len(), 1);
let batch = batches.pop().unwrap();
// id and id2
assert_eq!(batch.num_columns(), 2);
let id: &Int32Array = batch.column_by_name("id").unwrap().as_primitive();
let id2: &Int32Array = batch.column_by_name("id2").unwrap().as_primitive();
id.iter().zip(id2.iter()).for_each(|(id, id2)| {
let id = id.unwrap();
let id2 = id2.unwrap();
assert_eq!(id * 2, id2);
});
}
#[tokio::test]
async fn test_execute_no_vector() {
// TODO: Switch back to memory://foo after https://github.com/lancedb/lancedb/issues/1051
+4
View File
@@ -72,6 +72,10 @@ impl ServerVersion {
pub fn support_structural_fts(&self) -> bool {
self.0 >= semver::Version::new(0, 3, 0)
}
pub fn support_multipart_write(&self) -> bool {
self.0 >= semver::Version::new(0, 4, 0)
}
}
pub const OPT_REMOTE_PREFIX: &str = "remote_database_";
+811 -65
View File
@@ -10,6 +10,7 @@ use super::ARROW_STREAM_CONTENT_TYPE;
use super::client::RequestResultExt;
use super::client::{HttpSend, RestfulLanceDbClient, Sender};
use super::db::ServerVersion;
use crate::data::scannable::{PeekedScannable, Scannable, estimate_write_partitions};
use crate::index::Index;
use crate::index::IndexStatistics;
use crate::index::waiter::wait_for_index;
@@ -23,7 +24,7 @@ use crate::table::MergeResult;
use crate::table::Tags;
use crate::table::UpdateResult;
use crate::table::query::create_multi_vector_plan;
use crate::table::{AnyQuery, Filter, TableStatistics};
use crate::table::{AnyQuery, Filter, PreprocessingOutput, TableStatistics};
use crate::utils::background_cache::BackgroundCache;
use crate::utils::{supported_btree_data_type, supported_vector_data_type};
use crate::{DistanceType, Error};
@@ -43,7 +44,7 @@ use async_trait::async_trait;
use datafusion_common::DataFusionError;
use datafusion_physical_plan::stream::RecordBatchStreamAdapter;
use datafusion_physical_plan::{ExecutionPlan, RecordBatchStream, SendableRecordBatchStream};
use futures::TryStreamExt;
use futures::{StreamExt, TryStreamExt};
use http::header::CONTENT_TYPE;
use http::{HeaderName, StatusCode};
use lance::arrow::json::{JsonDataType, JsonSchema};
@@ -477,6 +478,16 @@ impl<S: HttpSend> RemoteTable<S> {
}));
body["columns"] = alias_map.into();
}
Select::Expr(pairs) => {
let alias_map: Result<serde_json::Map<String, serde_json::Value>> = pairs
.iter()
.map(|(name, expr)| {
expr_to_sql_string(expr)
.map(|sql| (name.clone(), serde_json::Value::String(sql)))
})
.collect();
body["columns"] = alias_map?.into();
}
}
if params.fast_search {
@@ -604,6 +615,66 @@ impl<S: HttpSend> RemoteTable<S> {
Ok(bodies)
}
async fn create_multipart_write(&self) -> Result<String> {
let request = self.client.post(&format!(
"/v1/table/{}/multipart_write/create",
self.identifier
));
let (request_id, response) = self.send(request, true).await?;
let response = self.check_table_response(&request_id, response).await?;
let body = response.text().await.err_to_http(request_id.clone())?;
let parsed: serde_json::Value = serde_json::from_str(&body).map_err(|e| Error::Http {
source: format!("Failed to parse multipart create response: {}", e).into(),
request_id,
status_code: None,
})?;
parsed["upload_id"]
.as_str()
.map(|s| s.to_string())
.ok_or_else(|| Error::Http {
source: "Missing upload_id in multipart create response".into(),
request_id: String::new(),
status_code: None,
})
}
async fn complete_multipart_write(&self, upload_id: &str) -> Result<AddResult> {
let request = self
.client
.post(&format!(
"/v1/table/{}/multipart_write/complete",
self.identifier
))
.query(&[("upload_id", upload_id)]);
let (request_id, response) = self.send(request, true).await?;
let response = self.check_table_response(&request_id, response).await?;
let body = response.text().await.err_to_http(request_id.clone())?;
let parsed: serde_json::Value = serde_json::from_str(&body).map_err(|e| Error::Http {
source: format!("Failed to parse multipart complete response: {}", e).into(),
request_id,
status_code: None,
})?;
let version = parsed["version"].as_u64().ok_or_else(|| Error::Http {
source: "Missing version in multipart complete response".into(),
request_id: String::new(),
status_code: None,
})?;
Ok(AddResult { version })
}
async fn abort_multipart_write(&self, upload_id: &str) -> Result<()> {
let request = self
.client
.post(&format!(
"/v1/table/{}/multipart_write/abort",
self.identifier
))
.query(&[("upload_id", upload_id)]);
let (request_id, response) = self.send(request, true).await?;
self.check_table_response(&request_id, response).await?;
Ok(())
}
async fn check_mutable(&self) -> Result<()> {
let read_guard = self.version.read().await;
match *read_guard {
@@ -807,6 +878,19 @@ mod test_utils {
}
pub fn new_mock_with_config<F, T>(name: String, handler: F, config: ClientConfig) -> Self
where
F: Fn(reqwest::Request) -> http::Response<T> + Send + Sync + 'static,
T: Into<reqwest::Body>,
{
Self::new_mock_with_version_and_config(name, handler, None, config)
}
pub fn new_mock_with_version_and_config<F, T>(
name: String,
handler: F,
version: Option<semver::Version>,
config: ClientConfig,
) -> Self
where
F: Fn(reqwest::Request) -> http::Response<T> + Send + Sync + 'static,
T: Into<reqwest::Body>,
@@ -817,7 +901,7 @@ mod test_utils {
name: name.clone(),
namespace: vec![],
identifier: name,
server_version: ServerVersion::default(),
server_version: version.map(ServerVersion).unwrap_or_default(),
version: RwLock::new(None),
location: RwLock::new(None),
schema_cache: BackgroundCache::new(SCHEMA_CACHE_TTL, SCHEMA_CACHE_REFRESH_WINDOW),
@@ -826,6 +910,185 @@ mod test_utils {
}
}
impl<S: HttpSend + 'static> RemoteTable<S> {
fn is_retryable_write_error(&self, err: &Error) -> bool {
match err {
Error::Http {
source,
status_code,
..
} => {
// Don't retry read errors (is_body/is_decode): the
// server may have committed the write already, and
// without an idempotency key we'd duplicate data.
source
.downcast_ref::<reqwest::Error>()
.is_some_and(|e| e.is_connect())
|| status_code.is_some_and(|s| self.client.retry_config.statuses.contains(&s))
}
// send_with_retry exhausted its internal retries on a retryable
// status. The outer loop can still retry the whole operation with
// a fresh session.
Error::Retry { status_code, .. } => {
status_code.is_some_and(|s| self.client.retry_config.statuses.contains(&s))
}
_ => false,
}
}
async fn add_single_partition(&self, output: PreprocessingOutput) -> Result<AddResult> {
use crate::remote::retry::RetryCounter;
let mut insert: Arc<dyn ExecutionPlan> = Arc::new(RemoteInsertExec::new(
self.name.clone(),
self.identifier.clone(),
self.client.clone(),
output.plan,
output.overwrite,
));
let mut retry_counter =
RetryCounter::new(&self.client.retry_config, uuid::Uuid::new_v4().to_string());
loop {
let stream = execute_plan(insert.clone(), Default::default())?;
let result: Result<Vec<_>> = stream.try_collect().await.map_err(Error::from);
match result {
Ok(_) => {
let add_result = insert
.as_any()
.downcast_ref::<RemoteInsertExec<S>>()
.and_then(|i| i.add_result())
.unwrap_or(AddResult { version: 0 });
if output.overwrite {
self.invalidate_schema_cache();
}
return Ok(add_result);
}
Err(err) if output.rescannable && self.is_retryable_write_error(&err) => {
retry_counter.increment_from_error(err)?;
tokio::time::sleep(retry_counter.next_sleep_time()).await;
insert = insert.reset_state()?;
continue;
}
Err(err) => return Err(err),
}
}
}
async fn add_multipart(
&self,
output: PreprocessingOutput,
num_partitions: usize,
) -> Result<AddResult> {
use crate::remote::retry::RetryCounter;
let mut retry_counter =
RetryCounter::new(&self.client.retry_config, uuid::Uuid::new_v4().to_string());
loop {
let upload_id = self.create_multipart_write().await?;
let result = self
.execute_multipart_inserts(&upload_id, &output, num_partitions)
.await;
match result {
Ok(()) => match self.complete_multipart_write(&upload_id).await {
Ok(result) => {
if output.overwrite {
self.invalidate_schema_cache();
}
return Ok(result);
}
Err(e) => {
if let Err(abort_err) = self.abort_multipart_write(&upload_id).await {
log::warn!(
"Failed to abort multipart write {}: {}",
upload_id,
abort_err
);
}
if output.rescannable && self.is_retryable_write_error(&e) {
retry_counter.increment_from_error(e)?;
tokio::time::sleep(retry_counter.next_sleep_time()).await;
continue;
}
return Err(e);
}
},
Err(e) => {
if let Err(abort_err) = self.abort_multipart_write(&upload_id).await {
log::warn!(
"Failed to abort multipart write {}: {}",
upload_id,
abort_err
);
}
if output.rescannable && self.is_retryable_write_error(&e) {
retry_counter.increment_from_error(e)?;
tokio::time::sleep(retry_counter.next_sleep_time()).await;
continue;
}
return Err(e);
}
}
}
}
async fn execute_multipart_inserts(
&self,
upload_id: &str,
output: &PreprocessingOutput,
num_partitions: usize,
) -> Result<()> {
let plan = Arc::new(
datafusion_physical_plan::repartition::RepartitionExec::try_new(
output.plan.clone(),
datafusion_physical_plan::Partitioning::RoundRobinBatch(num_partitions),
)?,
) as Arc<dyn ExecutionPlan>;
let insert = Arc::new(RemoteInsertExec::new_multipart(
self.name.clone(),
self.identifier.clone(),
self.client.clone(),
plan,
output.overwrite,
upload_id.to_string(),
));
let task_ctx = Arc::new(datafusion_execution::TaskContext::default());
let mut join_set = tokio::task::JoinSet::new();
for partition in 0..num_partitions {
let exec = insert.clone();
let ctx = task_ctx.clone();
join_set.spawn(async move {
let mut stream = exec
.execute(partition, ctx)
.map_err(|e| -> Error { e.into() })?;
while let Some(batch) = stream.next().await {
batch.map_err(|e| -> Error { e.into() })?;
}
Ok::<_, Error>(())
});
}
// JoinSet aborts all remaining tasks when dropped, so if we return
// early on error the orphaned tasks are automatically cancelled.
while let Some(result) = join_set.join_next().await {
result.map_err(|e| Error::Runtime {
message: format!("Insert task panicked: {}", e),
})??;
}
Ok(())
}
}
#[async_trait]
impl<S: HttpSend> BaseTable for RemoteTable<S> {
fn as_any(&self) -> &dyn std::any::Any {
@@ -976,74 +1239,44 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
status_code: None,
})
}
async fn add(&self, add: AddDataBuilder) -> Result<AddResult> {
use crate::remote::retry::RetryCounter;
async fn add(&self, mut add: AddDataBuilder) -> Result<AddResult> {
self.check_mutable().await?;
let table_schema = self.schema().await?;
let table_def = TableDefinition::try_from_rich_schema(table_schema.clone())?;
let num_partitions = if let Some(parallelism) = add.write_parallelism {
if parallelism > 1 && self.server_version.support_multipart_write() {
parallelism
} else {
1
}
} else if self.server_version.support_multipart_write() {
// Peek at the first batch to estimate write partitions, same as NativeTable.
let mut peeked = PeekedScannable::new(add.data);
let n = if let Some(first_batch) = peeked.peek().await {
let max_partitions = lance_core::utils::tokio::get_num_compute_intensive_cpus();
estimate_write_partitions(
first_batch.get_array_memory_size(),
first_batch.num_rows(),
peeked.num_rows(),
max_partitions,
)
} else {
1
};
add.data = Box::new(peeked);
n
} else {
1
};
let output = add.into_plan(&table_schema, &table_def)?;
let mut insert: Arc<dyn ExecutionPlan> = Arc::new(RemoteInsertExec::new(
self.name.clone(),
self.identifier.clone(),
self.client.clone(),
output.plan,
output.overwrite,
));
let mut retry_counter =
RetryCounter::new(&self.client.retry_config, uuid::Uuid::new_v4().to_string());
loop {
let stream = execute_plan(insert.clone(), Default::default())?;
let result: Result<Vec<_>> = stream.try_collect().await.map_err(Error::from);
match result {
Ok(_) => {
let add_result = insert
.as_any()
.downcast_ref::<RemoteInsertExec<S>>()
.and_then(|i| i.add_result())
.unwrap_or(AddResult { version: 0 });
if output.overwrite {
self.invalidate_schema_cache();
}
return Ok(add_result);
}
Err(err) if output.rescannable => {
let retryable = match &err {
Error::Http {
source,
status_code,
..
} => {
// Don't retry read errors (is_body/is_decode): the
// server may have committed the write already, and
// without an idempotency key we'd duplicate data.
source
.downcast_ref::<reqwest::Error>()
.is_some_and(|e| e.is_connect())
|| status_code
.is_some_and(|s| self.client.retry_config.statuses.contains(&s))
}
_ => false,
};
if retryable {
retry_counter.increment_from_error(err)?;
tokio::time::sleep(retry_counter.next_sleep_time()).await;
insert = insert.reset_state()?;
continue;
}
return Err(err);
}
Err(err) => return Err(err),
}
if num_partitions > 1 {
self.add_multipart(output, num_partitions).await
} else {
self.add_single_partition(output).await
}
}
@@ -1801,6 +2034,7 @@ mod tests {
use super::*;
use crate::remote::client::{ClientConfig, RetryConfig};
use crate::table::AddDataMode;
use arrow::{array::AsArray, compute::concat_batches, datatypes::Int32Type};
@@ -4821,4 +5055,516 @@ mod tests {
assert_eq!(data.len(), 1);
assert_eq!(data[0].as_ref().unwrap(), &expected_data);
}
fn schema_json() -> &'static str {
r#"{"fields": [{"name": "id", "type": {"type": "int32"}, "nullable": true}]}"#
}
fn simple_describe_response() -> http::Response<String> {
http::Response::builder()
.status(200)
.body(format!(r#"{{"version": 1, "schema": {}}}"#, schema_json()))
.unwrap()
}
#[tokio::test]
async fn test_multipart_write_happy_path() {
use std::sync::Mutex;
let create_count = Arc::new(AtomicUsize::new(0));
let insert_count = Arc::new(AtomicUsize::new(0));
let complete_count = Arc::new(AtomicUsize::new(0));
let abort_count = Arc::new(AtomicUsize::new(0));
let upload_ids = Arc::new(Mutex::new(Vec::<String>::new()));
let create_count_c = create_count.clone();
let insert_count_c = insert_count.clone();
let complete_count_c = complete_count.clone();
let abort_count_c = abort_count.clone();
let upload_ids_c = upload_ids.clone();
let table = Table::new_with_handler_version(
"my_table",
semver::Version::new(0, 4, 0),
move |request| {
let path = request.url().path();
let query = request.url().query().unwrap_or("");
if path == "/v1/table/my_table/describe/" {
return simple_describe_response();
}
if path == "/v1/table/my_table/multipart_write/create" {
create_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(r#"{"upload_id": "test-upload-123"}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/insert/" {
insert_count_c.fetch_add(1, Ordering::SeqCst);
let uid = url::form_urlencoded::parse(query.as_bytes())
.find(|(k, _)| k == "upload_id")
.map(|(_, v)| v.to_string());
upload_ids_c
.lock()
.unwrap()
.push(uid.expect("missing upload_id on insert"));
return http::Response::builder()
.status(200)
.body(r#"{"version": 1}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/complete" {
complete_count_c.fetch_add(1, Ordering::SeqCst);
let uid = url::form_urlencoded::parse(query.as_bytes())
.find(|(k, _)| k == "upload_id")
.map(|(_, v)| v.to_string());
upload_ids_c
.lock()
.unwrap()
.push(uid.expect("missing upload_id on complete"));
return http::Response::builder()
.status(200)
.body(r#"{"version": 5}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/abort" {
abort_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(String::new())
.unwrap();
}
panic!("Unexpected request path: {}", path);
},
);
let batch = record_batch!(("id", Int32, [1, 2, 3])).unwrap();
let result = table
.add(vec![batch])
.write_parallelism(2)
.execute()
.await
.unwrap();
assert_eq!(result.version, 5);
assert_eq!(create_count.load(Ordering::SeqCst), 1);
assert!(
insert_count.load(Ordering::SeqCst) > 1,
"Expected multiple insert calls, got {}",
insert_count.load(Ordering::SeqCst)
);
assert_eq!(complete_count.load(Ordering::SeqCst), 1);
assert_eq!(abort_count.load(Ordering::SeqCst), 0);
let ids = upload_ids.lock().unwrap();
assert!(
ids.iter().all(|id| id == "test-upload-123"),
"All requests should use the same upload_id, got: {:?}",
*ids
);
}
#[tokio::test]
async fn test_multipart_write_fallback_old_server() {
let insert_count = Arc::new(AtomicUsize::new(0));
let create_count = Arc::new(AtomicUsize::new(0));
let insert_count_c = insert_count.clone();
let create_count_c = create_count.clone();
// Server version 0.3.0 does not support multipart writes
let table = Table::new_with_handler_version(
"my_table",
semver::Version::new(0, 3, 0),
move |request| {
let path = request.url().path();
if path == "/v1/table/my_table/describe/" {
return simple_describe_response();
}
if path.contains("multipart_write") {
create_count_c.fetch_add(1, Ordering::SeqCst);
panic!("Should not call multipart write endpoints on old server");
}
if path == "/v1/table/my_table/insert/" {
let query = request.url().query().unwrap_or("");
assert!(
!query.contains("upload_id"),
"Should not have upload_id for old server"
);
insert_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(r#"{"version": 2}"#.to_string())
.unwrap();
}
panic!("Unexpected request path: {}", path);
},
);
let batch = record_batch!(("id", Int32, [1, 2, 3])).unwrap();
let result = table
.add(vec![batch])
.write_parallelism(2)
.execute()
.await
.unwrap();
assert_eq!(result.version, 2);
assert_eq!(create_count.load(Ordering::SeqCst), 0);
assert_eq!(insert_count.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_multipart_write_small_data_single_partition() {
let insert_count = Arc::new(AtomicUsize::new(0));
let create_count = Arc::new(AtomicUsize::new(0));
let insert_count_c = insert_count.clone();
let create_count_c = create_count.clone();
let table = Table::new_with_handler_version(
"my_table",
semver::Version::new(0, 4, 0),
move |request| {
let path = request.url().path();
if path == "/v1/table/my_table/describe/" {
return simple_describe_response();
}
if path.contains("multipart_write") {
create_count_c.fetch_add(1, Ordering::SeqCst);
panic!("Should not call multipart write endpoints for small data");
}
if path == "/v1/table/my_table/insert/" {
let query = request.url().query().unwrap_or("");
assert!(
!query.contains("upload_id"),
"Should not have upload_id for small data"
);
insert_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(r#"{"version": 2}"#.to_string())
.unwrap();
}
panic!("Unexpected request path: {}", path);
},
);
// Small data: only 3 rows
let batch = record_batch!(("id", Int32, [1, 2, 3])).unwrap();
let result = table.add(vec![batch]).execute().await.unwrap();
assert_eq!(result.version, 2);
assert_eq!(create_count.load(Ordering::SeqCst), 0);
assert_eq!(insert_count.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_multipart_write_abort_on_insert_failure() {
let create_count = Arc::new(AtomicUsize::new(0));
let insert_count = Arc::new(AtomicUsize::new(0));
let complete_count = Arc::new(AtomicUsize::new(0));
let abort_count = Arc::new(AtomicUsize::new(0));
let create_count_c = create_count.clone();
let insert_count_c = insert_count.clone();
let complete_count_c = complete_count.clone();
let abort_count_c = abort_count.clone();
let table = Table::new_with_handler_version(
"my_table",
semver::Version::new(0, 4, 0),
move |request| {
let path = request.url().path();
if path == "/v1/table/my_table/describe/" {
return simple_describe_response();
}
if path == "/v1/table/my_table/multipart_write/create" {
create_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(r#"{"upload_id": "test-upload-456"}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/insert/" {
let count = insert_count_c.fetch_add(1, Ordering::SeqCst);
// Fail on the first insert with non-retryable status
if count == 0 {
return http::Response::builder()
.status(400)
.body("Bad Request".to_string())
.unwrap();
}
return http::Response::builder()
.status(200)
.body(r#"{"version": 1}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/complete" {
complete_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(r#"{"version": 5}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/abort" {
abort_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(String::new())
.unwrap();
}
panic!("Unexpected request path: {}", path);
},
);
let batch = record_batch!(("id", Int32, [1, 2, 3])).unwrap();
let result = table.add(vec![batch]).write_parallelism(2).execute().await;
assert!(result.is_err());
assert_eq!(create_count.load(Ordering::SeqCst), 1);
assert_eq!(complete_count.load(Ordering::SeqCst), 0);
assert_eq!(abort_count.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_multipart_write_abort_on_complete_failure() {
let abort_count = Arc::new(AtomicUsize::new(0));
let abort_count_c = abort_count.clone();
let table = Table::new_with_handler_version(
"my_table",
semver::Version::new(0, 4, 0),
move |request| {
let path = request.url().path();
if path == "/v1/table/my_table/describe/" {
return simple_describe_response();
}
if path == "/v1/table/my_table/multipart_write/create" {
return http::Response::builder()
.status(200)
.body(r#"{"upload_id": "test-upload-789"}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/insert/" {
return http::Response::builder()
.status(200)
.body(r#"{"version": 1}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/complete" {
return http::Response::builder()
.status(400)
.body("Bad Request".to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/abort" {
abort_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(String::new())
.unwrap();
}
panic!("Unexpected request path: {}", path);
},
);
let batch = record_batch!(("id", Int32, [1, 2, 3])).unwrap();
let result = table.add(vec![batch]).write_parallelism(2).execute().await;
assert!(result.is_err());
assert_eq!(abort_count.load(Ordering::SeqCst), 1);
}
fn retry_config_no_backoff() -> ClientConfig {
ClientConfig {
retry_config: RetryConfig {
retries: Some(3),
connect_retries: Some(3),
read_retries: Some(3),
backoff_factor: Some(0.0),
backoff_jitter: Some(0.0),
statuses: Some(vec![502, 503]),
},
..Default::default()
}
}
#[tokio::test]
async fn test_multipart_write_retry_on_partition_failure() {
// All inserts for the first upload session return 503 (retryable).
// After exhausting internal retries, the outer loop retries with a
// new session and succeeds.
let create_count = Arc::new(AtomicUsize::new(0));
let complete_count = Arc::new(AtomicUsize::new(0));
let abort_count = Arc::new(AtomicUsize::new(0));
let create_count_c = create_count.clone();
let complete_count_c = complete_count.clone();
let abort_count_c = abort_count.clone();
let table = Table::new_with_handler_version_and_config(
"my_table",
semver::Version::new(0, 4, 0),
move |request| {
let path = request.url().path();
let query = request.url().query().unwrap_or("");
if path == "/v1/table/my_table/describe/" {
return simple_describe_response();
}
if path == "/v1/table/my_table/multipart_write/create" {
let n = create_count_c.fetch_add(1, Ordering::SeqCst);
let body = format!(r#"{{"upload_id": "upload-{}"}}"#, n + 1);
return http::Response::builder().status(200).body(body).unwrap();
}
if path == "/v1/table/my_table/insert/" {
// Fail all inserts for the first session
if query.contains("upload_id=upload-1") {
return http::Response::builder()
.status(503)
.body("Service Unavailable".to_string())
.unwrap();
}
return http::Response::builder()
.status(200)
.body(r#"{"version": 1}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/complete" {
complete_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(r#"{"version": 7}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/abort" {
abort_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(String::new())
.unwrap();
}
panic!("Unexpected request path: {}", path);
},
retry_config_no_backoff(),
);
let batch = record_batch!(("id", Int32, [1, 2, 3])).unwrap();
let result = table
.add(vec![batch])
.write_parallelism(2)
.execute()
.await
.unwrap();
assert_eq!(result.version, 7);
assert_eq!(create_count.load(Ordering::SeqCst), 2);
assert_eq!(abort_count.load(Ordering::SeqCst), 1);
assert_eq!(complete_count.load(Ordering::SeqCst), 1);
}
#[tokio::test]
async fn test_multipart_write_retry_on_complete_failure() {
// Complete returns 503 for the first session, succeeds for the second.
let create_count = Arc::new(AtomicUsize::new(0));
let abort_count = Arc::new(AtomicUsize::new(0));
let create_count_c = create_count.clone();
let abort_count_c = abort_count.clone();
let table = Table::new_with_handler_version_and_config(
"my_table",
semver::Version::new(0, 4, 0),
move |request| {
let path = request.url().path();
let query = request.url().query().unwrap_or("");
if path == "/v1/table/my_table/describe/" {
return simple_describe_response();
}
if path == "/v1/table/my_table/multipart_write/create" {
let n = create_count_c.fetch_add(1, Ordering::SeqCst);
let body = format!(r#"{{"upload_id": "upload-{}"}}"#, n + 1);
return http::Response::builder().status(200).body(body).unwrap();
}
if path == "/v1/table/my_table/insert/" {
return http::Response::builder()
.status(200)
.body(r#"{"version": 1}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/complete" {
// Fail complete for first session
if query.contains("upload_id=upload-1") {
return http::Response::builder()
.status(503)
.body("Service Unavailable".to_string())
.unwrap();
}
return http::Response::builder()
.status(200)
.body(r#"{"version": 9}"#.to_string())
.unwrap();
}
if path == "/v1/table/my_table/multipart_write/abort" {
abort_count_c.fetch_add(1, Ordering::SeqCst);
return http::Response::builder()
.status(200)
.body(String::new())
.unwrap();
}
panic!("Unexpected request path: {}", path);
},
retry_config_no_backoff(),
);
let batch = record_batch!(("id", Int32, [1, 2, 3])).unwrap();
let result = table
.add(vec![batch])
.write_parallelism(2)
.execute()
.await
.unwrap();
assert_eq!(result.version, 9);
assert_eq!(create_count.load(Ordering::SeqCst), 2);
assert_eq!(abort_count.load(Ordering::SeqCst), 1);
}
}
+86 -30
View File
@@ -12,7 +12,9 @@ use datafusion_common::{DataFusionError, Result as DataFusionResult};
use datafusion_execution::{SendableRecordBatchStream, TaskContext};
use datafusion_physical_expr::EquivalenceProperties;
use datafusion_physical_plan::stream::RecordBatchStreamAdapter;
use datafusion_physical_plan::{DisplayAs, DisplayFormatType, ExecutionPlan, PlanProperties};
use datafusion_physical_plan::{
DisplayAs, DisplayFormatType, ExecutionPlan, ExecutionPlanProperties, PlanProperties,
};
use futures::StreamExt;
use http::header::CONTENT_TYPE;
@@ -25,10 +27,12 @@ use crate::table::datafusion::insert::COUNT_SCHEMA;
/// ExecutionPlan for inserting data into a remote LanceDB table.
///
/// This plan:
/// 1. Requires single partition (no parallel remote inserts yet)
/// 2. Streams data as Arrow IPC to `/v1/table/{id}/insert/` endpoint
/// 3. Stores AddResult for retrieval after execution
/// Streams data as Arrow IPC to `/v1/table/{id}/insert/` endpoint.
///
/// When `upload_id` is set, inserts are staged as part of a multipart write
/// session and the plan supports multiple partitions for parallel uploads.
/// Without `upload_id`, the plan requires a single partition and commits
/// immediately.
#[derive(Debug)]
pub struct RemoteInsertExec<S: HttpSend = Sender> {
table_name: String,
@@ -38,10 +42,11 @@ pub struct RemoteInsertExec<S: HttpSend = Sender> {
overwrite: bool,
properties: PlanProperties,
add_result: Arc<Mutex<Option<AddResult>>>,
upload_id: Option<String>,
}
impl<S: HttpSend + 'static> RemoteInsertExec<S> {
/// Create a new RemoteInsertExec.
/// Create a new single-partition RemoteInsertExec.
pub fn new(
table_name: String,
identifier: String,
@@ -49,10 +54,49 @@ impl<S: HttpSend + 'static> RemoteInsertExec<S> {
input: Arc<dyn ExecutionPlan>,
overwrite: bool,
) -> Self {
Self::new_inner(table_name, identifier, client, input, overwrite, None)
}
/// Create a multi-partition RemoteInsertExec for use with multipart writes.
///
/// Each partition's insert is staged under the given `upload_id` without
/// committing. The caller is responsible for calling the complete (or abort)
/// endpoint after all partitions finish.
pub fn new_multipart(
table_name: String,
identifier: String,
client: RestfulLanceDbClient<S>,
input: Arc<dyn ExecutionPlan>,
overwrite: bool,
upload_id: String,
) -> Self {
Self::new_inner(
table_name,
identifier,
client,
input,
overwrite,
Some(upload_id),
)
}
fn new_inner(
table_name: String,
identifier: String,
client: RestfulLanceDbClient<S>,
input: Arc<dyn ExecutionPlan>,
overwrite: bool,
upload_id: Option<String>,
) -> Self {
let num_partitions = if upload_id.is_some() {
input.output_partitioning().partition_count()
} else {
1
};
let schema = COUNT_SCHEMA.clone();
let properties = PlanProperties::new(
EquivalenceProperties::new(schema),
datafusion_physical_plan::Partitioning::UnknownPartitioning(1),
datafusion_physical_plan::Partitioning::UnknownPartitioning(num_partitions),
datafusion_physical_plan::execution_plan::EmissionType::Final,
datafusion_physical_plan::execution_plan::Boundedness::Bounded,
);
@@ -65,6 +109,7 @@ impl<S: HttpSend + 'static> RemoteInsertExec<S> {
overwrite,
properties,
add_result: Arc::new(Mutex::new(None)),
upload_id,
}
}
@@ -174,8 +219,11 @@ impl<S: HttpSend + 'static> ExecutionPlan for RemoteInsertExec<S> {
}
fn required_input_distribution(&self) -> Vec<datafusion_physical_plan::Distribution> {
// Until we have a separate commit endpoint, we need to do all inserts in a single partition
vec![datafusion_physical_plan::Distribution::SinglePartition]
if self.upload_id.is_some() {
vec![datafusion_physical_plan::Distribution::UnspecifiedDistribution]
} else {
vec![datafusion_physical_plan::Distribution::SinglePartition]
}
}
fn benefits_from_input_partitioning(&self) -> Vec<bool> {
@@ -191,12 +239,13 @@ impl<S: HttpSend + 'static> ExecutionPlan for RemoteInsertExec<S> {
"RemoteInsertExec requires exactly one child".to_string(),
));
}
Ok(Arc::new(Self::new(
Ok(Arc::new(Self::new_inner(
self.table_name.clone(),
self.identifier.clone(),
self.client.clone(),
children[0].clone(),
self.overwrite,
self.upload_id.clone(),
)))
}
@@ -205,18 +254,20 @@ impl<S: HttpSend + 'static> ExecutionPlan for RemoteInsertExec<S> {
partition: usize,
context: Arc<TaskContext>,
) -> DataFusionResult<SendableRecordBatchStream> {
if partition != 0 {
if self.upload_id.is_none() && partition != 0 {
return Err(DataFusionError::Internal(
"RemoteInsertExec only supports single partition execution".to_string(),
"RemoteInsertExec only supports single partition execution without upload_id"
.to_string(),
));
}
let input_stream = self.input.execute(0, context)?;
let input_stream = self.input.execute(partition, context)?;
let client = self.client.clone();
let identifier = self.identifier.clone();
let overwrite = self.overwrite;
let add_result = self.add_result.clone();
let table_name = self.table_name.clone();
let upload_id = self.upload_id.clone();
let stream = futures::stream::once(async move {
let mut request = client
@@ -226,6 +277,9 @@ impl<S: HttpSend + 'static> ExecutionPlan for RemoteInsertExec<S> {
if overwrite {
request = request.query(&[("mode", "overwrite")]);
}
if let Some(ref uid) = upload_id {
request = request.query(&[("upload_id", uid.as_str())]);
}
let (error_tx, mut error_rx) = tokio::sync::oneshot::channel();
let body = Self::stream_as_http_body(input_stream, error_tx)?;
@@ -262,28 +316,30 @@ impl<S: HttpSend + 'static> ExecutionPlan for RemoteInsertExec<S> {
let (request_id, response) = result?;
let body_text = response.text().await.map_err(|e| {
DataFusionError::External(Box::new(Error::Http {
source: Box::new(e),
request_id: request_id.clone(),
status_code: None,
}))
})?;
let parsed_result = if body_text.trim().is_empty() {
// Backward compatible with old servers
AddResult { version: 0 }
} else {
serde_json::from_str(&body_text).map_err(|e| {
// For multipart writes, the staging response is not the final
// version. Only parse AddResult for non-multipart inserts.
if upload_id.is_none() {
let body_text = response.text().await.map_err(|e| {
DataFusionError::External(Box::new(Error::Http {
source: format!("Failed to parse add response: {}", e).into(),
source: Box::new(e),
request_id: request_id.clone(),
status_code: None,
}))
})?
};
})?;
let parsed_result = if body_text.trim().is_empty() {
// Backward compatible with old servers
AddResult { version: 0 }
} else {
serde_json::from_str(&body_text).map_err(|e| {
DataFusionError::External(Box::new(Error::Http {
source: format!("Failed to parse add response: {}", e).into(),
request_id: request_id.clone(),
status_code: None,
}))
})?
};
{
let mut res_lock = add_result.lock().map_err(|_| {
DataFusionError::Execution("Failed to acquire lock for add_result".to_string())
})?;
+49 -24
View File
@@ -75,6 +75,8 @@ pub mod query;
pub mod schema_evolution;
pub mod update;
use crate::index::waiter::wait_for_index;
#[cfg(feature = "remote")]
pub(crate) use add_data::PreprocessingOutput;
pub use add_data::{AddDataBuilder, AddDataMode, AddResult, NaNVectorBehavior};
pub use chrono::Duration;
pub use delete::DeleteResult;
@@ -440,6 +442,34 @@ mod test_utils {
embedding_registry: Arc::new(MemoryRegistry::new()),
}
}
pub fn new_with_handler_version_and_config<T>(
name: impl Into<String>,
version: semver::Version,
handler: impl Fn(reqwest::Request) -> http::Response<T> + Clone + Send + Sync + 'static,
config: crate::remote::ClientConfig,
) -> Self
where
T: Into<reqwest::Body>,
{
let inner = Arc::new(
crate::remote::table::RemoteTable::new_mock_with_version_and_config(
name.into(),
handler.clone(),
Some(version),
config.clone(),
),
);
let database = Arc::new(crate::remote::db::RemoteDatabase::new_mock_with_config(
handler, config,
));
Self {
inner,
database: Some(database),
// Registry is unused.
embedding_registry: Arc::new(MemoryRegistry::new()),
}
}
}
}
@@ -951,17 +981,7 @@ impl Table {
/// * Prune: Removes old versions of the dataset
/// * Index: Optimizes the indices, adding new data to existing indices
///
/// <section class="warning">Experimental API</section>
///
/// The optimization process is undergoing active development and may change.
/// Our goal with these changes is to improve the performance of optimization and
/// reduce the complexity.
///
/// That being said, it is essential today to run optimize if you want the best
/// performance. It should be stable and safe to use in production, but it our
/// hope that the API may be simplified (or not even need to be called) in the future.
///
/// The frequency an application shoudl call optimize is based on the frequency of
/// 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
/// you have added or modified 100,000 or more records or run more than 20 data
@@ -2208,21 +2228,26 @@ impl BaseTable for NativeTable {
let table_schema = Schema::from(&ds.schema().clone());
// Peek at the first batch to estimate a good partition count for
// write parallelism.
let mut peeked = PeekedScannable::new(add.data);
let num_partitions = if let Some(first_batch) = peeked.peek().await {
let max_partitions = lance_core::utils::tokio::get_num_compute_intensive_cpus();
estimate_write_partitions(
first_batch.get_array_memory_size(),
first_batch.num_rows(),
peeked.num_rows(),
max_partitions,
)
let num_partitions = if let Some(parallelism) = add.write_parallelism {
parallelism
} else {
1
// Peek at the first batch to estimate a good partition count for
// write parallelism.
let mut peeked = PeekedScannable::new(add.data);
let n = if let Some(first_batch) = peeked.peek().await {
let max_partitions = lance_core::utils::tokio::get_num_compute_intensive_cpus();
estimate_write_partitions(
first_batch.get_array_memory_size(),
first_batch.num_rows(),
peeked.num_rows(),
max_partitions,
)
} else {
1
};
add.data = Box::new(peeked);
n
};
add.data = Box::new(peeked);
let output = add.into_plan(&table_schema, &table_def)?;
+17
View File
@@ -52,6 +52,7 @@ pub struct AddDataBuilder {
pub(crate) write_options: WriteOptions,
pub(crate) on_nan_vectors: NaNVectorBehavior,
pub(crate) embedding_registry: Option<Arc<dyn EmbeddingRegistry>>,
pub(crate) write_parallelism: Option<usize>,
}
impl std::fmt::Debug for AddDataBuilder {
@@ -77,6 +78,7 @@ impl AddDataBuilder {
write_options: WriteOptions::default(),
on_nan_vectors: NaNVectorBehavior::default(),
embedding_registry,
write_parallelism: None,
}
}
@@ -101,7 +103,22 @@ impl AddDataBuilder {
self
}
/// Set the number of parallel write streams.
///
/// By default, the number of streams is estimated from the data size.
/// Setting this to `1` disables parallel writes.
pub fn write_parallelism(mut self, parallelism: usize) -> Self {
self.write_parallelism = Some(parallelism);
self
}
pub async fn execute(self) -> Result<AddResult> {
if self.write_parallelism.map(|p| p == 0).unwrap_or(false) {
return Err(Error::InvalidInput {
message: "write_parallelism must be greater than 0".to_string(),
});
}
self.parent.clone().add(self).await
}
+7
View File
@@ -64,6 +64,9 @@ pub enum OptimizeAction {
older_than: Option<Duration>,
/// Because they may be part of an in-progress transaction, files newer than 7 days old are not deleted by default.
/// If you are sure that there are no in-progress transactions, then you can set this to True to delete all files older than `older_than`.
///
/// **WARNING**: This should only be set to true if you can guarantee that no other process is
/// currently working on this dataset. Otherwise the dataset could be put into a corrupted state.
delete_unverified: Option<bool>,
/// If true, an error will be returned if there are any old versions that are still tagged.
error_if_tagged_old_versions: Option<bool>,
@@ -117,6 +120,10 @@ pub(crate) async fn optimize_indices(table: &NativeTable, options: &OptimizeOpti
/// If you are sure that there are no in-progress transactions, then you
/// can set this to True to delete all files older than `older_than`.
///
/// **WARNING**: This should only be set to true if you can guarantee that
/// no other process is currently working on this dataset. Otherwise the
/// dataset could be put into a corrupted state.
///
/// This calls into [lance::dataset::Dataset::cleanup_old_versions] and
/// returns the result.
pub(crate) async fn cleanup_old_versions(
+29
View File
@@ -186,6 +186,13 @@ pub async fn create_plan(
Select::Dynamic(ref select_with_transform) => {
scanner.project_with_transform(select_with_transform.as_slice())?;
}
Select::Expr(ref expr_pairs) => {
let sql_pairs: crate::Result<Vec<(String, String)>> = expr_pairs
.iter()
.map(|(name, expr)| expr_to_sql_string(expr).map(|sql| (name.clone(), sql)))
.collect();
scanner.project_with_transform(sql_pairs?.as_slice())?;
}
Select::All => {}
}
@@ -340,6 +347,17 @@ fn convert_to_namespace_query(query: &AnyQuery) -> Result<NsQueryTableRequest> {
.to_string(),
});
}
Select::Expr(pairs) => {
let sql_pairs: crate::Result<Vec<(String, String)>> = pairs
.iter()
.map(|(name, expr)| expr_to_sql_string(expr).map(|sql| (name.clone(), sql)))
.collect();
let sql_pairs = sql_pairs?;
Some(Box::new(QueryTableRequestColumns {
column_names: None,
column_aliases: Some(sql_pairs.into_iter().collect()),
}))
}
};
// Check for unsupported features
@@ -411,6 +429,17 @@ fn convert_to_namespace_query(query: &AnyQuery) -> Result<NsQueryTableRequest> {
.to_string(),
});
}
Select::Expr(pairs) => {
let sql_pairs: crate::Result<Vec<(String, String)>> = pairs
.iter()
.map(|(name, expr)| expr_to_sql_string(expr).map(|sql| (name.clone(), sql)))
.collect();
let sql_pairs = sql_pairs?;
Some(Box::new(QueryTableRequestColumns {
column_names: None,
column_aliases: Some(sql_pairs.into_iter().collect()),
}))
}
};
// Handle full text search if present