Compare commits

..

6 Commits

Author SHA1 Message Date
Lance Release b89f87f206 Bump version: 0.37.1-beta.3 → 0.37.1 2026-08-10 07:54:42 +00:00
Lance Release a022d3bcb3 Bump version: 0.37.1-beta.2 → 0.37.1-beta.3 2026-08-10 07:53:53 +00:00
Lance Release 8268532d64 Bump version: 0.37.1-beta.1 → 0.37.1-beta.2 2026-08-10 06:30:45 +00:00
Xuanwo f933ef9b21 fix(ci): tag releases after updating lockfiles 2026-08-10 14:23:41 +08:00
Xuanwo e885e5dd00 fix(ci): validate stable Lance dependencies 2026-08-10 14:23:33 +08:00
Xuanwo 7c7efa9743 feat: update lance dependency to v10.0.0 2026-08-10 14:23:29 +08:00
55 changed files with 837 additions and 2490 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.37.1-beta.1"
current_version = "0.37.1"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
+49 -70
View File
@@ -36,9 +36,7 @@ jobs:
permissions:
contents: read
outputs:
checker_outcome: ${{ steps.lychee.outcome }}
exit_code: ${{ steps.lychee.outputs.exit_code }}
status: ${{ steps.validate.outputs.status }}
steps:
- name: Checkout
uses: actions/checkout@v6
@@ -52,7 +50,6 @@ jobs:
- name: Check links
id: lychee
continue-on-error: true
uses: lycheeverse/lychee-action@e7477775783ea5526144ba13e8db5eec57747ce8 # v2.9.0
with:
# Restricted to http(s) on purpose. Much of docs/src is generated
@@ -71,50 +68,38 @@ jobs:
format: json
output: ./lychee/out.json
jobSummary: false
# The report issue, not a red workflow run, is the signal for link
# findings and checker failures alike.
# The report, not a red build, is the signal for broken links. The
# validation step below still fails the run if the check itself
# breaks.
fail: false
- name: Validate report
id: validate
# lychee does not reserve exit code 2 for broken links: its CLI
# parser also exits 2 on an invalid option, before any link was
# checked or any report written. Only a parseable report whose
# counts agree with a completed exit code (0 or 2) counts as a link
# verdict. Everything else becomes a checker-error report instead of
# failing the workflow. Exit 2 covers timeouts as well as errors, and a
# timed-out host is exactly the transient unavailability this report
# exists to surface, so both count as findings. Requiring total > 0
# also catches a glob that silently stopped matching any file.
if: always()
# counts agree with the exit code counts as a link verdict; anything
# else fails here, and the report job below is skipped entirely, so
# the tracking issue is never touched. Exit 2 covers timeouts as
# well as errors, and a timed-out host is exactly the transient
# unavailability this report exists to surface, so both count as
# findings. Requiring total > 0 also catches a glob that silently
# stopped matching any file.
if: steps.lychee.outputs.exit_code == 0 || steps.lychee.outputs.exit_code == 2
env:
CHECKER_OUTCOME: ${{ steps.lychee.outcome }}
EXIT_CODE: ${{ steps.lychee.outputs.exit_code }}
run: |
status=checker-error
if [[ "$CHECKER_OUTCOME" == success ]] &&
[[ "$EXIT_CODE" == 0 || "$EXIT_CODE" == 2 ]] &&
jq -e --argjson code "$EXIT_CODE" '
(.total > 0) and
(if $code == 0
then .errors == 0 and .timeouts == 0
and (.error_map | length == 0) and (.timeout_map | length == 0)
else (.errors + .timeouts) > 0
and ((.error_map | length) + (.timeout_map | length)) > 0
end)
' ./lychee/out.json
then
if [[ "$EXIT_CODE" == 0 ]]; then
status=healthy
else
status=findings
fi
fi
echo "status=$status" >> "$GITHUB_OUTPUT"
echo "Validated link check as $status"
jq -e --argjson code "$EXIT_CODE" '
(.total > 0) and
(if $code == 0
then .errors == 0 and .timeouts == 0
and (.error_map | length == 0) and (.timeout_map | length == 0)
else (.errors + .timeouts) > 0
and ((.error_map | length) + (.timeout_map | length)) > 0
end)
' ./lychee/out.json
- name: Upload report
if: steps.validate.outputs.status == 'findings'
if: steps.lychee.outputs.exit_code == 2
uses: actions/upload-artifact@v7
with:
name: link-report
@@ -130,11 +115,26 @@ jobs:
permissions:
issues: write
env:
CHECKER_OUTCOME: ${{ needs.scan.outputs.checker_outcome }}
EXIT_CODE: ${{ needs.scan.outputs.exit_code }}
STATUS: ${{ needs.scan.outputs.status }}
GH_TOKEN: ${{ secrets.GITHUB_TOKEN }}
steps:
- name: Classify checker result
# lychee exits 0 when every link resolves and 2 when links fail,
# both already cross-checked against the report by the scan job's
# validation step. Anything else (1 runtime, 3 bad config) means the
# check never produced a link verdict, which must surface as a failed
# run rather than be published as "broken documentation links".
run: |
case "$EXIT_CODE" in
0|2)
echo "lychee exit code $EXIT_CODE"
;;
*)
echo "::error::lychee exited with '$EXIT_CODE': the link check did not complete. Leaving the report issue untouched."
exit 1
;;
esac
- name: Find existing report issue
id: report
# Matched on title alone, and through search rather than a listing:
@@ -144,7 +144,7 @@ jobs:
# Closed issues are included because a healthy run closes the report:
# an open-only lookup would forget that identity and the next failing
# run would open a duplicate. The oldest match stays the canonical
# report and is reopened below when a problem recurs.
# report and is reopened below when links break again.
run: |
match=$(gh issue list --repo "$GITHUB_REPOSITORY" --state all \
--search "in:title \"$REPORT_TITLE\" author:app/github-actions" \
@@ -154,14 +154,14 @@ jobs:
echo "state=$(jq -r '.state // empty' <<<"$match")" >> "$GITHUB_OUTPUT"
- name: Download report
if: env.STATUS == 'findings'
if: env.EXIT_CODE == 2
uses: actions/download-artifact@v8
with:
name: link-report
path: ./lychee
- name: Compose report
if: env.STATUS == 'findings'
if: env.EXIT_CODE == 2
run: |
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
{
@@ -185,41 +185,22 @@ jobs:
' ./lychee/out.json
} > ./lychee/issue.md
- name: Compose checker error report
if: env.STATUS == 'checker-error'
run: |
mkdir -p ./lychee
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
{
echo "The documentation link check did not complete in [the latest run]($run_url)."
echo
echo "This issue is rewritten by every scheduled run and closed automatically once a trustworthy run finds that all links resolve."
echo
echo "The checker did not produce a trustworthy link verdict. Treat the previous result, if any, as stale until a later run completes."
echo
echo "* Action outcome: \`$CHECKER_OUTCOME\`"
echo "* Exit code: \`${EXIT_CODE:-not reported}\`"
echo "* Verdict validation: \`failed\`"
} > ./lychee/issue.md
- name: Reopen report issue
# A healthy run closes the report, and the issue action below only
# rewrites the body of whatever number it is given. Without an
# explicit reopen, a later finding or checker error would rewrite a
# closed issue. A CLOSED state implies the lookup found a canonical
# issue, so no separate emptiness check.
if: >-
env.STATUS != 'healthy' &&
steps.report.outputs.state == 'CLOSED'
# explicit reopen, the 2 -> 0 -> 2 sequence would keep rewriting a
# closed issue while links are broken. A CLOSED state implies the
# lookup found a canonical issue, so no separate emptiness check.
if: env.EXIT_CODE == 2 && steps.report.outputs.state == 'CLOSED'
env:
ISSUE_NUMBER: ${{ steps.report.outputs.number }}
run: |
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
gh issue reopen "$ISSUE_NUMBER" --repo "$GITHUB_REPOSITORY" \
--comment "The documentation link checker reported a problem again in [the latest run]($run_url)."
--comment "Broken documentation links found again in [the latest run]($run_url)."
- name: Report link-check problem
if: env.STATUS != 'healthy'
- name: Report broken links
if: env.EXIT_CODE == 2
uses: peter-evans/create-issue-from-file@fca9117c27cdc29c6c4db3b86c48e4115a786710 # v6.0.0
with:
# Empty on the first failing run, which creates the issue; afterwards
@@ -232,9 +213,7 @@ jobs:
- name: Close report issue once links are healthy
# An OPEN state implies the lookup found a canonical issue; a report
# that is already closed needs nothing.
if: >-
env.STATUS == 'healthy' &&
steps.report.outputs.state == 'OPEN'
if: env.EXIT_CODE == 0 && steps.report.outputs.state == 'OPEN'
env:
ISSUE_NUMBER: ${{ steps.report.outputs.number }}
run: |
+4 -2
View File
@@ -1,7 +1,8 @@
name: Create release commit
# This workflow increments the version, tags it, and pushes it. All SDKs share
# a single version, so one tag releases all of them.
# This workflow increments the version, updates lockfiles, tags the final
# commit, and pushes it. All SDKs share a single version, so one tag releases
# all of them.
# When a tag is pushed, another workflow is triggered that creates a GH release
# and uploads the binaries. This workflow is only for creating the tag.
@@ -63,6 +64,7 @@ jobs:
pip install bump-my-version PyGithub packaging
bash ci/bump_version.sh ${{ inputs.type }} ${{ inputs.bump-minor }}
bash ci/update_lockfiles.sh --amend
bash ci/create_release_tag.sh
- name: Push new version tag
if: ${{ !inputs.dry_run }}
uses: ad-m/github-push-action@881a6320fdb16eb5318c5054f31c218aec2b324c # v1.3.0
-10
View File
@@ -69,16 +69,6 @@ jobs:
uses: actions/setup-python@v6
with:
python-version: "3.10"
- name: Add swap for Arm fat LTO
if: matrix.config.platform == 'aarch64'
shell: bash
run: |
swap_file="$RUNNER_TEMP/lancedb-swap"
sudo fallocate --length 16G "$swap_file"
sudo chmod 600 "$swap_file"
sudo mkswap "$swap_file"
sudo swapon "$swap_file"
free -h
- uses: ./.github/workflows/build_linux_wheel
with:
python-minor-version: 10
Generated
+165 -138
View File
@@ -775,7 +775,7 @@ dependencies = [
"http 0.2.12",
"http 1.5.0",
"http-body 1.1.0",
"lru 0.16.4",
"lru",
"percent-encoding",
"regex-lite",
"sha2 0.11.0",
@@ -1970,6 +1970,15 @@ dependencies = [
"spin 0.10.1",
]
[[package]]
name = "crc32c"
version = "0.6.8"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "3a47af21622d091a8f0fb295b88bc886ac74efcc613efc19f5d0b21de5c89e47"
dependencies = [
"rustc_version",
]
[[package]]
name = "crc32fast"
version = "1.5.0"
@@ -3441,12 +3450,6 @@ dependencies = [
"percent-encoding",
]
[[package]]
name = "frostem"
version = "1.20260804.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "82eb03a32a1d50555353c85a7b9d3279a6f1e91af9890b789acdf544ed57c8d7"
[[package]]
name = "fs_extra"
version = "1.3.0"
@@ -3455,8 +3458,9 @@ checksum = "42703706b716c37f96a77aea830392ad231f44c9e9a67872fa5548707e11b11c"
[[package]]
name = "fsst"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d0981ce90521824089f3cd68a48b1c4c89e9ad96d0eb74f58d6e550a3d604545"
dependencies = [
"arrow-array",
"rand 0.9.5",
@@ -3813,23 +3817,14 @@ dependencies = [
[[package]]
name = "goosefs-sdk"
version = "0.1.9"
version = "0.1.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e1ea4eee6dcbc31b25ab4fd577adc55b677d2bed3aa3016c44c58fbe1b2298a5"
checksum = "9ae079b88ffe7772d12cfc5c40a5a324babb357893d95b5e3a22ae857f236c5f"
dependencies = [
"arc-swap",
"async-trait",
"bytes",
"dashmap",
"fastrand",
"futures",
"hostname",
"io-uring",
"itoa",
"libc",
"lru 0.18.2",
"memmap2 0.9.10",
"moka",
"prost",
"prost-types",
"rand 0.9.5",
@@ -3842,7 +3837,6 @@ dependencies = [
"tonic-prost",
"tracing",
"uuid",
"xxhash-rust",
]
[[package]]
@@ -4815,8 +4809,9 @@ checksum = "e037a2e1d8d5fdbd49b16a4ea09d5d6401c1f29eca5ff29d03d3824dba16256a"
[[package]]
name = "lance"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "830548f9fc92ae74b848a504282d4da06f016f2ee94b663773c2738b9cb61c42"
dependencies = [
"arc-swap",
"arrow",
@@ -4890,8 +4885,9 @@ dependencies = [
[[package]]
name = "lance-arrow"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "24187f972374567bb3573cffd8728f2f4f7f06d644c3c5d4430b79d6b0b67300"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4913,7 +4909,8 @@ dependencies = [
[[package]]
name = "lance-arrow-scalar"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "771f68b04b47f3addf781116f65061808de94b05e1e9411c23c18f32d14ebe79"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4927,17 +4924,20 @@ dependencies = [
[[package]]
name = "lance-arrow-stats"
version = "58.0.0"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "dd47ec33c90bf29f688fd02118e37d3a5ad5c339caa3163f89e417dc0867001f"
dependencies = [
"arrow-array",
"arrow-schema",
"half",
"lance-arrow-scalar",
]
[[package]]
name = "lance-bitpacking"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "55fa50ad941e25298afb54eec107c3584f80c16db4adf7a6e4e9c4b6066f4281"
dependencies = [
"arrayref",
"crunchy",
@@ -4947,8 +4947,9 @@ dependencies = [
[[package]]
name = "lance-core"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "21f4fd872bfe948150a6878d983327c4f150dc1b629e94ae7677310fa6a3f35f"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -4988,8 +4989,9 @@ dependencies = [
[[package]]
name = "lance-datafusion"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "230770734cd5f6fe1f1cb4a66750acad5c5a862c23c7f587d7def57c9f2c5f6b"
dependencies = [
"arrow",
"arrow-array",
@@ -5019,8 +5021,9 @@ dependencies = [
[[package]]
name = "lance-datagen"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c0c9edcea9dbf154a3baf470595bdcdb7ef6b44dcf2b5a0c28e8815c2d4837c5"
dependencies = [
"arrow",
"arrow-array",
@@ -5037,8 +5040,9 @@ dependencies = [
[[package]]
name = "lance-derive"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "03ac280ef94d66c2e7a4b0104e60d548f45156f114f02cd0ac57423721ec96ec"
dependencies = [
"proc-macro2",
"quote",
@@ -5047,8 +5051,9 @@ dependencies = [
[[package]]
name = "lance-encoding"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "65b345fecf1792d4147982e9ded1cf211992a442bcf73dd17a598b1c8c9b53a5"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5074,6 +5079,7 @@ dependencies = [
"prost",
"prost-build",
"rand 0.9.5",
"strum 0.26.3",
"tokio",
"tracing",
"xxhash-rust",
@@ -5082,8 +5088,9 @@ dependencies = [
[[package]]
name = "lance-file"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "203400160ecd4ac6f67f6bfdb3c186f94758ff70a6ca76d4f300c327f5c14a11"
dependencies = [
"arrow-arith",
"arrow-array",
@@ -5114,8 +5121,9 @@ dependencies = [
[[package]]
name = "lance-index"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "db67792e22332a7be08d35e60ce9e4c6fb7f75f70ab1566b67b462d1aac21b13"
dependencies = [
"arc-swap",
"arrow",
@@ -5182,8 +5190,9 @@ dependencies = [
[[package]]
name = "lance-index-core"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "49b5700d71f43e5da2057908f3243a6bc01247b5dd12247313f789c280155aae"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5205,8 +5214,9 @@ dependencies = [
[[package]]
name = "lance-io"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4b8a6872b9bff95e75e223e04011feebd990d8784878eaf5df73c86d6acebb47"
dependencies = [
"arrow",
"arrow-array",
@@ -5218,6 +5228,7 @@ dependencies = [
"bytes",
"chrono",
"futures",
"goosefs-sdk",
"http 1.5.0",
"io-uring",
"lance-arrow",
@@ -5242,8 +5253,9 @@ dependencies = [
[[package]]
name = "lance-linalg"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5a1dd49fdf61da86650dc335fbe1f239803f9d73bc2bc5ffbc0c905dd8809892"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5259,8 +5271,9 @@ dependencies = [
[[package]]
name = "lance-namespace"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "1449221f599147e570bec6e09eaea6ce6b9286f632e00a27f945dd0c6b70ccd8"
dependencies = [
"arrow",
"async-trait",
@@ -5272,8 +5285,9 @@ dependencies = [
[[package]]
name = "lance-namespace-impls"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c30bff40ab9a9a58313818ab97241cff7bab2b90b8d65039908425a065a797de"
dependencies = [
"arrow",
"arrow-ipc",
@@ -5303,6 +5317,7 @@ dependencies = [
"serde",
"serde_json",
"sha2 0.10.9",
"time",
"tokio",
"tower",
"tower-http 0.5.2",
@@ -5326,8 +5341,9 @@ dependencies = [
[[package]]
name = "lance-select"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "c66c5cb12ed194c7c8e427a0e90e990cdbc7159edba2bfc7b7e8db2c90c747af"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5342,8 +5358,9 @@ dependencies = [
[[package]]
name = "lance-table"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "6cd2a39f2e288d1acecf887cb9cc655711fb67bcf3b24a2034158dbac7a6f60e"
dependencies = [
"arrow",
"arrow-array",
@@ -5353,7 +5370,6 @@ dependencies = [
"async-trait",
"aws-credential-types",
"aws-sdk-dynamodb",
"blake3",
"byteorder",
"bytes",
"chrono",
@@ -5383,8 +5399,9 @@ dependencies = [
[[package]]
name = "lance-testing"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0847352699e869a465b49bd49ecbc467bc5a3dcb168757bf39368362aad05f99"
dependencies = [
"arrow-array",
"arrow-schema",
@@ -5397,13 +5414,14 @@ dependencies = [
[[package]]
name = "lance-tokenizer"
version = "11.0.0-beta.7"
source = "git+https://github.com/lance-format/lance.git?tag=v11.0.0-beta.7#e581c49338bc83baf1ea50c5e235bd702f3fbeea"
version = "10.0.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "254c1e9287789c8d744f93f2010a8389a5ea3fcedc0dddf7649c4117ec830c57"
dependencies = [
"frostem",
"icu_segmenter",
"jieba-rs",
"lindera",
"rust-stemmers",
"serde",
"stop-words",
"unicode-normalization",
@@ -5411,7 +5429,7 @@ dependencies = [
[[package]]
name = "lancedb"
version = "0.37.1-beta.1"
version = "0.37.1"
dependencies = [
"ahash",
"anyhow",
@@ -5495,12 +5513,11 @@ dependencies = [
"urlencoding",
"uuid",
"walkdir",
"windows-sys 0.61.2",
]
[[package]]
name = "lancedb-nodejs"
version = "0.37.1-beta.1"
version = "0.37.1"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5525,7 +5542,7 @@ dependencies = [
[[package]]
name = "lancedb-python"
version = "0.37.1-beta.1"
version = "0.37.1"
dependencies = [
"arrow",
"async-trait",
@@ -5682,7 +5699,7 @@ dependencies = [
"serde",
"serde_json",
"serde_yaml_ng",
"strum",
"strum 0.28.0",
"strum_macros 0.28.0",
"unicode-blocks",
"unicode-normalization",
@@ -5712,7 +5729,7 @@ dependencies = [
"rkyv",
"serde",
"serde_json",
"strum",
"strum 0.28.0",
"strum_macros 0.28.0",
"thiserror 2.0.18",
]
@@ -5784,15 +5801,6 @@ dependencies = [
"hashbrown 0.16.1",
]
[[package]]
name = "lru"
version = "0.18.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5d2f2f9b4ba7e6b24d95e7e899329d35be83bcded72c8540cdd5368932d1d90a"
dependencies = [
"hashbrown 0.17.1",
]
[[package]]
name = "lru-slab"
version = "0.1.2"
@@ -6430,9 +6438,9 @@ dependencies = [
[[package]]
name = "object_store_opendal"
version = "0.58.0"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "88f165780495c17aa3ce86846600504198c3fffd99073521552751c2430fa6ac"
checksum = "0eb12a624a41fce745838d0ef3701ff6c47797c13cd18ad3612fd2a3134fdbd8"
dependencies = [
"async-trait",
"bytes",
@@ -6493,13 +6501,12 @@ checksum = "d6790f58c7ff633d8771f42965289203411a5e5c68388703c06e14f24770b41e"
[[package]]
name = "opendal"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4f20562cc7447fcc915fc5c23df305a412ea80a733c9f2fd9e2d267e2815be6d"
checksum = "96c9c85ce253ff87225e7669979d877a20c98a06604ec9d6dd5f4473e08f1ae1"
dependencies = [
"ctor 1.0.12",
"opendal-core",
"opendal-http-transport-reqwest",
"opendal-layer-concurrent-limit",
"opendal-layer-logging",
"opendal-layer-retry",
@@ -6516,22 +6523,24 @@ dependencies = [
[[package]]
name = "opendal-core"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ec75551ff4cf3e57da98979f6a937aaa9ddb3915bf68cc17d03df733be6646ed"
checksum = "c4f8607c90e2c963a91467f50fb49fbc7fb3d573f88cea219ca59ccd3740b309"
dependencies = [
"anyhow",
"base64 0.23.1",
"base64 0.22.1",
"bytes",
"futures",
"http 1.5.0",
"http-body 1.1.0",
"jiff",
"log",
"md-5 0.11.0",
"mea",
"percent-encoding",
"quick-xml 0.41.0",
"quick-xml 0.39.4",
"reqsign-core",
"reqwest 0.13.4",
"serde",
"serde_json",
"tokio",
@@ -6540,25 +6549,11 @@ dependencies = [
"web-time",
]
[[package]]
name = "opendal-http-transport-reqwest"
version = "0.58.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "ad4d4f19c3ce01126a30611f8e544eaa217104a278c889ac17c9374fe4f9e4ef"
dependencies = [
"bytes",
"futures",
"http 1.5.0",
"http-body 1.1.0",
"opendal-core",
"reqwest 0.13.4",
]
[[package]]
name = "opendal-layer-concurrent-limit"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "249ac5b0aa5a7a6c3737342d10456067937f9c9a6f3f02544271f7908ab91081"
checksum = "0d6f81ba6960e3fae1882f253b114b21d7e444e1534f209c7737a79f6243eb6f"
dependencies = [
"futures",
"http 1.5.0",
@@ -6568,9 +6563,9 @@ dependencies = [
[[package]]
name = "opendal-layer-logging"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "5c75411ab00f77851ff086b686c1e9ca8175ac18c15afa2cb75b9036436cb06c"
checksum = "58ada45c6d81d1aa4c9305d0c7d4bc317c59c85866a0908a2d75a7a978aa5ee2"
dependencies = [
"log",
"opendal-core",
@@ -6578,9 +6573,9 @@ dependencies = [
[[package]]
name = "opendal-layer-retry"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "80b7738bd5f233ad8da39af9b9316b9b7a4eaddd91e8e32a1e19b7030688121d"
checksum = "7b2a25a718afb81fad81cb9a0580a1cb989221fa2317f888c6a37f8dad408eb7"
dependencies = [
"backon",
"log",
@@ -6589,9 +6584,9 @@ dependencies = [
[[package]]
name = "opendal-layer-timeout"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "a704141924500f3803c05ed871b53305d2a2f11cb5ef20160c3ee688a1857f66"
checksum = "1e91f731724c213af81e9d03517859c8fc47b4578e64ad61ae4f099f10fe36e3"
dependencies = [
"opendal-core",
"tokio",
@@ -6599,17 +6594,17 @@ dependencies = [
[[package]]
name = "opendal-service-azblob"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b3310fbbb48f111c6f590473c2cd15e1b7f8e384444b0d4e328f0464c864d767"
checksum = "0030644366ef5d8cbe3a4a5822bf99a4aafddc1666e9d24b44d158d9062fc76a"
dependencies = [
"base64 0.23.1",
"base64 0.22.1",
"bytes",
"http 1.5.0",
"log",
"opendal-core",
"opendal-service-azure-common",
"quick-xml 0.41.0",
"quick-xml 0.39.4",
"reqsign-azure-storage",
"reqsign-core",
"reqsign-file-read-tokio",
@@ -6620,18 +6615,17 @@ dependencies = [
[[package]]
name = "opendal-service-azdls"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "2e3c406729935fe214ce574d68681a1ff7e0b322548f14094912bdbfe50e5c53"
checksum = "6dea4908d490143a9b0b7f7a790e139ff829b06a023f670455ed3d44f664b361"
dependencies = [
"base64 0.23.1",
"base64 0.22.1",
"bytes",
"http 1.5.0",
"log",
"mea",
"opendal-core",
"opendal-service-azure-common",
"quick-xml 0.41.0",
"quick-xml 0.39.4",
"reqsign-azure-storage",
"reqsign-core",
"reqsign-file-read-tokio",
@@ -6641,9 +6635,9 @@ dependencies = [
[[package]]
name = "opendal-service-azure-common"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7348c88edf15af435b7be930077746b569fac5e738c1bf6a363b675e7317c9df"
checksum = "9b489f13c42e69d69bdd72952b634356ec43a7881a20259b38b540fcecdf4051"
dependencies = [
"http 1.5.0",
"opendal-core",
@@ -6651,15 +6645,15 @@ dependencies = [
[[package]]
name = "opendal-service-cos"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "d533d4582105d269c8aebeee5f0e8bcf960f41b8aab6197df7012254d9f39bf0"
checksum = "aa8cafe9729213375c7331019b0cb756ad3e1aff7f45cd32c45eae91ebde8901"
dependencies = [
"bytes",
"http 1.5.0",
"log",
"opendal-core",
"quick-xml 0.41.0",
"quick-xml 0.39.4",
"reqsign-core",
"reqsign-file-read-tokio",
"reqsign-tencent-cos",
@@ -6668,9 +6662,9 @@ dependencies = [
[[package]]
name = "opendal-service-gcs"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "007f3fba63c21e516c956b891e96ff9892d8175662bfb781cdada9d3766a11e6"
checksum = "48de101aac565ed06af4b47903c24eafd249075553ec1fb18256751c45148d47"
dependencies = [
"async-trait",
"bytes",
@@ -6678,7 +6672,7 @@ dependencies = [
"log",
"opendal-core",
"percent-encoding",
"quick-xml 0.41.0",
"quick-xml 0.39.4",
"reqsign-core",
"reqsign-file-read-tokio",
"reqsign-google",
@@ -6689,9 +6683,9 @@ dependencies = [
[[package]]
name = "opendal-service-goosefs"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "60871e6386f04d831e6a5bdbc032af4a91aeba49963252d0ef456a2cf36a9b78"
checksum = "69e43048bde419947ba826fbdc2f134d6c03f44ebf48bd33a03b72f9fc45fcb4"
dependencies = [
"bytes",
"goosefs-sdk",
@@ -6703,9 +6697,9 @@ dependencies = [
[[package]]
name = "opendal-service-hf"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "b41fd41eb7ed03c5e66cefda61e8e117808ffd2908f2916737cb020a6beb02c7"
checksum = "c4922661976a1d40794a2adfbdb888cc3c23097690f825a92f773af38908a848"
dependencies = [
"bytes",
"hf-xet",
@@ -6713,21 +6707,22 @@ dependencies = [
"log",
"opendal-core",
"percent-encoding",
"reqwest 0.13.4",
"serde",
"serde_json",
]
[[package]]
name = "opendal-service-oss"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "cd528ec2d49c5ca69e674ffed7b3e0686fb9cfcfea0596870de381467fda4f1b"
checksum = "328fa55e8888cbdfe00826bfea2a79042422b720e8369e9e021e46121dea5ace"
dependencies = [
"bytes",
"http 1.5.0",
"log",
"opendal-core",
"quick-xml 0.41.0",
"quick-xml 0.39.4",
"reqsign-aliyun-oss",
"reqsign-core",
"reqsign-file-read-tokio",
@@ -6736,18 +6731,18 @@ dependencies = [
[[package]]
name = "opendal-service-s3"
version = "0.58.1"
version = "0.57.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "58e80cdf192d7eff05feed747894d64f81905ac4eaf132edf7ea270abdd2d663"
checksum = "313d46c9f5ae70bca26b7c3e3fbb9b639292625f28af73aa016f47e788af9deb"
dependencies = [
"base64 0.23.1",
"base64 0.22.1",
"bytes",
"crc-fast",
"crc32c",
"http 1.5.0",
"log",
"md-5 0.11.0",
"opendal-core",
"quick-xml 0.41.0",
"quick-xml 0.39.4",
"reqsign-aws-v4",
"reqsign-core",
"reqsign-file-read-tokio",
@@ -8680,6 +8675,16 @@ dependencies = [
"ordered-multimap",
]
[[package]]
name = "rust-stemmers"
version = "1.2.0"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e46a2036019fdb888131db7a4c847a1063a7493f971ed94ea82c67eada63ca54"
dependencies = [
"serde",
"serde_derive",
]
[[package]]
name = "rustc-demangle"
version = "0.1.27"
@@ -9553,6 +9558,15 @@ version = "0.11.1"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "7da8b5736845d9f2fcb837ea5d9e2628564b3b043a70948a3f0b778838c5fb4f"
[[package]]
name = "strum"
version = "0.26.3"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "8fec0f0aef304996cf250b31b5a10dee7980c85da9d759361292b8bca5a18f06"
dependencies = [
"strum_macros 0.26.4",
]
[[package]]
name = "strum"
version = "0.28.0"
@@ -9575,6 +9589,19 @@ dependencies = [
"syn 2.0.117",
]
[[package]]
name = "strum_macros"
version = "0.26.4"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "4c6bee85a5a24955dc440386795aa378cd9cf82acd5f764469152d2270e581be"
dependencies = [
"heck 0.5.0",
"proc-macro2",
"quote",
"rustversion",
"syn 2.0.117",
]
[[package]]
name = "strum_macros"
version = "0.28.0"
+14 -14
View File
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
lance = { "version" = "=11.0.0-beta.7", default-features = false, "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-core = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-datagen = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-file = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-io = { "version" = "=11.0.0-beta.7", default-features = false, "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-index = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-linalg = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-namespace-impls = { "version" = "=11.0.0-beta.7", default-features = false, "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-table = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-testing = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-datafusion = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-encoding = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance-arrow = { "version" = "=11.0.0-beta.7", "tag" = "v11.0.0-beta.7", "git" = "https://github.com/lance-format/lance.git" }
lance = { "version" = "=10.0.0", default-features = false }
lance-core = "=10.0.0"
lance-datagen = "=10.0.0"
lance-file = "=10.0.0"
lance-io = { "version" = "=10.0.0", default-features = false }
lance-index = "=10.0.0"
lance-linalg = "=10.0.0"
lance-namespace = "=10.0.0"
lance-namespace-impls = { "version" = "=10.0.0", default-features = false }
lance-table = "=10.0.0"
lance-testing = "=10.0.0"
lance-datafusion = "=10.0.0"
lance-encoding = "=10.0.0"
lance-arrow = "=10.0.0"
ahash = "0.8"
# Note that this one does not include pyarrow
arrow = { version = "58.0.0", optional = false }
+4 -9
View File
@@ -6,16 +6,11 @@ HEAD_SHA=$(git rev-parse HEAD)
readonly TAG_PREFIX="v"
readonly SELF_DIR=$(cd "$( dirname "${BASH_SOURCE[0]}" )" && pwd )
readonly BUMP_ARGS="--no-tag"
PREV_TAG=$(git tag --sort='version:refname' | grep ^$TAG_PREFIX | python $SELF_DIR/semver_sort.py $TAG_PREFIX | tail -n 1)
echo "Found previous tag $PREV_TAG"
# Initially, we don't want to tag if we are doing stable, because we will bump
# again later. See comment at end for why.
if [[ "$RELEASE_TYPE" == 'stable' ]]; then
BUMP_ARGS="--no-tag"
fi
# If last is stable and not bumping minor
if [[ $PREV_TAG != *beta* ]]; then
if [[ "$BUMP_MINOR" != "false" ]]; then
@@ -39,12 +34,12 @@ fi
# a stable version, bump the pre-release level ("pre_l") to make it stable.
if [[ $RELEASE_TYPE == 'stable' ]]; then
# X.Y.Z-beta.N -> X.Y.Z
bump-my-version bump -vv pre_l
bump-my-version bump -vv $BUMP_ARGS pre_l
fi
# Validate that we have incremented version appropriately for breaking changes
NEW_TAG=$(git describe --tags --exact-match HEAD)
NEW_VERSION=$(echo $NEW_TAG | sed "s/^$TAG_PREFIX//")
NEW_VERSION=$(python -c 'import tomllib; print(tomllib.load(open(".bumpversion.toml", "rb"))["tool"]["bumpversion"]["current_version"])')
NEW_TAG="$TAG_PREFIX$NEW_VERSION"
LAST_STABLE_RELEASE=$(git tag --sort='version:refname' | grep ^$TAG_PREFIX | grep -v beta | grep -vF "$NEW_TAG" | python $SELF_DIR/semver_sort.py $TAG_PREFIX | tail -n 1)
LAST_STABLE_VERSION=$(echo $LAST_STABLE_RELEASE | sed "s/^$TAG_PREFIX//")
+21
View File
@@ -0,0 +1,21 @@
#!/usr/bin/env bash
set -euo pipefail
RELEASE_VERSION=$(python -c 'import tomllib; print(tomllib.load(open(".bumpversion.toml", "rb"))["tool"]["bumpversion"]["current_version"])')
RELEASE_TAG="v${RELEASE_VERSION}"
if git rev-parse --quiet --verify "refs/tags/${RELEASE_TAG}" >/dev/null; then
echo "Release tag ${RELEASE_TAG} already exists" >&2
exit 1
fi
git tag --annotate "$RELEASE_TAG" --message "Release ${RELEASE_TAG}"
HEAD_SHA=$(git rev-parse HEAD)
TAG_SHA=$(git rev-parse "refs/tags/${RELEASE_TAG}^{}")
if [[ "$TAG_SHA" != "$HEAD_SHA" ]]; then
echo "Release tag ${RELEASE_TAG} points to ${TAG_SHA}, expected ${HEAD_SHA}" >&2
exit 1
fi
echo "Created ${RELEASE_TAG} at ${HEAD_SHA}"
+72
View File
@@ -0,0 +1,72 @@
import tempfile
import unittest
from pathlib import Path
from validate_stable_lance import validate
class ValidateStableLanceTest(unittest.TestCase):
def write_fixture(
self,
root: Path,
*,
rust: str = "=10.0.0",
python: str = "10.0.0",
java: str = "10.0.0",
) -> None:
(root / "python").mkdir()
(root / "java").mkdir()
(root / "Cargo.toml").write_text(
f'[workspace.dependencies]\nlance = "{rust}"\nlance-core = "{rust}"\n'
)
(root / "python" / "pyproject.toml").write_text(
'[project.optional-dependencies]\ntests = ["pylance==' + python + '"]\n'
)
(root / "java" / "pom.xml").write_text(
"<project><properties><lance-core.version>"
+ java
+ "</lance-core.version></properties></project>"
)
def test_accepts_matching_stable_versions(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root)
self.assertEqual(validate(root), "10.0.0")
def test_rejects_prerelease(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root, python="10.0.0rc1")
with self.assertRaisesRegex(ValueError, "not stable"):
validate(root)
def test_rejects_sdk_mismatch(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root, java="9.0.0")
with self.assertRaisesRegex(ValueError, "do not match across SDKs"):
validate(root)
def test_rejects_non_exact_rust_dependency(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root, rust="10.0.0")
with self.assertRaisesRegex(ValueError, "not exact"):
validate(root)
def test_rejects_unpublished_rust_source(self) -> None:
with tempfile.TemporaryDirectory() as temp_dir:
root = Path(temp_dir)
self.write_fixture(root)
(root / "Cargo.toml").write_text(
"[workspace.dependencies]\n"
'lance = { version = "=10.0.0", git = "https://example.com/lance", '
'tag = "v10.0.0" }\n'
)
with self.assertRaisesRegex(ValueError, "unpublished source fields"):
validate(root)
if __name__ == "__main__":
unittest.main()
Regular → Executable
+104 -26
View File
@@ -1,34 +1,112 @@
#!/usr/bin/env python3
"""Validate that every SDK uses the same published stable Lance release."""
from __future__ import annotations
import re
import xml.etree.ElementTree as ET
from pathlib import Path
import tomllib
found_preview_lance = False
STABLE_VERSION = re.compile(r"[0-9]+\.[0-9]+\.[0-9]+")
with open("Cargo.toml", "rb") as f:
cargo_data = tomllib.load(f)
for name, dep in cargo_data["workspace"]["dependencies"].items():
if name == "lance" or name.startswith("lance-"):
if isinstance(dep, str):
version = dep
elif isinstance(dep, dict):
# Version doesn't have the beta tag in it, so we instead look
# at the git tag.
version = dep.get('tag', dep.get('version'))
else:
raise ValueError("Unexpected type for dependency: " + str(dep))
def _stable_version(raw: str, *, dependency: str, exact: bool = False) -> str:
value = raw.strip()
if exact and not value.startswith("="):
raise ValueError(f"Dependency '{dependency}' is not exact: {raw}")
value = value.removeprefix("=").removeprefix("v")
if STABLE_VERSION.fullmatch(value) is None:
raise ValueError(f"Dependency '{dependency}' is not stable: {raw}")
return value
if "beta" in version:
found_preview_lance = True
print(f"Dependency '{name}' is a preview version: {version}")
with open("python/pyproject.toml", "rb") as f:
py_proj_data = tomllib.load(f)
def rust_lance_version(repo_root: Path) -> str:
with (repo_root / "Cargo.toml").open("rb") as cargo_file:
dependencies = tomllib.load(cargo_file)["workspace"]["dependencies"]
for dep in py_proj_data["project"]["dependencies"]:
if dep.startswith("pylance"):
if "b" in dep:
found_preview_lance = True
print(f"Dependency '{dep}' is a preview version")
break # Only one pylance dependency
versions: dict[str, str] = {}
for name, dependency in dependencies.items():
if name != "lance" and not name.startswith("lance-"):
continue
if found_preview_lance:
raise ValueError("Found preview version of Lance in dependencies")
if isinstance(dependency, str):
raw_version = dependency
elif isinstance(dependency, dict):
forbidden_sources = [
source
for source in ("git", "path", "branch", "rev", "tag")
if source in dependency
]
if forbidden_sources:
joined = ", ".join(forbidden_sources)
raise ValueError(
f"Dependency '{name}' uses unpublished source fields: {joined}"
)
raw_version = dependency.get("version")
if raw_version is None:
raise ValueError(f"Dependency '{name}' has no version")
else:
raise TypeError(f"Dependency '{name}' has an unexpected definition")
versions[name] = _stable_version(raw_version, dependency=name, exact=True)
if not versions:
raise ValueError("No Rust Lance dependencies found")
unique_versions = set(versions.values())
if len(unique_versions) != 1:
details = ", ".join(f"{name}={version}" for name, version in versions.items())
raise ValueError(f"Rust Lance dependency versions do not match: {details}")
return unique_versions.pop()
def python_lance_version(repo_root: Path) -> str:
with (repo_root / "python" / "pyproject.toml").open("rb") as pyproject_file:
pyproject = tomllib.load(pyproject_file)
requirements = pyproject["project"]["optional-dependencies"]["tests"]
pylance_requirements = [
requirement for requirement in requirements if requirement.startswith("pylance")
]
if len(pylance_requirements) != 1:
raise ValueError(
"Expected exactly one pylance requirement in the Python test dependencies"
)
requirement = pylance_requirements[0]
prefix = "pylance=="
if not requirement.startswith(prefix):
raise ValueError(f"Python test dependency is not exact: {requirement}")
return _stable_version(requirement[len(prefix) :], dependency="pylance")
def java_lance_version(repo_root: Path) -> str:
pom_root = ET.parse(repo_root / "java" / "pom.xml").getroot()
versions = [
element.text
for element in pom_root.iter()
if element.tag.rsplit("}", 1)[-1] == "lance-core.version"
]
if len(versions) != 1 or versions[0] is None:
raise ValueError("Expected exactly one Java lance-core.version property")
return _stable_version(versions[0], dependency="Java lance-core")
def validate(repo_root: Path) -> str:
versions = {
"Rust": rust_lance_version(repo_root),
"Python": python_lance_version(repo_root),
"Java": java_lance_version(repo_root),
}
if len(set(versions.values())) != 1:
details = ", ".join(f"{sdk}={version}" for sdk, version in versions.items())
raise ValueError(
f"Lance dependency versions do not match across SDKs: {details}"
)
return versions["Rust"]
if __name__ == "__main__":
version = validate(Path(__file__).resolve().parents[1])
print(f"Validated published stable Lance v{version} across Rust, Python, and Java")
-7
View File
@@ -101,13 +101,6 @@ ignore = [
# https://rustsec.org/advisories/RUSTSEC-2026-0195
{ id = "RUSTSEC-2026-0194", reason = "transitive via inferno/lance/opendal; XML from trusted cloud endpoints, not attacker-controlled" },
{ id = "RUSTSEC-2026-0195", reason = "transitive via inferno/lance/opendal; XML from trusted cloud endpoints, not attacker-controlled" },
# smartstring: unmaintained — the repository was archived by its author on
# 2026-05-03. Not a vulnerability. Reached only transitively through polars
# (polars-core/-io/-ops/-time/-utils); nothing in LanceDB depends on it directly.
# The advisory states no safe upgrade is available: upstream recommends
# compact_str/smol_str, so clearing this requires polars to migrate.
# https://rustsec.org/advisories/RUSTSEC-2026-0249
{ id = "RUSTSEC-2026-0249", reason = "smartstring unmaintained via polars; no fixed upstream release" },
]
# ---------------------------------------------------------------------------
+1 -1
View File
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
<dependency>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-core</artifactId>
<version>0.37.1-beta.1</version>
<version>0.37.1</version>
</dependency>
```
+3 -9
View File
@@ -431,10 +431,9 @@ Read the [LsmWriteSpec](../interfaces/LsmWriteSpec.md) currently installed on th
Resolves to `undefined` when the MemWAL LSM write path is not enabled (no
spec has been set, or it was removed with [Table#unsetLsmWriteSpec](Table.md#unsetlsmwritespec)).
The returned spec mirrors what was passed to
[Table#setLsmWriteSpec](Table.md#setlsmwritespec), except that `maintainedIndexes` always
reports the concrete list resolved when the spec was set — `undefined`
never round-trips.
The returned spec — including its `maintainedIndexes` and
`writerConfigDefaults` — mirrors what was passed to
[Table#setLsmWriteSpec](Table.md#setlsmwritespec).
#### Returns
@@ -807,11 +806,6 @@ All variants require the table to have an unenforced primary key
([Table#setUnenforcedPrimaryKey](Table.md#setunenforcedprimarykey)); bucket sharding additionally
requires it to be the single column being bucketed.
Omitting `maintainedIndexes` maintains every index on the table, resolved
here, failing if one cannot be maintained — name them to install anyway.
Naming them pins an exact set, and a still-building index is rejected
rather than quietly omitted.
#### Parameters
* **spec**: [`LsmWriteSpec`](../interfaces/LsmWriteSpec.md)
+1 -3
View File
@@ -34,9 +34,7 @@ Bucket and identity variants: the sharding column.
optional maintainedIndexes: string[];
```
Indexes the MemWAL keeps up to date. Omit to maintain every supported
index, resolved on install — a snapshot, so indexes created later are not
maintained. Pass `[]` for none.
Names of indexes the MemWAL should keep up to date during writes.
***
+1 -1
View File
@@ -8,7 +8,7 @@
<parent>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.37.1-beta.1</version>
<version>0.37.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.37.1-beta.1</version>
<version>0.37.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>11.0.0-beta.7</lance-core.version>
<lance-core.version>10.0.0</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.37.1-beta.1"
version = "0.37.1"
publish = false
license.workspace = true
description.workspace = true
+4 -14
View File
@@ -197,11 +197,7 @@ export interface LsmWriteSpec {
column?: string;
/** Bucket variant: the number of buckets, in `[1, 1024]`. */
numBuckets?: number;
/**
* Indexes the MemWAL keeps up to date. Omit to maintain every supported
* index, resolved on install — a snapshot, so indexes created later are not
* maintained. Pass `[]` for none.
*/
/** Names of indexes the MemWAL should keep up to date during writes. */
maintainedIndexes?: string[];
/** Default `ShardWriter` configuration recorded in the MemWAL index. */
writerConfigDefaults?: Record<string, string>;
@@ -599,11 +595,6 @@ export abstract class Table {
* All variants require the table to have an unenforced primary key
* ({@link Table#setUnenforcedPrimaryKey}); bucket sharding additionally
* requires it to be the single column being bucketed.
*
* Omitting `maintainedIndexes` maintains every index on the table, resolved
* here, failing if one cannot be maintained — name them to install anyway.
* Naming them pins an exact set, and a still-building index is rejected
* rather than quietly omitted.
* @param {LsmWriteSpec} spec The sharding spec to install.
* @returns {Promise<void>}
* @example
@@ -631,10 +622,9 @@ export abstract class Table {
*
* Resolves to `undefined` when the MemWAL LSM write path is not enabled (no
* spec has been set, or it was removed with {@link Table#unsetLsmWriteSpec}).
* The returned spec mirrors what was passed to
* {@link Table#setLsmWriteSpec}, except that `maintainedIndexes` always
* reports the concrete list resolved when the spec was set — `undefined`
* never round-trips.
* The returned spec — including its `maintainedIndexes` and
* `writerConfigDefaults` — mirrors what was passed to
* {@link Table#setLsmWriteSpec}.
* @returns {Promise<LsmWriteSpec | undefined>}
*/
abstract getLsmWriteSpec(): Promise<LsmWriteSpec | undefined>;
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-darwin-arm64",
"version": "0.37.1-beta.1",
"version": "0.37.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.37.1-beta.1",
"version": "0.37.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.37.1-beta.1",
"version": "0.37.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.37.1-beta.1",
"version": "0.37.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.37.1-beta.1",
"version": "0.37.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.37.1-beta.1",
"version": "0.37.1",
"os": [
"win32"
],
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-x64-msvc",
"version": "0.37.1-beta.1",
"version": "0.37.1",
"os": ["win32"],
"cpu": ["x64"],
"main": "lancedb.win32-x64-msvc.node",
+2 -2
View File
@@ -1,12 +1,12 @@
{
"name": "@lancedb/lancedb",
"version": "0.37.1-beta.1",
"version": "0.37.1",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@lancedb/lancedb",
"version": "0.37.1-beta.1",
"version": "0.37.1",
"cpu": [
"x64",
"arm64"
+1 -1
View File
@@ -11,7 +11,7 @@
"ann"
],
"private": false,
"version": "0.37.1-beta.1",
"version": "0.37.1",
"main": "dist/index.js",
"exports": {
".": "./dist/index.js",
+6 -6
View File
@@ -772,8 +772,7 @@ pub struct LsmWriteSpec {
pub column: Option<String>,
/// Bucket variant: the number of buckets, in `[1, 1024]`.
pub num_buckets: Option<u32>,
/// Indexes the MemWAL keeps up to date. Omitted resolves every
/// maintainable index on install; an empty array means none.
/// Names of indexes the MemWAL should keep up to date during writes.
pub maintained_indexes: Option<Vec<String>>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
pub writer_config_defaults: Option<HashMap<String, String>>,
@@ -783,6 +782,7 @@ impl TryFrom<LsmWriteSpec> for lancedb::table::LsmWriteSpec {
type Error = napi::Error;
fn try_from(value: LsmWriteSpec) -> napi::Result<Self> {
let maintained = value.maintained_indexes.unwrap_or_default();
let writer_config_defaults = value.writer_config_defaults.unwrap_or_default();
let spec = match value.spec_type.as_str() {
"bucket" => {
@@ -809,7 +809,7 @@ impl TryFrom<LsmWriteSpec> for lancedb::table::LsmWriteSpec {
}
};
Ok(spec
.with_maintained_indexes(value.maintained_indexes)
.with_maintained_indexes(maintained)
.with_writer_config_defaults(writer_config_defaults))
}
}
@@ -827,7 +827,7 @@ impl From<lancedb::table::LsmWriteSpec> for LsmWriteSpec {
spec_type: "bucket".to_string(),
column: Some(column),
num_buckets: Some(num_buckets),
maintained_indexes,
maintained_indexes: Some(maintained_indexes),
writer_config_defaults: Some(writer_config_defaults),
},
Native::Identity {
@@ -838,7 +838,7 @@ impl From<lancedb::table::LsmWriteSpec> for LsmWriteSpec {
spec_type: "identity".to_string(),
column: Some(column),
num_buckets: None,
maintained_indexes,
maintained_indexes: Some(maintained_indexes),
writer_config_defaults: Some(writer_config_defaults),
},
Native::Unsharded {
@@ -848,7 +848,7 @@ impl From<lancedb::table::LsmWriteSpec> for LsmWriteSpec {
spec_type: "unsharded".to_string(),
column: None,
num_buckets: None,
maintained_indexes,
maintained_indexes: Some(maintained_indexes),
writer_config_defaults: Some(writer_config_defaults),
},
}
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb-python"
version = "0.37.1-beta.1"
version = "0.37.1"
publish = false
edition.workspace = true
description = "Python bindings for LanceDB"
+1 -1
View File
@@ -63,7 +63,7 @@ tests = [
"polars>=0.19, <=1.32.3",
"pyarrow<25",
"pyarrow-stubs>=16.0",
"pylance==9.0.0rc1",
"pylance==10.0.0",
"requests>=2.31.0",
"datafusion>=54,<55",
"opentelemetry-sdk>=1.30.0",
+4 -7
View File
@@ -653,10 +653,9 @@ class LsmWriteSpec:
def identity(column: str) -> "LsmWriteSpec": ...
@staticmethod
def unsharded() -> "LsmWriteSpec": ...
def with_maintained_indexes(self, indexes: Optional[List[str]]) -> "LsmWriteSpec":
"""Set which indexes the MemWAL keeps up to date. None resolves every
index on the table at install, failing if one cannot be maintained;
a list is verbatim, empty means none."""
def with_maintained_indexes(self, indexes: List[str]) -> "LsmWriteSpec":
"""Return a copy of this spec asking the MemWAL to keep the named
indexes up to date as rows are appended."""
...
def with_writer_config_defaults(self, defaults: Dict[str, str]) -> "LsmWriteSpec":
"""Return a copy of this spec recording the given default
@@ -671,9 +670,7 @@ class LsmWriteSpec:
@property
def num_buckets(self) -> Optional[int]: ...
@property
def maintained_indexes(self) -> Optional[List[str]]:
"""Indexes the MemWAL keeps up to date, or None for every supported one."""
...
def maintained_indexes(self) -> List[str]: ...
@property
def writer_config_defaults(self) -> Dict[str, str]: ...
+27 -315
View File
@@ -11,11 +11,6 @@ Provides StreamingDataset, a PyTorch IterableDataset that guarantees:
- **Resumability**: state_dict / load_state_dict capture per-split consumption
counts so training can resume from an exact mid-epoch position even when the
distributed topology changes between runs.
Transform failures on bad rows (e.g. nulls or NaNs from incomplete data) can
be tolerated with ``on_transform_error="skip"``; see the parameter
documentation on StreamingDataset for how this interacts with the guarantees
above.
"""
import ctypes
@@ -27,7 +22,7 @@ import time
from collections import deque
from concurrent.futures import ThreadPoolExecutor
from multiprocessing import RawArray
from typing import Any, Callable, Iterator, Optional, Union
from typing import Any, Callable, Iterator, Optional
from torch.utils.data import IterableDataset, get_worker_info
@@ -132,49 +127,6 @@ class StreamingDataset(IterableDataset):
Maximum number of transforms to run concurrently. Must be greater
than zero. When ``None`` (the default), uses ``os.cpu_count()`` or 1
when the CPU count is unavailable.
on_transform_error:
What to do when the transform raises an exception:
- ``"raise"`` (the default): the exception propagates and iteration
aborts.
- ``"skip"``: the failing rows are dropped and iteration continues.
- ``"warn"``: like ``"skip"``, but a warning is logged for each
failing batch.
- a callable ``handler(exc) -> bool``: called with the exception;
return ``True`` to skip the failing rows or ``False`` to re-raise.
Useful to skip only expected error types (compatible with
``webdataset.handlers`` style handlers).
When a batch fails, the transform is re-invoked on each single-row
slice of the batch so that only the rows that actually fail are
dropped. Transforms should therefore be deterministic and accept
batches of any size (including one row). Skipped rows are counted in
``rows_skipped``.
Skipping weakens the elastic-determinism guarantee at the end of the
epoch: splits that lose more rows than others run dry earlier, and
each rank's iterator ends at the last cycle where every split *it
owns* still has a row. Because bad rows are not distributed evenly
across splits, this means one rank's iterator can yield noticeably
fewer or more steps than another rank's *in the same run* — there is
no cross-rank coordination that stops every rank at the same global
step. This is generally safe for asynchronous or single-rank use,
but synchronous distributed training (e.g. ranks that call
``all_reduce`` every step) can hang or deadlock if one rank's
iterator is exhausted while others are still stepping; callers doing
synchronous multi-rank training with ``on_transform_error != "raise"``
are responsible for their own cross-rank stopping mechanism (e.g.
broadcasting a stop signal on ``StopIteration``). The final few
global steps can also differ across topologies (bounded by the skew
in bad-row counts across splits). The sequence of samples yielded
from each split remains deterministic. Mid-epoch
checkpoints remain exact provided the transform fails
deterministically; in multi-rank training each rank must save its
own ``state_dict`` and the states must be combined with
``merge_state_dicts`` before resuming on a different topology.
Prefer the ``filter`` parameter when bad rows can be expressed as a
SQL predicate (e.g. ``"col IS NOT NULL"``) — filtering happens before
splits are built, so every guarantee is fully preserved.
worker_info_override:
If set, used in place of ``torch.utils.data.get_worker_info()`` to
determine the DataLoader worker assignment. Intended for unit tests
@@ -200,7 +152,6 @@ class StreamingDataset(IterableDataset):
filter: Optional[str] = None,
transform: Optional[Callable] = None,
transform_parallelism: Optional[int] = None,
on_transform_error: Union[str, Callable[[Exception], bool]] = "raise",
connection_factory: Optional[Callable[[str], Any]] = None,
worker_info_override=None,
):
@@ -216,13 +167,6 @@ class StreamingDataset(IterableDataset):
)
if transform_parallelism is not None and transform_parallelism <= 0:
raise ValueError("transform_parallelism must be greater than 0")
if on_transform_error not in ("raise", "skip", "warn") and not callable(
on_transform_error
):
raise ValueError(
"on_transform_error must be 'raise', 'skip', 'warn', or a "
f"callable, got {on_transform_error!r}"
)
self._table = table
self._num_splits = num_splits
@@ -238,7 +182,6 @@ class StreamingDataset(IterableDataset):
self._filter = filter
self._transform = transform
self._transform_parallelism = transform_parallelism
self._on_transform_error = on_transform_error
self._connection_factory = connection_factory
self._worker_info_override = worker_info_override
@@ -256,28 +199,19 @@ class StreamingDataset(IterableDataset):
# in the main process. RawArray is picklable via the forkserver
# reduction protocol so it survives the dataset pickle round-trip.
# Layout: [unscanned_rows, raw_rows, cooked_rows, consumed_rows,
# bytes_loaded, fetch_time_us, transform_time_us,
# rows_skipped]
self._worker_stats: RawArray = RawArray(ctypes.c_int64, 8)
# bytes_loaded, fetch_time_us, transform_time_us]
self._worker_stats: RawArray = RawArray(ctypes.c_int64, 7)
# Cumulative bytes of Arrow buffer data fetched across all iterations.
self._bytes_loaded: int = 0
# Cumulative seconds spent in LanceDB I/O and in transform functions.
self._fetch_time: float = 0.0
self._transform_time: float = 0.0
# Cumulative rows dropped by on_transform_error across all iterations.
self._rows_skipped: int = 0
# Number of samples each split has already been consumed. At global
# step boundaries all splits have consumed this many samples, so a
# single scalar captures the topology-independent checkpoint state.
self._resume_offset: int = 0
# Permutation position each split has consumed through, keyed by
# global split index. Equal to _resume_offset for every split unless
# on_transform_error skipped rows, in which case skipped positions
# push the watermark of the affected splits further ahead. Splits
# this instance has never iterated have no entry.
self._resume_positions: dict[int, int] = {}
# Build the permutation table once, deterministically.
builder = permutation_builder(table)
@@ -341,7 +275,6 @@ class StreamingDataset(IterableDataset):
# Set identity transform on each Permutation so __getitems__ returns
# the raw RecordBatch. Stage 2 applies the real transform.
permutations: list[Permutation] = []
initial_positions: list[int] = []
for split_idx in my_splits:
perm = Permutation.from_tables(
self._table, self._perm_table, split=split_idx
@@ -349,20 +282,14 @@ class StreamingDataset(IterableDataset):
if self._columns is not None:
perm = perm.select_columns(self._columns)
perm = perm.with_transform(lambda batch: batch)
start_pos = self._resume_positions.get(split_idx, self._resume_offset)
if start_pos > 0:
perm = perm.with_skip(start_pos)
initial_positions.append(start_pos)
if self._resume_offset > 0:
perm = perm.with_skip(self._resume_offset)
permutations.append(perm)
n = len(permutations)
split_sizes = [perm.num_rows for perm in permutations]
initial_offset = self._resume_offset
local_consumed = [0] * n
# Permutation position each split has consumed through (absolute,
# i.e. counted from the start of the unskipped split). Runs ahead of
# initial + local_consumed when rows are skipped.
pos_consumed = list(initial_positions)
batch_size = self._read_batch_size
max_prefetch = self._prefetch_batches
@@ -375,14 +302,12 @@ class StreamingDataset(IterableDataset):
self._transform if self._transform is not None else Transforms.arrow2python
)
# Per-split pipeline state. Batches are paired with the absolute
# permutation position of their first row so that skipped rows can be
# accounted for in pos_consumed.
# Per-split pipeline state.
fetch_head = [0] * n
io_pending = [deque() for _ in range(n)] # (abs_start, Future[RecordBatch])
raw_batches = [deque() for _ in range(n)] # (abs_start, RecordBatch)
tx_pending = [deque() for _ in range(n)] # Future[list[(abs_pos, row)]]
cooked = [deque() for _ in range(n)] # (abs_pos, row) ready to yield
io_pending = [deque() for _ in range(n)] # Future[RecordBatch]
raw_batches = [deque() for _ in range(n)] # RecordBatch — fetched, awaiting tx
tx_pending = [deque() for _ in range(n)] # Future[list[Any]]
cooked = [deque() for _ in range(n)] # rows ready to yield
# Limit simultaneous transforms to transform_workers across all splits.
tx_semaphore = threading.Semaphore(transform_workers)
@@ -405,8 +330,7 @@ class StreamingDataset(IterableDataset):
fetch_head[i] += fetch
perm_i = permutations[i]
indices = list(range(start, start + fetch))
abs_start = initial_positions[i] + start
io_pending[i].append((abs_start, io_pool.submit(_io_call, perm_i, indices)))
io_pending[i].append(io_pool.submit(_io_call, perm_i, indices))
def _fill_io(i: int) -> None:
while len(io_pending[i]) < max_prefetch and fetch_head[i] < split_sizes[i]:
@@ -414,72 +338,15 @@ class StreamingDataset(IterableDataset):
def _drain_io(i: int) -> None:
"""Move completed I/O futures into raw_batches non-blockingly."""
while io_pending[i] and io_pending[i][0][1].done():
abs_start, fut = io_pending[i].popleft()
raw_batches[i].append((abs_start, fut.result()))
while io_pending[i] and io_pending[i][0].done():
raw_batches[i].append(io_pending[i].popleft().result())
# ── Stage 2 helpers ───────────────────────────────────────────────────
on_error = self._on_transform_error
def _should_skip(exc: Exception) -> bool:
if on_error == "raise":
return False
if callable(on_error):
return bool(on_error(exc))
return True # "skip" or "warn"
def _check_row_count(rows: list, num_rows: int) -> None:
if len(rows) != num_rows:
raise ValueError(
f"transform returned {len(rows)} rows for a batch of "
f"{num_rows}; transforms must return exactly one output "
"row per input row. To drop bad rows, raise inside the "
"transform and pass on_transform_error='skip'."
)
def _transform_isolated(abs_start, batch, batch_exc):
"""Re-run the transform on single-row slices, dropping failures."""
out = []
skipped = 0
first_exc = None
for j in range(batch.num_rows):
try:
rows = list(final_transform(batch.slice(j, 1)))
except Exception as exc:
if not _should_skip(exc):
raise
skipped += 1
if first_exc is None:
first_exc = exc
continue
_check_row_count(rows, 1)
out.append((abs_start + j, rows[0]))
self._rows_skipped += skipped
if skipped and on_error == "warn":
logger.warning(
"Skipped %d of %d rows whose transform failed (first error: %r)",
skipped,
batch.num_rows,
first_exc if first_exc is not None else batch_exc,
)
return out
def _transform_batch(abs_start, batch):
"""Apply the transform, returning [(abs_pos, row), ...]."""
try:
rows = list(final_transform(batch))
except Exception as exc:
if not _should_skip(exc):
raise
return _transform_isolated(abs_start, batch, exc)
_check_row_count(rows, batch.num_rows)
return [(abs_start + j, row) for j, row in enumerate(rows)]
def _tx_call_guarded(abs_start, batch):
def _tx_call_guarded(batch):
try:
t0 = time.perf_counter()
result = _transform_batch(abs_start, batch)
result = final_transform(batch)
self._transform_time += time.perf_counter() - t0
return result
finally:
@@ -488,8 +355,8 @@ class StreamingDataset(IterableDataset):
def _try_submit_tx(i: int) -> None:
"""Submit transforms for raw_batches[i] up to available capacity."""
while raw_batches[i] and tx_semaphore.acquire(blocking=False):
abs_start, batch = raw_batches[i].popleft()
tx_pending[i].append(tx_pool.submit(_tx_call_guarded, abs_start, batch))
batch = raw_batches[i].popleft()
tx_pending[i].append(tx_pool.submit(_tx_call_guarded, batch))
def _drain_tx(i: int) -> None:
"""Move completed transform futures into cooked non-blockingly."""
@@ -517,14 +384,11 @@ class StreamingDataset(IterableDataset):
# Acquire a transform slot (may block briefly if all
# transform_workers are busy with other splits).
tx_semaphore.acquire()
abs_start, batch = raw_batches[i].popleft()
tx_pending[i].append(
tx_pool.submit(_tx_call_guarded, abs_start, batch)
)
batch = raw_batches[i].popleft()
tx_pending[i].append(tx_pool.submit(_tx_call_guarded, batch))
elif io_pending[i]:
# Block on the oldest in-flight I/O fetch.
abs_start, fut = io_pending[i].popleft()
raw_batches[i].append((abs_start, fut.result()))
raw_batches[i].append(io_pending[i].popleft().result())
_advance(i)
else:
break # split exhausted
@@ -543,28 +407,15 @@ class StreamingDataset(IterableDataset):
_fill_io(i)
while True:
# A cycle only runs if every split can still produce a
# row. Without skips all splits exhaust simultaneously
# (equal split sizes + round-robin); when
# on_transform_error drops rows a split can run dry
# early, ending the epoch at the last complete cycle.
# This check only sees splits owned by this rank/worker
# (my_splits) — there is no cross-rank coordination, so
# a different rank with fewer skipped rows keeps going;
# see the on_transform_error docstring.
exhausted = False
for i in range(n):
_ensure_cooked(i)
if not cooked[i]:
exhausted = True
break
if exhausted:
# Stop when any split is exhausted (all exhaust
# simultaneously: equal split sizes + round-robin).
if any(local_consumed[i] >= split_sizes[i] for i in range(n)):
break
for i in range(n):
pos, row = cooked[i].popleft()
_ensure_cooked(i)
row = cooked[i].popleft()
local_consumed[i] += 1
pos_consumed[i] = pos + 1
_advance(i)
# After the last split in each cycle: update the
@@ -573,39 +424,21 @@ class StreamingDataset(IterableDataset):
# even when __iter__ runs in a worker process.
if i == n - 1:
self._resume_offset = initial_offset + local_consumed[i]
for j, split_idx in enumerate(my_splits):
self._resume_positions[split_idx] = pos_consumed[j]
ws = self._worker_stats
ws[0] = sum(
split_sizes[j] - fetch_head[j] for j in range(n)
)
ws[1] = sum(
batch.num_rows
for q in raw_batches
for _, batch in q
batch.num_rows for q in raw_batches for batch in q
)
ws[2] = sum(len(q) for q in cooked)
ws[3] = sum(local_consumed)
ws[4] = self._bytes_loaded
ws[5] = int(self._fetch_time * 1_000_000)
ws[6] = int(self._transform_time * 1_000_000)
ws[7] = self._rows_skipped
yield row
finally:
# Final stats flush: the per-cycle write above never runs
# when iteration ends mid-cycle (e.g. a split whose rows
# were all skipped before completing a single cycle), so
# counters like rows_skipped would otherwise be stale.
ws = self._worker_stats
ws[0] = sum(split_sizes[j] - fetch_head[j] for j in range(n))
ws[1] = 0 # queue-depth properties document 0 when idle
ws[2] = 0
ws[3] = sum(local_consumed)
ws[4] = self._bytes_loaded
ws[5] = int(self._fetch_time * 1_000_000)
ws[6] = int(self._transform_time * 1_000_000)
ws[7] = self._rows_skipped
self._raw_batches_ref = None
self._cooked_ref = None
self._fetch_head_ref = None
@@ -659,7 +492,7 @@ class StreamingDataset(IterableDataset):
batches. Returns 0 when not iterating.
"""
if self._raw_batches_ref is not None:
return sum(batch.num_rows for q in self._raw_batches_ref for _, batch in q)
return sum(batch.num_rows for q in self._raw_batches_ref for batch in q)
return int(self._worker_stats[1])
@property
@@ -689,19 +522,6 @@ class StreamingDataset(IterableDataset):
)
return int(self._worker_stats[0])
@property
def rows_skipped(self) -> int:
"""Number of rows dropped because their transform raised an exception.
Only ever non-zero when ``on_transform_error`` is set to ``"skip"``,
``"warn"``, or a callable that returned ``True``. Accumulates across
multiple iterations of the same dataset instance and is never reset
automatically.
"""
if self._raw_batches_ref is not None:
return self._rows_skipped
return int(self._worker_stats[7])
@property
def consumed_rows(self) -> int:
"""Number of rows already yielded to the caller across all splits.
@@ -767,27 +587,12 @@ class StreamingDataset(IterableDataset):
every split has been consumed the same number of times (by the
round-robin design), so the per-split count is a single uniform value
that is identical across all ranks and DataLoader workers.
``positions_consumed_per_split`` records how far into each split's
permutation iteration has advanced. It only differs from
``samples_consumed_per_split`` when ``on_transform_error`` skipped
rows, in which case entries are exact for the splits this instance
iterated and a lower bound (the sample count) for splits owned by
other ranks or workers. Combine the state dicts from all ranks with
[merge_state_dicts][lancedb.streaming.StreamingDataset.merge_state_dicts]
to recover the exact value for every split before resuming on a
different topology.
"""
positions = [
self._resume_positions.get(split, self._resume_offset)
for split in range(self._num_splits)
]
return {
"shuffle_seed": self._shuffle_seed,
"num_splits": self._num_splits,
"epoch": self._epoch,
"samples_consumed_per_split": [self._resume_offset] * self._num_splits,
"positions_consumed_per_split": positions,
}
def load_state_dict(self, state: dict) -> None:
@@ -813,96 +618,3 @@ class StreamingDataset(IterableDataset):
self._resume_offset = consumed[0] if consumed else 0
else:
self._resume_offset = int(consumed)
# Older checkpoints predate positions_consumed_per_split; without
# skipped rows positions equal sample counts, so falling back to
# _resume_offset (the .get default in __iter__) is exact.
positions = state.get("positions_consumed_per_split")
if positions is None:
self._resume_positions = {}
else:
self._resume_positions = {
split: int(pos) for split, pos in enumerate(positions)
}
@staticmethod
def merge_state_dicts(states: list[dict]) -> dict:
"""Merge state dicts saved by different ranks into one exact state.
Only needed when ``on_transform_error`` skips rows in multi-rank
training: each rank then knows the exact permutation position only for
its own splits, and records a lower bound for the rest. Because
exactly one rank owns each split, the elementwise maximum across all
ranks' ``positions_consumed_per_split`` recovers the exact position of
every split. Without skipped rows every rank's state is already
identical and merging is a no-op.
Raises ``ValueError`` if the states are empty or were not produced by
the same run (mismatched seed, split count, epoch, or sample counts).
The merge is always all-to-all and topology-agnostic: collect the
``state_dict()`` from every rank of the *previous* run into one list,
merge that whole list, and hand the identical merged result to every
rank of the *next* run — regardless of whether the rank count grew,
shrank, or stayed the same. There is no pairwise or subset merging
step, because each split's exact position is only known to whichever
rank owned that split, and the elementwise maximum needs every rank's
contribution to be correct.
For example, checkpointing 8 ranks and resuming on 4 (the same
pattern applies when growing, e.g. 4 ranks resuming on 8)::
states = [ds.state_dict() for ds in previous_run_datasets] # 8
merged = StreamingDataset.merge_state_dicts(states)
for ds in resumed_datasets: # now only 4 ranks
ds.load_state_dict(merged) # same dict on every rank
The rank count on either side never affects the merge itself, since
``merge_state_dicts`` only cares about the list of states it is
given. Each split's position is recovered by elementwise maximum;
here rank 0 owned split 0 (and skipped two rows there) while rank 1
owned split 1 (and skipped one row):
>>> rank0 = {
... "shuffle_seed": 0, "num_splits": 2, "epoch": 0,
... "samples_consumed_per_split": [3, 3],
... "positions_consumed_per_split": [5, 3],
... }
>>> rank1 = {
... "shuffle_seed": 0, "num_splits": 2, "epoch": 0,
... "samples_consumed_per_split": [3, 3],
... "positions_consumed_per_split": [3, 4],
... }
>>> merged = StreamingDataset.merge_state_dicts([rank0, rank1])
>>> merged["positions_consumed_per_split"]
[5, 4]
"""
if not states:
raise ValueError("merge_state_dicts requires at least one state dict")
first = states[0]
for state in states[1:]:
for key in ("shuffle_seed", "num_splits", "epoch"):
if state[key] != first[key]:
raise ValueError(
f"{key} mismatch across state dicts: "
f"{state[key]} != {first[key]}"
)
if (
state["samples_consumed_per_split"]
!= first["samples_consumed_per_split"]
):
raise ValueError(
"samples_consumed_per_split mismatch across state dicts; "
"state_dict() must be called at the same global step "
"boundary on every rank"
)
merged = dict(first)
all_positions = [
state.get(
"positions_consumed_per_split", state["samples_consumed_per_split"]
)
for state in states
]
merged["positions_consumed_per_split"] = [
max(per_split) for per_split in zip(*all_positions)
]
return merged
+3 -10
View File
@@ -4676,13 +4676,6 @@ class AsyncTable:
via [`set_unenforced_primary_key`]; bucket sharding additionally
requires it to be the single column being bucketed.
By default the MemWAL maintains every index on the table, resolved
here — a snapshot, so an index created afterwards needs the spec unset
and set again. This fails if one cannot be maintained; name the set
with ``with_maintained_indexes`` to install anyway. That pins an exact
set (a still-building index is rejected, not omitted); ``[]`` maintains
none.
Parameters
----------
spec : LsmWriteSpec
@@ -4709,9 +4702,9 @@ class AsyncTable:
Returns ``None`` when the MemWAL LSM write path is not enabled (no
spec has been set, or it was removed with `unset_lsm_write_spec`).
The returned spec mirrors what was passed to `set_lsm_write_spec`,
except that ``maintained_indexes`` always reports the concrete list
resolved when the spec was set — ``None`` never round-trips.
The returned spec — including its ``maintained_indexes`` and
``writer_config_defaults`` — mirrors what was passed to
`set_lsm_write_spec`.
"""
return await self._inner.get_lsm_write_spec()
@@ -1456,408 +1456,6 @@ def test_shuffle_clump_size_yields_all_rows(lance_table):
)
# ---------------------------------------------------------------------------
# on_transform_error tests
# ---------------------------------------------------------------------------
class BadRowError(ValueError):
"""Raised by the failing transforms below when a batch contains a bad id."""
def _failing_transform(bad_ids: set):
"""A transform that raises BadRowError whenever the batch has a bad id.
Raises on the full batch and on any single-row slice containing a bad id,
so per-row isolation drops exactly the bad rows.
"""
def transform(batch: pa.RecordBatch) -> list:
ids = batch.column("id").to_pylist()
bad = sorted(set(ids) & bad_ids)
if bad:
raise BadRowError(f"bad ids in batch: {bad}")
return [{"id": i} for i in ids]
return transform
def _sequential_split_members(table) -> list[list[int]]:
"""Return each split's ids in yield order for shuffle=False.
With a single rank and no workers the round-robin yields one row per split
per cycle, so item k of a clean run belongs to split k % NUM_SPLITS.
"""
ds = StreamingDataset(table, num_splits=NUM_SPLITS, shuffle=False)
members: list[list[int]] = [[] for _ in range(NUM_SPLITS)]
for k, row in enumerate(ds):
members[k % NUM_SPLITS].append(row["id"])
return members
def test_on_transform_error_default_raises(lance_table):
"""By default a transform exception propagates and aborts iteration."""
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle_seed=SHUFFLE_SEED,
transform=_failing_transform({7}),
)
with pytest.raises(BadRowError):
list(ds)
def test_on_transform_error_invalid_value(lance_table):
with pytest.raises(ValueError, match="on_transform_error"):
StreamingDataset(lance_table, num_splits=NUM_SPLITS, on_transform_error="bogus")
def test_on_transform_error_skip_drops_bad_rows(lance_table):
"""With one bad row per split, 'skip' yields every good row exactly once
and counts the dropped rows in rows_skipped."""
members = _sequential_split_members(lance_table)
bad_ids = {members[i][4] for i in range(NUM_SPLITS)}
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle=False,
transform=_failing_transform(bad_ids),
on_transform_error="skip",
)
assert ds.rows_skipped == 0
ids = [row["id"] for row in ds]
assert sorted(ids) == sorted(set(range(NUM_ROWS)) - bad_ids)
assert ds.rows_skipped == NUM_SPLITS
def test_on_transform_error_skip_uneven_ends_at_last_complete_cycle(lance_table):
"""When one split loses more rows than the others, the epoch ends at the
last cycle where every split still has a row no crash, no bad rows, and
every step remains one sample per split."""
members = _sequential_split_members(lance_table)
bad_ids = set(members[0][:3]) # all 3 bad rows in split 0
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle=False,
transform=_failing_transform(bad_ids),
on_transform_error="skip",
)
items = [row["id"] for row in ds]
rows_per_split = NUM_ROWS // NUM_SPLITS
expected_cycles = rows_per_split - len(bad_ids)
assert len(items) == expected_cycles * NUM_SPLITS
assert len(set(items)) == len(items), "duplicate samples yielded"
assert not set(items) & bad_ids, "a bad row was yielded"
# Split 0 contributed exactly its surviving rows, in order, one per cycle.
survivors = [i for i in members[0] if i not in bad_ids]
assert items[0::NUM_SPLITS] == survivors[:expected_cycles]
def test_on_transform_error_warn_logs(lance_table, caplog):
"""'warn' skips like 'skip' but logs a warning for the failing batch."""
members = _sequential_split_members(lance_table)
bad_ids = {members[i][3] for i in range(NUM_SPLITS)}
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle=False,
transform=_failing_transform(bad_ids),
on_transform_error="warn",
)
with caplog.at_level(logging.WARNING, logger="lancedb.streaming"):
items = list(ds)
assert len(items) == NUM_ROWS - NUM_SPLITS
assert ds.rows_skipped == NUM_SPLITS
assert "Skipped" in caplog.text
assert "BadRowError" in caplog.text
def test_on_transform_error_callable_selective(lance_table):
"""A callable handler can skip expected errors and re-raise the rest."""
members = _sequential_split_members(lance_table)
bad_ids = {members[i][0] for i in range(NUM_SPLITS)}
handled: list[Exception] = []
def handler(exc: Exception) -> bool:
handled.append(exc)
return isinstance(exc, BadRowError)
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle=False,
transform=_failing_transform(bad_ids),
on_transform_error=handler,
)
items = list(ds)
assert len(items) == NUM_ROWS - NUM_SPLITS
assert handled and all(isinstance(exc, BadRowError) for exc in handled)
def broken_transform(batch: pa.RecordBatch) -> list:
raise TypeError("boom")
ds2 = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle=False,
transform=broken_transform,
on_transform_error=handler,
)
with pytest.raises(TypeError, match="boom"):
list(ds2)
def test_transform_wrong_row_count_raises(lance_table):
"""A transform that returns the wrong number of rows is an error even with
on_transform_error='skip' silent shrinkage would corrupt accounting."""
def drops_rows(batch: pa.RecordBatch) -> list:
return batch.column("id").to_pylist()[:-1]
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle_seed=SHUFFLE_SEED,
transform=drops_rows,
on_transform_error="skip",
)
with pytest.raises(ValueError, match="one output row per input row"):
list(ds)
def test_skip_deterministic_across_runs(lance_table):
"""With a fixed seed, skipping produces the identical sample sequence on
every run skips are data-dependent, not run-dependent."""
bad_ids = {5, 17, 46}
def run() -> tuple[list[int], int]:
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle_seed=SHUFFLE_SEED,
transform=_failing_transform(bad_ids),
on_transform_error="skip",
)
return [row["id"] for row in ds], ds.rows_skipped
ids_a, skipped_a = run()
ids_b, skipped_b = run()
assert ids_a == ids_b
assert skipped_a == skipped_b
assert not set(ids_a) & bad_ids
def test_skip_elastic_det_across_world_sizes(lance_table):
"""With equal bad-row counts per split, skipping preserves the full
elastic-determinism guarantee: identical global batches at every step for
every compatible world_size."""
members = _sequential_split_members(lance_table)
bad_ids = {members[i][6] for i in range(NUM_SPLITS)}
def collect(world_size: int) -> list[frozenset[int]]:
micro = GLOBAL_BATCH_SIZE // world_size
iters = [
iter(
StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle=False,
rank=rank,
world_size=world_size,
transform=_failing_transform(bad_ids),
on_transform_error="skip",
)
)
for rank in range(world_size)
]
_STOP = object()
batches: list[frozenset[int]] = []
while True:
step_samples: set[int] = set()
exhausted = 0
for it in iters:
for _ in range(micro):
val = next(it, _STOP)
if val is _STOP:
exhausted += 1
break
step_samples.add(val["id"])
if exhausted == len(iters):
break
assert exhausted == 0, (
"Rank iterators exhausted at different steps despite equal "
"bad-row counts per split"
)
batches.append(frozenset(step_samples))
return batches
reference = collect(1)
assert len(reference) == NUM_ROWS // NUM_SPLITS - 1
for ws in (2, 3, 4):
assert collect(ws) == reference, f"world_size={ws} diverged"
def test_resumability_with_skips_same_topology(lance_table):
"""Checkpointing mid-epoch with skipped rows resumes exactly: no sample
repeated, no sample lost, skipped rows stay skipped."""
members = _sequential_split_members(lance_table)
# Uneven skips: positions diverge across splits (2 bad in split 0, 1 in
# split 5), which only a position-based checkpoint can resume exactly.
bad_ids = {members[0][2], members[0][3], members[5][7]}
kwargs = dict(
num_splits=NUM_SPLITS,
shuffle=False,
transform=_failing_transform(bad_ids),
on_transform_error="skip",
)
reference = [row["id"] for row in StreamingDataset(lance_table, **kwargs)]
rows_per_split = NUM_ROWS // NUM_SPLITS
assert len(reference) == (rows_per_split - 2) * NUM_SPLITS
steps = 3
ds = StreamingDataset(lance_table, **kwargs)
it = iter(ds)
consumed = [next(it)["id"] for _ in range(steps * NUM_SPLITS)]
checkpoint = ds.state_dict()
it.close()
# Split 0 skipped positions 2 and 3 within its first 3 yields; split 5's
# bad row is beyond the checkpoint. Everything else is at 3 = the sample
# count.
positions = checkpoint["positions_consumed_per_split"]
assert positions[0] == 5
assert positions[1:] == [3] * (NUM_SPLITS - 1)
assert checkpoint["samples_consumed_per_split"] == [3] * NUM_SPLITS
ds2 = StreamingDataset(lance_table, **kwargs)
ds2.load_state_dict(checkpoint)
resumed = [row["id"] for row in ds2]
assert consumed == reference[: steps * NUM_SPLITS]
assert resumed == reference[steps * NUM_SPLITS :]
def test_resumability_with_skips_elastic_merge(lance_table):
"""Elastic resume with skips: each rank's checkpoint knows exact positions
only for its own splits; merge_state_dicts recovers the global state, and
a run on a different world_size continues exactly."""
members = _sequential_split_members(lance_table)
# Bad rows early in split 0 (rank 0) and split 6 (rank 1 of a ws=2 run) so
# both ranks' position vectors diverge before the checkpoint.
bad_ids = {members[0][0], members[0][2], members[6][1]}
kwargs = dict(
num_splits=NUM_SPLITS,
shuffle=False,
transform=_failing_transform(bad_ids),
on_transform_error="skip",
)
reference = [row["id"] for row in StreamingDataset(lance_table, **kwargs)]
steps = 3
world_size = 2
micro = GLOBAL_BATCH_SIZE // world_size
datasets = [
StreamingDataset(lance_table, rank=rank, world_size=world_size, **kwargs)
for rank in range(world_size)
]
iters = [iter(ds) for ds in datasets]
seen: list[frozenset[int]] = []
for _ in range(steps):
step_samples = set()
for it in iters:
for _ in range(micro):
step_samples.add(next(it)["id"])
seen.append(frozenset(step_samples))
states = [ds.state_dict() for ds in datasets]
for it in iters:
it.close()
merged = StreamingDataset.merge_state_dicts(states)
expected_positions = [3] * NUM_SPLITS
expected_positions[0] = 5 # skipped positions 0 and 2
expected_positions[6] = 4 # skipped position 1
assert merged["positions_consumed_per_split"] == expected_positions
# The first 3 global batches match the world_size=1 reference.
ref_batches = [
frozenset(reference[s * NUM_SPLITS : (s + 1) * NUM_SPLITS])
for s in range(len(reference) // NUM_SPLITS)
]
assert seen == ref_batches[:steps]
# Resume on world_size=1 from the merged state.
ds_resume = StreamingDataset(lance_table, **kwargs)
ds_resume.load_state_dict(merged)
resumed = [row["id"] for row in ds_resume]
assert resumed == reference[steps * NUM_SPLITS :]
def test_rows_skipped_flushed_when_split_entirely_bad(lance_table):
"""A split whose rows all fail never completes a cycle, so the epoch ends
immediately but rows_skipped must still report the drops after the
iterator exits (the shared-memory counter is flushed on exhaustion)."""
members = _sequential_split_members(lance_table)
bad_ids = set(members[0]) # every row of split 0 is bad
ds = StreamingDataset(
lance_table,
num_splits=NUM_SPLITS,
shuffle=False,
transform=_failing_transform(bad_ids),
on_transform_error="skip",
)
assert list(ds) == []
assert ds.rows_skipped == len(bad_ids)
def test_merge_state_dicts_validates_consistency(lance_table):
ds = StreamingDataset(lance_table, num_splits=NUM_SPLITS, shuffle_seed=SHUFFLE_SEED)
state = ds.state_dict()
other = dict(state, shuffle_seed=SHUFFLE_SEED + 1)
with pytest.raises(ValueError, match="shuffle_seed mismatch"):
StreamingDataset.merge_state_dicts([state, other])
with pytest.raises(ValueError, match="at least one"):
StreamingDataset.merge_state_dicts([])
def test_load_state_dict_without_positions_key(lance_table):
"""Checkpoints from before positions_consumed_per_split existed still
resume exactly (positions equal sample counts when nothing is skipped)."""
reference = [
row["id"]
for row in StreamingDataset(
lance_table, num_splits=NUM_SPLITS, shuffle_seed=SHUFFLE_SEED
)
]
steps = 4
ds = StreamingDataset(lance_table, num_splits=NUM_SPLITS, shuffle_seed=SHUFFLE_SEED)
it = iter(ds)
for _ in range(steps * NUM_SPLITS):
next(it)
checkpoint = ds.state_dict()
it.close()
del checkpoint["positions_consumed_per_split"]
ds2 = StreamingDataset(
lance_table, num_splits=NUM_SPLITS, shuffle_seed=SHUFFLE_SEED
)
ds2.load_state_dict(checkpoint)
resumed = [row["id"] for row in ds2]
assert resumed == reference[steps * NUM_SPLITS :]
def test_num_splits_defaults_to_world_size(lance_table):
"""Omitting num_splits gives world_size splits (one per rank)."""
ds = StreamingDataset(
+4 -11
View File
@@ -83,9 +83,7 @@ def test_lsm_write_spec_repr():
assert s.spec_type == "bucket"
assert s.column == "id"
assert s.num_buckets == 4
# A fresh spec defers its maintained set to install time.
assert s.maintained_indexes is None
assert s.with_maintained_indexes([]).maintained_indexes == []
assert s.maintained_indexes == []
assert "bucket" in repr(s)
assert "id" in repr(s)
assert "4" in repr(s)
@@ -171,23 +169,18 @@ def test_get_lsm_write_spec(tmp_path):
table.unset_lsm_write_spec()
assert table.get_lsm_write_spec() is None
# Identity round-trips (column recovered from the schema). Leaving the
# maintained set to be inferred picks up the index on the table, so the
# spec reads back naming it rather than as "infer".
# Identity round-trips (column recovered from the schema).
table.set_lsm_write_spec(LsmWriteSpec.identity("id"))
spec = table.get_lsm_write_spec()
assert spec.spec_type == "identity"
assert spec.column == "id"
assert spec.maintained_indexes == [idx_name]
table.unset_lsm_write_spec()
# Unsharded round-trips (no routing column). Opting out is distinct from
# the inferred default.
table.set_lsm_write_spec(LsmWriteSpec.unsharded().with_maintained_indexes([]))
# Unsharded round-trips (no routing column).
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
spec = table.get_lsm_write_spec()
assert spec.spec_type == "unsharded"
assert spec.column is None
assert spec.maintained_indexes == []
@pytest.mark.asyncio
+2 -2
View File
@@ -544,7 +544,7 @@ def test_lsm_read_fts_unmaintained_index_errors(tmp_path):
table.create_index("text", config=FTS())
# No maintained indexes: the active memtable FTS arm cannot serve un-compacted
# docs, so the search would silently omit them — reject instead.
table.set_lsm_write_spec(LsmWriteSpec.unsharded().with_maintained_indexes([]))
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
with pytest.raises(Exception, match="maintained"):
table.search("fox", query_type="fts", fts_columns="text").to_arrow()
@@ -631,7 +631,7 @@ def test_lsm_read_vector_unmaintained_index_errors(tmp_path):
)
# Spec with NO maintained indexes: the base vector index's catch-up is untracked,
# so the scanner rejects rather than risk dropping compacted-but-unindexed rows.
table.set_lsm_write_spec(LsmWriteSpec.unsharded().with_maintained_indexes([]))
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
with pytest.raises(Exception, match="maintained"):
table.search([1.0] * VECTOR_DIM).to_arrow()
+1 -1
View File
@@ -289,7 +289,7 @@ struct IvfHnswFlatParams {
target_partition_size: Option<u32>,
}
#[pyclass(module = "lancedb._lancedb", get_all)]
#[pyclass(get_all)]
/// A description of an index currently configured on a column
pub struct IndexConfig {
/// The type of the index
+1 -1
View File
@@ -11,7 +11,7 @@ use pyo3::{PyResult, pyclass, pymethods};
/// Sessions allow you to configure cache sizes for index and metadata caches,
/// which can significantly impact memory use and performance. They can
/// also be re-used across multiple connections to share the same cache state.
#[pyclass(module = "lancedb._lancedb", from_py_object)]
#[pyclass(from_py_object)]
#[derive(Clone)]
pub struct Session {
pub(crate) inner: Arc<LanceSession>,
+16 -32
View File
@@ -246,22 +246,12 @@ impl From<lancedb::table::MergeResult> for MergeResult {
}
}
/// Render for `__repr__`, so the default reads as Python's `None` rather than
/// Rust's `Some([..])`.
fn fmt_maintained(maintained: &Option<Vec<String>>) -> String {
match maintained {
Some(names) => format!("{:?}", names),
None => "None".to_string(),
}
}
/// Specification selecting Lance's MemWAL LSM-style write path for
/// `merge_insert`.
///
/// Constructed via the `bucket(...)`, `identity(...)`, or `unsharded()`
/// classmethods, then optionally chain `with_maintained_indexes(...)` and
/// `with_writer_config_defaults(...)`. A fresh spec maintains every index the
/// MemWAL supports, resolved on install.
/// `with_writer_config_defaults(...)`.
#[pyclass(from_py_object)]
#[derive(Clone, Debug)]
pub struct LsmWriteSpec {
@@ -301,11 +291,11 @@ impl LsmWriteSpec {
}
}
/// Set which indexes the MemWAL maintains. `None` (the default)
/// resolves every supported index on install; a list is verbatim,
/// and an empty list maintains nothing.
#[pyo3(signature = (indexes))]
pub fn with_maintained_indexes(&self, indexes: Option<Vec<String>>) -> Self {
/// Replace the list of indexes the MemWAL should keep up to date as
/// rows are appended. Each name must reference an index that
/// already exists on the table at the time `set_lsm_write_spec`
/// is called.
pub fn with_maintained_indexes(&self, indexes: Vec<String>) -> Self {
Self {
inner: self.inner.clone().with_maintained_indexes(indexes),
}
@@ -327,29 +317,23 @@ impl LsmWriteSpec {
maintained_indexes,
writer_config_defaults,
} => format!(
"LsmWriteSpec.bucket(column={:?}, num_buckets={}, maintained_indexes={}, writer_config_defaults={:?})",
column,
num_buckets,
fmt_maintained(maintained_indexes),
writer_config_defaults,
"LsmWriteSpec.bucket(column={:?}, num_buckets={}, maintained_indexes={:?}, writer_config_defaults={:?})",
column, num_buckets, maintained_indexes, writer_config_defaults,
),
lancedb::table::LsmWriteSpec::Identity {
column,
maintained_indexes,
writer_config_defaults,
} => format!(
"LsmWriteSpec.identity(column={:?}, maintained_indexes={}, writer_config_defaults={:?})",
column,
fmt_maintained(maintained_indexes),
writer_config_defaults,
"LsmWriteSpec.identity(column={:?}, maintained_indexes={:?}, writer_config_defaults={:?})",
column, maintained_indexes, writer_config_defaults,
),
lancedb::table::LsmWriteSpec::Unsharded {
maintained_indexes,
writer_config_defaults,
} => format!(
"LsmWriteSpec.unsharded(maintained_indexes={}, writer_config_defaults={:?})",
fmt_maintained(maintained_indexes),
writer_config_defaults,
"LsmWriteSpec.unsharded(maintained_indexes={:?}, writer_config_defaults={:?})",
maintained_indexes, writer_config_defaults,
),
}
}
@@ -384,10 +368,10 @@ impl LsmWriteSpec {
}
}
/// Indexes the MemWAL keeps up to date, or `None` for every supported one.
/// Names of indexes the MemWAL should keep up to date during writes.
#[getter]
pub fn maintained_indexes(&self) -> Option<Vec<String>> {
self.inner.maintained_indexes().map(<[String]>::to_vec)
pub fn maintained_indexes(&self) -> Vec<String> {
self.inner.maintained_indexes().to_vec()
}
/// Default `ShardWriter` configuration recorded by this spec.
@@ -579,7 +563,7 @@ impl PyBlobFile {
}
}
#[pyclass(module = "lancedb._lancedb", get_all, from_py_object)]
#[pyclass(get_all, from_py_object)]
#[derive(Clone, Debug)]
pub struct FtsToken {
pub text: String,
+95 -95
View File
@@ -799,7 +799,7 @@ name = "cuda-bindings"
version = "13.3.1"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "cuda-pathfinder", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "cuda-pathfinder" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/a9/21/8464d133752951c154feafb3b65c297e7d80f301183d220bec4c830f1441/cuda_bindings-13.3.1-cp310-cp310-manylinux_2_24_aarch64.manylinux_2_28_aarch64.whl", hash = "sha256:120fcc53d57903df529c3486962c56528cba5b7d6c57c99537320ed9922c8b86", size = 6073403, upload-time = "2026-05-29T23:11:36.22Z" },
@@ -834,37 +834,37 @@ wheels = [
[package.optional-dependencies]
cublas = [
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-cublas" },
]
cudart = [
{ name = "nvidia-cuda-runtime", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-cuda-runtime" },
]
cufft = [
{ name = "nvidia-cufft", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-cufft" },
]
cufile = [
{ name = "nvidia-cufile", marker = "sys_platform == 'linux'" },
{ name = "nvidia-cufile" },
]
cupti = [
{ name = "nvidia-cuda-cupti", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-cuda-cupti" },
]
curand = [
{ name = "nvidia-curand", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-curand" },
]
cusolver = [
{ name = "nvidia-cusolver", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-cusolver" },
]
cusparse = [
{ name = "nvidia-cusparse", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-cusparse" },
]
nvjitlink = [
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-nvjitlink" },
]
nvrtc = [
{ name = "nvidia-cuda-nvrtc", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-cuda-nvrtc" },
]
nvtx = [
{ name = "nvidia-nvtx", marker = "(python_full_version < '3.14' and sys_platform == 'win32') or sys_platform == 'linux'" },
{ name = "nvidia-nvtx" },
]
[[package]]
@@ -1023,7 +1023,7 @@ name = "exceptiongroup"
version = "1.3.1"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "typing-extensions", marker = "python_full_version < '3.11'" },
{ name = "typing-extensions" },
]
sdist = { url = "https://files.pythonhosted.org/packages/50/79/66800aadf48771f6b62f7eb014e352e5d06856655206165d775e675a02c9/exceptiongroup-1.3.1.tar.gz", hash = "sha256:8b412432c6055b0b7d14c310000ae93352ed6754f70fa8f7c34141f91c4e3219", size = 30371, upload-time = "2025-11-21T23:01:54.787Z" }
wheels = [
@@ -1440,16 +1440,16 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "cachetools", marker = "python_full_version < '3.11'" },
{ name = "certifi", marker = "python_full_version < '3.11'" },
{ name = "httpx", marker = "python_full_version < '3.11'" },
{ name = "ibm-cos-sdk", marker = "python_full_version < '3.11'" },
{ name = "lomond", marker = "python_full_version < '3.11'" },
{ name = "packaging", marker = "python_full_version < '3.11'" },
{ name = "pandas", version = "2.2.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "requests", marker = "python_full_version < '3.11'" },
{ name = "tabulate", marker = "python_full_version < '3.11'" },
{ name = "urllib3", marker = "python_full_version < '3.11'" },
{ name = "cachetools" },
{ name = "certifi" },
{ name = "httpx" },
{ name = "ibm-cos-sdk" },
{ name = "lomond" },
{ name = "packaging" },
{ name = "pandas", version = "2.2.3", source = { registry = "https://pypi.org/simple" } },
{ name = "requests" },
{ name = "tabulate" },
{ name = "urllib3" },
]
sdist = { url = "https://files.pythonhosted.org/packages/c7/56/2e3df38a1f13062095d7bde23c87a92f3898982993a15186b1bfecbd206f/ibm_watsonx_ai-1.3.42.tar.gz", hash = "sha256:ee5be59009004245d957ce97d1227355516df95a2640189749487614fef674ff", size = 688651, upload-time = "2025-10-01T13:35:41.527Z" }
wheels = [
@@ -1468,17 +1468,17 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "cachetools", marker = "python_full_version >= '3.11'" },
{ name = "certifi", marker = "python_full_version >= '3.11'" },
{ name = "httpx", marker = "python_full_version >= '3.11'" },
{ name = "ibm-cos-sdk", marker = "python_full_version >= '3.11'" },
{ name = "lomond", marker = "python_full_version >= '3.11'" },
{ name = "packaging", marker = "python_full_version >= '3.11'" },
{ name = "pandas", version = "2.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
{ name = "cachetools" },
{ name = "certifi" },
{ name = "httpx" },
{ name = "ibm-cos-sdk" },
{ name = "lomond" },
{ name = "packaging" },
{ name = "pandas", version = "2.3.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.14'" },
{ name = "pandas", version = "3.0.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.14'" },
{ name = "requests", marker = "python_full_version >= '3.11'" },
{ name = "tabulate", marker = "python_full_version >= '3.11'" },
{ name = "urllib3", marker = "python_full_version >= '3.11'" },
{ name = "requests" },
{ name = "tabulate" },
{ name = "urllib3" },
]
sdist = { url = "https://files.pythonhosted.org/packages/29/a3/c756b534696ab2f3f29882fdb7ca7198b7a5c94e10c0a3a327853d6d6b79/ibm_watsonx_ai-1.5.14.tar.gz", hash = "sha256:a756488bd57e87c0fc51be42dcba871143cfe0ac1e805c497c5047e1e4f13e9d", size = 735804, upload-time = "2026-06-22T12:32:43.85Z" }
wheels = [
@@ -1554,17 +1554,17 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "colorama", marker = "python_full_version < '3.11' and sys_platform == 'win32'" },
{ name = "decorator", marker = "python_full_version < '3.11'" },
{ name = "exceptiongroup", marker = "python_full_version < '3.11'" },
{ name = "jedi", marker = "python_full_version < '3.11'" },
{ name = "matplotlib-inline", marker = "python_full_version < '3.11'" },
{ name = "pexpect", marker = "python_full_version < '3.11' and sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit", marker = "python_full_version < '3.11'" },
{ name = "pygments", marker = "python_full_version < '3.11'" },
{ name = "stack-data", marker = "python_full_version < '3.11'" },
{ name = "traitlets", marker = "python_full_version < '3.11'" },
{ name = "typing-extensions", marker = "python_full_version < '3.11'" },
{ name = "colorama", marker = "sys_platform == 'win32'" },
{ name = "decorator" },
{ name = "exceptiongroup" },
{ name = "jedi" },
{ name = "matplotlib-inline" },
{ name = "pexpect", marker = "sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit" },
{ name = "pygments" },
{ name = "stack-data" },
{ name = "traitlets" },
{ name = "typing-extensions" },
]
sdist = { url = "https://files.pythonhosted.org/packages/40/18/f8598d287006885e7136451fdea0755af4ebcbfe342836f24deefaed1164/ipython-8.39.0.tar.gz", hash = "sha256:4110ae96012c379b8b6db898a07e186c40a2a1ef5d57a7fa83166047d9da7624", size = 5513971, upload-time = "2026-03-27T10:02:13.94Z" }
wheels = [
@@ -1583,18 +1583,18 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "colorama", marker = "python_full_version >= '3.11' and sys_platform == 'win32'" },
{ name = "decorator", marker = "python_full_version >= '3.11'" },
{ name = "ipython-pygments-lexers", marker = "python_full_version >= '3.11'" },
{ name = "jedi", marker = "python_full_version >= '3.11'" },
{ name = "matplotlib-inline", marker = "python_full_version >= '3.11'" },
{ name = "pexpect", marker = "python_full_version >= '3.11' and sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit", marker = "python_full_version >= '3.11'" },
{ name = "psutil", marker = "python_full_version >= '3.11' and sys_platform != 'cygwin' and sys_platform != 'emscripten'" },
{ name = "pygments", marker = "python_full_version >= '3.11'" },
{ name = "stack-data", marker = "python_full_version >= '3.11'" },
{ name = "traitlets", marker = "python_full_version >= '3.11'" },
{ name = "typing-extensions", marker = "python_full_version == '3.11.*'" },
{ name = "colorama", marker = "sys_platform == 'win32'" },
{ name = "decorator" },
{ name = "ipython-pygments-lexers" },
{ name = "jedi" },
{ name = "matplotlib-inline" },
{ name = "pexpect", marker = "sys_platform != 'emscripten' and sys_platform != 'win32'" },
{ name = "prompt-toolkit" },
{ name = "psutil", marker = "sys_platform != 'cygwin' and sys_platform != 'emscripten'" },
{ name = "pygments" },
{ name = "stack-data" },
{ name = "traitlets" },
{ name = "typing-extensions", marker = "python_full_version < '3.12'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/53/59/165d3b4d75cc34add3122c4417ecb229085140ac573103c223cd01dde96f/ipython-9.15.0.tar.gz", hash = "sha256:da2819ce2aa83135257df830660b1176d986c3d2876db24df01974fa955b2756", size = 4442580, upload-time = "2026-06-26T11:03:35.913Z" }
wheels = [
@@ -1606,7 +1606,7 @@ name = "ipython-pygments-lexers"
version = "1.1.1"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "pygments", marker = "python_full_version >= '3.11'" },
{ name = "pygments" },
]
sdist = { url = "https://files.pythonhosted.org/packages/ef/4c/5dd1d8af08107f88c7f741ead7a40854b8ac24ddf9ae850afbcf698aa552/ipython_pygments_lexers-1.1.1.tar.gz", hash = "sha256:09c0138009e56b6854f9535736f4171d855c8c08a563a0dcd8022f78355c7e81", size = 8393, upload-time = "2025-01-17T11:24:34.505Z" }
wheels = [
@@ -2005,7 +2005,7 @@ requires-dist = [
{ name = "pyarrow-stubs", marker = "extra == 'tests'", specifier = ">=16.0" },
{ name = "pydantic", specifier = ">=1.10" },
{ name = "pylance", marker = "extra == 'pylance'", specifier = ">=5.0.0b5" },
{ name = "pylance", marker = "extra == 'tests'", specifier = "==9.0.0rc1" },
{ name = "pylance", marker = "extra == 'tests'", specifier = "==10.0.0" },
{ name = "pyright", marker = "extra == 'dev'", specifier = ">=1.1.350" },
{ name = "pytest", marker = "extra == 'tests'", specifier = ">=7.0" },
{ name = "pytest-asyncio", marker = "extra == 'tests'", specifier = ">=0.21" },
@@ -2858,7 +2858,7 @@ name = "nvidia-cudnn-cu13"
version = "9.19.0.56"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-cublas" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/f1/84/26025437c1e6b61a707442184fa0c03d083b661adf3a3eecfd6d21677740/nvidia_cudnn_cu13-9.19.0.56-py3-none-manylinux_2_27_aarch64.whl", hash = "sha256:6ed29ffaee1176c612daf442e4dd6cfeb6a0caa43ddcbeb59da94953030b1be4", size = 433781201, upload-time = "2026-02-03T20:40:53.805Z" },
@@ -2870,7 +2870,7 @@ name = "nvidia-cufft"
version = "12.0.0.61"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-nvjitlink" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/8b/ae/f417a75c0259e85c1d2f83ca4e960289a5f814ed0cea74d18c353d3e989d/nvidia_cufft-12.0.0.61-py3-none-manylinux2014_aarch64.manylinux_2_17_aarch64.whl", hash = "sha256:2708c852ef8cd89d1d2068bdbece0aa188813a0c934db3779b9b1faa8442e5f5", size = 214053554, upload-time = "2025-09-04T08:31:38.196Z" },
@@ -2900,9 +2900,9 @@ name = "nvidia-cusolver"
version = "12.0.4.66"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-cublas", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-cusparse", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-cublas" },
{ name = "nvidia-cusparse" },
{ name = "nvidia-nvjitlink" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/c8/c3/b30c9e935fc01e3da443ec0116ed1b2a009bb867f5324d3f2d7e533e776b/nvidia_cusolver-12.0.4.66-py3-none-manylinux_2_27_aarch64.whl", hash = "sha256:02c2457eaa9e39de20f880f4bd8820e6a1cfb9f9a34f820eb12a155aa5bc92d2", size = 223467760, upload-time = "2025-09-04T08:33:04.222Z" },
@@ -2914,7 +2914,7 @@ name = "nvidia-cusparse"
version = "12.6.3.3"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "nvidia-nvjitlink", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "nvidia-nvjitlink" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/f8/94/5c26f33738ae35276672f12615a64bd008ed5be6d1ebcb23579285d960a9/nvidia_cusparse-12.6.3.3-py3-none-manylinux2014_aarch64.manylinux_2_17_aarch64.whl", hash = "sha256:80bcc4662f23f1054ee334a15c72b8940402975e0eab63178fc7e670aa59472c", size = 162155568, upload-time = "2025-09-04T08:33:42.864Z" },
@@ -3091,10 +3091,10 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "python-dateutil", marker = "python_full_version < '3.11'" },
{ name = "pytz", marker = "python_full_version < '3.11'" },
{ name = "tzdata", marker = "python_full_version < '3.11'" },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
{ name = "python-dateutil" },
{ name = "pytz" },
{ name = "tzdata" },
]
sdist = { url = "https://files.pythonhosted.org/packages/9c/d6/9f8431bacc2e19dca897724cd097b1bb224a6ad5433784a44b587c7c13af/pandas-2.2.3.tar.gz", hash = "sha256:4f18ba62b61d7e192368b84517265a99b4d7ee8912f8708660fb4a366cc82667", size = 4399213, upload-time = "2024-09-20T13:10:04.827Z" }
wheels = [
@@ -3143,11 +3143,11 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12' or python_full_version >= '3.14'" },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12' and python_full_version < '3.14'" },
{ name = "python-dateutil", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
{ name = "pytz", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
{ name = "tzdata", marker = "python_full_version >= '3.11' and python_full_version < '3.14'" },
{ name = "python-dateutil" },
{ name = "pytz" },
{ name = "tzdata" },
]
sdist = { url = "https://files.pythonhosted.org/packages/33/01/d40b85317f86cf08d853a4f495195c73815fdf205eef3993821720274518/pandas-2.3.3.tar.gz", hash = "sha256:e05e1af93b977f7eafa636d043f9f94c7ee3ac81af99c13508215942e64c993b", size = 4495223, upload-time = "2025-09-29T23:34:51.853Z" }
wheels = [
@@ -3210,9 +3210,9 @@ resolution-markers = [
"python_full_version >= '3.14' and sys_platform != 'emscripten' and sys_platform != 'win32'",
]
dependencies = [
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.14'" },
{ name = "python-dateutil", marker = "python_full_version >= '3.14'" },
{ name = "tzdata", marker = "(python_full_version >= '3.14' and sys_platform == 'emscripten') or (python_full_version >= '3.14' and sys_platform == 'win32')" },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" } },
{ name = "python-dateutil" },
{ name = "tzdata", marker = "sys_platform == 'emscripten' or sys_platform == 'win32'" },
]
sdist = { url = "https://files.pythonhosted.org/packages/f8/87/4341c6252d1c47b08768c3d25ac487362bf403f0313ddae4a2a26c9b1b4c/pandas-3.0.3.tar.gz", hash = "sha256:696a4a00a2a2a35d4e5deb3fc946641b96c944f02230e4f76137fe35d806c4fc", size = 4651414, upload-time = "2026-05-11T18:54:29.21Z" }
wheels = [
@@ -3320,7 +3320,7 @@ name = "pexpect"
version = "4.9.0"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "ptyprocess", marker = "(python_full_version < '3.14' and sys_platform == 'emscripten') or (python_full_version < '3.14' and sys_platform == 'win32') or (sys_platform != 'emscripten' and sys_platform != 'win32')" },
{ name = "ptyprocess" },
]
sdist = { url = "https://files.pythonhosted.org/packages/42/92/cc564bf6381ff43ce1f4d06852fc19a2f11d180f23dc32d9588bee2f149d/pexpect-4.9.0.tar.gz", hash = "sha256:ee7d41123f3c9911050ea2c2dac107568dc43b2d3b0c7557a33212c398ead30f", size = 166450, upload-time = "2023-11-25T09:07:26.339Z" }
wheels = [
@@ -3912,7 +3912,7 @@ crypto = [
[[package]]
name = "pylance"
version = "7.0.0"
version = "10.0.0"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "lance-namespace" },
@@ -3922,12 +3922,12 @@ dependencies = [
{ name = "pyarrow" },
]
wheels = [
{ url = "https://files.pythonhosted.org/packages/ac/ad/2f64921bf346e7075aef24a72595db44821724a3d89a9a92dd24e79632aa/pylance-7.0.0-cp39-abi3-macosx_11_0_arm64.whl", hash = "sha256:98422021975be76e72b1572f41b8c9abb3bee5bdc9bfa5e9ce731110a65ed4d1", size = 62134146, upload-time = "2026-05-27T21:59:37.459Z" },
{ url = "https://files.pythonhosted.org/packages/73/1c/c5a01bee0160b55d9a98895cbd33091d038f0a0995b121ab72e629008d02/pylance-7.0.0-cp39-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:4bec86ee5b6fbd8bfc493e653f0a1fba0303cfe5492b9b46fc25ab908edc7183", size = 65373684, upload-time = "2026-05-27T22:04:01.584Z" },
{ url = "https://files.pythonhosted.org/packages/eb/da/1fe8b8f7dbfe734d76af76acc994fc360a0d0c79a4874ef69f5a72a58fe3/pylance-7.0.0-cp39-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:881491432c53184e52f8d1db8d5f872f39a03f36fb104bec77b33d379519d8b5", size = 69458555, upload-time = "2026-05-27T22:16:50.567Z" },
{ url = "https://files.pythonhosted.org/packages/76/f0/dd505cf3fd0226ab9d94759acd713125af1d3bfacfd80bbd52e3b9f89509/pylance-7.0.0-cp39-abi3-manylinux_2_28_aarch64.whl", hash = "sha256:18453999e7fff4f76b16d6b7882c9df0628bd142ff95e2461bd7dd5ee3fe0af3", size = 65394430, upload-time = "2026-05-27T22:05:30.923Z" },
{ url = "https://files.pythonhosted.org/packages/17/ba/2357b81034f28eb00790e258ed140289a6a887a7468ca9df6349fd186b27/pylance-7.0.0-cp39-abi3-manylinux_2_28_x86_64.whl", hash = "sha256:04a58051d408c60fe76d41a220dcaf8fea8fb6d1aa0ca78a709b60bc3cc8d19a", size = 69473470, upload-time = "2026-05-27T22:17:18.935Z" },
{ url = "https://files.pythonhosted.org/packages/1f/ec/5c00b6303a67d787f9475141832cbdc513d674ac3dcaeef8a7b169905e65/pylance-7.0.0-cp39-abi3-win_amd64.whl", hash = "sha256:467d4864af047eaab4e1370e2f1e88e2c6f507c079874421116cb41d78bc3629", size = 74792863, upload-time = "2026-05-27T22:19:23.875Z" },
{ url = "https://files.pythonhosted.org/packages/82/b2/c81de196076c4c8d768f485324a0043e113ce0950a339977d83fee9783ef/pylance-10.0.0-cp310-abi3-macosx_11_0_arm64.whl", hash = "sha256:d4bba56ae829202b7e9cdc82c92c00a4b52f03e679dc3edfad1db09d1285be2e", size = 69279797, upload-time = "2026-08-07T18:25:24.812Z" },
{ url = "https://files.pythonhosted.org/packages/d2/e9/af671a6225740bd70c1e70bd83fa08091d628b960397cbecc477f16aaaf1/pylance-10.0.0-cp310-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl", hash = "sha256:489b944827c0271e16a62b4006f8c75acc3bd7fbc381f35388a2afaaae38d438", size = 72775384, upload-time = "2026-08-07T18:31:29.686Z" },
{ url = "https://files.pythonhosted.org/packages/5d/22/e07194195bb3bbdf062b0c31690fc92fccb2686d5e768efa1dc379a93350/pylance-10.0.0-cp310-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl", hash = "sha256:018efe7d437d326b9049c1223458bd955e85e48509b4b0bbf8b2d3d94075dbff", size = 76632194, upload-time = "2026-08-07T18:44:29.001Z" },
{ url = "https://files.pythonhosted.org/packages/f2/71/ed9956cf657e86a5fee5d5ddcfa7bfed0d7c2968bee39ed85d12629d8c93/pylance-10.0.0-cp310-abi3-manylinux_2_28_aarch64.whl", hash = "sha256:254bad5d765c14db6c4eddd5323133dee3c76ff6d40c4777afe0b0e91d4c04d2", size = 72804222, upload-time = "2026-08-07T18:31:24.264Z" },
{ url = "https://files.pythonhosted.org/packages/94/96/de449c246b2892df9d5a0af775b79a061e9897e6942bcee5c4ab6d6071f8/pylance-10.0.0-cp310-abi3-manylinux_2_28_x86_64.whl", hash = "sha256:9f0c089e30389d9b7765a4c16f3d9325b05ed4fbf6cc2bbc9418292983472e54", size = 76602015, upload-time = "2026-08-07T18:46:27.724Z" },
{ url = "https://files.pythonhosted.org/packages/53/81/4a5a9072b6d68c4dbb8c7b5530a381c7a6bea4ba92f9eca58ac885142722/pylance-10.0.0-cp310-abi3-win_amd64.whl", hash = "sha256:9fecf46e6835dc64b2d71767d5eab4f73787543a852432e9ba8ecf8527b3cdab", size = 82799247, upload-time = "2026-08-07T18:47:46.855Z" },
]
[[package]]
@@ -4683,10 +4683,10 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "joblib", marker = "python_full_version < '3.11'" },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "threadpoolctl", marker = "python_full_version < '3.11'" },
{ name = "joblib" },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
{ name = "scipy", version = "1.15.3", source = { registry = "https://pypi.org/simple" } },
{ name = "threadpoolctl" },
]
sdist = { url = "https://files.pythonhosted.org/packages/98/c2/a7855e41c9d285dfe86dc50b250978105dce513d6e459ea66a6aeb0e1e0c/scikit_learn-1.7.2.tar.gz", hash = "sha256:20e9e49ecd130598f1ca38a1d85090e1a600147b9c02fa6f15d69cb53d968fda", size = 7193136, upload-time = "2025-09-09T08:21:29.075Z" }
wheels = [
@@ -4734,13 +4734,13 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "joblib", marker = "python_full_version >= '3.11'" },
{ name = "narwhals", marker = "python_full_version >= '3.11'" },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
{ name = "joblib" },
{ name = "narwhals" },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
{ name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
{ name = "scipy", version = "1.17.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.12'" },
{ name = "scipy", version = "1.18.0", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
{ name = "threadpoolctl", marker = "python_full_version >= '3.11'" },
{ name = "threadpoolctl" },
]
sdist = { url = "https://files.pythonhosted.org/packages/fa/6f/37092bdb25f712817231799fc5674d8e704066a8a70c1d2d40517e18b4ab/scikit_learn-1.9.0.tar.gz", hash = "sha256:8833266989d3a5110178a9fae30783675460724d0e1efb13b14901d2c660c557", size = 7750767, upload-time = "2026-06-02T11:54:32.706Z" }
wheels = [
@@ -4784,7 +4784,7 @@ resolution-markers = [
"python_full_version < '3.11'",
]
dependencies = [
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version < '3.11'" },
{ name = "numpy", version = "2.2.6", source = { registry = "https://pypi.org/simple" } },
]
sdist = { url = "https://files.pythonhosted.org/packages/0f/37/6964b830433e654ec7485e45a00fc9a27cf868d622838f6b6d9c5ec0d532/scipy-1.15.3.tar.gz", hash = "sha256:eae3cf522bc7df64b42cad3925c876e1b0b6c35c1337c93e12c0f366f55b0eaf", size = 59419214, upload-time = "2025-05-08T16:13:05.955Z" }
wheels = [
@@ -4843,7 +4843,7 @@ resolution-markers = [
"python_full_version == '3.11.*'",
]
dependencies = [
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version == '3.11.*'" },
{ name = "numpy", version = "2.4.6", source = { registry = "https://pypi.org/simple" } },
]
sdist = { url = "https://files.pythonhosted.org/packages/7a/97/5a3609c4f8d58b039179648e62dd220f89864f56f7357f5d4f45c29eb2cc/scipy-1.17.1.tar.gz", hash = "sha256:95d8e012d8cb8816c226aef832200b1d45109ed4464303e997c5b13122b297c0", size = 30573822, upload-time = "2026-02-23T00:26:24.851Z" }
wheels = [
@@ -4920,7 +4920,7 @@ resolution-markers = [
"python_full_version >= '3.12' and python_full_version < '3.14'",
]
dependencies = [
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" }, marker = "python_full_version >= '3.12'" },
{ name = "numpy", version = "2.5.1", source = { registry = "https://pypi.org/simple" } },
]
sdist = { url = "https://files.pythonhosted.org/packages/a7/25/c2700dfaf6442b4effaa91af24ebce5dc9d31bb4a69706313aae70d72cd0/scipy-1.18.0.tar.gz", hash = "sha256:67b2ad2ad54c72ca6d04975a9b2df8c3638c34ddd5b28738e94fc2b57929d378", size = 30774447, upload-time = "2026-06-19T15:01:43.456Z" }
wheels = [
+1 -9
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb"
version = "0.37.1-beta.1"
version = "0.37.1"
edition.workspace = true
description = "LanceDB: A serverless, low-latency vector database for AI applications"
license.workspace = true
@@ -115,11 +115,6 @@ serial_test = "3"
[target.'cfg(unix)'.dev-dependencies]
pprof = { version = "0.14", features = ["flamegraph"] }
[target.'cfg(windows)'.dependencies]
windows-sys = { version = "0.61", features = [
"Win32_Foundation",
"Win32_Storage_FileSystem",
] }
[features]
default = []
@@ -193,9 +188,6 @@ required-features = ["bedrock"]
[[example]]
name = "bench_streaming_dataloader"
[[example]]
name = "bench_open_missing_table"
[[example]]
name = "simple"
@@ -1,150 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
// Release benchmark for opening a missing table as sibling-table cardinality grows.
//
// The fixture uses real `.lance` directories and marker files. Fixture creation is
// outside the timed section. Defaults intentionally cover 1k, 10k, and 100k siblings
// with 10 warmups and 100 distinct missing-table opens per scale:
//
// ```text
// cargo run --release -p lancedb --example bench_open_missing_table
// ```
//
// `BENCH_SIBLINGS`, `BENCH_WARMUPS`, and `BENCH_TRIALS` override those defaults.
// Reduced settings are useful only as a smoke test. Performance comparisons require
// the same machine, filesystem, fixture sizes, settings, lockfile, and alternating
// baseline/candidate execution order.
use std::time::{Duration, Instant};
use anyhow::{Context, Result, bail};
use lancedb::connection::Connection;
use lancedb::{Error, connect};
use object_store::ObjectStoreExt as _;
use object_store::path::Path;
const MAX_SIBLINGS: usize = 1_000_000;
const MAX_WARMUPS: usize = 10_000;
const MAX_TRIALS: usize = 100_000;
fn env_usize(key: &str, default: usize, max: usize) -> Result<usize> {
let value = match std::env::var(key) {
Ok(value) => value
.parse()
.with_context(|| format!("invalid {key} value: {value}"))?,
Err(std::env::VarError::NotPresent) => default,
Err(error) => return Err(error).with_context(|| format!("reading {key}")),
};
if value == 0 || value > max {
bail!("{key} must be between 1 and {max}");
}
Ok(value)
}
fn sibling_counts() -> Result<Vec<usize>> {
let raw = std::env::var("BENCH_SIBLINGS").unwrap_or_else(|_| "1000,10000,100000".into());
let mut counts = raw
.split(',')
.map(|value| {
value
.trim()
.parse::<usize>()
.with_context(|| format!("invalid BENCH_SIBLINGS value: {value}"))
})
.collect::<Result<Vec<_>>>()?;
counts.sort_unstable();
counts.dedup();
if counts.is_empty() || counts[0] == 0 || counts[counts.len() - 1] > MAX_SIBLINGS {
bail!("BENCH_SIBLINGS values must be between 1 and {MAX_SIBLINGS}");
}
Ok(counts)
}
async fn add_siblings(
store: &object_store::local::LocalFileSystem,
start: usize,
end: usize,
) -> Result<()> {
for index in start..end {
let marker = Path::from(format!("sibling_{index:06}.lance/_marker"));
store
.put(&marker, bytes::Bytes::new().into())
.await
.with_context(|| format!("creating benchmark marker {marker}"))?;
}
Ok(())
}
async fn time_missing_open(db: &Connection, name: &str) -> Result<Duration> {
let started = Instant::now();
let result = db.open_table(name).execute().await;
let elapsed = started.elapsed();
match result {
Err(Error::TableNotFound { .. }) => Ok(elapsed),
Err(error) => bail!("expected TableNotFound for {name}, got {error:?}"),
Ok(_) => bail!("benchmark missing-table name unexpectedly exists: {name}"),
}
}
fn percentile(sorted: &[Duration], percentile: usize) -> Duration {
let rank = (sorted.len() * percentile).div_ceil(100).saturating_sub(1);
sorted[rank]
}
#[tokio::main]
async fn main() -> Result<()> {
let counts = sibling_counts()?;
let warmups = env_usize("BENCH_WARMUPS", 10, MAX_WARMUPS)?;
let trials = env_usize("BENCH_TRIALS", 100, MAX_TRIALS)?;
let fixture = tempfile::tempdir().context("creating benchmark fixture")?;
let database_path = fixture.path();
let fixture_store = object_store::local::LocalFileSystem::new_with_prefix(database_path)
.context("creating benchmark object store")?;
let db = connect(database_path.to_str().context("non-UTF-8 fixture path")?)
.execute()
.await?;
println!(
"config: siblings={counts:?} warmups={warmups} trials={trials} profile={} os={} arch={}",
if cfg!(debug_assertions) {
"debug"
} else {
"release"
},
std::env::consts::OS,
std::env::consts::ARCH,
);
println!("lower is better; fixture setup and teardown are excluded");
println!("| siblings | samples | p50 | p95 | max |");
println!("| ---: | ---: | ---: | ---: | ---: |");
let mut created = 0;
for sibling_count in counts {
add_siblings(&fixture_store, created, sibling_count).await?;
created = sibling_count;
for index in 0..warmups {
let name = format!("__missing_warmup_{sibling_count}_{index}");
let _ = time_missing_open(&db, &name).await?;
}
let mut samples = Vec::with_capacity(trials);
for index in 0..trials {
let name = format!("__missing_trial_{sibling_count}_{index}");
samples.push(time_missing_open(&db, &name).await?);
}
samples.sort_unstable();
println!(
"| {sibling_count} | {} | {:?} | {:?} | {:?} |",
samples.len(),
percentile(&samples, 50),
percentile(&samples, 95),
samples[samples.len() - 1],
);
}
Ok(())
}
+4 -7
View File
@@ -17,7 +17,7 @@ use arrow_array::builder::LargeBinaryBuilder;
use arrow_schema::{DataType, Field, Schema};
use lance::dataset::{BlobRangeRequest as LanceBlobRangeRequest, Dataset, WriteParams};
use lance_arrow::FieldExt;
use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
use lance_file::version::LanceFileVersion;
use lance_io::object_store::ObjectStore;
use object_store::path::Path;
@@ -333,10 +333,7 @@ pub(crate) fn ensure_blob_storage_version(schema: &Schema, params: &mut WritePar
.data_storage_version
.unwrap_or(LanceFileVersion::Stable)
.resolve();
if matches!(
resolved,
ConcreteFileVersion::V1 | ConcreteFileVersion::V2_0 | ConcreteFileVersion::V2_1
) {
if resolved < LanceFileVersion::V2_2 {
params.data_storage_version = Some(LanceFileVersion::V2_2);
}
}
@@ -502,7 +499,7 @@ mod tests {
ensure_blob_storage_version(&blob_schema(), &mut params);
assert_eq!(
params.data_storage_version.unwrap().resolve(),
ConcreteFileVersion::V2_2
LanceFileVersion::V2_2
);
}
@@ -515,7 +512,7 @@ mod tests {
ensure_blob_storage_version(&blob_schema(), &mut params);
assert_eq!(
params.data_storage_version.unwrap().resolve(),
ConcreteFileVersion::V2_2
LanceFileVersion::V2_2
);
}
+4 -8
View File
@@ -409,11 +409,6 @@ impl Connection {
///
/// The names will be returned in lexicographical order (ascending)
///
/// Listing databases discover physical `*.lance` entries without opening every
/// dataset. The result is a point-in-time discovery snapshot: an entry may still be
/// under creation, may contain only uncommitted storage, or may be concurrently
/// dropped before it is opened.
///
/// The parameters `page_token` and `limit` can be used to paginate the results
pub fn table_names(&self) -> TableNamesBuilder {
TableNamesBuilder::new(self.internal.clone())
@@ -461,9 +456,10 @@ impl Connection {
///
/// # Returns
/// Created [`TableRef`], or [`Error::TableNotFound`] if the table does not exist.
/// On listing databases, a committed Lance manifest is authoritative for table
/// existence. Uncommitted files or a physical `<name>.lance` directory alone do not
/// make a table openable.
/// If the table's storage is present but holds no readable dataset (for example a
/// `<name>.lance` directory left behind by an interrupted drop and re-create, which
/// [`Self::table_names`] still lists) this returns [`Error::TableCorrupted`]
/// instead.
pub fn open_table(&self, name: impl Into<String>) -> OpenTableBuilder {
OpenTableBuilder::new(
self.internal.clone(),
+3 -2
View File
@@ -438,9 +438,10 @@ mod tests {
.await
.unwrap()
.data_storage_format
.lance_file_format();
.lance_file_version()
.unwrap();
// Compare resolved versions since Stable/Next are aliases that resolve at storage time
assert_eq!(storage_format, data_storage_version.resolve());
assert_eq!(storage_format.resolve(), data_storage_version.resolve());
}
#[tokio::test]
+5 -234
View File
@@ -25,7 +25,7 @@ use crate::database::namespace::LanceNamespaceDatabase;
use crate::error::{CreateDirSnafu, Error, Result};
use crate::io::object_store::MirroringObjectStoreWrapper;
use crate::table::NativeTable;
use crate::utils::{PatchStoreParam, validate_table_name};
use crate::utils::validate_table_name;
use lance_namespace::models::{
CreateNamespaceRequest, CreateNamespaceResponse, DescribeNamespaceRequest,
@@ -355,14 +355,6 @@ impl ListingDatabase {
url.to_string()
}
#[cfg(any(windows, test))]
fn uses_local_file_provider(object_store: &ObjectStore) -> bool {
matches!(
object_store.scheme(),
"file" | "file-object-store" | "file+uring"
)
}
async fn prepare_namespace_root(
uri: &str,
storage_options: &HashMap<String, String>,
@@ -589,17 +581,6 @@ impl ListingDatabase {
}
None => None,
};
#[cfg(windows)]
let write_store_wrapper = if Self::uses_local_file_provider(&object_store) {
// Local manifest commits need create-only rename semantics,
// including on filesystems that do not support hard links.
Some(
Arc::new(crate::io::object_store::windows::WindowsLocalFileSystemWrapper)
as Arc<dyn WrappingObjectStore>,
)
} else {
write_store_wrapper
};
let namespace_database = Self::connect_namespace_database(
&storage_base_uri,
@@ -664,22 +645,12 @@ impl ListingDatabase {
)
.await?;
#[cfg(windows)]
let write_store_wrapper = Self::uses_local_file_provider(&object_store).then(|| {
// Local manifest commits need create-only rename semantics,
// including on filesystems that do not support hard links.
Arc::new(crate::io::object_store::windows::WindowsLocalFileSystemWrapper)
as Arc<dyn WrappingObjectStore>
});
#[cfg(not(windows))]
let write_store_wrapper = None;
Ok(Self {
uri: path.to_string(),
query_string: None,
base_path,
object_store,
store_wrapper: write_store_wrapper,
store_wrapper: None,
read_consistency_interval,
storage_options: HashMap::new(),
storage_options_provider: None,
@@ -1141,12 +1112,6 @@ impl Database for ListingDatabase {
},
..Default::default()
};
let storage_params = match self.store_wrapper.clone() {
Some(wrapper) => Some(storage_params)
.patch_with_store_wrapper(wrapper)?
.expect("patching store params always returns parameters"),
None => storage_params,
};
let read_params = ReadParams {
store_options: Some(storage_params.clone()),
session: Some(self.session.clone()),
@@ -1326,37 +1291,16 @@ impl Database for ListingDatabase {
mod tests {
use super::*;
use crate::Table;
use crate::arrow::{SendableRecordBatchStream, SimpleRecordBatchStream};
use crate::connection::ConnectRequest;
use crate::data::scannable::Scannable;
use crate::database::{CreateTableMode, CreateTableRequest};
use crate::io::object_store::io_tracking::IoStatsHolder;
use crate::query::QueryRequest;
use crate::table::{AnyQuery, WriteOptions};
use arrow_array::{Int32Array, RecordBatch, StringArray};
use arrow_schema::{DataType, Field, Schema, SchemaRef};
use futures::{TryStreamExt, stream::once};
use arrow_schema::{DataType, Field, Schema};
use futures::TryStreamExt;
use std::path::PathBuf;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::time::Duration;
use tempfile::tempdir;
use tokio::sync::Barrier;
use tokio::time::timeout;
#[derive(Debug)]
struct PassthroughStoreWrapper(Arc<AtomicUsize>);
impl WrappingObjectStore for PassthroughStoreWrapper {
fn wrap(
&self,
_store_prefix: &str,
target: Arc<dyn object_store::ObjectStore>,
) -> Arc<dyn object_store::ObjectStore> {
self.0.fetch_add(1, Ordering::Relaxed);
target
}
}
async fn setup_database() -> (tempfile::TempDir, ListingDatabase) {
let tempdir = tempdir().unwrap();
@@ -1380,114 +1324,6 @@ mod tests {
(tempdir, db)
}
struct BarrierScannable {
batch: RecordBatch,
barrier: Arc<Barrier>,
}
impl Scannable for BarrierScannable {
fn schema(&self) -> SchemaRef {
self.batch.schema()
}
fn scan_as_stream(&mut self) -> SendableRecordBatchStream {
let batch = self.batch.clone();
let schema = batch.schema();
let barrier = self.barrier.clone();
Box::pin(SimpleRecordBatchStream {
schema,
stream: once(async move {
barrier.wait().await;
Ok(batch)
}),
})
}
}
fn create_request(name: &str, data: Box<dyn Scannable>) -> CreateTableRequest {
CreateTableRequest {
name: name.to_string(),
namespace_path: vec![],
data,
mode: CreateTableMode::Create,
write_options: Default::default(),
location: None,
namespace_client: None,
}
}
#[tokio::test]
async fn test_create_ignores_uncommitted_storage_without_manifest() {
let (tmp_dir, db) = setup_database().await;
let data_dir = tmp_dir.path().join("test.lance/data");
std::fs::create_dir_all(&data_dir).unwrap();
std::fs::write(data_dir.join("orphan.lance"), b"uncommitted").unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
let batch =
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1]))]).unwrap();
let table = db
.create_table(create_request("test", Box::new(batch)))
.await
.unwrap();
assert_eq!(table.count_rows(None).await.unwrap(), 1);
}
#[tokio::test]
async fn test_concurrent_create_is_arbitrated_by_manifest_commit() {
let uri = format!("memory:///concurrent-create-{}", uuid::Uuid::new_v4());
let db = crate::connect(&uri).execute().await.unwrap();
let store: Arc<dyn object_store::ObjectStore> =
Arc::new(object_store::memory::InMemory::new());
let table_url = url::Url::parse("memory:///database/test.lance").unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
let batch =
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1]))]).unwrap();
let barrier = Arc::new(Barrier::new(2));
#[allow(deprecated)]
let request = |batch, barrier| {
let mut request = create_request("test", Box::new(BarrierScannable { batch, barrier }));
request.write_options = WriteOptions {
lance_write_params: Some(lance::dataset::WriteParams {
store_params: Some(ObjectStoreParams {
object_store: Some((store.clone(), table_url.clone())),
..Default::default()
}),
commit_handler: Some(Arc::new(
lance_table::io::commit::ConditionalPutCommitHandler,
)),
..Default::default()
}),
};
request
};
let left = db
.database()
.create_table(request(batch.clone(), barrier.clone()));
let right = db.database().create_table(request(batch, barrier));
let (left, right) = timeout(Duration::from_secs(30), async { tokio::join!(left, right) })
.await
.expect("concurrent creates deadlocked");
let results = [left, right];
assert_eq!(
results.iter().filter(|result| result.is_ok()).count(),
1,
"expected one successful create, got {results:?}"
);
assert_eq!(
results
.iter()
.filter(|result| matches!(result, Err(Error::TableAlreadyExists { .. })))
.count(),
1,
"expected one manifest conflict, got {results:?}"
);
}
#[tokio::test]
async fn test_listing_database_root_ops_do_not_create_manifest() {
let tempdir = tempdir().unwrap();
@@ -1542,25 +1378,6 @@ mod tests {
assert!(!tempdir.path().join("__manifest").exists());
}
#[tokio::test]
async fn test_file_object_store_uses_local_file_provider() {
let tempdir = tempdir().unwrap();
let path = tempdir.path().to_string_lossy().replace('\\', "/");
let uri = if path.starts_with('/') {
format!("file-object-store://{path}")
} else {
format!("file-object-store:///{path}")
};
let registry = Arc::new(lance_io::object_store::ObjectStoreRegistry::default());
let (store, _) =
ObjectStore::from_uri_and_params(registry, &uri, &ObjectStoreParams::default())
.await
.unwrap();
assert_eq!(store.scheme(), "file-object-store");
assert!(ListingDatabase::uses_local_file_provider(&store));
}
/// Regression test for https://github.com/lancedb/lancedb/issues/1600.
///
/// Opening a table used to create a separate object-store client instead of
@@ -1584,13 +1401,9 @@ mod tests {
read_consistency_interval: None,
session: Some(session),
};
let mut db = ListingDatabase::connect_with_options(&request)
let db = ListingDatabase::connect_with_options(&request)
.await
.unwrap();
// A connection-level write wrapper must not prevent table opens from
// reusing the connection's registered object store.
let wrapper_calls = Arc::new(AtomicUsize::new(0));
db.store_wrapper = Some(Arc::new(PassthroughStoreWrapper(wrapper_calls.clone())));
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
db.create_table(CreateTableRequest {
@@ -1606,7 +1419,6 @@ mod tests {
.unwrap();
let before_open = registry.stats();
let wrapper_calls_before_open = wrapper_calls.load(Ordering::Relaxed);
for _ in 0..3 {
let table = db
.open_table(OpenTableRequest {
@@ -1626,7 +1438,6 @@ mod tests {
let after_open = registry.stats();
assert_eq!(after_open.misses, before_open.misses);
assert!(after_open.hits >= before_open.hits + 3);
assert!(wrapper_calls.load(Ordering::Relaxed) >= wrapper_calls_before_open + 3);
}
/// Regression test for https://github.com/lancedb/lancedb/issues/3197.
@@ -1770,46 +1581,6 @@ mod tests {
);
}
#[tokio::test]
async fn test_clone_table_uses_connection_store_wrapper() {
let (_tempdir, mut db) = setup_database().await;
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
db.create_table(CreateTableRequest {
name: "source_table".to_string(),
namespace_path: vec![],
data: Box::new(RecordBatch::new_empty(schema)) as Box<dyn Scannable>,
mode: CreateTableMode::Create,
write_options: Default::default(),
location: None,
namespace_client: None,
})
.await
.unwrap();
let source_uri = db.table_uri("source_table").unwrap();
let tracker = IoStatsHolder::default();
db.store_wrapper = Some(Arc::new(tracker.clone()));
let _ = tracker.incremental_stats();
db.clone_table(CloneTableRequest {
target_table_name: "cloned_table".to_string(),
target_namespace_path: vec![],
source_uri,
source_version: None,
source_tag: None,
is_shallow: true,
namespace_client: None,
})
.await
.unwrap();
let stats = tracker.incremental_stats();
assert!(
stats.write_iops > 0,
"clone bypassed the wrapper: {stats:?}"
);
}
#[tokio::test]
async fn test_clone_table_with_data() {
let (_tempdir, db) = setup_database().await;
@@ -12,8 +12,7 @@ use lance_encoding::decoder::{DecoderPlugins, FilterExpression};
use lance_file::{
reader::{FileReader, FileReaderOptions},
version::ConcreteFileVersion,
versions,
writer::FileWriterOptions,
writer::{FileWriter, FileWriterOptions},
};
use lance_io::{
ReadBatchParams,
@@ -154,11 +153,13 @@ impl Shuffler {
source: None,
})?;
let object_writer = object_store.create(&path).await?;
let writer = versions::create_writer(
ConcreteFileVersion::V2_1,
let writer = FileWriter::try_new(
object_writer,
schema.clone(),
FileWriterOptions::default(),
FileWriterOptions {
format_version: Some(ConcreteFileVersion::V2_1.into()),
..Default::default()
},
)?;
file_writers.push(writer);
}
-3
View File
@@ -18,9 +18,6 @@ use async_trait::async_trait;
#[cfg(test)]
pub mod io_tracking;
#[cfg(windows)]
pub(crate) mod windows;
#[derive(Debug)]
struct MirroringObjectStore {
primary: Arc<dyn ObjectStore>,
-223
View File
@@ -1,223 +0,0 @@
// SPDX-License-Identifier: Apache-2.0
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
//! Windows local filesystem compatibility for atomic manifest commits.
use std::ffi::OsStr;
use std::fmt::{Display, Formatter};
use std::os::windows::ffi::OsStrExt;
use std::path::{Path as StdPath, PathBuf};
use std::sync::Arc;
use bytes::Bytes;
use futures::stream::BoxStream;
use lance::io::WrappingObjectStore;
use object_store::{
CopyOptions, Error, GetOptions, GetResult, ListResult, MultipartUpload, ObjectMeta,
ObjectStore, PutMultipartOptions, PutOptions, PutPayload, PutResult, RenameOptions,
RenameTargetMode, Result, UploadPart, path::Path,
};
use windows_sys::Win32::Foundation::{ERROR_ALREADY_EXISTS, ERROR_FILE_EXISTS};
use windows_sys::Win32::Storage::FileSystem::MoveFileExW;
const STORE_NAME: &str = "WindowsLocalFileSystem";
/// Uses the Windows move primitive for create-only renames on local stores.
///
/// `object_store` implements create-only local renames with a hard link followed
/// by a delete. Some Windows filesystems do not support hard links, but
/// `MoveFileExW` without `MOVEFILE_REPLACE_EXISTING` provides the same atomic
/// create-only rename semantics without requiring them.
#[derive(Debug, Default)]
pub struct WindowsLocalFileSystemWrapper;
impl WrappingObjectStore for WindowsLocalFileSystemWrapper {
fn wrap(&self, _store_prefix: &str, target: Arc<dyn ObjectStore>) -> Arc<dyn ObjectStore> {
Arc::new(WindowsLocalFileSystem { target })
}
}
#[derive(Debug)]
struct WindowsLocalFileSystem {
target: Arc<dyn ObjectStore>,
}
impl Display for WindowsLocalFileSystem {
fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
write!(f, "{STORE_NAME}({})", self.target)
}
}
#[async_trait::async_trait]
#[deny(clippy::missing_trait_methods)]
impl ObjectStore for WindowsLocalFileSystem {
async fn put_opts(
&self,
location: &Path,
bytes: PutPayload,
opts: PutOptions,
) -> Result<PutResult> {
self.target.put_opts(location, bytes, opts).await
}
async fn put_multipart_opts(
&self,
location: &Path,
opts: PutMultipartOptions,
) -> Result<Box<dyn MultipartUpload>> {
self.target.put_multipart_opts(location, opts).await
}
async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
self.target.get_opts(location, options).await
}
async fn get_ranges(
&self,
location: &Path,
ranges: &[std::ops::Range<u64>],
) -> Result<Vec<Bytes>> {
self.target.get_ranges(location, ranges).await
}
fn delete_stream(
&self,
locations: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
self.target.delete_stream(locations)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
self.target.list(prefix)
}
fn list_with_offset(
&self,
prefix: Option<&Path>,
offset: &Path,
) -> BoxStream<'static, Result<ObjectMeta>> {
self.target.list_with_offset(prefix, offset)
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
self.target.list_with_delimiter(prefix).await
}
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
self.target.copy_opts(from, to, options).await
}
async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
if options.target_mode != RenameTargetMode::Create {
return self.target.rename_opts(from, to, options).await;
}
let from = PathBuf::from(from.as_ref());
let to = PathBuf::from(to.as_ref());
tokio::task::spawn_blocking(move || move_file_if_not_exists(&from, &to))
.await
.map_err(|source| Error::Generic {
store: STORE_NAME,
source: Box::new(source),
})?
}
}
fn move_file_if_not_exists(from: &StdPath, to: &StdPath) -> Result<()> {
let from_wide = null_terminated_wide(from.as_os_str());
let to_wide = null_terminated_wide(to.as_os_str());
// SAFETY: both pointers reference null-terminated UTF-16 buffers that remain
// alive for the duration of this call. A zero flag value deliberately omits
// MOVEFILE_REPLACE_EXISTING, giving this operation create-only semantics.
if unsafe { MoveFileExW(from_wide.as_ptr(), to_wide.as_ptr(), 0) } != 0 {
return Ok(());
}
let source = std::io::Error::last_os_error();
let path = to.to_string_lossy().into_owned();
match source.raw_os_error().map(|code| code as u32) {
Some(ERROR_ALREADY_EXISTS | ERROR_FILE_EXISTS) => Err(Error::AlreadyExists {
path,
source: Box::new(source),
}),
_ if source.kind() == std::io::ErrorKind::NotFound => Err(Error::NotFound {
path,
source: Box::new(source),
}),
_ => Err(Error::Generic {
store: STORE_NAME,
source: Box::new(source),
}),
}
}
fn null_terminated_wide(value: &OsStr) -> Vec<u16> {
value.encode_wide().chain(Some(0)).collect()
}
#[cfg(test)]
mod tests {
use object_store::memory::InMemory;
use super::*;
#[tokio::test]
async fn create_only_rename_does_not_use_hard_links() {
let tempdir = tempfile::tempdir().unwrap();
let source_path = tempdir.path().join("staged.manifest");
let destination_path = tempdir.path().join("1.manifest");
std::fs::write(&source_path, b"manifest").unwrap();
let source = Path::from_absolute_path(&source_path).unwrap();
let destination = Path::from_absolute_path(&destination_path).unwrap();
let store = WindowsLocalFileSystem {
// The source does not exist in this inner store. Delegating the
// rename would fail, proving the wrapper uses the native move path.
target: Arc::new(InMemory::new()),
};
store
.rename_opts(
&source,
&destination,
RenameOptions::new().with_target_mode(RenameTargetMode::Create),
)
.await
.unwrap();
assert!(!source_path.exists());
assert_eq!(std::fs::read(destination_path).unwrap(), b"manifest");
}
#[tokio::test]
async fn create_only_rename_preserves_existing_destination() {
let tempdir = tempfile::tempdir().unwrap();
let source_path = tempdir.path().join("staged.manifest");
let destination_path = tempdir.path().join("1.manifest");
std::fs::write(&source_path, b"new manifest").unwrap();
std::fs::write(&destination_path, b"existing manifest").unwrap();
let source = Path::from_absolute_path(&source_path).unwrap();
let destination = Path::from_absolute_path(&destination_path).unwrap();
let store = WindowsLocalFileSystem {
target: Arc::new(InMemory::new()),
};
let error = store
.rename_opts(
&source,
&destination,
RenameOptions::new().with_target_mode(RenameTargetMode::Create),
)
.await
.unwrap_err();
assert!(matches!(error, Error::AlreadyExists { .. }));
assert_eq!(std::fs::read(source_path).unwrap(), b"new manifest");
assert_eq!(
std::fs::read(destination_path).unwrap(),
b"existing manifest"
);
}
}
+9 -28
View File
@@ -2520,9 +2520,9 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
self.check_mutable().await?;
// Map the spec onto the server's request DTO. `sharding` is internally
// tagged on `mode` to mirror sophon's `Sharding` enum. A null
// `maintained_indexes` asks the server to resolve every maintainable
// index at HEAD; a list is verbatim, an empty one meaning none.
// tagged on `mode` to mirror sophon's `Sharding` enum; `maintained_indexes`
// and `writer_config_defaults` are sent verbatim (an empty list means "no
// maintained indexes", not "default to all").
let sharding = match &spec {
LsmWriteSpec::Bucket {
column,
@@ -5955,18 +5955,17 @@ mod tests {
.await
.unwrap();
// Positions are relative to the first retained token, so dropping the
// leading "hello" stop word does not shift the remaining tokens.
// Lance 10 retains original token positions after stop-word removal.
assert_eq!(
tokens,
vec![
FtsToken {
text: "こんにちは".to_string(),
position: 0,
position: 1,
},
FtsToken {
text: "世界".to_string(),
position: 1,
position: 2,
},
]
);
@@ -6599,7 +6598,7 @@ mod tests {
.unwrap()
});
let spec = crate::table::LsmWriteSpec::unsharded()
.with_maintained_indexes(vec!["id_idx".to_string()])
.with_maintained_indexes(["id_idx"])
.with_writer_config_defaults([("max_memtable_rows", "1000")]);
table.set_lsm_write_spec(spec).await.unwrap();
}
@@ -6618,29 +6617,11 @@ mod tests {
body["sharding"],
serde_json::json!({ "mode": "bucket", "column": "id", "num_buckets": 16 })
);
// An unpinned maintained set sends null: resolve server-side.
assert_eq!(body["maintained_indexes"], serde_json::Value::Null);
http::Response::builder().status(200).body("{}").unwrap()
});
table
.set_lsm_write_spec(crate::table::LsmWriteSpec::bucket("id", 16))
.await
.unwrap();
}
/// `[]` (none) must stay distinguishable on the wire from null (all).
#[tokio::test]
async fn test_set_lsm_write_spec_no_maintained_indexes() {
let table = Table::new_with_handler("my_table", |request| {
let body = request.body().unwrap().as_bytes().unwrap();
let body: serde_json::Value = serde_json::from_slice(body).unwrap();
assert_eq!(body["maintained_indexes"], serde_json::json!([]));
http::Response::builder().status(200).body("{}").unwrap()
});
table
.set_lsm_write_spec(
crate::table::LsmWriteSpec::bucket("id", 16).with_maintained_indexes(Vec::new()),
)
.set_lsm_write_spec(crate::table::LsmWriteSpec::bucket("id", 16))
.await
.unwrap();
}
@@ -6719,7 +6700,7 @@ mod tests {
} => {
assert_eq!(column, "id");
assert_eq!(num_buckets, 4);
assert_eq!(maintained_indexes, Some(vec!["id_idx".to_string()]));
assert_eq!(maintained_indexes, vec!["id_idx".to_string()]);
assert_eq!(
writer_config_defaults
.get("durable_write")
+157 -471
View File
@@ -50,6 +50,7 @@ use crate::DistanceType;
use crate::blob::BlobRangeRequest;
use crate::data::scannable::{PeekedScannable, Scannable, estimate_write_partitions};
use crate::database::Database;
use crate::database::listing::LANCE_FILE_EXTENSION;
use crate::database::read_freshness::TableFreshness;
use crate::embeddings::{EmbeddingDefinition, EmbeddingRegistry, MemoryRegistry};
use crate::error::{Error, Result};
@@ -151,6 +152,55 @@ pub(crate) fn map_namespace_lance_error(err: lance::Error, table_name: &str) ->
}
}
/// Map a `lance::Error::DatasetNotFound` for the table at `uri` into a `lancedb::Error`.
///
/// Lance reports "there is nothing at this location" and "there is a table directory
/// here but nothing loadable inside it" with the same error. Only the first is a
/// `TableNotFound`: a `<name>.lance` directory left behind by an interrupted drop and
/// re-create is still reported by `Connection::table_names`, so callers need to be able
/// to tell "never existed" from "exists but is broken".
///
/// See <https://github.com/lancedb/lancedb/issues/3127>.
async fn map_dataset_not_found(
uri: &str,
name: &str,
params: ReadParams,
err: lance::Error,
) -> Error {
let name = name.to_string();
let source = Box::new(err);
if table_dir_exists(uri, params).await.unwrap_or(false) {
Error::TableCorrupted { name, source }
} else {
Error::TableNotFound { name, source }
}
}
/// Whether a table directory is present at `uri`, even though no dataset could be
/// loaded from it.
///
/// This looks for a `<name>.lance` entry in the parent directory, which is exactly what
/// `ListingDatabase::table_names` lists, so the two APIs agree on whether a table is
/// present. Probing `uri` itself would not work: object stores have no empty
/// directories to probe, and on a local filesystem the interesting case is precisely an
/// empty directory.
async fn table_dir_exists(uri: &str, params: ReadParams) -> Result<bool> {
let (object_store, path, _) = DatasetBuilder::from_uri(uri)
.with_read_params(params)
.build_object_store()
.await?;
// Only `*.lance` entries are ever reported as tables, so nothing else can produce
// the list-then-open mismatch this guards against.
if path.extension() != Some(LANCE_FILE_EXTENSION) {
return Ok(false);
}
let (Some(parent), Some(dir_name)) = (path.parent(), path.filename()) else {
return Ok(false);
};
let entries = object_store.read_dir(parent).await?;
Ok(entries.iter().any(|entry| entry.as_str() == dir_name))
}
/// Defines the type of column
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum ColumnKind {
@@ -320,8 +370,6 @@ pub use self::merge::MergeResult;
/// date) and [`LsmWriteSpec::with_writer_config_defaults`] (default
/// `ShardWriter` configuration recorded in the MemWAL index).
///
/// A fresh spec maintains every index on the table, resolved on install.
///
/// Install a spec with [`Table::set_lsm_write_spec`] and remove it with
/// [`Table::unset_lsm_write_spec`]. The actual `merge_insert` dispatch
/// onto the MemWAL writer is a follow-up.
@@ -336,12 +384,9 @@ pub enum LsmWriteSpec {
Bucket {
column: String,
num_buckets: u32,
/// Indexes the MemWAL maintains in-memory as rows are appended.
///
/// `None` means every index it can maintain, resolved on install — a
/// snapshot, so indexes created later need the spec unset and re-set.
/// `Some([])` maintains nothing.
maintained_indexes: Option<Vec<String>>,
/// Names of indexes (already created on the table) that the
/// MemWAL should maintain in-memory as rows are appended.
maintained_indexes: Vec<String>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
writer_config_defaults: HashMap<String, String>,
},
@@ -351,41 +396,35 @@ pub enum LsmWriteSpec {
/// distinct value of `column` becomes its own shard.
Identity {
column: String,
/// Indexes the MemWAL maintains in-memory as rows are appended.
///
/// `None` means every index it can maintain, resolved on install — a
/// snapshot, so indexes created later need the spec unset and re-set.
/// `Some([])` maintains nothing.
maintained_indexes: Option<Vec<String>>,
/// Names of indexes (already created on the table) that the
/// MemWAL should maintain in-memory as rows are appended.
maintained_indexes: Vec<String>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
writer_config_defaults: HashMap<String, String>,
},
/// No sharding — every `merge_insert` call writes to a single MemWAL shard.
Unsharded {
/// Indexes the MemWAL maintains in-memory as rows are appended.
///
/// `None` means every index it can maintain, resolved on install — a
/// snapshot, so indexes created later need the spec unset and re-set.
/// `Some([])` maintains nothing.
maintained_indexes: Option<Vec<String>>,
/// Names of indexes (already created on the table) that the
/// MemWAL should maintain in-memory as rows are appended.
maintained_indexes: Vec<String>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
writer_config_defaults: HashMap<String, String>,
},
}
impl LsmWriteSpec {
/// Construct a hash-bucket sharding spec maintaining every index on the table.
/// Construct a hash-bucket sharding spec with no maintained indexes.
pub fn bucket(column: impl Into<String>, num_buckets: u32) -> Self {
Self::Bucket {
column: column.into(),
num_buckets,
maintained_indexes: None,
maintained_indexes: Vec::new(),
writer_config_defaults: HashMap::new(),
}
}
/// Construct an identity-sharding spec (shard by the raw value of
/// `column`) maintaining every index on the table.
/// `column`) with no maintained indexes.
///
/// `column` must be a deterministic function of the unenforced primary
/// key: every row with a given primary key must always produce the same
@@ -397,37 +436,28 @@ impl LsmWriteSpec {
pub fn identity(column: impl Into<String>) -> Self {
Self::Identity {
column: column.into(),
maintained_indexes: None,
maintained_indexes: Vec::new(),
writer_config_defaults: HashMap::new(),
}
}
/// Construct an unsharded spec maintaining every index on the table.
/// Construct an unsharded spec with no maintained indexes.
pub fn unsharded() -> Self {
Self::Unsharded {
maintained_indexes: None,
maintained_indexes: Vec::new(),
writer_config_defaults: HashMap::new(),
}
}
/// Set which indexes the MemWAL maintains.
///
/// `None` (the default) resolves to every index on the table at install,
/// failing if one cannot be maintained — name the set to install anyway. A
/// list is verbatim: each name must already exist and be maintainable, and
/// an empty list maintains nothing.
///
/// ```
/// # use lancedb::table::LsmWriteSpec;
/// // Every index the table has when the spec is installed:
/// LsmWriteSpec::unsharded().with_maintained_indexes(None);
/// // Exactly these:
/// LsmWriteSpec::unsharded().with_maintained_indexes(vec!["id_idx".to_string()]);
/// // None at all:
/// LsmWriteSpec::unsharded().with_maintained_indexes(Vec::new());
/// ```
pub fn with_maintained_indexes(mut self, indexes: impl Into<Option<Vec<String>>>) -> Self {
let indexes = indexes.into();
/// Replace the list of indexes the MemWAL should keep up to date as
/// rows are appended. Each name must reference an index that already
/// exists on the table at the time `set_lsm_write_spec` is called.
pub fn with_maintained_indexes<I, S>(mut self, indexes: I) -> Self
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
let v: Vec<String> = indexes.into_iter().map(Into::into).collect();
match &mut self {
Self::Bucket {
maintained_indexes, ..
@@ -437,7 +467,7 @@ impl LsmWriteSpec {
}
| Self::Unsharded {
maintained_indexes, ..
} => *maintained_indexes = indexes,
} => *maintained_indexes = v,
}
self
}
@@ -473,9 +503,8 @@ impl LsmWriteSpec {
self
}
/// Borrow the list of index names this spec asks MemWAL to maintain, or
/// `None` when it asks for every index on the table.
pub fn maintained_indexes(&self) -> Option<&[String]> {
/// Borrow the list of index names this spec asks MemWAL to maintain.
pub fn maintained_indexes(&self) -> &[String] {
match self {
Self::Bucket {
maintained_indexes, ..
@@ -485,7 +514,7 @@ impl LsmWriteSpec {
}
| Self::Unsharded {
maintained_indexes, ..
} => maintained_indexes.as_deref(),
} => maintained_indexes,
}
}
@@ -1683,7 +1712,7 @@ impl Table {
/// # async fn example(table: &Table) -> Result<(), Box<dyn std::error::Error>> {
/// table
/// .set_lsm_write_spec(
/// LsmWriteSpec::bucket("id", 16).with_maintained_indexes(vec!["id_idx".to_string()]),
/// LsmWriteSpec::bucket("id", 16).with_maintained_indexes(["id_idx"]),
/// )
/// .await?;
/// # Ok(())
@@ -1705,10 +1734,9 @@ impl Table {
///
/// Returns `Ok(None)` when the MemWAL LSM write path is not enabled (no
/// spec has been set, or it was removed with [`Table::unset_lsm_write_spec`]).
/// The returned spec mirrors what was passed to
/// [`Table::set_lsm_write_spec`], except that
/// [`LsmWriteSpec::maintained_indexes`] always reports the concrete list
/// resolved when the spec was set — `None` never round-trips.
/// The returned spec — including its [`LsmWriteSpec::maintained_indexes`] and
/// [`LsmWriteSpec::writer_config_defaults`] — mirrors what was passed to
/// [`Table::set_lsm_write_spec`].
///
/// # Example
///
@@ -2335,19 +2363,10 @@ impl NativeTable {
managed_versioning: Option<bool>,
) -> Result<Self> {
let params = params.unwrap_or_default();
let has_caller_store_wrapper = params
.store_options
.as_ref()
.and_then(|options| options.object_store_wrapper.as_ref())
.is_some();
// A caller wrapper must remain outside connection-level compatibility
// behavior. When there is no caller wrapper, apply the compatibility
// layer after loading so the session's registered store can be reused.
let (params, wrapper_after_load) = match write_store_wrapper {
Some(wrapper) if has_caller_store_wrapper => {
(params.patch_with_store_wrapper(wrapper)?, None)
}
wrapper => (params, wrapper),
// patch the params if we have a write store wrapper
let params = match write_store_wrapper.clone() {
Some(wrapper) => params.patch_with_store_wrapper(wrapper)?,
None => params,
};
// Build table_id from namespace + name
@@ -2379,6 +2398,8 @@ impl NativeTable {
None => false,
};
// Kept so that a `DatasetNotFound` can be re-checked against storage below.
let recovery_params = params.clone();
let mut builder = DatasetBuilder::from_uri(uri).with_read_params(params);
// Set up commit handler when managed_versioning is enabled
@@ -2397,23 +2418,10 @@ impl NativeTable {
let dataset = match builder.load().await {
Ok(dataset) => dataset,
Err(e @ lance::Error::DatasetNotFound { .. }) => {
// The manifest load is the existence check. A physical prefix may be
// from a concurrent or abandoned create, so it cannot refine this error.
return Err(Error::TableNotFound {
name: name.to_string(),
source: Box::new(e),
});
return Err(map_dataset_not_found(uri, name, recovery_params, e).await);
}
Err(e) => return Err(e.into()),
};
// Resolve the store from the session registry before applying a
// connection-level write wrapper. Wrapper identity is part of the
// registry key, so including it in ReadParams prevents reuse when the
// opened table (and its wrapped store) is short-lived.
let dataset = match wrapper_after_load {
Some(wrapper) => dataset.with_object_store_wrappers([wrapper]),
None => dataset,
};
let dataset = DatasetConsistencyWrapper::new_latest(dataset, read_consistency_interval);
let id = Self::build_id(&namespace, name);
@@ -2515,16 +2523,11 @@ impl NativeTable {
if let Some(sess) = session {
params.session(sess);
}
let has_caller_store_wrapper = params
.store_options
.as_ref()
.and_then(|options| options.object_store_wrapper.as_ref())
.is_some();
let (params, wrapper_after_load) = match write_store_wrapper {
Some(wrapper) if has_caller_store_wrapper => {
(params.patch_with_store_wrapper(wrapper)?, None)
}
wrapper => (params, wrapper),
// patch the params if we have a write store wrapper
let params = match write_store_wrapper.clone() {
Some(wrapper) => params.patch_with_store_wrapper(wrapper)?,
None => params,
};
// Build table_id from namespace + name
@@ -2548,13 +2551,6 @@ impl NativeTable {
},
e => e.into(),
})?;
// Apply the write wrapper after the session registry has resolved the
// shared store. The cloned dataset retains the wrapper for subsequent
// reads, manifest commits, and any additional base stores.
let dataset = match wrapper_after_load {
Some(wrapper) => dataset.with_object_store_wrappers([wrapper]),
None => dataset,
};
let uri = dataset.uri().to_string();
let dataset = DatasetConsistencyWrapper::new_latest(dataset, read_consistency_interval);
@@ -3689,8 +3685,8 @@ pub struct FragmentSummaryStats {
#[cfg(test)]
#[allow(deprecated)]
mod tests {
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
use std::time::Duration;
use arrow_array::{
@@ -3772,50 +3768,73 @@ mod tests {
);
}
#[tokio::test]
async fn test_open_not_found_when_empty_directory_exists() {
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
std::fs::create_dir(&dataset_path).unwrap();
/// Write a table and then break it, leaving the `<name>.lance` directory in place.
///
/// `remove_all` reproduces an interrupted drop + re-create (the directory is left
/// empty); otherwise only the manifests are removed, leaving the data files behind.
async fn write_then_corrupt_table(dir: &std::path::Path, remove_all: bool) -> String {
let dataset_path = dir.join("test.lance");
let uri = dataset_path.to_str().unwrap().to_string();
let err = NativeTable::open(dataset_path.to_str().unwrap())
.await
.unwrap_err();
let batch = make_test_batches();
let reader = RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
Dataset::write(reader, &uri, None).await.unwrap();
if remove_all {
for entry in std::fs::read_dir(&dataset_path).unwrap() {
let entry = entry.unwrap();
if entry.file_type().unwrap().is_dir() {
std::fs::remove_dir_all(entry.path()).unwrap();
} else {
std::fs::remove_file(entry.path()).unwrap();
}
}
assert_eq!(std::fs::read_dir(&dataset_path).unwrap().count(), 0);
} else {
let versions = dataset_path.join("_versions");
assert!(versions.is_dir(), "expected manifests under {versions:?}");
std::fs::remove_dir_all(&versions).unwrap();
assert!(std::fs::read_dir(&dataset_path).unwrap().count() > 0);
}
uri
}
#[tokio::test]
async fn test_open_corrupt_empty_dir() {
let tmp_dir = tempdir().unwrap();
let uri = write_then_corrupt_table(tmp_dir.path(), true).await;
let err = NativeTable::open(&uri).await.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "test"),
matches!(&err, Error::TableCorrupted { name, .. } if name == "test"),
"got {err:?}"
);
}
#[tokio::test]
async fn test_open_not_found_when_only_uncommitted_storage_exists() {
async fn test_open_corrupt_missing_manifest() {
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let data_dir = dataset_path.join("data");
std::fs::create_dir_all(&data_dir).unwrap();
std::fs::write(data_dir.join("orphan.lance"), b"uncommitted").unwrap();
let uri = write_then_corrupt_table(tmp_dir.path(), false).await;
let err = NativeTable::open(dataset_path.to_str().unwrap())
.await
.unwrap_err();
let err = NativeTable::open(&uri).await.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "test"),
matches!(&err, Error::TableCorrupted { name, .. } if name == "test"),
"got {err:?}"
);
}
/// Listing databases discover physical `*.lance` entries. That snapshot is not an
/// authoritative table-existence check: only a committed manifest makes a table
/// openable, and the entry could also be concurrently created or dropped.
/// A table listed by `table_names()` must not be reported as missing by
/// `open_table()`. See <https://github.com/lancedb/lancedb/issues/3127>.
#[tokio::test]
async fn test_table_names_may_include_uncommitted_storage() {
async fn test_open_table_corrupt_is_still_listed() {
let tmp_dir = tempdir().unwrap();
let db = connect(tmp_dir.path().to_str().unwrap())
.execute()
.await
.unwrap();
std::fs::create_dir(tmp_dir.path().join("test.lance")).unwrap();
write_then_corrupt_table(tmp_dir.path(), true).await;
assert_eq!(
db.table_names().execute().await.unwrap(),
@@ -3823,177 +3842,12 @@ mod tests {
);
let err = db.open_table("test").execute().await.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "test"),
"physical storage without a committed manifest is not a table: {err:?}"
);
}
#[derive(Debug)]
struct ParentListGuardStore {
inner: Arc<dyn object_store::ObjectStore>,
parent: object_store::path::Path,
parent_list_calls: Arc<AtomicUsize>,
}
impl std::fmt::Display for ParentListGuardStore {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("ParentListGuardStore")
}
}
#[async_trait::async_trait]
#[deny(clippy::missing_trait_methods)]
impl object_store::ObjectStore for ParentListGuardStore {
async fn put_opts(
&self,
location: &object_store::path::Path,
payload: object_store::PutPayload,
opts: object_store::PutOptions,
) -> object_store::Result<object_store::PutResult> {
self.inner.put_opts(location, payload, opts).await
}
async fn put_multipart_opts(
&self,
location: &object_store::path::Path,
opts: object_store::PutMultipartOptions,
) -> object_store::Result<Box<dyn object_store::MultipartUpload>> {
self.inner.put_multipart_opts(location, opts).await
}
async fn get_opts(
&self,
location: &object_store::path::Path,
options: object_store::GetOptions,
) -> object_store::Result<object_store::GetResult> {
self.inner.get_opts(location, options).await
}
async fn get_ranges(
&self,
location: &object_store::path::Path,
ranges: &[std::ops::Range<u64>],
) -> object_store::Result<Vec<bytes::Bytes>> {
self.inner.get_ranges(location, ranges).await
}
fn delete_stream(
&self,
locations: futures::stream::BoxStream<
'static,
object_store::Result<object_store::path::Path>,
>,
) -> futures::stream::BoxStream<'static, object_store::Result<object_store::path::Path>>
{
self.inner.delete_stream(locations)
}
fn list(
&self,
prefix: Option<&object_store::path::Path>,
) -> futures::stream::BoxStream<'static, object_store::Result<object_store::ObjectMeta>>
{
if prefix == Some(&self.parent) {
self.parent_list_calls.fetch_add(1, Ordering::Relaxed);
}
self.inner.list(prefix)
}
fn list_with_offset(
&self,
prefix: Option<&object_store::path::Path>,
offset: &object_store::path::Path,
) -> futures::stream::BoxStream<'static, object_store::Result<object_store::ObjectMeta>>
{
if prefix == Some(&self.parent) {
self.parent_list_calls.fetch_add(1, Ordering::Relaxed);
}
self.inner.list_with_offset(prefix, offset)
}
async fn list_with_delimiter(
&self,
prefix: Option<&object_store::path::Path>,
) -> object_store::Result<object_store::ListResult> {
if prefix == Some(&self.parent) {
self.parent_list_calls.fetch_add(1, Ordering::Relaxed);
}
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(
&self,
from: &object_store::path::Path,
to: &object_store::path::Path,
options: object_store::CopyOptions,
) -> object_store::Result<()> {
self.inner.copy_opts(from, to, options).await
}
async fn rename_opts(
&self,
from: &object_store::path::Path,
to: &object_store::path::Path,
options: object_store::RenameOptions,
) -> object_store::Result<()> {
self.inner.rename_opts(from, to, options).await
}
}
#[derive(Debug)]
struct ParentListGuardWrapper {
parent_list_calls: Arc<AtomicUsize>,
}
impl WrappingObjectStore for ParentListGuardWrapper {
fn wrap(
&self,
_store_prefix: &str,
inner: Arc<dyn object_store::ObjectStore>,
) -> Arc<dyn object_store::ObjectStore> {
Arc::new(ParentListGuardStore {
inner,
parent: object_store::path::Path::from("database"),
parent_list_calls: self.parent_list_calls.clone(),
})
}
}
#[tokio::test]
async fn test_open_missing_never_lists_database_parent() {
let parent_list_calls = Arc::new(AtomicUsize::new(0));
let params = ReadParams {
store_options: Some(ObjectStoreParams {
object_store_wrapper: Some(Arc::new(ParentListGuardWrapper {
parent_list_calls: parent_list_calls.clone(),
})),
..Default::default()
}),
..Default::default()
};
let err = NativeTable::open_with_params(
"memory:///database/missing.lance",
"missing",
Vec::new(),
None,
Some(params),
None,
None,
HashSet::new(),
None,
)
.await
.unwrap_err();
assert!(
matches!(&err, Error::TableNotFound { name, .. } if name == "missing"),
matches!(&err, Error::TableCorrupted { name, .. } if name == "test"),
"got {err:?}"
);
assert_eq!(
parent_list_calls.load(Ordering::Relaxed),
0,
"opening one missing table must not enumerate sibling tables"
assert!(
err.to_string().contains("exists but could not be loaded"),
"got {err}"
);
}
@@ -4062,66 +3916,6 @@ mod tests {
}
}
#[derive(Debug)]
struct OrderedStoreWrapper {
name: &'static str,
order: Arc<Mutex<Vec<&'static str>>>,
}
impl WrappingObjectStore for OrderedStoreWrapper {
fn wrap(
&self,
_store_prefix: &str,
original: Arc<dyn object_store::ObjectStore>,
) -> Arc<dyn object_store::ObjectStore> {
self.order.lock().unwrap().push(self.name);
original
}
}
#[tokio::test]
async fn test_open_with_params_keeps_caller_store_wrapper_outermost() {
let tmp_dir = tempdir().unwrap();
let dataset_path = tmp_dir.path().join("test.lance");
let uri = dataset_path.to_str().unwrap();
let batch = make_test_batches();
let reader = RecordBatchIterator::new(vec![Ok(batch.clone())], batch.schema());
Dataset::write(reader, uri, None).await.unwrap();
let order = Arc::new(Mutex::new(Vec::new()));
let caller_wrapper = Arc::new(OrderedStoreWrapper {
name: "caller",
order: order.clone(),
});
let compatibility_wrapper = Arc::new(OrderedStoreWrapper {
name: "compatibility",
order: order.clone(),
});
let params = ReadParams {
store_options: Some(ObjectStoreParams {
object_store_wrapper: Some(caller_wrapper),
..Default::default()
}),
..Default::default()
};
NativeTable::open_with_params(
uri,
"test",
vec![],
Some(compatibility_wrapper),
Some(params),
None,
None,
HashSet::new(),
None,
)
.await
.unwrap();
assert_eq!(*order.lock().unwrap(), vec!["compatibility", "caller"]);
}
#[tokio::test]
async fn test_open_table_options() {
let tmp_dir = tempdir().unwrap();
@@ -5291,7 +5085,7 @@ mod tests {
// Bucket spec round-trips exactly, including the routing column (recovered
// from its field id), maintained indexes, and writer config defaults.
let spec = LsmWriteSpec::bucket("id", 4)
.with_maintained_indexes(vec![idx_name.clone()])
.with_maintained_indexes([idx_name])
.with_writer_config_defaults([("durable_write", "false")]);
table.set_lsm_write_spec(spec.clone()).await.unwrap();
assert_eq!(table.get_lsm_write_spec().await.unwrap(), Some(spec));
@@ -5301,125 +5095,15 @@ mod tests {
assert_eq!(table.get_lsm_write_spec().await.unwrap(), None);
// Identity sharding round-trips (column recovered from the schema).
// A spec left at its default maintains every index on the table, so it
// reads back naming the one on the table rather than as "infer".
let spec = LsmWriteSpec::identity("region");
table.set_lsm_write_spec(spec.clone()).await.unwrap();
assert_eq!(
table.get_lsm_write_spec().await.unwrap(),
Some(spec.with_maintained_indexes(vec![idx_name.clone()]))
);
assert_eq!(table.get_lsm_write_spec().await.unwrap(), Some(spec));
table.unset_lsm_write_spec().await.unwrap();
// Unsharded round-trips (no routing column).
let spec = LsmWriteSpec::unsharded();
table.set_lsm_write_spec(spec.clone()).await.unwrap();
assert_eq!(
table.get_lsm_write_spec().await.unwrap(),
Some(spec.with_maintained_indexes(vec![idx_name]))
);
}
/// The maintained set defaults to every index on the table, resolved at
/// install. An index the memtable cannot build fails the install rather
/// than being dropped: maintaining it would take the table offline for
/// writes, dropping it would hide that from the caller.
#[tokio::test]
async fn test_set_lsm_write_spec_infers_maintained_indexes() {
let tmp_dir = tempdir().unwrap();
let uri = tmp_dir.path().to_str().unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new("tag", DataType::Utf8, true),
]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(arrow_array::Int64Array::from(vec![1, 2, 3])),
Arc::new(StringArray::from(vec!["a", "b", "c"])),
],
)
.unwrap();
let reader: Box<dyn arrow_array::RecordBatchReader + Send> =
Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema.clone()));
let conn = ConnectBuilder::new(uri)
.read_consistency_interval(Duration::from_secs(0))
.execute()
.await
.unwrap();
let table = conn.create_table("t", reader).execute().await.unwrap();
table
.create_index(&["id"], Index::BTree(Default::default()))
.name("id_btree".to_string())
.execute()
.await
.unwrap();
table
.create_index(&["tag"], Index::Bitmap(Default::default()))
.name("tag_bitmap".to_string())
.execute()
.await
.unwrap();
// Explicitly naming the bitmap index fails before anything commits.
let err = table
.set_lsm_write_spec(
LsmWriteSpec::unsharded().with_maintained_indexes(vec!["tag_bitmap".to_string()]),
)
.await
.unwrap_err();
assert!(
matches!(err, Error::InvalidInput { ref message } if message.contains("tag_bitmap")),
"expected the bitmap index to be rejected, got {err:?}"
);
assert_eq!(table.get_lsm_write_spec().await.unwrap(), None);
// The default covers every index, so the bitmap fails it too.
let err = table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap_err();
assert!(
matches!(err, Error::InvalidInput { ref message }
if message.contains("tag_bitmap") && message.contains("maintained_indexes")),
"expected the inferred set to be rejected, got {err:?}"
);
assert_eq!(table.get_lsm_write_spec().await.unwrap(), None);
// Naming the maintainable subset installs.
table
.set_lsm_write_spec(
LsmWriteSpec::unsharded().with_maintained_indexes(vec!["id_btree".to_string()]),
)
.await
.unwrap();
assert_eq!(
table
.get_lsm_write_spec()
.await
.unwrap()
.unwrap()
.maintained_indexes(),
Some(["id_btree".to_string()].as_slice())
);
// Opting out entirely is distinct from the default.
table.unset_lsm_write_spec().await.unwrap();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes(Vec::new()))
.await
.unwrap();
assert_eq!(
table
.get_lsm_write_spec()
.await
.unwrap()
.unwrap()
.maintained_indexes(),
Some([].as_slice())
);
assert_eq!(table.get_lsm_write_spec().await.unwrap(), Some(spec));
}
#[tokio::test]
@@ -5523,8 +5207,8 @@ mod tests {
pub async fn test_stats_includes_index_and_overlay_files() {
use lance::dataset::WriteDestination;
use lance::dataset::transaction::{DataOverlayGroup, Operation};
use lance_file::version::stable_file_version;
use lance_file::writer::FileWriterOptions;
use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
use lance_file::writer::{FileWriter, FileWriterOptions};
use lance_io::utils::CachedFileSize;
use lance_table::format::DataFile;
use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage};
@@ -5589,17 +5273,19 @@ mod tests {
let fragment_id = dataset.get_fragments()[0].id() as u64;
let foo_field_id = dataset.schema().field("foo").unwrap().id;
let overlay_schema = dataset.schema().project_by_ids(&[foo_field_id], true);
let file_version = stable_file_version();
let file_version = ConcreteFileVersion::from(LanceFileVersion::Stable);
let filename = "overlay.lance".to_string();
let store = dataset.object_store(None).await.unwrap();
let path = dataset.data_dir().child(filename.clone());
let obj_writer = store.create(&path).await.unwrap();
let mut writer = lance_file::versions::create_writer(
file_version,
let mut writer = FileWriter::try_new(
obj_writer,
overlay_schema,
FileWriterOptions::default(),
FileWriterOptions {
format_version: Some(file_version.into()),
..Default::default()
},
)
.unwrap();
writer
+2 -2
View File
@@ -1161,7 +1161,7 @@ mod lsm_tests {
.unwrap();
let fts_index = table.list_indices().await.unwrap()[0].name.clone();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes(vec![fts_index]))
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes([fts_index]))
.await
.unwrap();
@@ -1254,7 +1254,7 @@ mod lsm_tests {
.unwrap();
let vec_index = table.list_indices().await.unwrap()[0].name.clone();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes(vec![vec_index]))
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes([vec_index]))
.await
.unwrap();
+8 -74
View File
@@ -29,7 +29,6 @@ use arrow_schema::{DataType, Schema as ArrowSchema, SchemaRef};
use lance::Dataset;
use lance::dataset::mem_wal::{
DatasetMemWalExt, ShardWriter, ShardWriterConfig, evaluate_sharding_spec,
validate_maintained_indexes,
};
use lance::index::DatasetIndexExt;
use lance_core::datatypes::Schema as LanceSchema;
@@ -38,9 +37,8 @@ use tokio::sync::RwLock;
use uuid::Uuid;
use crate::error::{Error, Result};
use crate::index::IndexConfig;
use crate::table::merge::{MergeInsertBuilder, MergeResult};
use crate::table::{BaseTable, LsmWriteSpec, NativeTable};
use crate::table::{LsmWriteSpec, NativeTable};
/// Spec id of the sole sharding spec installed by [`set_lsm_write_spec`].
/// Must match Lance's `InitializeMemWalBuilder` (`SHARDING_SPEC_ID`).
@@ -82,44 +80,32 @@ pub(crate) async fn set_lsm_write_spec(table: &NativeTable, spec: LsmWriteSpec)
}
}
// Before the builder borrows the dataset clone. `list_indices` merges an
// index's segments into one entry, so the result needs no dedup.
let maintained_indexes = {
let dataset = table.dataset.get().await?;
resolve_maintained_indexes(
&dataset,
&table.list_indices().await?,
spec.maintained_indexes(),
)
.await?
};
let mut dataset = (*table.dataset.get().await?).clone();
let mut builder = dataset.initialize_mem_wal();
let writer_config_defaults = match spec {
let (maintained_indexes, writer_config_defaults) = match spec {
LsmWriteSpec::Bucket {
column,
num_buckets,
maintained_indexes,
writer_config_defaults,
..
} => {
builder = builder.bucket_sharding(column, num_buckets);
writer_config_defaults
(maintained_indexes, writer_config_defaults)
}
LsmWriteSpec::Identity {
column,
maintained_indexes,
writer_config_defaults,
..
} => {
builder = builder.identity_sharding(column);
writer_config_defaults
(maintained_indexes, writer_config_defaults)
}
LsmWriteSpec::Unsharded {
maintained_indexes,
writer_config_defaults,
..
} => {
builder = builder.unsharded();
writer_config_defaults
(maintained_indexes, writer_config_defaults)
}
};
builder = builder.maintained_indexes(maintained_indexes);
@@ -131,58 +117,6 @@ pub(crate) async fn set_lsm_write_spec(table: &NativeTable, spec: LsmWriteSpec)
Ok(())
}
/// Resolve a spec's maintained-index selection against `indices`, as reported
/// by [`Table::list_indices`](crate::Table::list_indices).
///
/// `None` means every index on the table, snapshotted now. Lance validates
/// either selection against its shard-writer rules, so a spec that installs is
/// one the MemWAL can open.
///
/// An unmaintainable index fails an inferred set rather than being dropped from
/// it — dropping would leave the caller believing it is maintained.
async fn resolve_maintained_indexes(
dataset: &Dataset,
indices: &[IndexConfig],
requested: Option<&[String]>,
) -> Result<Vec<String>> {
let Some(requested) = requested else {
let all: Vec<String> = indices.iter().map(|index| index.name.clone()).collect();
validate_maintained_indexes(dataset, &all)
.await
.map_err(|source| Error::InvalidInput {
message: format!(
"cannot maintain every index on this table: {source}. Set \
maintained_indexes explicitly to choose from {}",
index_name_list(indices),
),
})?;
return Ok(all);
};
for name in requested {
if !indices.iter().any(|index| &index.name == name) {
return Err(Error::InvalidInput {
message: format!(
"maintained index '{}' does not exist on this table; it has {}",
name,
index_name_list(indices),
),
});
}
}
validate_maintained_indexes(dataset, requested).await?;
Ok(requested.to_vec())
}
/// Index names for an error message.
fn index_name_list(indices: &[IndexConfig]) -> String {
if indices.is_empty() {
return "no indexes".to_string();
}
let mut names: Vec<&str> = indices.iter().map(|index| index.name.as_str()).collect();
names.sort_unstable();
format!("[{}]", names.join(", "))
}
// =============================================================================
// unset_lsm_write_spec
// =============================================================================
+7 -54
View File
@@ -14,7 +14,6 @@ use lance::arrow::json::JsonDataType;
use lance::dataset::{ReadParams, WriteParams};
use lance::index::vector::utils::infer_vector_dim;
use lance::io::{ObjectStoreParams, WrappingObjectStore};
use lance_io::object_store::ChainedWrappingObjectStore;
use std::pin::Pin;
use crate::error::{Error, Result};
@@ -38,13 +37,13 @@ impl PatchStoreParam for Option<ObjectStoreParams> {
wrapper: Arc<dyn WrappingObjectStore>,
) -> Result<Option<ObjectStoreParams>> {
let mut params = self.unwrap_or_default();
params.object_store_wrapper = Some(match params.object_store_wrapper.take() {
// The wrapper being patched in is connection-level compatibility
// behavior. Keep it closest to the target store so an existing
// caller wrapper remains outermost and can observe every operation.
Some(existing) => Arc::new(ChainedWrappingObjectStore::new(vec![wrapper, existing])),
None => wrapper,
});
if params.object_store_wrapper.is_some() {
return Err(Error::Other {
message: "can not patch param because object store is already set".into(),
source: None,
});
}
params.object_store_wrapper = Some(wrapper);
Ok(Some(params))
}
@@ -473,60 +472,14 @@ impl Stream for MaxBatchLengthStream {
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use arrow_array::Int32Array;
use arrow_schema::Field;
use datafusion_physical_plan::stream::RecordBatchStreamAdapter;
use futures::{StreamExt, stream};
use object_store::{ObjectStore, memory::InMemory};
use tokio::time::sleep;
use super::*;
#[derive(Debug)]
struct OrderedStoreWrapper {
name: &'static str,
order: Arc<Mutex<Vec<&'static str>>>,
}
impl WrappingObjectStore for OrderedStoreWrapper {
fn wrap(
&self,
_store_prefix: &str,
original: Arc<dyn ObjectStore>,
) -> Arc<dyn ObjectStore> {
self.order.lock().unwrap().push(self.name);
original
}
}
#[test]
fn test_patch_store_param_keeps_caller_wrapper_outermost() {
let order = Arc::new(Mutex::new(Vec::new()));
let params = Some(ObjectStoreParams {
object_store_wrapper: Some(Arc::new(OrderedStoreWrapper {
name: "caller",
order: order.clone(),
})),
..Default::default()
});
let params = params
.patch_with_store_wrapper(Arc::new(OrderedStoreWrapper {
name: "compatibility",
order: order.clone(),
}))
.unwrap()
.unwrap();
params
.object_store_wrapper
.unwrap()
.wrap("memory", Arc::new(InMemory::new()) as Arc<dyn ObjectStore>);
assert_eq!(*order.lock().unwrap(), vec!["compatibility", "caller"]);
}
#[test]
fn test_guess_default_column() {
let schema_no_vector = Schema::new(vec![
+13 -18
View File
@@ -10,7 +10,7 @@ use arrow_array::{
use arrow_schema::{DataType, Field, Fields, Schema};
use futures::TryStreamExt;
use lance::Dataset;
use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
use lance_file::version::LanceFileVersion;
use lancedb::{
Connection, Error, Result, Table,
blob::{BlobRangeRequest, blob},
@@ -61,7 +61,7 @@ async fn create_inline_blob_table(
Ok(table)
}
async fn storage_format_version(table: &Table) -> ConcreteFileVersion {
async fn storage_format_version(table: &Table) -> LanceFileVersion {
table
.as_native()
.unwrap()
@@ -69,14 +69,9 @@ async fn storage_format_version(table: &Table) -> ConcreteFileVersion {
.await
.unwrap()
.data_storage_format
.lance_file_format()
}
fn supports_blob_v2(version: ConcreteFileVersion) -> bool {
matches!(
version,
ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3
)
.lance_file_version()
.unwrap()
.resolve()
}
async fn uses_stable_row_ids(table: &Table) -> bool {
@@ -117,7 +112,7 @@ async fn declaring_blob_column_bumps_format_and_enables_stable_row_ids() -> Resu
.execute()
.await?;
assert!(supports_blob_v2(storage_format_version(&table).await));
assert!(storage_format_version(&table).await >= LanceFileVersion::V2_2);
assert!(uses_stable_row_ids(&table).await);
Ok(())
}
@@ -132,7 +127,7 @@ async fn explicit_stable_row_id_setting_wins_over_blob_default() -> Result<()> {
.execute()
.await?;
assert!(supports_blob_v2(storage_format_version(&table).await));
assert!(storage_format_version(&table).await >= LanceFileVersion::V2_2);
assert!(!uses_stable_row_ids(&table).await);
Ok(())
}
@@ -144,7 +139,7 @@ async fn non_blob_table_keeps_default_format_and_row_id_setting() -> Result<()>
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int64, false)]));
let table = db.create_empty_table("t", schema).execute().await?;
assert!(!supports_blob_v2(storage_format_version(&table).await));
assert!(storage_format_version(&table).await < LanceFileVersion::V2_2);
assert!(!uses_stable_row_ids(&table).await);
Ok(())
}
@@ -176,7 +171,7 @@ async fn creating_with_blob_data_bumps_format() -> Result<()> {
.unwrap();
let table = db.create_table("t", batch).execute().await?;
assert!(supports_blob_v2(storage_format_version(&table).await));
assert!(storage_format_version(&table).await >= LanceFileVersion::V2_2);
assert!(uses_stable_row_ids(&table).await);
assert_eq!(table.count_rows(None).await?, 1);
Ok(())
@@ -286,7 +281,7 @@ async fn connection_level_stable_row_id_setting_wins_over_blob_default() -> Resu
.execute()
.await?;
assert!(supports_blob_v2(storage_format_version(&table).await));
assert!(storage_format_version(&table).await >= LanceFileVersion::V2_2);
assert!(!uses_stable_row_ids(&table).await);
Ok(())
}
@@ -302,7 +297,7 @@ async fn namespace_create_applies_blob_defaults() -> Result<()> {
.execute()
.await?;
assert!(supports_blob_v2(storage_format_version(&table).await));
assert!(storage_format_version(&table).await >= LanceFileVersion::V2_2);
assert!(uses_stable_row_ids(&table).await);
Ok(())
}
@@ -479,7 +474,7 @@ async fn fetch_blobs_round_trips_nested_blob_column() -> Result<()> {
let batch = RecordBatch::try_new(schema, vec![Arc::new(info_array) as ArrayRef]).unwrap();
let table = db.create_table("t", batch).execute().await?;
assert!(supports_blob_v2(storage_format_version(&table).await));
assert!(storage_format_version(&table).await >= LanceFileVersion::V2_2);
assert!(uses_stable_row_ids(&table).await);
let ids = collect_row_ids(&table).await?;
@@ -1310,7 +1305,7 @@ async fn optimize_preserves_blob_v2_null_and_empty_distinction() -> Result<()> {
.await?;
table.add(null_empty_input_batch()).execute().await?;
assert!(
supports_blob_v2(storage_format_version(&table).await),
storage_format_version(&table).await >= LanceFileVersion::V2_2,
"blob v2 columns require storage >= 2.2"
);