mirror of
https://github.com/lancedb/lancedb.git
synced 2026-08-30 18:08:24 +00:00
Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 1825f6df9e | |||
| 6a1dacb1fe | |||
| 62dea8acd8 | |||
| ae9e8e8f8d | |||
| 36e44ab7a9 | |||
| feccabd739 | |||
| 1f3093a51f |
+1
-1
@@ -1,5 +1,5 @@
|
||||
[tool.bumpversion]
|
||||
current_version = "0.37.1-beta.1"
|
||||
current_version = "0.37.1-beta.0"
|
||||
parse = """(?x)
|
||||
(?P<major>0|[1-9]\\d*)\\.
|
||||
(?P<minor>0|[1-9]\\d*)\\.
|
||||
|
||||
@@ -1,222 +0,0 @@
|
||||
name: Check doc links
|
||||
|
||||
# Checking external links is inherently noisy: third-party sites rate-limit
|
||||
# automated clients, reject non-browser user agents, and go down temporarily.
|
||||
# Blocking pull requests on that trades a lot of false failures for very little
|
||||
# signal, so this runs on a schedule and reports findings in a single tracking
|
||||
# issue instead of failing anyone's build.
|
||||
on:
|
||||
schedule:
|
||||
- cron: "0 7 * * *"
|
||||
workflow_dispatch:
|
||||
|
||||
# The report lives in one repository-global issue, so runs must not overlap: a
|
||||
# lookup racing a create produces duplicate issues, and a healthy run closing
|
||||
# the issue while a failing run only rewrites its body would leave a broken
|
||||
# report closed. The group is deliberately ref-independent so that a manual
|
||||
# dispatch serializes against the scheduled run.
|
||||
concurrency:
|
||||
group: docs-link-check
|
||||
cancel-in-progress: false
|
||||
|
||||
permissions: {}
|
||||
|
||||
env:
|
||||
REPORT_TITLE: "Docs link checker report"
|
||||
|
||||
jobs:
|
||||
scan:
|
||||
name: Scan links
|
||||
runs-on: ubuntu-24.04
|
||||
# lychee-action is pinned by SHA, but its wrapper downloads the lychee
|
||||
# release tarball at run time without verifying a digest, and hands the
|
||||
# resulting binary a GitHub token. Release assets remain replaceable, so
|
||||
# that binary is confined to a job whose token can only read public
|
||||
# content; everything that writes runs in the report job below.
|
||||
permissions:
|
||||
contents: read
|
||||
outputs:
|
||||
exit_code: ${{ steps.lychee.outputs.exit_code }}
|
||||
steps:
|
||||
- name: Checkout
|
||||
uses: actions/checkout@v6
|
||||
with:
|
||||
# workflow_dispatch can run from any ref, but the report is
|
||||
# repository-global. Always measure the default branch so a manual
|
||||
# run from a topic branch cannot close a report that main warrants,
|
||||
# or overwrite it with branch-only findings.
|
||||
ref: ${{ github.event.repository.default_branch }}
|
||||
persist-credentials: false
|
||||
|
||||
- name: Check links
|
||||
id: lychee
|
||||
uses: lycheeverse/lychee-action@e7477775783ea5526144ba13e8db5eec57747ce8 # v2.9.0
|
||||
with:
|
||||
# Restricted to http(s) on purpose. Much of docs/src is generated
|
||||
# API reference (the js/ tree comes from `npm run docs` in nodejs)
|
||||
# and the hand-written pages use mkdocstrings cross-references and
|
||||
# nav-relative paths that only resolve in the site mkdocs builds,
|
||||
# not in this checkout, so relative links would be reported as
|
||||
# broken on every run.
|
||||
args: >-
|
||||
--scheme https
|
||||
--scheme http
|
||||
--no-progress
|
||||
--max-retries 3
|
||||
--timeout 20
|
||||
'docs/src/**/*.md'
|
||||
format: json
|
||||
output: ./lychee/out.json
|
||||
jobSummary: false
|
||||
# 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
|
||||
# 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 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:
|
||||
EXIT_CODE: ${{ steps.lychee.outputs.exit_code }}
|
||||
run: |
|
||||
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.lychee.outputs.exit_code == 2
|
||||
uses: actions/upload-artifact@v7
|
||||
with:
|
||||
name: link-report
|
||||
path: ./lychee/out.json
|
||||
retention-days: 7
|
||||
|
||||
report:
|
||||
name: Update report issue
|
||||
needs: scan
|
||||
runs-on: ubuntu-24.04
|
||||
# Deliberately no checkout: this job needs the report artifact and the
|
||||
# issues API, not the repository contents.
|
||||
permissions:
|
||||
issues: write
|
||||
env:
|
||||
EXIT_CODE: ${{ needs.scan.outputs.exit_code }}
|
||||
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:
|
||||
# the issue action applies labels in a separate call after creating the
|
||||
# issue, so a label filter misses a half-created report, and this
|
||||
# repository has far more open issues than one listing page holds.
|
||||
# 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 links break again.
|
||||
run: |
|
||||
match=$(gh issue list --repo "$GITHUB_REPOSITORY" --state all \
|
||||
--search "in:title \"$REPORT_TITLE\" author:app/github-actions" \
|
||||
--limit 50 --json number,title,state \
|
||||
--jq "[.[] | select(.title == \"$REPORT_TITLE\")] | sort_by(.number) | first // empty")
|
||||
echo "number=$(jq -r '.number // empty' <<<"$match")" >> "$GITHUB_OUTPUT"
|
||||
echo "state=$(jq -r '.state // empty' <<<"$match")" >> "$GITHUB_OUTPUT"
|
||||
|
||||
- name: Download report
|
||||
if: env.EXIT_CODE == 2
|
||||
uses: actions/download-artifact@v8
|
||||
with:
|
||||
name: link-report
|
||||
path: ./lychee
|
||||
|
||||
- name: Compose report
|
||||
if: env.EXIT_CODE == 2
|
||||
run: |
|
||||
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
|
||||
{
|
||||
echo "Broken documentation links found by [\`$GITHUB_WORKFLOW\`]($run_url)."
|
||||
echo
|
||||
echo "This issue is rewritten by every scheduled run and closed automatically once all links resolve."
|
||||
echo
|
||||
echo "Entries can be false positives: some sites rate-limit or block automated clients while working fine in a browser. Confirm before editing the docs, and add persistent offenders to \`--exclude\` in \`.github/workflows/docs-link-check.yml\`."
|
||||
echo
|
||||
# Timeouts are reported alongside errors: entries land in
|
||||
# timeout_map with a status text instead of an HTTP code.
|
||||
jq -r '
|
||||
"\(.errors) of \(.total) links failed, \(.timeouts) timed out.",
|
||||
"",
|
||||
([(.error_map | to_entries[]), (.timeout_map | to_entries[])]
|
||||
| group_by(.key)[] |
|
||||
"### Errors in \(.[0].key)",
|
||||
"",
|
||||
(map(.value[])[] | "* [\(.status.code // .status.text // "ERR")] <\(.url)> — \(.status.details // .status.text // "unknown error")"),
|
||||
"")
|
||||
' ./lychee/out.json
|
||||
} > ./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, 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 "Broken documentation links found again in [the latest run]($run_url)."
|
||||
|
||||
- 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
|
||||
# the same issue is updated in place.
|
||||
issue-number: ${{ steps.report.outputs.number }}
|
||||
title: ${{ env.REPORT_TITLE }}
|
||||
content-filepath: ./lychee/issue.md
|
||||
labels: documentation
|
||||
|
||||
- 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.EXIT_CODE == 0 && steps.report.outputs.state == 'OPEN'
|
||||
env:
|
||||
ISSUE_NUMBER: ${{ steps.report.outputs.number }}
|
||||
run: |
|
||||
run_url="$GITHUB_SERVER_URL/$GITHUB_REPOSITORY/actions/runs/$GITHUB_RUN_ID"
|
||||
gh issue close "$ISSUE_NUMBER" --repo "$GITHUB_REPOSITORY" \
|
||||
--comment "All documentation links resolved in [the latest run]($run_url)."
|
||||
@@ -296,18 +296,16 @@ jobs:
|
||||
cargo update -p aws-types --precise 1.3.9
|
||||
cargo update -p aws-sigv4 --precise 1.3.5
|
||||
cargo update -p aws-credential-types --precise 1.2.8
|
||||
# aws-smithy-checksums must stay at or above 0.63.13: OpenDAL's S3
|
||||
# service needs crc-fast ~1.9, and older releases pin it to ~1.3.
|
||||
cargo update -p aws-smithy-checksums --precise 0.63.13
|
||||
cargo update -p aws-smithy-checksums --precise 0.63.9
|
||||
cargo update -p aws-smithy-runtime --precise 1.9.3
|
||||
cargo update -p aws-smithy-http --precise 0.62.6
|
||||
cargo update -p aws-smithy-eventstream --precise 0.60.14
|
||||
cargo update -p aws-smithy-http --precise 0.62.4
|
||||
cargo update -p aws-smithy-eventstream --precise 0.60.12
|
||||
cargo update -p aws-smithy-http-client --precise 1.1.3
|
||||
cargo update -p aws-smithy-observability --precise 0.1.4
|
||||
cargo update -p aws-smithy-query --precise 0.60.8
|
||||
cargo update -p aws-smithy-runtime-api --precise 1.9.3
|
||||
cargo update -p aws-smithy-async --precise 1.2.7
|
||||
cargo update -p aws-smithy-types --precise 1.3.6
|
||||
cargo update -p aws-smithy-runtime-api --precise 1.9.1
|
||||
cargo update -p aws-smithy-async --precise 1.2.6
|
||||
cargo update -p aws-smithy-types --precise 1.3.5
|
||||
cargo update -p aws-smithy-xml --precise 0.60.11
|
||||
cargo update -p home --precise 0.5.9
|
||||
- name: cargo +${{ matrix.msrv }} check
|
||||
|
||||
Generated
+237
-262
File diff suppressed because it is too large
Load Diff
+15
-15
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
lance = { "version" = "=11.0.0-beta.3", default-features = false, "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=11.0.0-beta.3", default-features = false, "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=11.0.0-beta.3", default-features = false, "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=11.0.0-beta.3", "tag" = "v11.0.0-beta.3", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
ahash = "0.8"
|
||||
# Note that this one does not include pyarrow
|
||||
arrow = { version = "58.0.0", optional = false }
|
||||
@@ -52,7 +52,7 @@ env_logger = "0.11"
|
||||
half = { "version" = "2.7.1", default-features = false, features = [
|
||||
"num-traits",
|
||||
] }
|
||||
futures = "0.3"
|
||||
futures = "0"
|
||||
log = "0.4"
|
||||
metrics = "0.24"
|
||||
metrics-util = "0.19"
|
||||
|
||||
@@ -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-beta.0</version>
|
||||
</dependency>
|
||||
```
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
***
|
||||
|
||||
|
||||
@@ -44,7 +44,4 @@ The number of rows in the table
|
||||
totalBytes: number;
|
||||
```
|
||||
|
||||
The total size, in bytes, of the table's data files, index files, and
|
||||
overlay files
|
||||
|
||||
Read from the manifest, so this excludes deletion files and manifests.
|
||||
The total number of bytes in the table
|
||||
|
||||
@@ -31,7 +31,7 @@ is also an [asynchronous API client](#connections-asynchronous).
|
||||
## Namespaces (Synchronous)
|
||||
|
||||
A namespace-backed connection resolves tables through a
|
||||
[Lance namespace](https://lance-format.github.io/lance-namespace/) service instead of
|
||||
[Lance namespace](https://lancedb.github.io/lance-namespace/) service instead of
|
||||
listing a storage directory.
|
||||
|
||||
::: lancedb.connect_namespace
|
||||
|
||||
@@ -8,7 +8,7 @@
|
||||
<parent>
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-parent</artifactId>
|
||||
<version>0.37.1-beta.1</version>
|
||||
<version>0.37.1-beta.0</version>
|
||||
<relativePath>../pom.xml</relativePath>
|
||||
</parent>
|
||||
|
||||
|
||||
+2
-2
@@ -6,7 +6,7 @@
|
||||
|
||||
<groupId>com.lancedb</groupId>
|
||||
<artifactId>lancedb-parent</artifactId>
|
||||
<version>0.37.1-beta.1</version>
|
||||
<version>0.37.1-beta.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.3</lance-core.version>
|
||||
<lance-core.version>10.1.0-beta.1</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
@@ -1,7 +1,7 @@
|
||||
[package]
|
||||
name = "lancedb-nodejs"
|
||||
edition.workspace = true
|
||||
version = "0.37.1-beta.1"
|
||||
version = "0.37.1-beta.0"
|
||||
publish = false
|
||||
license.workspace = true
|
||||
description.workspace = true
|
||||
|
||||
@@ -69,33 +69,6 @@ describe("given a connection", () => {
|
||||
await expect(tbl.countRows()).resolves.toBe(1);
|
||||
});
|
||||
|
||||
it("should isolate object-form table creation across databases", async () => {
|
||||
const otherTmpDir = tmp.dirSync({ unsafeCleanup: true });
|
||||
const otherDb = await connect(otherTmpDir.name);
|
||||
|
||||
try {
|
||||
const firstTable = await db.createTable({
|
||||
name: "defaultTable",
|
||||
data: [{ rowId: "id1", vector: Array(384).fill(0) }],
|
||||
});
|
||||
const secondTable = await otherDb.createTable({
|
||||
name: "defaultTable",
|
||||
data: [{ rowId: "id2", vector: Array(384).fill(0) }],
|
||||
});
|
||||
|
||||
await expect(db.tableNames()).resolves.toEqual(["defaultTable"]);
|
||||
await expect(otherDb.tableNames()).resolves.toEqual(["defaultTable"]);
|
||||
|
||||
const firstRows = await firstTable.query().select(["rowId"]).toArray();
|
||||
const secondRows = await secondTable.query().select(["rowId"]).toArray();
|
||||
expect(firstRows.map((row) => row.rowId)).toEqual(["id1"]);
|
||||
expect(secondRows.map((row) => row.rowId)).toEqual(["id2"]);
|
||||
} finally {
|
||||
otherDb.close();
|
||||
otherTmpDir.removeCallback();
|
||||
}
|
||||
});
|
||||
|
||||
it("should be able to drop tables`", async () => {
|
||||
await db.createTable("test", [{ id: 1 }, { id: 2 }]);
|
||||
await db.createTable("test2", [{ id: 1 }, { id: 2 }]);
|
||||
|
||||
@@ -11,11 +11,8 @@ import {
|
||||
Float16,
|
||||
Float32,
|
||||
Float64,
|
||||
Int32,
|
||||
Schema,
|
||||
Utf8,
|
||||
fromDataToBuffer,
|
||||
tableFromIPC,
|
||||
} from "../lancedb/arrow";
|
||||
import { EmbeddingFunction, LanceSchema } from "../lancedb/embedding";
|
||||
import { getRegistry, register } from "../lancedb/embedding/registry";
|
||||
@@ -187,63 +184,6 @@ describe("embedding functions", () => {
|
||||
const vector0 = JSON.parse(JSON.stringify(arr[0].vector));
|
||||
expect(vector0).toEqual([1, 2, 3]);
|
||||
});
|
||||
|
||||
it("should append generated vectors to a non-nullable schema", async () => {
|
||||
@register("non_nullable_schema_test")
|
||||
class MockEmbeddingFunction extends EmbeddingFunction<string> {
|
||||
ndims() {
|
||||
return 3;
|
||||
}
|
||||
embeddingDataType(): Float {
|
||||
return new Float64();
|
||||
}
|
||||
async computeSourceEmbeddings(data: string[]) {
|
||||
return data.map(() => [1, 2, 3]);
|
||||
}
|
||||
}
|
||||
|
||||
const schema = new Schema([
|
||||
new Field("id", new Int32()),
|
||||
new Field("text", new Utf8()),
|
||||
new Field("type", new Utf8()),
|
||||
new Field(
|
||||
"vector",
|
||||
new FixedSizeList(3, new Field("item", new Float64())),
|
||||
),
|
||||
]);
|
||||
const func = new MockEmbeddingFunction();
|
||||
const db = await connect(tmpDir.name);
|
||||
const table = await db.createEmptyTable("test_non_nullable", schema, {
|
||||
embeddingFunction: {
|
||||
function: func,
|
||||
sourceColumn: "text",
|
||||
},
|
||||
});
|
||||
|
||||
const data = [
|
||||
{ id: 1, text: "Carrot", type: "vegetable" },
|
||||
{ id: 2, text: "Apple", type: "fruit" },
|
||||
];
|
||||
const buffer = await fromDataToBuffer(
|
||||
data,
|
||||
undefined,
|
||||
await table.schema(),
|
||||
);
|
||||
const generatedTable = tableFromIPC(buffer);
|
||||
const vectorField = generatedTable.schema.fields.find(
|
||||
(field) => field.name === "vector",
|
||||
);
|
||||
expect(vectorField?.nullable).toBe(false);
|
||||
|
||||
await table.add(data);
|
||||
|
||||
const rows = await table.query().toArray();
|
||||
expect(rows).toHaveLength(2);
|
||||
for (const row of rows) {
|
||||
expect([...row.vector]).toEqual([1, 2, 3]);
|
||||
}
|
||||
});
|
||||
|
||||
it("should error when appending to a table with an unregistered embedding function", async () => {
|
||||
@register("mock")
|
||||
class MockEmbeddingFunction extends EmbeddingFunction<string> {
|
||||
|
||||
@@ -1,14 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
import packageJson = require("../package.json");
|
||||
|
||||
describe("package metadata", () => {
|
||||
it("requires Node.js type declarations compatible with the runtime", () => {
|
||||
expect(packageJson.engines.node).toBe(">= 18");
|
||||
expect(packageJson.peerDependencies["@types/node"]).toBe(">=18");
|
||||
expect(packageJson.peerDependenciesMeta["@types/node"]).toEqual({
|
||||
optional: true,
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -110,81 +110,6 @@ describe("Query outputSchema", () => {
|
||||
});
|
||||
});
|
||||
|
||||
describe("Search pagination", () => {
|
||||
let tmpDir: tmp.DirResult;
|
||||
let table: Table;
|
||||
|
||||
beforeEach(async () => {
|
||||
tmpDir = tmp.dirSync({ unsafeCleanup: true });
|
||||
const db = await connect(tmpDir.name);
|
||||
const schema = new Schema([
|
||||
new Field("id", new Int64(), false),
|
||||
new Field("text", new Utf8(), false),
|
||||
new Field(
|
||||
"vector",
|
||||
new FixedSizeList(2, new Field("item", new Float32())),
|
||||
false,
|
||||
),
|
||||
]);
|
||||
const data = makeArrowTable(
|
||||
[
|
||||
{ id: 1n, text: "common", vector: [0, 0] },
|
||||
{ id: 2n, text: "common common", vector: [1, 1] },
|
||||
{ id: 3n, text: "common common common", vector: [2, 2] },
|
||||
{ id: 4n, text: "common common common common", vector: [3, 3] },
|
||||
],
|
||||
{ schema },
|
||||
);
|
||||
table = await db.createTable("test", data);
|
||||
});
|
||||
|
||||
afterEach(() => {
|
||||
tmpDir.removeCallback();
|
||||
});
|
||||
|
||||
it("applies offset after the vector search limit", async () => {
|
||||
const allResults = await table
|
||||
.vectorSearch([0, 0])
|
||||
.select(["id"])
|
||||
.limit(4)
|
||||
.toArray();
|
||||
const secondPage = await table
|
||||
.vectorSearch([0, 0])
|
||||
.select(["id"])
|
||||
.limit(2)
|
||||
.offset(2)
|
||||
.toArray();
|
||||
|
||||
expect(allResults).toHaveLength(4);
|
||||
expect(secondPage).toHaveLength(2);
|
||||
expect(secondPage.map((row) => row.id)).toEqual(
|
||||
allResults.slice(2, 4).map((row) => row.id),
|
||||
);
|
||||
});
|
||||
|
||||
it("applies offset after the full-text search limit", async () => {
|
||||
await table.createIndex("text", { config: Index.fts() });
|
||||
|
||||
const allResults = await table
|
||||
.search("common", "fts")
|
||||
.select(["id"])
|
||||
.limit(4)
|
||||
.toArray();
|
||||
const secondPage = await table
|
||||
.search("common", "fts")
|
||||
.select(["id"])
|
||||
.limit(2)
|
||||
.offset(2)
|
||||
.toArray();
|
||||
|
||||
expect(allResults).toHaveLength(4);
|
||||
expect(secondPage).toHaveLength(2);
|
||||
expect(secondPage.map((row) => row.id)).toEqual(
|
||||
allResults.slice(2, 4).map((row) => row.id),
|
||||
);
|
||||
});
|
||||
});
|
||||
|
||||
describe("Query orderBy", () => {
|
||||
let tmpDir: tmp.DirResult;
|
||||
let table: Table;
|
||||
|
||||
@@ -277,16 +277,8 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
|
||||
},
|
||||
numIndices: 0,
|
||||
numRows: 3,
|
||||
// Full on-disk size of the two data files, footers and metadata included.
|
||||
totalBytes: 684,
|
||||
totalBytes: 44,
|
||||
});
|
||||
|
||||
// Index files count toward totalBytes too (only deletion files and
|
||||
// manifests are excluded).
|
||||
await table.createIndex("id", { config: Index.btree() });
|
||||
const statsWithIndex = await table.stats();
|
||||
expect(statsWithIndex.numIndices).toBe(1);
|
||||
expect(statsWithIndex.totalBytes).toBeGreaterThan(684);
|
||||
});
|
||||
|
||||
it("should overwrite data if asked", async () => {
|
||||
|
||||
+4
-14
@@ -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,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-darwin-arm64",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"os": ["darwin"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.darwin-arm64.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-arm64-gnu",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"os": ["linux"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.linux-arm64-gnu.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-arm64-musl",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"os": ["linux"],
|
||||
"cpu": ["arm64"],
|
||||
"main": "lancedb.linux-arm64-musl.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-x64-gnu",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"os": ["linux"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.linux-x64-gnu.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-linux-x64-musl",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"os": ["linux"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.linux-x64-musl.node",
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-win32-arm64-msvc",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"os": [
|
||||
"win32"
|
||||
],
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-win32-x64-msvc",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"os": ["win32"],
|
||||
"cpu": ["x64"],
|
||||
"main": "lancedb.win32-x64-msvc.node",
|
||||
|
||||
Generated
+2
-8
@@ -1,12 +1,12 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"lockfileVersion": 3,
|
||||
"requires": true,
|
||||
"packages": {
|
||||
"": {
|
||||
"name": "@lancedb/lancedb",
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"cpu": [
|
||||
"x64",
|
||||
"arm64"
|
||||
@@ -55,13 +55,7 @@
|
||||
"openai": "4.29.2"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@types/node": ">=18",
|
||||
"apache-arrow": ">=15.0.0 <=18.1.0"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
"@types/node": {
|
||||
"optional": true
|
||||
}
|
||||
}
|
||||
},
|
||||
"node_modules/@aws-crypto/crc32": {
|
||||
|
||||
+1
-7
@@ -11,7 +11,7 @@
|
||||
"ann"
|
||||
],
|
||||
"private": false,
|
||||
"version": "0.37.1-beta.1",
|
||||
"version": "0.37.1-beta.0",
|
||||
"main": "dist/index.js",
|
||||
"exports": {
|
||||
".": "./dist/index.js",
|
||||
@@ -101,12 +101,6 @@
|
||||
"openai": "4.29.2"
|
||||
},
|
||||
"peerDependencies": {
|
||||
"@types/node": ">=18",
|
||||
"apache-arrow": ">=15.0.0 <=18.1.0"
|
||||
},
|
||||
"peerDependenciesMeta": {
|
||||
"@types/node": {
|
||||
"optional": true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+7
-10
@@ -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),
|
||||
},
|
||||
}
|
||||
@@ -1043,10 +1043,7 @@ impl From<lancedb::index::IndexStatistics> for IndexStatistics {
|
||||
|
||||
#[napi(object)]
|
||||
pub struct TableStatistics {
|
||||
/// The total size, in bytes, of the table's data files, index files, and
|
||||
/// overlay files
|
||||
///
|
||||
/// Read from the manifest, so this excludes deletion files and manifests.
|
||||
/// The total number of bytes in the table
|
||||
pub total_bytes: i64,
|
||||
|
||||
/// The number of rows in the table
|
||||
|
||||
+1
-1
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "lancedb-python"
|
||||
version = "0.37.1-beta.1"
|
||||
version = "0.37.1-beta.0"
|
||||
publish = false
|
||||
edition.workspace = true
|
||||
description = "Python bindings for LanceDB"
|
||||
|
||||
@@ -60,7 +60,7 @@ tests = [
|
||||
"pytest-asyncio>=0.21",
|
||||
"duckdb>=0.9.0",
|
||||
"pytz>=2023.3",
|
||||
"polars>=0.19, <=1.32.3",
|
||||
"polars>=0.19, <=1.3.0",
|
||||
"pyarrow<25",
|
||||
"pyarrow-stubs>=16.0",
|
||||
"pylance==9.0.0rc1",
|
||||
@@ -140,7 +140,6 @@ include = [
|
||||
"python/lancedb/remote/errors.py",
|
||||
"python/lancedb/embeddings/__init__.py",
|
||||
"python/lancedb/_lancedb.pyi",
|
||||
"python/type_tests/connect.py",
|
||||
]
|
||||
exclude = ["python/tests/"]
|
||||
pythonVersion = "3.13"
|
||||
|
||||
@@ -355,10 +355,6 @@ class Table:
|
||||
async def set_lsm_write_spec(self, spec: LsmWriteSpec) -> None: ...
|
||||
async def unset_lsm_write_spec(self) -> None: ...
|
||||
async def get_lsm_write_spec(self) -> Optional[LsmWriteSpec]: ...
|
||||
async def checkpoint_lsm(self) -> None: ...
|
||||
async def flush_lsm(self) -> None: ...
|
||||
async def compact_lsm(self) -> None: ...
|
||||
async def get_lsm_stats(self, include_generation_rows: bool) -> Optional[dict]: ...
|
||||
async def close_lsm_writers(self) -> None: ...
|
||||
@property
|
||||
def tags(self) -> Tags: ...
|
||||
@@ -653,10 +649,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 +666,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]: ...
|
||||
|
||||
|
||||
@@ -804,7 +804,7 @@ class LanceDBConnection(DBConnection):
|
||||
"manifest_enabled": self._manifest_enabled,
|
||||
"namespace_client_properties": self._namespace_client_properties,
|
||||
"read_consistency_interval_seconds": (
|
||||
rci.total_seconds() if rci else None
|
||||
rci.total_seconds() if rci is not None else None
|
||||
),
|
||||
}
|
||||
)
|
||||
|
||||
@@ -87,13 +87,12 @@ class JinaEmbeddings(EmbeddingFunction):
|
||||
if isinstance(image, bytes):
|
||||
image_dict = {"image": base64.b64encode(image).decode("utf-8")}
|
||||
elif isinstance(image, (str, Path)):
|
||||
parsed = urlparse(str(image))
|
||||
parsed = urlparse.urlparse(image)
|
||||
# TODO handle drive letter on windows.
|
||||
PIL_Image = attempt_import_or_raise("PIL.Image", "pillow")
|
||||
if parsed.scheme == "file":
|
||||
pil_image = PIL_Image.open(parsed.path)
|
||||
elif parsed.scheme == "" or (os.name == "nt" and len(parsed.scheme) == 1):
|
||||
# A Windows drive letter parses as a one-character scheme
|
||||
# ("C:\\img.png" -> scheme="c"), so treat it as a local path.
|
||||
elif parsed.scheme == "":
|
||||
pil_image = PIL_Image.open(image if os.name == "nt" else parsed.path)
|
||||
elif parsed.scheme.startswith("http"):
|
||||
pil_image = PIL_Image.open(io.BytesIO(url_retrieve(image)))
|
||||
|
||||
@@ -482,6 +482,16 @@ class LanceNamespaceDBConnection(DBConnection):
|
||||
def serialize(self) -> str:
|
||||
import json
|
||||
|
||||
if (
|
||||
self._namespace_client_impl is None
|
||||
or self._namespace_client_properties is None
|
||||
):
|
||||
raise ValueError(
|
||||
"Cannot serialize a namespace connection constructed from an "
|
||||
"opaque namespace client. Pass namespace_client_impl and "
|
||||
"namespace_client_properties when constructing the connection."
|
||||
)
|
||||
|
||||
return json.dumps(
|
||||
{
|
||||
"connection_type": "namespace",
|
||||
@@ -493,7 +503,7 @@ class LanceNamespaceDBConnection(DBConnection):
|
||||
"storage_options": self.storage_options or None,
|
||||
"read_consistency_interval_seconds": (
|
||||
self.read_consistency_interval.total_seconds()
|
||||
if self.read_consistency_interval
|
||||
if self.read_consistency_interval is not None
|
||||
else None
|
||||
),
|
||||
}
|
||||
@@ -569,6 +579,7 @@ class LanceNamespaceDBConnection(DBConnection):
|
||||
self,
|
||||
name,
|
||||
namespace_path=namespace_path,
|
||||
storage_options=storage_options,
|
||||
namespace_client=self._namespace_client,
|
||||
pushdown_operations=self._namespace_client_pushdown_operations,
|
||||
route_pushdown_to_rust=self._route_pushdown_to_rust,
|
||||
@@ -607,6 +618,8 @@ class LanceNamespaceDBConnection(DBConnection):
|
||||
self,
|
||||
name,
|
||||
namespace_path=namespace_path,
|
||||
storage_options=storage_options,
|
||||
index_cache_size=index_cache_size,
|
||||
namespace_client=self._namespace_client,
|
||||
pushdown_operations=self._namespace_client_pushdown_operations,
|
||||
route_pushdown_to_rust=self._route_pushdown_to_rust,
|
||||
@@ -899,10 +912,13 @@ class LanceNamespaceDBConnection(DBConnection):
|
||||
self,
|
||||
name,
|
||||
namespace_path=namespace_path,
|
||||
storage_options=storage_options,
|
||||
index_cache_size=index_cache_size,
|
||||
location=table_uri,
|
||||
namespace_client=namespace_client,
|
||||
managed_versioning=managed_versioning,
|
||||
pushdown_operations=self._namespace_client_pushdown_operations,
|
||||
route_pushdown_to_rust=self._route_pushdown_to_rust,
|
||||
_async=async_table,
|
||||
)
|
||||
|
||||
|
||||
@@ -591,6 +591,13 @@ class Permutation:
|
||||
then the first split will be used.
|
||||
"""
|
||||
assert base_table is not None, "base_table is required"
|
||||
# A PyTorch fork worker may construct its Permutation lazily from a
|
||||
# table opened in the parent process. Reopen that table before the
|
||||
# Rust reader clones its object-store clients and connection pools.
|
||||
if hasattr(base_table, "_ensure_open"):
|
||||
base_table._ensure_open()
|
||||
if permutation_table is not None and hasattr(permutation_table, "_ensure_open"):
|
||||
permutation_table._ensure_open()
|
||||
if split is not None:
|
||||
if permutation_table is None:
|
||||
raise ValueError(
|
||||
|
||||
@@ -1 +0,0 @@
|
||||
|
||||
@@ -153,16 +153,6 @@ def Vector(
|
||||
return FixedSizeList
|
||||
|
||||
|
||||
def _raise_bare_vector_error(*_args):
|
||||
raise TypeError("Vector must be parameterized with a dimension, e.g. Vector(128).")
|
||||
|
||||
|
||||
# Pydantic v1 and v2 otherwise treat the bare Vector factory as a field validator
|
||||
# and inspect its signature, which produces misleading errors about internal types.
|
||||
setattr(Vector, "__get_validators__", _raise_bare_vector_error)
|
||||
setattr(Vector, "__get_pydantic_core_schema__", _raise_bare_vector_error)
|
||||
|
||||
|
||||
def MultiVector(
|
||||
dim: int, value_type: pa.DataType = pa.float32(), nullable: bool = True
|
||||
) -> Type:
|
||||
|
||||
+310
-128
@@ -6,6 +6,8 @@ from __future__ import annotations
|
||||
import asyncio
|
||||
import inspect
|
||||
import deprecation
|
||||
import os
|
||||
import threading
|
||||
import warnings
|
||||
from abc import ABC, abstractmethod
|
||||
from dataclasses import dataclass
|
||||
@@ -108,11 +110,6 @@ def _should_push_down_query_table(
|
||||
return namespace_client is not None and "QueryTable" in pushdown_operations
|
||||
|
||||
|
||||
def _polars_predicate_pushdown_barrier(frame: Any) -> Any:
|
||||
"""Return a Polars frame unchanged while blocking predicate pushdown."""
|
||||
return frame
|
||||
|
||||
|
||||
_MODEL_BACKED_TOKENIZER_PREFIXES = ("jieba", "lindera")
|
||||
_MODEL_BACKED_TOKENIZER_ERRORS = (
|
||||
"unknown base tokenizer",
|
||||
@@ -168,7 +165,7 @@ def _maybe_add_fts_error_note(
|
||||
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from .db import LanceDBConnection
|
||||
from .db import DBConnection, LanceDBConnection
|
||||
from ._lancedb import (
|
||||
Table as LanceDBTable,
|
||||
OptimizeStats,
|
||||
@@ -869,18 +866,12 @@ class Table(ABC):
|
||||
"""
|
||||
raise NotImplementedError
|
||||
|
||||
def to_polars(self, **kwargs) -> "pl.LazyFrame":
|
||||
"""Return the table as a Polars LazyFrame.
|
||||
|
||||
Note
|
||||
----
|
||||
The Polars streaming engine is not supported because it does not currently
|
||||
implement Python PyArrow dataset scans. Use the default engine when collecting
|
||||
this LazyFrame.
|
||||
def to_polars(self, **kwargs) -> "pl.DataFrame":
|
||||
"""Return the table as a polars.DataFrame.
|
||||
|
||||
Returns
|
||||
-------
|
||||
polars.LazyFrame
|
||||
polars.DataFrame
|
||||
"""
|
||||
raise NotImplementedError
|
||||
|
||||
@@ -2116,6 +2107,23 @@ class Table(ABC):
|
||||
"""
|
||||
|
||||
|
||||
@dataclass
|
||||
class _LanceTableReopenState:
|
||||
"""Process-independent coordinates for reopening a native table."""
|
||||
|
||||
connection_state: Optional[str]
|
||||
can_reopen_after_fork: bool
|
||||
fork_reopen_error: Optional[str]
|
||||
name: str
|
||||
namespace_path: List[str]
|
||||
storage_options: Optional[Dict[str, str]]
|
||||
index_cache_size: Optional[int]
|
||||
location: Optional[str]
|
||||
managed_versioning: Optional[bool]
|
||||
branch: Optional[str]
|
||||
checkout_version: Optional[int]
|
||||
|
||||
|
||||
class LanceTable(Table):
|
||||
"""
|
||||
A table in a LanceDB database.
|
||||
@@ -2150,7 +2158,10 @@ class LanceTable(Table):
|
||||
namespace_path = []
|
||||
self._conn = connection
|
||||
self._namespace_path = namespace_path
|
||||
self._storage_options = storage_options
|
||||
self._index_cache_size = index_cache_size
|
||||
self._location = location # Store location for use in _dataset_path
|
||||
self._managed_versioning = managed_versioning
|
||||
self._namespace_client = namespace_client
|
||||
self._pushdown_operations = pushdown_operations or set()
|
||||
# When the connection built the namespace client natively (e.g. an
|
||||
@@ -2175,9 +2186,206 @@ class LanceTable(Table):
|
||||
managed_versioning=managed_versioning,
|
||||
)
|
||||
)
|
||||
self._initialize_reopen_state(name)
|
||||
|
||||
def _initialize_reopen_state(self, name: str) -> None:
|
||||
"""Capture the state needed to replace inherited native handles."""
|
||||
self._name = name
|
||||
self._pid = os.getpid()
|
||||
self._native_state_guard = (self._pid, threading.RLock())
|
||||
|
||||
# A native table owns object-store clients and connection pools. Those
|
||||
# handles must not be used after fork, so retain a process-independent
|
||||
# connection description while it is still safe to inspect the parent
|
||||
# connection. Connections without reconstructible metadata are not
|
||||
# safe to reuse in a forked child, so retain a clear diagnostic rather
|
||||
# than advertising them as reopenable based on JSON encoding alone.
|
||||
try:
|
||||
connection_uri: Optional[str] = self._conn.uri
|
||||
except Exception:
|
||||
connection_uri = None
|
||||
|
||||
fork_reopen_error: Optional[str] = None
|
||||
try:
|
||||
connection_state: Optional[str] = self._conn.serialize()
|
||||
can_reopen_after_fork = connection_uri is not None and not (
|
||||
connection_uri.startswith("memory://")
|
||||
)
|
||||
except Exception as error:
|
||||
connection_state = None
|
||||
can_reopen_after_fork = False
|
||||
if connection_uri is not None and not connection_uri.startswith(
|
||||
"memory://"
|
||||
):
|
||||
fork_reopen_error = (
|
||||
f"Cannot reopen table {name!r} in a forked process: {error}"
|
||||
)
|
||||
|
||||
self._reopen_state = _LanceTableReopenState(
|
||||
connection_state=connection_state,
|
||||
can_reopen_after_fork=can_reopen_after_fork,
|
||||
fork_reopen_error=fork_reopen_error,
|
||||
name=name,
|
||||
namespace_path=list(self._namespace_path),
|
||||
storage_options=(
|
||||
dict(self._storage_options)
|
||||
if self._storage_options is not None
|
||||
else None
|
||||
),
|
||||
index_cache_size=self._index_cache_size,
|
||||
location=self._location,
|
||||
managed_versioning=self._managed_versioning,
|
||||
branch=self._table.current_branch(),
|
||||
checkout_version=None,
|
||||
)
|
||||
|
||||
@property
|
||||
def _connection_state(self) -> Optional[str]:
|
||||
"""Serialized connection retained for worker reconstruction."""
|
||||
return self._reopen_state.connection_state
|
||||
|
||||
@property
|
||||
def _can_reopen_after_fork(self) -> bool:
|
||||
return self._reopen_state.can_reopen_after_fork
|
||||
|
||||
@property
|
||||
def _branch(self) -> Optional[str]:
|
||||
state = getattr(self, "_reopen_state", None)
|
||||
if state is not None:
|
||||
return state.branch
|
||||
return getattr(self, "_legacy_branch", None)
|
||||
|
||||
@_branch.setter
|
||||
def _branch(self, value: Optional[str]) -> None:
|
||||
state = getattr(self, "_reopen_state", None)
|
||||
if state is not None:
|
||||
state.branch = value
|
||||
else:
|
||||
self._legacy_branch = value
|
||||
|
||||
@property
|
||||
def _checkout_version(self) -> Optional[int]:
|
||||
state = getattr(self, "_reopen_state", None)
|
||||
if state is not None:
|
||||
return state.checkout_version
|
||||
return getattr(self, "_legacy_checkout_version", None)
|
||||
|
||||
@_checkout_version.setter
|
||||
def _checkout_version(self, value: Optional[int]) -> None:
|
||||
state = getattr(self, "_reopen_state", None)
|
||||
if state is not None:
|
||||
state.checkout_version = value
|
||||
else:
|
||||
self._legacy_checkout_version = value
|
||||
|
||||
def _native_state_lock(self):
|
||||
"""Return the per-process lock coordinating native mode and reopen state."""
|
||||
pid = os.getpid()
|
||||
guard = getattr(self, "_native_state_guard", None)
|
||||
if guard is None:
|
||||
candidate = (pid, threading.RLock())
|
||||
guard = self.__dict__.setdefault("_native_state_guard", candidate)
|
||||
elif guard[0] != pid:
|
||||
# A lock inherited while another parent thread held it cannot be
|
||||
# safely acquired in the child. Child state starts single-threaded,
|
||||
# so replace it before coordinating the first reopen.
|
||||
guard = (pid, threading.RLock())
|
||||
self._native_state_guard = guard
|
||||
return guard[1]
|
||||
|
||||
@classmethod
|
||||
def _open_from_reopen_state(
|
||||
cls,
|
||||
connection: "DBConnection",
|
||||
state: "_LanceTableReopenState",
|
||||
) -> "LanceTable":
|
||||
"""Open a table from its complete process-independent descriptor."""
|
||||
async_connection = getattr(connection, "_conn", None)
|
||||
if async_connection is None:
|
||||
async_connection = connection._inner
|
||||
|
||||
namespace_client = getattr(connection, "_namespace_client", None)
|
||||
async_table = LOOP.run(
|
||||
async_connection.open_table(
|
||||
state.name,
|
||||
namespace_path=state.namespace_path,
|
||||
storage_options=state.storage_options,
|
||||
index_cache_size=state.index_cache_size,
|
||||
location=state.location,
|
||||
namespace_client=namespace_client,
|
||||
managed_versioning=state.managed_versioning,
|
||||
)
|
||||
)
|
||||
table = cls(
|
||||
connection,
|
||||
state.name,
|
||||
namespace_path=state.namespace_path,
|
||||
storage_options=state.storage_options,
|
||||
index_cache_size=state.index_cache_size,
|
||||
location=state.location,
|
||||
namespace_client=namespace_client,
|
||||
managed_versioning=state.managed_versioning,
|
||||
pushdown_operations=getattr(
|
||||
connection, "_namespace_client_pushdown_operations", None
|
||||
),
|
||||
route_pushdown_to_rust=getattr(
|
||||
connection, "_route_pushdown_to_rust", False
|
||||
),
|
||||
_async=async_table,
|
||||
)
|
||||
if state.branch is not None:
|
||||
table = table.branches.checkout(state.branch, state.checkout_version)
|
||||
elif state.checkout_version is not None:
|
||||
table.checkout(state.checkout_version)
|
||||
return table
|
||||
|
||||
def _ensure_open(self) -> None:
|
||||
"""Reopen native table handles inherited from another process."""
|
||||
with self._native_state_lock():
|
||||
pid = os.getpid()
|
||||
if getattr(self, "_pid", pid) == pid:
|
||||
return
|
||||
|
||||
state = getattr(self, "_reopen_state", None)
|
||||
fork_reopen_error = getattr(state, "fork_reopen_error", None)
|
||||
if fork_reopen_error is not None:
|
||||
raise RuntimeError(fork_reopen_error)
|
||||
if (
|
||||
state is None
|
||||
or not state.can_reopen_after_fork
|
||||
or state.connection_state is None
|
||||
):
|
||||
# In-memory and opaque Rust-only connections cannot be recreated
|
||||
# from connection metadata. Their local handles retain the prior
|
||||
# best-effort fork behavior.
|
||||
self._pid = pid
|
||||
return
|
||||
|
||||
from lancedb import deserialize_conn
|
||||
|
||||
connection = deserialize_conn(state.connection_state, for_worker=True)
|
||||
reopened = self._open_from_reopen_state(
|
||||
connection,
|
||||
state,
|
||||
)
|
||||
|
||||
# Keep this Python object stable because user datasets commonly retain
|
||||
# it across fork. Replace every process-bound component with the fresh
|
||||
# child's equivalent.
|
||||
self._conn = reopened._conn
|
||||
self._table = reopened._table
|
||||
self._namespace_client = reopened._namespace_client
|
||||
self._pushdown_operations = reopened._pushdown_operations
|
||||
self._route_pushdown_to_rust = reopened._route_pushdown_to_rust
|
||||
self._reopen_state = reopened._reopen_state
|
||||
self._pid = pid
|
||||
|
||||
@property
|
||||
def name(self) -> str:
|
||||
if hasattr(self, "_name"):
|
||||
return self._name
|
||||
# Preserve compatibility with lightweight / legacy instances that
|
||||
# were constructed without running ``LanceTable.__init__``.
|
||||
return self._table.name
|
||||
|
||||
@property
|
||||
@@ -2394,18 +2602,75 @@ class LanceTable(Table):
|
||||
def _wrap_branch_handle(
|
||||
self, async_table: "AsyncTable", version: Optional[int] = None
|
||||
) -> "LanceTable":
|
||||
# version is unused locally: the pin already lives on async_table and a
|
||||
# local handle is not reopened via a serialized connection.
|
||||
return LanceTable(
|
||||
table = LanceTable(
|
||||
self._conn,
|
||||
async_table.name,
|
||||
namespace_path=self._namespace_path,
|
||||
storage_options=self._storage_options,
|
||||
index_cache_size=self._index_cache_size,
|
||||
namespace_client=self._namespace_client,
|
||||
pushdown_operations=self._pushdown_operations,
|
||||
route_pushdown_to_rust=self._route_pushdown_to_rust,
|
||||
location=self._location,
|
||||
managed_versioning=self._managed_versioning,
|
||||
_async=async_table,
|
||||
)
|
||||
table._checkout_version = version
|
||||
return table
|
||||
|
||||
def _resolve_checkout_version(self, version: Union[int, str]) -> int:
|
||||
if isinstance(version, int):
|
||||
return version
|
||||
try:
|
||||
return self.tags.get_version(version)
|
||||
except RuntimeError as err:
|
||||
# Native checkout historically exposes an unknown tag as ValueError.
|
||||
# Preserve that contract while resolving tags before mutating the table.
|
||||
if "Ref not found" in str(err) and "does not exist" in str(err):
|
||||
raise ValueError(str(err)) from err
|
||||
raise
|
||||
|
||||
async def _commit_native_state(
|
||||
self,
|
||||
transition,
|
||||
version: Optional[int],
|
||||
started: threading.Event,
|
||||
finished: threading.Event,
|
||||
):
|
||||
"""Commit a native transition and its fork coordinate as one task."""
|
||||
started.set()
|
||||
try:
|
||||
task = asyncio.ensure_future(transition)
|
||||
try:
|
||||
result = await asyncio.shield(task)
|
||||
except asyncio.CancelledError:
|
||||
# BackgroundEventLoop cancels its submitted task when the
|
||||
# waiting caller is interrupted. Let an already-started native
|
||||
# transition reach its authoritative terminal state before the
|
||||
# per-table boundary is released.
|
||||
result = await task
|
||||
self._checkout_version = version
|
||||
raise
|
||||
self._checkout_version = version
|
||||
return result
|
||||
finally:
|
||||
finished.set()
|
||||
|
||||
def _run_native_state_transition(self, transition, version: Optional[int]):
|
||||
started = threading.Event()
|
||||
finished = threading.Event()
|
||||
try:
|
||||
return LOOP.run(
|
||||
self._commit_native_state(transition, version, started, finished)
|
||||
)
|
||||
except BaseException:
|
||||
if started.is_set():
|
||||
while not finished.is_set():
|
||||
try:
|
||||
finished.wait()
|
||||
except BaseException: # noqa: PERF203
|
||||
continue
|
||||
raise
|
||||
|
||||
def checkout(self, version: Union[int, str]):
|
||||
"""Checkout a version of the table. This is an in-place operation.
|
||||
@@ -2443,7 +2708,14 @@ class LanceTable(Table):
|
||||
vector type
|
||||
0 [1.1, 0.9] vector
|
||||
"""
|
||||
LOOP.run(self._table.checkout(version))
|
||||
# Resolve tags before mutating the native handle. This leaves the live
|
||||
# handle and reopen descriptor aligned if tag lookup fails, and avoids a
|
||||
# second fallible version lookup after checkout succeeds.
|
||||
with self._native_state_lock():
|
||||
resolved_version = self._resolve_checkout_version(version)
|
||||
self._run_native_state_transition(
|
||||
self._table.checkout(resolved_version), resolved_version
|
||||
)
|
||||
|
||||
def checkout_latest(self):
|
||||
"""Checkout the latest version of the table. This is an in-place operation.
|
||||
@@ -2451,7 +2723,8 @@ class LanceTable(Table):
|
||||
The table will be set back into standard mode, and will track the latest
|
||||
version of the table.
|
||||
"""
|
||||
LOOP.run(self._table.checkout_latest())
|
||||
with self._native_state_lock():
|
||||
self._run_native_state_transition(self._table.checkout_latest(), None)
|
||||
|
||||
def restore(self, version: Optional[Union[int, str]] = None):
|
||||
"""Restore a version of the table. This is an in-place operation.
|
||||
@@ -2497,9 +2770,13 @@ class LanceTable(Table):
|
||||
>>> len(table.list_versions())
|
||||
4
|
||||
"""
|
||||
if version is not None:
|
||||
LOOP.run(self._table.checkout(version))
|
||||
LOOP.run(self._table.restore())
|
||||
with self._native_state_lock():
|
||||
if version is not None:
|
||||
resolved_version = self._resolve_checkout_version(version)
|
||||
self._run_native_state_transition(
|
||||
self._table.checkout(resolved_version), resolved_version
|
||||
)
|
||||
self._run_native_state_transition(self._table.restore(), None)
|
||||
|
||||
def count_rows(self, filter: Optional[str] = None) -> int:
|
||||
return LOOP.run(self._table.count_rows(filter))
|
||||
@@ -2580,9 +2857,6 @@ class LanceTable(Table):
|
||||
2. Currently we've disabled push-down of the filters from polars
|
||||
because polars pushdown into pyarrow uses pyarrow compute
|
||||
expressions rather than SQl strings (which LanceDB supports)
|
||||
3. The Polars streaming engine is not supported because it does not
|
||||
currently implement Python PyArrow dataset scans. Use the default
|
||||
engine when collecting this LazyFrame.
|
||||
|
||||
Returns
|
||||
-------
|
||||
@@ -2591,12 +2865,8 @@ class LanceTable(Table):
|
||||
from lancedb.integrations.pyarrow import PyarrowDatasetAdapter
|
||||
|
||||
dataset = PyarrowDatasetAdapter(self)
|
||||
# Polars 1.32's non-PyArrow callback path passes batch_size twice. Keep
|
||||
# the compatible PyArrow path, but block predicates because this adapter
|
||||
# cannot translate PyArrow expressions into LanceDB filters.
|
||||
return pl.scan_pyarrow_dataset(dataset, batch_size=batch_size).map_batches(
|
||||
_polars_predicate_pushdown_barrier,
|
||||
predicate_pushdown=False,
|
||||
return pl.scan_pyarrow_dataset(
|
||||
dataset, allow_pyarrow_filter=False, batch_size=batch_size
|
||||
)
|
||||
|
||||
# New unified API overload
|
||||
@@ -3617,7 +3887,9 @@ class LanceTable(Table):
|
||||
self = cls.__new__(cls)
|
||||
self._conn = db
|
||||
self._namespace_path = namespace_path
|
||||
self._index_cache_size = None
|
||||
self._location = location
|
||||
self._managed_versioning = None
|
||||
self._namespace_client = namespace_client
|
||||
self._pushdown_operations = pushdown_operations or set()
|
||||
self._route_pushdown_to_rust = route_pushdown_to_rust
|
||||
@@ -3645,6 +3917,7 @@ class LanceTable(Table):
|
||||
enable_v2_manifest_paths
|
||||
)
|
||||
|
||||
self._storage_options = storage_options
|
||||
self._table = LOOP.run(
|
||||
self._conn._conn.create_table(
|
||||
name,
|
||||
@@ -3661,6 +3934,7 @@ class LanceTable(Table):
|
||||
namespace_client=namespace_client,
|
||||
)
|
||||
)
|
||||
self._initialize_reopen_state(name)
|
||||
return self
|
||||
|
||||
def delete(self, where: Union[str, Expr]) -> DeleteResult:
|
||||
@@ -3976,28 +4250,6 @@ class LanceTable(Table):
|
||||
[`AsyncTable.get_lsm_write_spec`][lancedb.AsyncTable.get_lsm_write_spec]."""
|
||||
return LOOP.run(self._table.get_lsm_write_spec())
|
||||
|
||||
def checkpoint_lsm(self) -> None:
|
||||
"""Synchronous version of
|
||||
[`AsyncTable.checkpoint_lsm`][lancedb.AsyncTable.checkpoint_lsm]."""
|
||||
return LOOP.run(self._table.checkpoint_lsm())
|
||||
|
||||
def flush_lsm(self) -> None:
|
||||
"""Synchronous version of
|
||||
[`AsyncTable.flush_lsm`][lancedb.AsyncTable.flush_lsm]."""
|
||||
return LOOP.run(self._table.flush_lsm())
|
||||
|
||||
def compact_lsm(self) -> None:
|
||||
"""Synchronous version of
|
||||
[`AsyncTable.compact_lsm`][lancedb.AsyncTable.compact_lsm]."""
|
||||
return LOOP.run(self._table.compact_lsm())
|
||||
|
||||
def get_lsm_stats(self, *, include_generation_rows: bool = False) -> Optional[dict]:
|
||||
"""Synchronous version of
|
||||
[`AsyncTable.get_lsm_stats`][lancedb.AsyncTable.get_lsm_stats]."""
|
||||
return LOOP.run(
|
||||
self._table.get_lsm_stats(include_generation_rows=include_generation_rows)
|
||||
)
|
||||
|
||||
def close_lsm_writers(self) -> None:
|
||||
"""Close cached MemWAL shard writers. See
|
||||
[`AsyncTable.close_lsm_writers`][lancedb.AsyncTable.close_lsm_writers]."""
|
||||
@@ -4676,13 +4928,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,73 +4954,12 @@ 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()
|
||||
|
||||
async def checkpoint_lsm(self) -> None:
|
||||
"""Converge this table's LSM write path into its base table.
|
||||
|
||||
One flush, sealing every memtable into L0, then compaction triggers
|
||||
until every generation that existed at that moment has reached base.
|
||||
The loop runs client-side, reading progress from ``get_lsm_stats``.
|
||||
|
||||
Best-effort: generations created *while* it runs are deliberately not
|
||||
waited on, which is what lets it terminate on a table taking writes.
|
||||
Idempotent and safe on a cadence.
|
||||
|
||||
There is no deadline, and the caller owns that. It returns when the
|
||||
target generations are gone, raises on a terminal server fault, and
|
||||
otherwise waits however long the server takes. A slow table and a
|
||||
stuck one are the same picture from the client: the compactor pool is
|
||||
shared across every table on the node, so a checkpoint queued behind
|
||||
unrelated work looks exactly like one that is merging. Wrap this in
|
||||
``asyncio.wait_for`` for a wall-clock bound; abandoning it partway
|
||||
costs nothing.
|
||||
"""
|
||||
return await self._inner.checkpoint_lsm()
|
||||
|
||||
async def flush_lsm(self) -> None:
|
||||
"""Seal every bucket's active memtable into L0.
|
||||
|
||||
Does not touch the base table — moving L0 into base is
|
||||
`compact_lsm`. On a node that has not claimed this table, this claims
|
||||
it and replays its WAL log first.
|
||||
"""
|
||||
return await self._inner.flush_lsm()
|
||||
|
||||
async def compact_lsm(self) -> None:
|
||||
"""Trigger a background L0 to base compaction pass per bucket.
|
||||
|
||||
Returns once the passes are dispatched, not once they finish: watch
|
||||
``get_lsm_stats`` for progress, or use ``checkpoint_lsm`` to loop
|
||||
until the current L0 has reached base.
|
||||
"""
|
||||
return await self._inner.compact_lsm()
|
||||
|
||||
async def get_lsm_stats(
|
||||
self, *, include_generation_rows: bool = False
|
||||
) -> Optional[dict]:
|
||||
"""Read live per-bucket LSM state.
|
||||
|
||||
Answers "how far behind is my fresh tier", "which bucket is hot", and
|
||||
"why is my fresh-tier vector search brute-force". Mutates no table
|
||||
state, though on a node that has not claimed this table it claims it,
|
||||
exactly as a read would.
|
||||
|
||||
Returns ``None`` only when the LSM write path is not enabled.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
include_generation_rows
|
||||
Report a row count per L0 generation. Off by default: each count
|
||||
opens an uncached Lance dataset, and ``checkpoint_lsm`` polls this
|
||||
needing only generation numbers.
|
||||
"""
|
||||
return await self._inner.get_lsm_stats(include_generation_rows)
|
||||
|
||||
async def close_lsm_writers(self) -> None:
|
||||
"""Drain and close any cached MemWAL shard writers for this table.
|
||||
|
||||
@@ -6341,9 +6525,7 @@ class TableStatistics:
|
||||
Attributes
|
||||
----------
|
||||
total_bytes: int
|
||||
The total size, in bytes, of the table's data files, index files, and
|
||||
overlay files. Read from the manifest, so this excludes deletion files
|
||||
and manifests.
|
||||
The total number of bytes in the table.
|
||||
num_rows: int
|
||||
The total number of rows in the table.
|
||||
num_indices: int
|
||||
|
||||
@@ -2,11 +2,11 @@
|
||||
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
|
||||
import json
|
||||
import inspect
|
||||
import re
|
||||
import sys
|
||||
from datetime import timedelta
|
||||
from importlib import resources
|
||||
import os
|
||||
from types import SimpleNamespace
|
||||
|
||||
@@ -19,10 +19,6 @@ from lance_namespace.errors import NamespaceNotEmptyError, TableNotFoundError
|
||||
from lancedb.pydantic import LanceModel, Vector
|
||||
|
||||
|
||||
def test_package_includes_pep_561_marker():
|
||||
assert resources.files(lancedb).joinpath("py.typed").is_file()
|
||||
|
||||
|
||||
def test_basic(tmp_path):
|
||||
db = lancedb.connect(tmp_path)
|
||||
|
||||
@@ -106,6 +102,17 @@ def test_read_consistency_interval_does_not_use_background_loop(tmp_path, monkey
|
||||
assert db_from_inner.read_consistency_interval == consistency_interval
|
||||
|
||||
|
||||
def test_serialize_preserves_zero_read_consistency_interval(tmp_path):
|
||||
db = lancedb.connect(tmp_path, read_consistency_interval=timedelta(0))
|
||||
table = db.create_table("items", pa.table({"x": [1]}))
|
||||
|
||||
encoded = json.loads(table._connection_state)
|
||||
assert encoded["read_consistency_interval_seconds"] == 0.0
|
||||
|
||||
restored = lancedb.deserialize_conn(table._connection_state)
|
||||
assert restored.read_consistency_interval == timedelta(0)
|
||||
|
||||
|
||||
def test_ingest_pd(tmp_path):
|
||||
db = lancedb.connect(tmp_path)
|
||||
|
||||
|
||||
@@ -631,23 +631,3 @@ def test_url_retrieve_downloads_image():
|
||||
image_bytes = url_retrieve(image_url)
|
||||
img = Image.open(io.BytesIO(image_bytes))
|
||||
assert img.size[0] > 0 and img.size[1] > 0
|
||||
|
||||
|
||||
def test_jina_generate_image_input_dict_local_path(tmp_path):
|
||||
"""
|
||||
JinaEmbeddings._generate_image_input_dict must accept a local image path
|
||||
(str or Path), not just bytes. Previously it crashed with
|
||||
`AttributeError: 'function' object has no attribute 'urlparse'` on any
|
||||
str/Path input because it called `urlparse.urlparse(image)` instead of
|
||||
`urlparse(image)` (urlparse was imported as a function, not a module).
|
||||
"""
|
||||
Image = pytest.importorskip("PIL.Image")
|
||||
from lancedb.embeddings.jinaai import JinaEmbeddings
|
||||
|
||||
image_path = tmp_path / "test.png"
|
||||
Image.new("RGB", (4, 4), color="red").save(image_path, format="PNG")
|
||||
|
||||
for image in (str(image_path), image_path):
|
||||
image_dict = JinaEmbeddings._generate_image_input_dict(image)
|
||||
assert "image" in image_dict
|
||||
assert isinstance(image_dict["image"], str) and len(image_dict["image"]) > 0
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -6,9 +6,13 @@
|
||||
import tempfile
|
||||
import shutil
|
||||
import importlib
|
||||
import multiprocessing as mp
|
||||
import sys
|
||||
from datetime import timedelta
|
||||
import pytest
|
||||
import pyarrow as pa
|
||||
import lancedb
|
||||
from lance_namespace import connect as namespace_connect
|
||||
from lance_namespace.errors import NamespaceNotEmptyError, TableNotFoundError
|
||||
from lancedb.namespace import _MAX_QUERY_K
|
||||
from lancedb.table import AsyncTable, LanceTable
|
||||
@@ -72,6 +76,16 @@ def _namespace_lance_table(namespace_client: _NamespaceClient) -> LanceTable:
|
||||
return table
|
||||
|
||||
|
||||
def _direct_namespace_fork_child(table, result_queue):
|
||||
from lancedb.permutation import Permutation
|
||||
|
||||
try:
|
||||
permutation = Permutation.identity(table)
|
||||
result_queue.put(("ok", permutation.num_rows))
|
||||
except Exception as error:
|
||||
result_queue.put((type(error).__name__, str(error)))
|
||||
|
||||
|
||||
class TestNamespaceConnection:
|
||||
"""Test namespace-based LanceDB connection using DirectoryNamespace."""
|
||||
|
||||
@@ -419,7 +433,95 @@ class TestNamespaceConnection:
|
||||
pa.field("vector", pa.list_(pa.float32(), 2)),
|
||||
]
|
||||
)
|
||||
db.create_table("test_table", schema=schema, storage_options=table_opts)
|
||||
created = db.create_table(
|
||||
"test_table", schema=schema, storage_options=table_opts
|
||||
)
|
||||
assert created._storage_options == table_opts
|
||||
|
||||
opened = db.open_table(
|
||||
"test_table",
|
||||
storage_options={"allow_http": "true"},
|
||||
index_cache_size=17,
|
||||
)
|
||||
assert opened._storage_options == {"allow_http": "true"}
|
||||
assert opened._index_cache_size == 17
|
||||
opened._pid = -1
|
||||
opened._ensure_open()
|
||||
assert opened.count_rows() == 0
|
||||
|
||||
def test_serialize_preserves_zero_read_consistency_interval(self):
|
||||
db = lancedb.connect_namespace(
|
||||
"dir",
|
||||
{"root": self.temp_dir},
|
||||
read_consistency_interval=timedelta(0),
|
||||
)
|
||||
|
||||
restored = lancedb.deserialize_conn(db.serialize())
|
||||
assert restored.read_consistency_interval == timedelta(0)
|
||||
|
||||
@pytest.mark.skipif(
|
||||
sys.platform != "linux",
|
||||
reason="fork() is only supported safely for this test on Linux",
|
||||
)
|
||||
def test_direct_namespace_with_descriptor_reopens_after_fork(self):
|
||||
properties = {"root": self.temp_dir}
|
||||
namespace = namespace_connect("dir", properties)
|
||||
db = lancedb.LanceNamespaceDBConnection(
|
||||
namespace,
|
||||
namespace_client_impl="dir",
|
||||
namespace_client_properties=properties,
|
||||
)
|
||||
table = db.create_table("items", pa.table({"id": [1]}))
|
||||
|
||||
ctx = mp.get_context("fork")
|
||||
result_queue = ctx.Queue()
|
||||
process = ctx.Process(
|
||||
target=_direct_namespace_fork_child,
|
||||
args=(table, result_queue),
|
||||
)
|
||||
process.start()
|
||||
process.join(10)
|
||||
|
||||
if process.is_alive():
|
||||
process.terminate()
|
||||
process.join(5)
|
||||
pytest.fail("Direct namespace table hung while reopening after fork")
|
||||
|
||||
assert process.exitcode == 0
|
||||
assert result_queue.get(timeout=2) == ("ok", 1)
|
||||
|
||||
@pytest.mark.skipif(
|
||||
sys.platform != "linux",
|
||||
reason="fork() is only supported safely for this test on Linux",
|
||||
)
|
||||
def test_opaque_direct_namespace_reports_unsupported_fork(self):
|
||||
namespace = namespace_connect("dir", {"root": self.temp_dir})
|
||||
db = lancedb.LanceNamespaceDBConnection(namespace)
|
||||
table = db.create_table("items", pa.table({"id": [1]}))
|
||||
|
||||
with pytest.raises(ValueError, match="opaque namespace client"):
|
||||
db.serialize()
|
||||
assert not table._can_reopen_after_fork
|
||||
|
||||
ctx = mp.get_context("fork")
|
||||
result_queue = ctx.Queue()
|
||||
process = ctx.Process(
|
||||
target=_direct_namespace_fork_child,
|
||||
args=(table, result_queue),
|
||||
)
|
||||
process.start()
|
||||
process.join(10)
|
||||
|
||||
if process.is_alive():
|
||||
process.terminate()
|
||||
process.join(5)
|
||||
pytest.fail("Opaque namespace table hung after fork")
|
||||
|
||||
assert process.exitcode == 0
|
||||
error_type, message = result_queue.get(timeout=2)
|
||||
assert error_type == "RuntimeError"
|
||||
assert "Cannot reopen table 'items' in a forked process" in message
|
||||
assert "namespace_client_impl and namespace_client_properties" in message
|
||||
|
||||
def test_namespace_operations(self):
|
||||
"""Test namespace management operations."""
|
||||
|
||||
@@ -415,17 +415,6 @@ def test_nullable_vector():
|
||||
assert schema == pa.schema([pa.field("vec", pa.list_(pa.float32(), 16), True)])
|
||||
|
||||
|
||||
def test_bare_vector_raises_clear_error():
|
||||
namespace = {
|
||||
"__name__": "test_model_without_pyarrow",
|
||||
"LanceModel": LanceModel,
|
||||
"Vector": Vector,
|
||||
}
|
||||
|
||||
with pytest.raises(TypeError, match=r"Vector must be parameterized.*Vector\(128\)"):
|
||||
exec("class TestModel(LanceModel):\n vector: Vector", namespace)
|
||||
|
||||
|
||||
def test_fixed_size_list_field():
|
||||
class TestModel(pydantic.BaseModel):
|
||||
vec: Vector(16)
|
||||
|
||||
@@ -570,15 +570,6 @@ def test_query_builder(table):
|
||||
assert all(np.array(rs[0]["vector"]) == [1, 2])
|
||||
|
||||
|
||||
def test_query_multiple_vectors(table):
|
||||
results = table.search([np.array([1, 2]), np.array([4, 5])]).limit(1).to_list()
|
||||
|
||||
assert len(results) == 2
|
||||
results_by_query = {result["query_index"]: result for result in results}
|
||||
assert results_by_query[0]["id"] == 1
|
||||
assert results_by_query[1]["id"] == 2
|
||||
|
||||
|
||||
def test_with_row_id(table: lancedb.table.Table):
|
||||
rs = table.search().with_row_id(True).to_arrow()
|
||||
assert "_rowid" in rs.column_names
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
|
||||
import asyncio
|
||||
import ctypes
|
||||
import gc
|
||||
import os
|
||||
@@ -9,7 +10,7 @@ import sys
|
||||
import threading
|
||||
import warnings
|
||||
import weakref
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from concurrent.futures import CancelledError, ThreadPoolExecutor
|
||||
from datetime import date, datetime, timedelta
|
||||
from time import sleep
|
||||
from typing import List
|
||||
@@ -929,7 +930,6 @@ def test_polars(mem_db: DBConnection):
|
||||
|
||||
# enter table to polars dataframe
|
||||
result = table.to_polars()
|
||||
assert isinstance(result, pl.LazyFrame)
|
||||
assert np.allclose(result.collect()["vector"].to_list(), data["vector"])
|
||||
|
||||
# make sure filtering isn't broken
|
||||
@@ -1846,27 +1846,6 @@ def test_add_with_empty_fixed_size_list_drops_bad_rows(mem_db: DBConnection):
|
||||
assert np.allclose(data["embedding"].to_pylist()[0], np.array([0.1] * 16))
|
||||
|
||||
|
||||
def test_add_nullable_fixed_size_list_with_none(mem_db: DBConnection):
|
||||
"""Regression test for issue #2340."""
|
||||
table = mem_db.create_table(
|
||||
"test_nullable_fixed_size_list",
|
||||
schema=pa.schema(
|
||||
[
|
||||
pa.field("id", pa.string()),
|
||||
pa.field("feature", pa.list_(pa.float32(), 256)),
|
||||
pa.field("tags", pa.list_(pa.string())),
|
||||
]
|
||||
),
|
||||
)
|
||||
|
||||
table.add([{"id": "1", "feature": None, "tags": ["tag1", "tag2"]}])
|
||||
|
||||
result = table.to_arrow()
|
||||
assert result.to_pylist() == [
|
||||
{"id": "1", "feature": None, "tags": ["tag1", "tag2"]}
|
||||
]
|
||||
|
||||
|
||||
def test_add_nullable_struct_with_none(mem_db: DBConnection):
|
||||
"""Regression test for issue #2654: a nullable struct column whose
|
||||
first batch contains only None values must not crash in
|
||||
@@ -2167,6 +2146,260 @@ def test_restore(mem_db: DBConnection):
|
||||
table.restore(0)
|
||||
|
||||
|
||||
def test_restore_tracks_checkout_when_restore_fails():
|
||||
class FailingRestore:
|
||||
def __init__(self):
|
||||
self.live_version = None
|
||||
|
||||
async def checkout(self, version):
|
||||
self.live_version = version
|
||||
|
||||
async def restore(self):
|
||||
raise RuntimeError("injected restore failure")
|
||||
|
||||
inner = FailingRestore()
|
||||
table = LanceTable.__new__(LanceTable)
|
||||
table._table = inner
|
||||
table._checkout_version = None
|
||||
|
||||
with pytest.raises(RuntimeError, match="injected restore failure"):
|
||||
table.restore(7)
|
||||
|
||||
assert table._checkout_version == inner.live_version
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("operation", "expected_descriptor", "expected_restore_calls"),
|
||||
[("checkout", 11, 0), ("restore", None, 1)],
|
||||
)
|
||||
def test_string_tag_resolves_before_checkout(
|
||||
operation, expected_descriptor, expected_restore_calls
|
||||
):
|
||||
class Tags:
|
||||
async def get_version(self, tag):
|
||||
assert tag == "tag-v1"
|
||||
return 11
|
||||
|
||||
class NoPostCheckoutVersionLookup:
|
||||
def __init__(self):
|
||||
self.tags = Tags()
|
||||
self.checkout_versions = []
|
||||
self.restore_calls = 0
|
||||
|
||||
async def checkout(self, version):
|
||||
self.checkout_versions.append(version)
|
||||
|
||||
async def version(self):
|
||||
raise RuntimeError("post-checkout version lookup must not run")
|
||||
|
||||
async def restore(self):
|
||||
self.restore_calls += 1
|
||||
|
||||
inner = NoPostCheckoutVersionLookup()
|
||||
table = LanceTable.__new__(LanceTable)
|
||||
table._table = inner
|
||||
table._checkout_version = 3
|
||||
|
||||
getattr(table, operation)("tag-v1")
|
||||
|
||||
assert inner.checkout_versions == [11]
|
||||
assert table._checkout_version == expected_descriptor
|
||||
assert inner.restore_calls == expected_restore_calls
|
||||
|
||||
|
||||
@pytest.mark.parametrize("operation", ["checkout", "restore"])
|
||||
def test_string_tag_resolution_failure_does_not_mutate_handle(operation):
|
||||
class FailingTags:
|
||||
async def get_version(self, tag):
|
||||
assert tag == "missing-tag"
|
||||
raise RuntimeError("injected tag lookup failure")
|
||||
|
||||
class UnchangedTable:
|
||||
def __init__(self):
|
||||
self.tags = FailingTags()
|
||||
self.checkout_calls = 0
|
||||
self.restore_calls = 0
|
||||
|
||||
async def checkout(self, version):
|
||||
self.checkout_calls += 1
|
||||
|
||||
async def restore(self):
|
||||
self.restore_calls += 1
|
||||
|
||||
inner = UnchangedTable()
|
||||
table = LanceTable.__new__(LanceTable)
|
||||
table._table = inner
|
||||
table._checkout_version = 3
|
||||
|
||||
with pytest.raises(RuntimeError, match="injected tag lookup failure"):
|
||||
getattr(table, operation)("missing-tag")
|
||||
|
||||
assert table._checkout_version == 3
|
||||
assert inner.checkout_calls == 0
|
||||
assert inner.restore_calls == 0
|
||||
|
||||
|
||||
def test_native_state_transitions_are_serialized(monkeypatch):
|
||||
from lancedb.background_loop import LOOP
|
||||
|
||||
class Inner:
|
||||
def __init__(self):
|
||||
self.live_version = None
|
||||
|
||||
async def checkout(self, version):
|
||||
self.live_version = version
|
||||
|
||||
inner = Inner()
|
||||
table = LanceTable.__new__(LanceTable)
|
||||
table._table = inner
|
||||
table._checkout_version = None
|
||||
|
||||
first_native_done = threading.Event()
|
||||
release_first_call = threading.Event()
|
||||
second_call_started = threading.Event()
|
||||
second_call_done = threading.Event()
|
||||
errors = []
|
||||
original_run = LOOP.run
|
||||
|
||||
def delayed_delivery(awaitable):
|
||||
result = original_run(awaitable)
|
||||
if threading.current_thread().name == "checkout-1":
|
||||
first_native_done.set()
|
||||
assert release_first_call.wait(5)
|
||||
return result
|
||||
|
||||
monkeypatch.setattr(LOOP, "run", delayed_delivery)
|
||||
|
||||
def checkout(version):
|
||||
if version == 2:
|
||||
second_call_started.set()
|
||||
try:
|
||||
table.checkout(version)
|
||||
except BaseException as err:
|
||||
errors.append(err)
|
||||
finally:
|
||||
if version == 2:
|
||||
second_call_done.set()
|
||||
|
||||
first = threading.Thread(target=checkout, args=(1,), name="checkout-1")
|
||||
first.start()
|
||||
assert first_native_done.wait(5)
|
||||
|
||||
second = threading.Thread(target=checkout, args=(2,), name="checkout-2")
|
||||
second.start()
|
||||
assert second_call_started.wait(5)
|
||||
assert not second_call_done.wait(0.1)
|
||||
|
||||
release_first_call.set()
|
||||
first.join(5)
|
||||
second.join(5)
|
||||
|
||||
assert not first.is_alive()
|
||||
assert not second.is_alive()
|
||||
assert errors == []
|
||||
assert inner.live_version == 2
|
||||
assert table._checkout_version == 2
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("operation", "args", "initial_version", "expected_version"),
|
||||
[
|
||||
("checkout", (11,), 3, 11),
|
||||
("checkout_latest", (), 3, None),
|
||||
("restore", (), 11, None),
|
||||
],
|
||||
)
|
||||
def test_native_state_commits_before_success_delivery(
|
||||
monkeypatch, operation, args, initial_version, expected_version
|
||||
):
|
||||
from lancedb.background_loop import LOOP
|
||||
|
||||
class Inner:
|
||||
def __init__(self, live_version):
|
||||
self.live_version = live_version
|
||||
|
||||
async def checkout(self, version):
|
||||
self.live_version = version
|
||||
|
||||
async def checkout_latest(self):
|
||||
self.live_version = None
|
||||
|
||||
async def restore(self):
|
||||
self.live_version = None
|
||||
|
||||
inner = Inner(initial_version)
|
||||
table = LanceTable.__new__(LanceTable)
|
||||
table._table = inner
|
||||
table._checkout_version = initial_version
|
||||
|
||||
original_run = LOOP.run
|
||||
|
||||
def success_then_interrupt(awaitable):
|
||||
original_run(awaitable)
|
||||
raise KeyboardInterrupt("injected after native success")
|
||||
|
||||
monkeypatch.setattr(LOOP, "run", success_then_interrupt)
|
||||
|
||||
with pytest.raises(KeyboardInterrupt, match="injected after native success"):
|
||||
getattr(table, operation)(*args)
|
||||
|
||||
assert inner.live_version == expected_version
|
||||
assert table._checkout_version == expected_version
|
||||
|
||||
|
||||
def test_native_state_waits_for_cancelled_delivery(monkeypatch):
|
||||
from lancedb.background_loop import LOOP
|
||||
|
||||
class Inner:
|
||||
def __init__(self):
|
||||
self.live_version = 3
|
||||
|
||||
async def checkout(self, version):
|
||||
await asyncio.sleep(0.01)
|
||||
self.live_version = version
|
||||
|
||||
inner = Inner()
|
||||
table = LanceTable.__new__(LanceTable)
|
||||
table._table = inner
|
||||
table._checkout_version = 3
|
||||
|
||||
original_run = LOOP.run
|
||||
|
||||
def cancel_while_running(awaitable):
|
||||
async def cancel_after_start():
|
||||
task = asyncio.create_task(awaitable)
|
||||
await asyncio.sleep(0)
|
||||
task.cancel()
|
||||
return await task
|
||||
|
||||
return original_run(cancel_after_start())
|
||||
|
||||
monkeypatch.setattr(LOOP, "run", cancel_while_running)
|
||||
|
||||
with pytest.raises(CancelledError):
|
||||
table.checkout(11)
|
||||
|
||||
assert inner.live_version == 11
|
||||
assert table._checkout_version == 11
|
||||
|
||||
|
||||
def test_reopen_preserves_explicit_table_location(tmp_path):
|
||||
db = lancedb.connect(tmp_path / "db")
|
||||
location = str(tmp_path / "physical-table")
|
||||
table = LanceTable.create(
|
||||
db,
|
||||
"items",
|
||||
pa.table({"x": [1]}),
|
||||
location=location,
|
||||
)
|
||||
|
||||
table._pid = -1
|
||||
table._ensure_open()
|
||||
|
||||
assert table.count_rows() == 1
|
||||
assert table._location == location
|
||||
|
||||
|
||||
def test_restore_with_tags(mem_db: DBConnection):
|
||||
table = mem_db.create_table(
|
||||
"my_table",
|
||||
@@ -2218,45 +2451,6 @@ def test_merge(tmp_db: DBConnection, tmp_path):
|
||||
table.merge(other_dataset, left_on="id")
|
||||
|
||||
|
||||
@pytest.mark.parametrize("storage_version", ["legacy", "stable"])
|
||||
def test_search_after_merge(tmp_path, storage_version):
|
||||
pytest.importorskip("lance")
|
||||
pd = pytest.importorskip("pandas")
|
||||
|
||||
db = lancedb.connect(
|
||||
tmp_path,
|
||||
storage_options={"new_table_data_storage_version": storage_version},
|
||||
)
|
||||
rng = np.random.default_rng(42)
|
||||
row_count = 512
|
||||
vectors = rng.standard_normal((row_count, 8)).astype(np.float32)
|
||||
table = db.create_table(
|
||||
"search_after_merge",
|
||||
data=pd.DataFrame(
|
||||
{
|
||||
"id": [str(i) for i in range(row_count)],
|
||||
"vector": list(vectors),
|
||||
}
|
||||
),
|
||||
)
|
||||
table.create_index("vector", config=IvfPq(num_partitions=1, num_sub_vectors=2))
|
||||
|
||||
links = pd.DataFrame(
|
||||
{
|
||||
"id": [str(i) for i in range(row_count // 2)],
|
||||
"link": [f"https://example.com/{i}" for i in range(row_count // 2)],
|
||||
}
|
||||
)
|
||||
table.merge(links, left_on="id")
|
||||
|
||||
query = table.search(vectors[-1]).refine_factor(50).limit(10)
|
||||
assert "ANN" in query.explain_plan(verbose=True)
|
||||
|
||||
result = query.to_arrow()
|
||||
links_by_id = dict(zip(result["id"].to_pylist(), result["link"].to_pylist()))
|
||||
assert links_by_id[str(row_count - 1)] is None
|
||||
|
||||
|
||||
def test_delete(mem_db: DBConnection):
|
||||
table = mem_db.create_table(
|
||||
"my_table",
|
||||
@@ -2799,40 +2993,15 @@ def test_create_with_embedding_function(mem_db: DBConnection):
|
||||
assert actual == expected
|
||||
|
||||
|
||||
def test_create_f16_table_from_arrow_data(mem_db: DBConnection):
|
||||
dimension = 32
|
||||
num_rows = 512
|
||||
values = pa.array(
|
||||
np.random.default_rng(42)
|
||||
.standard_normal(num_rows * dimension)
|
||||
.astype(np.float16)
|
||||
)
|
||||
df = pa.table(
|
||||
{
|
||||
"text": [f"s-{i}" for i in range(num_rows)],
|
||||
"vector": pa.FixedSizeListArray.from_arrays(values, dimension),
|
||||
}
|
||||
)
|
||||
table = mem_db.create_table("f16_tbl", data=df)
|
||||
assert table.schema.field("vector").type == pa.list_(pa.float16(), dimension)
|
||||
table.create_index(num_partitions=2, num_sub_vectors=2)
|
||||
|
||||
query = df["vector"][2].as_py()
|
||||
expected = table.search(query).limit(2).to_arrow()
|
||||
|
||||
assert "s-2" in expected["text"].to_pylist()
|
||||
|
||||
|
||||
def test_create_f16_table(mem_db: DBConnection):
|
||||
class MyTable(LanceModel):
|
||||
text: str
|
||||
vector: Vector(32, value_type=pa.float16())
|
||||
|
||||
rng = np.random.default_rng(42)
|
||||
df = pa.table(
|
||||
{
|
||||
"text": [f"s-{i}" for i in range(512)],
|
||||
"vector": [rng.standard_normal(32).astype(np.float16) for _ in range(512)],
|
||||
"vector": [np.random.randn(32).astype(np.float16) for _ in range(512)],
|
||||
}
|
||||
)
|
||||
table = mem_db.create_table(
|
||||
@@ -3713,8 +3882,7 @@ def test_stats(mem_db: DBConnection):
|
||||
stats = table.stats()
|
||||
print(f"{stats=}")
|
||||
assert stats == {
|
||||
# Full on-disk size of the data file, footer and metadata included.
|
||||
"total_bytes": 633,
|
||||
"total_bytes": 60,
|
||||
"num_rows": 2,
|
||||
"num_indices": 0,
|
||||
"fragment_stats": {
|
||||
@@ -3732,13 +3900,6 @@ def test_stats(mem_db: DBConnection):
|
||||
},
|
||||
}
|
||||
|
||||
# Index files count toward total_bytes too (only deletion files and
|
||||
# manifests are excluded).
|
||||
table.create_index("id", config=BTree())
|
||||
stats_with_index = table.stats()
|
||||
assert stats_with_index["num_indices"] == 1
|
||||
assert stats_with_index["total_bytes"] > stats["total_bytes"]
|
||||
|
||||
|
||||
def test_create_table_empty_list_with_schema(mem_db: DBConnection):
|
||||
"""Test creating table with empty list data and schema
|
||||
|
||||
@@ -342,6 +342,42 @@ def _multiworker_dataloader_target(db_uri: str, result_queue):
|
||||
result_queue.put(count)
|
||||
|
||||
|
||||
class _LazyPermutationDataset(torch.utils.data.Dataset):
|
||||
"""Match applications that create their Permutation inside a fork worker."""
|
||||
|
||||
def __init__(self, table):
|
||||
self._table = table
|
||||
self._permutation = None
|
||||
self._length = table.count_rows()
|
||||
|
||||
def __len__(self):
|
||||
return self._length
|
||||
|
||||
def __getitems__(self, indices):
|
||||
if self._permutation is None:
|
||||
inherited_connection = self._table._conn
|
||||
self._permutation = Permutation.identity(self._table)
|
||||
if self._table._conn is inherited_connection:
|
||||
raise RuntimeError("Permutation reused a connection inherited by fork")
|
||||
return self._permutation.__getitems__(indices)
|
||||
|
||||
|
||||
def _lazy_multiworker_dataloader_target(db_uri: str, result_queue):
|
||||
table = lancedb.connect(db_uri).open_table("test_table")
|
||||
dataset = _LazyPermutationDataset(table)
|
||||
dataloader = torch.utils.data.DataLoader(
|
||||
dataset,
|
||||
batch_size=10,
|
||||
num_workers=2,
|
||||
multiprocessing_context="fork",
|
||||
)
|
||||
count = 0
|
||||
for batch in dataloader:
|
||||
assert batch["a"].size(0) == 10
|
||||
count += 1
|
||||
result_queue.put(count)
|
||||
|
||||
|
||||
def _remote_multiworker_dataloader_target(port: int, result_queue):
|
||||
import lancedb
|
||||
from lancedb.permutation import Permutation
|
||||
@@ -410,6 +446,46 @@ def test_permutation_dataloader_fork_workers(tmp_path):
|
||||
assert queue.get() == 100
|
||||
|
||||
|
||||
@pytest.mark.skipif(
|
||||
sys.platform != "linux",
|
||||
reason=(
|
||||
"fork() is unavailable on Windows and unsafe on macOS "
|
||||
"(Apple frameworks/TLS are not fork-safe)"
|
||||
),
|
||||
)
|
||||
def test_lazy_permutation_reopens_inherited_table_in_fork_worker(tmp_path):
|
||||
"""A lazily built Permutation must not reuse an inherited table client.
|
||||
|
||||
Object-store table handles contain HTTP connection pools that are unsafe
|
||||
after fork. The local table makes the handle replacement deterministic
|
||||
without requiring an S3 service in the unit-test environment.
|
||||
"""
|
||||
db_uri = str(tmp_path / "db")
|
||||
db = lancedb.connect(db_uri)
|
||||
db.create_table("test_table", pa.table({"a": list(range(1000))}))
|
||||
|
||||
ctx = mp.get_context("spawn")
|
||||
queue = ctx.Queue()
|
||||
proc = ctx.Process(
|
||||
target=_lazy_multiworker_dataloader_target,
|
||||
args=(db_uri, queue),
|
||||
)
|
||||
proc.start()
|
||||
proc.join(timeout=30)
|
||||
|
||||
if proc.is_alive():
|
||||
proc.terminate()
|
||||
proc.join(timeout=5)
|
||||
if proc.is_alive():
|
||||
proc.kill()
|
||||
proc.join()
|
||||
pytest.fail("Lazy Permutation hung in a fork-based DataLoader worker")
|
||||
|
||||
assert proc.exitcode == 0, f"child exited with code {proc.exitcode}"
|
||||
assert not queue.empty(), "child produced no batches"
|
||||
assert queue.get() == 100
|
||||
|
||||
|
||||
@pytest.mark.skipif(
|
||||
sys.platform != "linux",
|
||||
reason=(
|
||||
|
||||
@@ -1,15 +0,0 @@
|
||||
# SPDX-License-Identifier: Apache-2.0
|
||||
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
from typing import assert_type
|
||||
|
||||
import lancedb
|
||||
from lancedb import AsyncConnection, DBConnection
|
||||
|
||||
|
||||
def check_connect_type() -> None:
|
||||
assert_type(lancedb.connect("memory://"), DBConnection)
|
||||
|
||||
|
||||
async def check_connect_async_type() -> None:
|
||||
assert_type(await lancedb.connect_async("memory://"), AsyncConnection)
|
||||
+16
-138
@@ -28,72 +28,11 @@ use pyo3::{
|
||||
Bound, FromPyObject, Py, PyAny, PyRef, PyResult, Python,
|
||||
exceptions::{PyRuntimeError, PyValueError},
|
||||
pyclass, pyfunction, pymethods,
|
||||
types::{IntoPyDict, PyAnyMethods, PyBytes, PyDict, PyDictMethods, PyList, PyListMethods},
|
||||
types::{IntoPyDict, PyAnyMethods, PyBytes, PyDict, PyDictMethods},
|
||||
};
|
||||
|
||||
mod scannable;
|
||||
|
||||
/// Convert `LsmStats` to a Python dict, preserving the per-bucket list.
|
||||
///
|
||||
/// Deliberately not flattened to a table-level summary: a table is N
|
||||
/// buckets on one node, and the per-bucket detail is the reason the
|
||||
/// endpoint exists — flattening hides the single hot bucket someone opened
|
||||
/// it to find.
|
||||
fn lsm_stats_to_py(py: Python<'_>, stats: &lancedb::table::LsmStats) -> PyResult<Py<PyDict>> {
|
||||
let out = PyDict::new(py);
|
||||
let buckets = PyList::empty(py);
|
||||
for b in &stats.buckets {
|
||||
let e = PyDict::new(py);
|
||||
e.set_item("shard_id", &b.shard_id)?;
|
||||
e.set_item("status", &b.status)?;
|
||||
e.set_item("writer_epoch", b.writer_epoch)?;
|
||||
e.set_item("manifest_version", b.manifest_version)?;
|
||||
e.set_item("current_generation", b.current_generation)?;
|
||||
e.set_item(
|
||||
"replay_after_wal_entry_position",
|
||||
b.replay_after_wal_entry_position,
|
||||
)?;
|
||||
e.set_item(
|
||||
"wal_entry_position_last_seen",
|
||||
b.wal_entry_position_last_seen,
|
||||
)?;
|
||||
|
||||
let generations = PyList::empty(py);
|
||||
for g in &b.generations {
|
||||
let ge = PyDict::new(py);
|
||||
ge.set_item("generation", g.generation)?;
|
||||
ge.set_item("bytes", g.bytes)?;
|
||||
ge.set_item("rows", g.rows)?;
|
||||
generations.append(ge)?;
|
||||
}
|
||||
e.set_item("generations", generations)?;
|
||||
e.set_item("compacting", b.compacting)?;
|
||||
|
||||
e.set_item(
|
||||
"memtables",
|
||||
b.memtables
|
||||
.as_ref()
|
||||
.map(|ms| {
|
||||
let l = PyList::empty(py);
|
||||
for m in ms {
|
||||
let d = PyDict::new(py);
|
||||
d.set_item("generation", m.generation)?;
|
||||
d.set_item("rows", m.rows)?;
|
||||
d.set_item("bytes", m.bytes)?;
|
||||
d.set_item("batches", m.batches)?;
|
||||
d.set_item("indexes", m.indexes.clone())?;
|
||||
l.append(d)?;
|
||||
}
|
||||
PyResult::Ok(l.unbind())
|
||||
})
|
||||
.transpose()?,
|
||||
)?;
|
||||
buckets.append(e)?;
|
||||
}
|
||||
out.set_item("buckets", buckets)?;
|
||||
Ok(out.unbind())
|
||||
}
|
||||
|
||||
#[derive(FromPyObject)]
|
||||
enum PredicateArg {
|
||||
Expr(PyExpr),
|
||||
@@ -246,22 +185,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 +230,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 +256,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 +307,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.
|
||||
@@ -1416,51 +1339,6 @@ impl Table {
|
||||
})
|
||||
}
|
||||
|
||||
/// Converge the table's LSM write path into its base table.
|
||||
///
|
||||
/// Best-effort: with writes flowing, new rows may land after the last
|
||||
/// pass. Errors if the table stops making progress.
|
||||
pub fn checkpoint_lsm(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
inner.checkpoint_lsm().await.infer_error()
|
||||
})
|
||||
}
|
||||
|
||||
/// Seal every bucket's active memtable into L0.
|
||||
pub fn flush_lsm(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(
|
||||
self_.py(),
|
||||
async move { inner.flush_lsm().await.infer_error() },
|
||||
)
|
||||
}
|
||||
|
||||
/// Trigger a background L0 → base pass per bucket. Returns once the
|
||||
/// passes are dispatched, not once they finish — watch `get_lsm_stats`.
|
||||
pub fn compact_lsm(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
inner.compact_lsm().await.infer_error()
|
||||
})
|
||||
}
|
||||
|
||||
/// Live LSM state, or `None` when the LSM write path is not enabled.
|
||||
#[pyo3(signature = (include_generation_rows=false))]
|
||||
pub fn get_lsm_stats(
|
||||
self_: PyRef<'_, Self>,
|
||||
include_generation_rows: bool,
|
||||
) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
let stats = inner
|
||||
.get_lsm_stats(include_generation_rows)
|
||||
.await
|
||||
.infer_error()?;
|
||||
Python::attach(|py| stats.map(|s| lsm_stats_to_py(py, &s)).transpose())
|
||||
})
|
||||
}
|
||||
|
||||
pub fn close_lsm_writers(self_: PyRef<'_, Self>) -> PyResult<Bound<'_, PyAny>> {
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
|
||||
Generated
+1
-1
@@ -1998,7 +1998,7 @@ requires-dist = [
|
||||
{ name = "pillow", marker = "extra == 'clip'", specifier = ">=12.1.1" },
|
||||
{ name = "pillow", marker = "extra == 'embeddings'", specifier = ">=12.1.1" },
|
||||
{ name = "pillow", marker = "extra == 'siglip'", specifier = ">=12.1.1" },
|
||||
{ name = "polars", marker = "extra == 'tests'", specifier = ">=0.19,<=1.32.3" },
|
||||
{ name = "polars", marker = "extra == 'tests'", specifier = ">=0.19,<=1.3.0" },
|
||||
{ name = "pre-commit", marker = "extra == 'dev'", specifier = ">=3.5.0" },
|
||||
{ name = "pyarrow", specifier = ">=16" },
|
||||
{ name = "pyarrow", marker = "extra == 'tests'", specifier = "<25" },
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
[package]
|
||||
name = "lancedb"
|
||||
version = "0.37.1-beta.1"
|
||||
version = "0.37.1-beta.0"
|
||||
edition.workspace = true
|
||||
description = "LanceDB: A serverless, low-latency vector database for AI applications"
|
||||
license.workspace = true
|
||||
@@ -49,8 +49,8 @@ lance-namespace = { workspace = true }
|
||||
lance-namespace-impls = { workspace = true }
|
||||
metrics = { workspace = true, optional = true }
|
||||
metrics-util = { workspace = true, optional = true }
|
||||
# Pin the GooseFS SDK to the version required by Lance's OpenDAL dependency.
|
||||
goosefs-sdk = { version = "=0.1.9", optional = true }
|
||||
# Pin the transitive GooseFS SDK until the 0.1.6 compile break is fixed upstream.
|
||||
goosefs-sdk = { version = "=0.1.5", optional = true }
|
||||
moka = { workspace = true }
|
||||
pin-project = { workspace = true }
|
||||
tokio = { version = "1.23", features = ["rt-multi-thread", "sync"] }
|
||||
@@ -100,8 +100,7 @@ anyhow = "1"
|
||||
lance-testing = { workspace = true }
|
||||
tempfile = "3.5.0"
|
||||
random_word = { version = "0.4.3", features = ["en"] }
|
||||
roaring = "0.11.4"
|
||||
tokio = { version = "1.23", features = ["io-util", "macros", "net", "rt-multi-thread", "sync", "test-util"] }
|
||||
tokio = { version = "1.23", features = ["io-util", "macros", "net", "rt-multi-thread", "sync"] }
|
||||
uuid = { version = "1.7.0", features = ["v4"] }
|
||||
walkdir = "2"
|
||||
aws-sdk-dynamodb = { version = "1.55.0" }
|
||||
|
||||
@@ -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::LanceFileVersion;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lance_io::object_store::ObjectStore;
|
||||
use object_store::path::Path;
|
||||
|
||||
|
||||
@@ -34,7 +34,7 @@ use crate::remote::{
|
||||
db::{OPT_REMOTE_API_KEY, OPT_REMOTE_HOST_OVERRIDE, OPT_REMOTE_REGION},
|
||||
};
|
||||
use lance::io::ObjectStoreParams;
|
||||
pub use lance_file::version::LanceFileVersion;
|
||||
pub use lance_encoding::version::LanceFileVersion;
|
||||
#[cfg(feature = "remote")]
|
||||
use lance_io::object_store::StorageOptions;
|
||||
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
|
||||
|
||||
@@ -12,7 +12,7 @@ use lance::dataset::refs::Ref;
|
||||
use lance::dataset::{ReadParams, WriteMode, builder::DatasetBuilder};
|
||||
use lance::io::{ObjectStore, ObjectStoreParams, WrappingObjectStore};
|
||||
use lance_datafusion::utils::StreamingWriteSource;
|
||||
use lance_file::version::LanceFileVersion;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
|
||||
use lance_table::io::commit::commit_handler_from_url;
|
||||
use object_store::local::LocalFileSystem;
|
||||
@@ -1294,11 +1294,9 @@ mod tests {
|
||||
use crate::connection::ConnectRequest;
|
||||
use crate::data::scannable::Scannable;
|
||||
use crate::database::{CreateTableMode, CreateTableRequest};
|
||||
use crate::query::QueryRequest;
|
||||
use crate::table::{AnyQuery, WriteOptions};
|
||||
use crate::table::WriteOptions;
|
||||
use arrow_array::{Int32Array, RecordBatch, StringArray};
|
||||
use arrow_schema::{DataType, Field, Schema};
|
||||
use futures::TryStreamExt;
|
||||
use std::path::PathBuf;
|
||||
use tempfile::tempdir;
|
||||
|
||||
@@ -1440,94 +1438,6 @@ mod tests {
|
||||
assert!(after_open.hits >= before_open.hits + 3);
|
||||
}
|
||||
|
||||
/// Regression test for https://github.com/lancedb/lancedb/issues/3197.
|
||||
#[cfg(unix)]
|
||||
#[tokio::test]
|
||||
async fn test_open_table_follows_hugging_face_symlinks() {
|
||||
let (tempdir, db) = setup_database().await;
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
|
||||
db.create_table(CreateTableRequest {
|
||||
name: "test".to_string(),
|
||||
namespace_path: vec![],
|
||||
data: Box::new(
|
||||
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from(vec![1, 2, 3]))])
|
||||
.unwrap(),
|
||||
) as Box<dyn Scannable>,
|
||||
mode: CreateTableMode::Create,
|
||||
write_options: Default::default(),
|
||||
location: None,
|
||||
namespace_client: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let table_dir = tempdir.path().join("test.lance");
|
||||
let versions_dir = table_dir.join("_versions");
|
||||
let manifest_path = std::fs::read_dir(&versions_dir)
|
||||
.unwrap()
|
||||
.map(|entry| entry.unwrap().path())
|
||||
.find(|path| path.extension().is_some_and(|ext| ext == "manifest"))
|
||||
.unwrap();
|
||||
let data_path = std::fs::read_dir(table_dir.join("data"))
|
||||
.unwrap()
|
||||
.map(|entry| entry.unwrap().path())
|
||||
.find(|path| path.extension().is_some_and(|ext| ext == "lance"))
|
||||
.unwrap();
|
||||
|
||||
// Hugging Face snapshots keep dataset objects in a separate blob directory and
|
||||
// expose them through relative symlinks.
|
||||
let blobs_dir = tempdir.path().join("blobs");
|
||||
std::fs::create_dir(&blobs_dir).unwrap();
|
||||
let manifest_blob = "9b603c63d0e692e05d58be25605f2f2064cc781e5ff94fe983a405059547b816";
|
||||
let data_blob = "be64f20e5723bd0a27cfdbdb41cf7d6fad94cd572a71973b717fb8340f4310c5";
|
||||
std::fs::rename(&manifest_path, blobs_dir.join(manifest_blob)).unwrap();
|
||||
std::fs::rename(&data_path, blobs_dir.join(data_blob)).unwrap();
|
||||
std::os::unix::fs::symlink(Path::new("../../blobs").join(manifest_blob), &manifest_path)
|
||||
.unwrap();
|
||||
std::os::unix::fs::symlink(Path::new("../../blobs").join(data_blob), &data_path).unwrap();
|
||||
let symlink_len = std::fs::symlink_metadata(&manifest_path).unwrap().len();
|
||||
let target_len = std::fs::metadata(&manifest_path).unwrap().len();
|
||||
assert_ne!(symlink_len, target_len);
|
||||
|
||||
drop(db);
|
||||
let db = ListingDatabase::connect_with_options(&ConnectRequest {
|
||||
uri: tempdir.path().to_str().unwrap().to_string(),
|
||||
#[cfg(feature = "remote")]
|
||||
client_config: Default::default(),
|
||||
options: Default::default(),
|
||||
namespace_client_properties: Default::default(),
|
||||
manifest_enabled: false,
|
||||
read_consistency_interval: None,
|
||||
session: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let table = db
|
||||
.open_table(OpenTableRequest {
|
||||
name: "test".to_string(),
|
||||
namespace_path: vec![],
|
||||
index_cache_size: None,
|
||||
lance_read_params: None,
|
||||
location: None,
|
||||
namespace_client: None,
|
||||
managed_versioning: None,
|
||||
})
|
||||
.await
|
||||
.unwrap();
|
||||
let batches = table
|
||||
.query(
|
||||
&AnyQuery::Query(QueryRequest::default()),
|
||||
Default::default(),
|
||||
)
|
||||
.await
|
||||
.unwrap()
|
||||
.try_collect::<Vec<_>>()
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::<usize>(), 3);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_clone_table_basic() {
|
||||
let (_tempdir, db) = setup_database().await;
|
||||
@@ -2432,7 +2342,7 @@ mod tests {
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_table_uri() {
|
||||
let (_tempdir, mut db) = setup_database().await;
|
||||
let (_tempdir, db) = setup_database().await;
|
||||
|
||||
let mut pb = PathBuf::new();
|
||||
pb.push(db.uri.clone());
|
||||
@@ -2441,18 +2351,6 @@ mod tests {
|
||||
let expected = pb.to_str().unwrap();
|
||||
let uri = db.table_uri("test").ok().unwrap();
|
||||
assert_eq!(uri, expected);
|
||||
|
||||
// URI paths always use forward slashes, even on Windows. Using
|
||||
// `Path::join` here used to produce `az://container/prefix\\test.lance`,
|
||||
// which Azure treated as a different object from the table returned by
|
||||
// `table_names` (https://github.com/lancedb/lancedb/issues/1072).
|
||||
for base_uri in ["az://container/prefix", "az://container/prefix/"] {
|
||||
db.uri = base_uri.to_string();
|
||||
assert_eq!(
|
||||
db.table_uri("test").unwrap(),
|
||||
"az://container/prefix/test.lance"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// Regression: connecting via a URL-style URI (which goes through
|
||||
|
||||
@@ -201,7 +201,7 @@ impl LanceNamespaceDatabase {
|
||||
&self,
|
||||
request: &DbCreateTableRequest,
|
||||
) -> Result<(
|
||||
Option<lance_file::version::LanceFileVersion>,
|
||||
Option<lance_encoding::version::LanceFileVersion>,
|
||||
Option<bool>,
|
||||
Option<bool>,
|
||||
)> {
|
||||
@@ -214,7 +214,7 @@ impl LanceNamespaceDatabase {
|
||||
|
||||
let storage_version_override = storage_options
|
||||
.and_then(|opts| opts.get(OPT_NEW_TABLE_STORAGE_VERSION))
|
||||
.map(|s| s.parse::<lance_file::version::LanceFileVersion>())
|
||||
.map(|s| s.parse::<lance_encoding::version::LanceFileVersion>())
|
||||
.transpose()?;
|
||||
|
||||
let v2_manifest_override = storage_options
|
||||
|
||||
@@ -169,12 +169,6 @@ impl From<DataFusionError> for Error {
|
||||
|
||||
impl From<lance::Error> for Error {
|
||||
fn from(source: lance::Error) -> Self {
|
||||
if has_unsupported_local_filesystem_source(&source) {
|
||||
return Self::NotSupported {
|
||||
message: "the filesystem does not support an operation required for safe Lance commits (such as atomic rename). Object-storage mounts such as Mountpoint for Amazon S3 are not supported; use the native object-store URI (for example, s3://bucket/path) instead".to_string(),
|
||||
};
|
||||
}
|
||||
|
||||
// Try to unwrap external errors that were wrapped by lance
|
||||
match source {
|
||||
lance::Error::Wrapped { error, .. } => Self::from_box_error(error),
|
||||
@@ -187,27 +181,6 @@ impl From<lance::Error> for Error {
|
||||
}
|
||||
}
|
||||
|
||||
fn has_unsupported_local_filesystem_source(error: &(dyn std::error::Error + 'static)) -> bool {
|
||||
let mut current = Some(error);
|
||||
let mut is_local_filesystem = false;
|
||||
let mut is_unsupported = false;
|
||||
while let Some(error) = current {
|
||||
is_local_filesystem |= error
|
||||
.downcast_ref::<object_store::Error>()
|
||||
.is_some_and(|error| {
|
||||
matches!(error, object_store::Error::Generic { store, .. } if *store == "LocalFileSystem")
|
||||
});
|
||||
is_unsupported |= error
|
||||
.downcast_ref::<std::io::Error>()
|
||||
.is_some_and(|error| error.kind() == std::io::ErrorKind::Unsupported);
|
||||
if is_local_filesystem && is_unsupported {
|
||||
return true;
|
||||
}
|
||||
current = error.source();
|
||||
}
|
||||
false
|
||||
}
|
||||
|
||||
impl Error {
|
||||
fn from_box_error(mut source: Box<dyn std::error::Error + Send + Sync>) -> Self {
|
||||
source = match source.downcast::<Self>() {
|
||||
@@ -297,46 +270,3 @@ impl From<candle_core::Error> for Error {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn unsupported_filesystem_operations_have_actionable_error() {
|
||||
let object_store_error = object_store::Error::Generic {
|
||||
store: "LocalFileSystem",
|
||||
source: Box::new(std::io::Error::from(std::io::ErrorKind::Unsupported)),
|
||||
};
|
||||
let lance_error = lance::Error::io_source(Box::new(object_store_error));
|
||||
|
||||
let error = Error::from(lance_error);
|
||||
|
||||
assert!(matches!(
|
||||
error,
|
||||
Error::NotSupported { message }
|
||||
if message.contains("Mountpoint for Amazon S3")
|
||||
&& message.contains("s3://bucket/path")
|
||||
));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn other_io_errors_remain_lance_errors() {
|
||||
let object_store_error = object_store::Error::Generic {
|
||||
store: "LocalFileSystem",
|
||||
source: Box::new(std::io::Error::from(std::io::ErrorKind::PermissionDenied)),
|
||||
};
|
||||
let lance_error = lance::Error::io_source(Box::new(object_store_error));
|
||||
|
||||
assert!(matches!(Error::from(lance_error), Error::Lance { .. }));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unsupported_non_filesystem_errors_remain_lance_errors() {
|
||||
let lance_error = lance::Error::io_source(Box::new(std::io::Error::from(
|
||||
std::io::ErrorKind::Unsupported,
|
||||
)));
|
||||
|
||||
assert!(matches!(Error::from(lance_error), Error::Lance { .. }));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -23,13 +23,11 @@ use crate::table::AddResult;
|
||||
use crate::table::BranchDiff;
|
||||
use crate::table::DeleteResult;
|
||||
use crate::table::DropColumnsResult;
|
||||
use crate::table::LsmStats;
|
||||
use crate::table::LsmWriteSpec;
|
||||
use crate::table::MergeBranchResult;
|
||||
use crate::table::MergeResult;
|
||||
use crate::table::Tags;
|
||||
use crate::table::UpdateResult;
|
||||
use crate::table::lsm_stats::GetLsmStatsResponse;
|
||||
use crate::table::merge::MergeFilter;
|
||||
use crate::table::query::create_multi_vector_plan;
|
||||
use crate::table::write_progress::FinishOnDrop;
|
||||
@@ -993,18 +991,6 @@ impl<S: HttpSend> RemoteTable<S> {
|
||||
}
|
||||
}
|
||||
|
||||
/// Send an LSM operator request with the transport retry layer **off**.
|
||||
///
|
||||
/// Retry policy on these routes belongs to the checkpoint loop, which
|
||||
/// reads the status and can tell contention from a lost claim. Leaving the
|
||||
/// transport layer on would re-ask on its own schedule first, and surface
|
||||
/// an `Error::Retry` whose status the loop would then have to unwrap.
|
||||
async fn send_lsm_route(&self, request: RequestBuilder) -> Result<(String, reqwest::Response)> {
|
||||
let (request_id, response) = self.send(request, false).await?;
|
||||
let response = self.check_table_response(&request_id, response).await?;
|
||||
Ok((request_id, response))
|
||||
}
|
||||
|
||||
/// Build a POST request and attach the read-freshness headers
|
||||
/// (`x-lancedb-min-version`, `x-lancedb-min-timestamp`).
|
||||
fn post_read(&self, uri: &str) -> RequestBuilder {
|
||||
@@ -2482,47 +2468,13 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
|
||||
})
|
||||
}
|
||||
|
||||
async fn flush_lsm(&self) -> Result<()> {
|
||||
let request = self
|
||||
.client
|
||||
.post(&format!("/v1/table/{}/flush_lsm/", self.identifier));
|
||||
self.send_lsm_route(request).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn compact_lsm(&self) -> Result<()> {
|
||||
let request = self
|
||||
.client
|
||||
.post(&format!("/v1/table/{}/compact_lsm/", self.identifier));
|
||||
self.send_lsm_route(request).await?;
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn get_lsm_stats(&self, include_generation_rows: bool) -> Result<Option<LsmStats>> {
|
||||
// Read-semantics POST, like `get_lsm_write_spec`.
|
||||
let request = self
|
||||
.post_read(&format!("/v1/table/{}/get_lsm_stats/", self.identifier))
|
||||
.json(&serde_json::json!({
|
||||
"include_generation_rows": include_generation_rows,
|
||||
}));
|
||||
let (request_id, response) = self.send_lsm_route(request).await?;
|
||||
let body = response.text().await.err_to_http(request_id.clone())?;
|
||||
let parsed: GetLsmStatsResponse = serde_json::from_str(&body).map_err(|e| Error::Http {
|
||||
source: format!("Failed to parse get_lsm_stats response: {e}").into(),
|
||||
request_id,
|
||||
status_code: None,
|
||||
})?;
|
||||
// `null` — and only — when the table has no LSM write path.
|
||||
Ok(parsed.lsm_stats)
|
||||
}
|
||||
|
||||
async fn set_lsm_write_spec(&self, spec: LsmWriteSpec) -> Result<()> {
|
||||
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,
|
||||
@@ -2990,7 +2942,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
|
||||
}
|
||||
|
||||
#[derive(Serialize, Clone, Debug)]
|
||||
pub struct MergeInsertRequest {
|
||||
pub(crate) struct MergeInsertRequest {
|
||||
on: String,
|
||||
when_matched_update_all: bool,
|
||||
when_matched_update_all_filt: Option<String>,
|
||||
@@ -5955,18 +5907,16 @@ 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.
|
||||
assert_eq!(
|
||||
tokens,
|
||||
vec![
|
||||
FtsToken {
|
||||
text: "こんにちは".to_string(),
|
||||
position: 0,
|
||||
position: 1,
|
||||
},
|
||||
FtsToken {
|
||||
text: "世界".to_string(),
|
||||
position: 1,
|
||||
position: 2,
|
||||
},
|
||||
]
|
||||
);
|
||||
@@ -6599,7 +6549,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 +6568,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 +6651,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")
|
||||
@@ -6748,499 +6680,6 @@ mod tests {
|
||||
assert!(table.get_lsm_write_spec().await.unwrap().is_none());
|
||||
}
|
||||
|
||||
/// Build a `get_lsm_stats` body for one bucket holding `generations`.
|
||||
fn stats_body(generations: &[u64], compacting: bool) -> String {
|
||||
serde_json::json!({
|
||||
"lsm_stats": {
|
||||
"buckets": [{
|
||||
"shard_id": "b0",
|
||||
"status": "Active",
|
||||
"writer_epoch": 1,
|
||||
"manifest_version": 1,
|
||||
"current_generation": generations.iter().max().copied().unwrap_or(0) + 1,
|
||||
"replay_after_wal_entry_position": 0,
|
||||
"wal_entry_position_last_seen": 0,
|
||||
"generations": generations.iter()
|
||||
.map(|g| serde_json::json!({ "generation": g, "bytes": 1 }))
|
||||
.collect::<Vec<_>>(),
|
||||
"compacting": compacting,
|
||||
"memtables": [],
|
||||
}],
|
||||
}
|
||||
})
|
||||
.to_string()
|
||||
}
|
||||
|
||||
/// `flush_lsm` / `compact_lsm` answer 202 with no body at all.
|
||||
fn accepted() -> http::Response<String> {
|
||||
http::Response::builder()
|
||||
.status(202)
|
||||
.body(String::new())
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
fn ok_json(body: String) -> http::Response<String> {
|
||||
http::Response::builder().status(200).body(body).unwrap()
|
||||
}
|
||||
|
||||
/// A flush landing in an empty L0 finishes on the opening stats read
|
||||
/// alone. Asserting zero compacts is the point: "it returned Ok" is also
|
||||
/// true of a loop that ran a pointless pass.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_checkpoint_short_circuits_on_empty_l0() {
|
||||
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = compacts.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("compact_lsm") {
|
||||
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
panic!("an already-converged table must issue no compact calls");
|
||||
}
|
||||
if path.contains("flush_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
assert_eq!(path, "/v1/table/my_table/get_lsm_stats/");
|
||||
ok_json(stats_body(&[], false))
|
||||
});
|
||||
|
||||
table.checkpoint_lsm().await.unwrap();
|
||||
assert_eq!(compacts.load(std::sync::atomic::Ordering::SeqCst), 0);
|
||||
}
|
||||
|
||||
/// The loop triggers compaction until every generation that existed at
|
||||
/// the start is gone, one bounded prefix per pass.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_checkpoint_triggers_until_targets_are_drained() {
|
||||
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = compacts.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("flush_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
if path.contains("compact_lsm") {
|
||||
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
return accepted();
|
||||
}
|
||||
// Each pass drains the oldest generation.
|
||||
let drained = seen.load(std::sync::atomic::Ordering::SeqCst);
|
||||
let left: Vec<u64> = [1u64, 2, 3].into_iter().skip(drained).collect();
|
||||
ok_json(stats_body(&left, false))
|
||||
});
|
||||
|
||||
table.checkpoint_lsm().await.unwrap();
|
||||
assert_eq!(
|
||||
compacts.load(std::sync::atomic::Ordering::SeqCst),
|
||||
3,
|
||||
"one trigger per generation prefix, then stop"
|
||||
);
|
||||
}
|
||||
|
||||
/// Generations created *during* the checkpoint are not waited on, which
|
||||
/// is what lets the loop terminate on a table taking writes where "L0 is
|
||||
/// empty" never becomes true.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_checkpoint_ignores_generations_created_while_it_runs() {
|
||||
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = compacts.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("flush_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
if path.contains("compact_lsm") {
|
||||
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
return accepted();
|
||||
}
|
||||
// Target is 5. One pass drains it; a writer keeps adding above.
|
||||
let n = seen.load(std::sync::atomic::Ordering::SeqCst);
|
||||
let body = if n == 0 {
|
||||
stats_body(&[5], false)
|
||||
} else {
|
||||
stats_body(&[6, 7], false)
|
||||
};
|
||||
ok_json(body)
|
||||
});
|
||||
|
||||
table.checkpoint_lsm().await.unwrap();
|
||||
assert_eq!(
|
||||
compacts.load(std::sync::atomic::Ordering::SeqCst),
|
||||
1,
|
||||
"the loop must not chase generations written after it started"
|
||||
);
|
||||
}
|
||||
|
||||
/// Contention is a 429 and must be retried. The server keeps it off 503
|
||||
/// precisely so the client can act on the status alone — reading it as
|
||||
/// terminal stops the checkpoint early on a healthy node.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_checkpoint_retries_contention() {
|
||||
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = compacts.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("flush_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
if path.contains("compact_lsm") {
|
||||
// First two triggers: every bucket already latched.
|
||||
if seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 2 {
|
||||
return http::Response::builder()
|
||||
.status(429)
|
||||
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
|
||||
.unwrap();
|
||||
}
|
||||
return accepted();
|
||||
}
|
||||
let accepted_triggers = seen
|
||||
.load(std::sync::atomic::Ordering::SeqCst)
|
||||
.saturating_sub(2);
|
||||
let left: Vec<u64> = if accepted_triggers == 0 {
|
||||
vec![1]
|
||||
} else {
|
||||
vec![]
|
||||
};
|
||||
ok_json(stats_body(&left, false))
|
||||
});
|
||||
|
||||
table
|
||||
.checkpoint_lsm()
|
||||
.await
|
||||
.expect("contention must not abort the checkpoint");
|
||||
assert_eq!(
|
||||
compacts.load(std::sync::atomic::Ordering::SeqCst),
|
||||
3,
|
||||
"assert the retry count, not just the outcome"
|
||||
);
|
||||
}
|
||||
|
||||
/// A transient fault on the poll must not abort the checkpoint. This route
|
||||
/// meets the most contention — it runs every `POLL_INTERVAL` for the
|
||||
/// checkpoint's whole life, with the transport retry layer disabled — yet
|
||||
/// was the one call reached with a bare `?`.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_checkpoint_retries_a_contended_stats_poll() {
|
||||
let polls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = polls.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("flush_lsm") || path.contains("compact_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
// The opening read lands; the next two polls are latched out.
|
||||
let n = seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
if (1..3).contains(&n) {
|
||||
return http::Response::builder()
|
||||
.status(429)
|
||||
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
|
||||
.unwrap();
|
||||
}
|
||||
ok_json(stats_body(if n < 4 { &[1] } else { &[] }, false))
|
||||
});
|
||||
|
||||
table
|
||||
.checkpoint_lsm()
|
||||
.await
|
||||
.expect("a contended poll must be retried, not surfaced");
|
||||
assert_eq!(
|
||||
polls.load(std::sync::atomic::Ordering::SeqCst),
|
||||
5,
|
||||
"the two rejected polls must be re-issued, not skipped"
|
||||
);
|
||||
}
|
||||
|
||||
/// Contention and a lost claim draw on separate budgets: five straight
|
||||
/// 429s on `flush`, more than `MAX_REISSUES`, must still converge. On one
|
||||
/// shared counter this spent the re-issue cap and then reported a lost
|
||||
/// claim nothing had ever reported.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_contention_does_not_exhaust_the_reissue_budget() {
|
||||
let flushes = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = flushes.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("flush_lsm") {
|
||||
if seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst) < 5 {
|
||||
return http::Response::builder()
|
||||
.status(429)
|
||||
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
|
||||
.unwrap();
|
||||
}
|
||||
return accepted();
|
||||
}
|
||||
if path.contains("compact_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
ok_json(stats_body(&[], false))
|
||||
});
|
||||
|
||||
table
|
||||
.checkpoint_lsm()
|
||||
.await
|
||||
.expect("contention must not be reported as a lost claim");
|
||||
assert_eq!(
|
||||
flushes.load(std::sync::atomic::Ordering::SeqCst),
|
||||
6,
|
||||
"five retries against one seal, then it lands"
|
||||
);
|
||||
}
|
||||
|
||||
/// An exhausted retry budget surfaces the fault that consumed it, not a
|
||||
/// message the loop invented: "429, nine times" points an operator at a
|
||||
/// saturated pool, a generic runtime error points them nowhere.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_exhausted_retries_surface_the_underlying_fault() {
|
||||
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = calls.clone();
|
||||
let table = Table::new_with_handler("my_table", move |_request| {
|
||||
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
http::Response::builder()
|
||||
.status(429)
|
||||
.body(r#"{"code":21,"error":"Too many concurrent writes"}"#.to_string())
|
||||
.unwrap()
|
||||
});
|
||||
|
||||
let err = table.checkpoint_lsm().await.unwrap_err();
|
||||
assert!(
|
||||
matches!(&err, Error::Http { status_code: Some(s), .. } if s.as_u16() == 429),
|
||||
"the fault that spent the budget must be the one reported: {err:?}"
|
||||
);
|
||||
assert_eq!(
|
||||
calls.load(std::sync::atomic::Ordering::SeqCst),
|
||||
9,
|
||||
"one call plus MAX_RETRIES — the re-issue budget is not spent on top"
|
||||
);
|
||||
}
|
||||
|
||||
/// A draining node is terminal, but the client does not know that from the
|
||||
/// status: draining and a proxy blip are both 503, and telling them apart
|
||||
/// takes parsing the body for a namespace code. So it spends the retry
|
||||
/// budget and then reports what the server said — the drain gate never
|
||||
/// releases, so the answer does not change, and the operator still reads
|
||||
/// "WAL node draining" in the error.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_draining_surfaces_after_the_retry_budget() {
|
||||
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = calls.clone();
|
||||
let table = Table::new_with_handler("my_table", move |_request| {
|
||||
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
http::Response::builder()
|
||||
.status(503)
|
||||
.body(r#"{"code":19,"error":"WAL node draining"}"#.to_string())
|
||||
.unwrap()
|
||||
});
|
||||
|
||||
let err = table.checkpoint_lsm().await.unwrap_err();
|
||||
let message = err.to_string();
|
||||
assert!(
|
||||
matches!(&err, Error::Http { status_code: Some(s), .. } if s.as_u16() == 503),
|
||||
"the 503 must surface as itself: {err:?}"
|
||||
);
|
||||
assert!(
|
||||
message.contains("WAL node draining"),
|
||||
"the server's own diagnosis must survive to the caller: {message}"
|
||||
);
|
||||
assert_eq!(
|
||||
calls.load(std::sync::atomic::Ordering::SeqCst),
|
||||
9,
|
||||
"one call plus MAX_RETRIES, then it reports rather than spinning"
|
||||
);
|
||||
}
|
||||
|
||||
/// A long stall with nothing compacting must keep waiting, not fail. The
|
||||
/// client cannot judge this: a checkpoint queued behind unrelated tables
|
||||
/// on the pod-wide compactor pool reports exactly these numbers — flat
|
||||
/// generations, an idle latch — as one whose merges are failing. The
|
||||
/// deadline is the caller's.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_checkpoint_waits_out_a_long_stall_rather_than_failing() {
|
||||
let polls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = polls.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("flush_lsm") || path.contains("compact_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
// Flat for far longer than any bound this loop ever had, with
|
||||
// `compacting: false` throughout — then it drains.
|
||||
let n = seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
ok_json(stats_body(if n < 40 { &[1, 2] } else { &[] }, false))
|
||||
});
|
||||
|
||||
table
|
||||
.checkpoint_lsm()
|
||||
.await
|
||||
.expect("a stall is the server being slow, not the client's call to make");
|
||||
assert!(
|
||||
polls.load(std::sync::atomic::Ordering::SeqCst) > 40,
|
||||
"the loop must have kept polling well past the old ten-poll bound"
|
||||
);
|
||||
}
|
||||
|
||||
/// A pass already owns the latch on every outstanding bucket, so the loop
|
||||
/// waits rather than piling on triggers it would only refuse. This is the
|
||||
/// sole thing `compacting` is read for.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_checkpoint_waits_while_a_pass_is_running() {
|
||||
let polls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let compacts = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen_polls = polls.clone();
|
||||
let seen_compacts = compacts.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
if path.contains("flush_lsm") {
|
||||
return accepted();
|
||||
}
|
||||
if path.contains("compact_lsm") {
|
||||
seen_compacts.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
return accepted();
|
||||
}
|
||||
// Latched for many polls, then done.
|
||||
let n = seen_polls.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
ok_json(if n > 15 {
|
||||
stats_body(&[], false)
|
||||
} else {
|
||||
stats_body(&[1], true)
|
||||
})
|
||||
});
|
||||
|
||||
table
|
||||
.checkpoint_lsm()
|
||||
.await
|
||||
.expect("a running pass is progress, not a stall");
|
||||
assert_eq!(
|
||||
compacts.load(std::sync::atomic::Ordering::SeqCst),
|
||||
0,
|
||||
"never trigger against a bucket already compacting"
|
||||
);
|
||||
}
|
||||
|
||||
/// WAL off ⇒ `None`; WAL on ⇒ a fully populated `Some` with no field
|
||||
/// defaulting to a zero it did not measure. `include_generation_rows`
|
||||
/// rides in the body and is off unless asked for.
|
||||
#[tokio::test]
|
||||
async fn test_get_lsm_stats_round_trip() {
|
||||
let table = Table::new_with_handler("my_table", |request| {
|
||||
assert_eq!(request.url().path(), "/v1/table/my_table/get_lsm_stats/");
|
||||
let body = request.body().unwrap().as_bytes().unwrap();
|
||||
let body: serde_json::Value = serde_json::from_slice(body).unwrap();
|
||||
assert_eq!(
|
||||
body["include_generation_rows"], true,
|
||||
"the flag must reach the server, not be silently dropped"
|
||||
);
|
||||
let response = serde_json::json!({
|
||||
"lsm_stats": {
|
||||
"buckets": [{
|
||||
"shard_id": "b0",
|
||||
"status": "Active",
|
||||
"writer_epoch": 3,
|
||||
"manifest_version": 11,
|
||||
"current_generation": 9,
|
||||
"replay_after_wal_entry_position": 100,
|
||||
"wal_entry_position_last_seen": 140,
|
||||
"generations": [{ "generation": 8, "bytes": 4096, "rows": 30 }],
|
||||
"compacting": false,
|
||||
"memtables": [
|
||||
{ "generation": 9, "rows": 12, "bytes": 900, "batches": 2,
|
||||
"indexes": ["vec_idx"] }
|
||||
],
|
||||
}],
|
||||
}
|
||||
});
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(response.to_string())
|
||||
.unwrap()
|
||||
});
|
||||
|
||||
let stats = table
|
||||
.get_lsm_stats(true)
|
||||
.await
|
||||
.unwrap()
|
||||
.expect("a WAL-backed table reports Some");
|
||||
let bucket = &stats.buckets[0];
|
||||
assert_eq!(bucket.replay_after_wal_entry_position, 100);
|
||||
assert_eq!(bucket.wal_entry_position_last_seen, 140);
|
||||
assert!(!bucket.compacting);
|
||||
assert_eq!(bucket.generations[0].generation, 8);
|
||||
assert_eq!(bucket.generations[0].rows, Some(30));
|
||||
// The line that answers "why is my fresh-tier vector search
|
||||
// brute-force" — an absent index name is the whole explanation.
|
||||
let memtables = bucket.memtables.as_ref().unwrap();
|
||||
assert_eq!(memtables[0].indexes, vec!["vec_idx".to_string()]);
|
||||
}
|
||||
|
||||
/// A 404 arrives as `TableNotFound`, not as a lost claim the loop
|
||||
/// re-issues from flush until its cap. The two are distinguished by
|
||||
/// status: 404 is "no such table", 421 is "this node holds no claim".
|
||||
/// They shared 404 once, and the loop chased a name that never existed.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_missing_table_is_not_read_as_a_lost_claim() {
|
||||
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = calls.clone();
|
||||
let table = Table::new_with_handler("my_table", move |_request| {
|
||||
seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
http::Response::builder()
|
||||
.status(404)
|
||||
.body(r#"{"code":4,"error":"Not found: Table not found: my_table"}"#.to_string())
|
||||
.unwrap()
|
||||
});
|
||||
|
||||
let err = table.checkpoint_lsm().await.unwrap_err();
|
||||
assert!(
|
||||
matches!(err, Error::TableNotFound { .. }),
|
||||
"a missing table must say so: {err:?}"
|
||||
);
|
||||
assert_eq!(
|
||||
calls.load(std::sync::atomic::Ordering::SeqCst),
|
||||
1,
|
||||
"no point re-claiming a table that does not exist"
|
||||
);
|
||||
}
|
||||
|
||||
/// A lost claim — 421, not 404 — does re-issue from flush, the call that
|
||||
/// re-claims and replays.
|
||||
#[tokio::test(start_paused = true)]
|
||||
async fn test_registry_miss_reissues_from_flush() {
|
||||
let calls = Arc::new(std::sync::atomic::AtomicUsize::new(0));
|
||||
let seen = calls.clone();
|
||||
let table = Table::new_with_handler("my_table", move |request| {
|
||||
let path = request.url().path().to_string();
|
||||
let n = seen.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
|
||||
if path.contains("flush_lsm") {
|
||||
// First flush lands; the claim is then lost, and the
|
||||
// re-issued flush succeeds.
|
||||
return accepted();
|
||||
}
|
||||
if path.contains("compact_lsm") {
|
||||
if n < 4 {
|
||||
return http::Response::builder()
|
||||
.status(421)
|
||||
.body(r#"{"code":19,"error":"table not claimed"}"#.to_string())
|
||||
.unwrap();
|
||||
}
|
||||
return accepted();
|
||||
}
|
||||
ok_json(stats_body(if n < 6 { &[1] } else { &[] }, false))
|
||||
});
|
||||
|
||||
table
|
||||
.checkpoint_lsm()
|
||||
.await
|
||||
.expect("a lost claim must be recovered by re-flushing, not surfaced");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_get_lsm_stats_absent_when_wal_off() {
|
||||
let table = Table::new_with_handler("my_table", |_request| {
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(serde_json::json!({ "lsm_stats": null }).to_string())
|
||||
.unwrap()
|
||||
});
|
||||
assert!(table.get_lsm_stats(false).await.unwrap().is_none());
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_wait_for_index() {
|
||||
let table = _make_table_with_indices(0);
|
||||
|
||||
@@ -90,7 +90,7 @@ struct RemoteBlobState {
|
||||
|
||||
/// Seekable Cloud blob handle over HTTP Range.
|
||||
#[derive(Debug)]
|
||||
pub struct RemoteBlobFile {
|
||||
pub(crate) struct RemoteBlobFile {
|
||||
requester: Arc<dyn BlobRangeRequester>,
|
||||
state: Mutex<RemoteBlobState>,
|
||||
closed: AtomicBool,
|
||||
|
||||
@@ -33,7 +33,7 @@ use crate::table::{AddResult, MergeResult};
|
||||
/// same Arrow-IPC streaming body and error side-channel; only the target
|
||||
/// endpoint, query parameters, and parsed result type differ.
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum WriteOp {
|
||||
pub(crate) enum WriteOp {
|
||||
/// `add`: stream to `/v1/table/{id}/insert/`, optionally overwriting.
|
||||
Insert { overwrite: bool },
|
||||
/// `merge_insert`: stream to `/v1/table/{id}/merge_insert/` with the merge
|
||||
@@ -49,7 +49,7 @@ pub enum WriteOp {
|
||||
/// The parsed server response for a completed write, discriminated by the
|
||||
/// operation that produced it.
|
||||
#[derive(Debug, Clone)]
|
||||
pub enum WriteResult {
|
||||
pub(crate) enum WriteResult {
|
||||
Add(AddResult),
|
||||
Merge(MergeResult),
|
||||
}
|
||||
|
||||
+41
-496
@@ -68,12 +68,10 @@ use self::merge::MergeInsertBuilder;
|
||||
pub mod add_columns;
|
||||
mod add_data;
|
||||
pub mod branch_merge;
|
||||
pub mod checkpoint;
|
||||
mod create_index;
|
||||
pub mod datafusion;
|
||||
pub(crate) mod dataset;
|
||||
pub mod delete;
|
||||
pub mod lsm_stats;
|
||||
pub mod merge;
|
||||
pub mod optimize;
|
||||
mod primary_key;
|
||||
@@ -95,8 +93,8 @@ pub use delete::DeleteResult;
|
||||
use futures::future::join_all;
|
||||
pub use lance::dataset::refs::{BranchContents, Ref, TagContents, Tags as LanceTags};
|
||||
pub use lance::dataset::scanner::DatasetRecordBatchStream;
|
||||
use lance::dataset::statistics::DatasetStatisticsExt;
|
||||
pub use lance_index::optimize::OptimizeOptions;
|
||||
pub use lsm_stats::{BucketStats, GenerationStats, LsmStats, MemtableStats};
|
||||
pub use optimize::{CompactionOptions, OptimizeAction, OptimizeStats};
|
||||
pub use schema_evolution::{
|
||||
AddColumnsResult, AlterColumnsResult, DropColumnsResult, FieldMetadataUpdate,
|
||||
@@ -370,8 +368,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.
|
||||
@@ -386,12 +382,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>,
|
||||
},
|
||||
@@ -401,41 +394,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
|
||||
@@ -447,37 +434,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, ..
|
||||
@@ -487,7 +465,7 @@ impl LsmWriteSpec {
|
||||
}
|
||||
| Self::Unsharded {
|
||||
maintained_indexes, ..
|
||||
} => *maintained_indexes = indexes,
|
||||
} => *maintained_indexes = v,
|
||||
}
|
||||
self
|
||||
}
|
||||
@@ -523,9 +501,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, ..
|
||||
@@ -535,7 +512,7 @@ impl LsmWriteSpec {
|
||||
}
|
||||
| Self::Unsharded {
|
||||
maintained_indexes, ..
|
||||
} => maintained_indexes.as_deref(),
|
||||
} => maintained_indexes,
|
||||
}
|
||||
}
|
||||
|
||||
@@ -708,31 +685,6 @@ pub trait BaseTable: std::fmt::Display + std::fmt::Debug + Send + Sync {
|
||||
message: "get_lsm_write_spec is not supported on this table type".into(),
|
||||
})
|
||||
}
|
||||
/// Seal every bucket's active memtable into L0.
|
||||
///
|
||||
/// The default implementation returns `NotSupported`.
|
||||
async fn flush_lsm(&self) -> Result<()> {
|
||||
Err(Error::NotSupported {
|
||||
message: "flush_lsm is not supported on this table type".into(),
|
||||
})
|
||||
}
|
||||
/// Trigger a background L0 → base compaction pass per bucket.
|
||||
///
|
||||
/// The default implementation returns `NotSupported`.
|
||||
async fn compact_lsm(&self) -> Result<()> {
|
||||
Err(Error::NotSupported {
|
||||
message: "compact_lsm is not supported on this table type".into(),
|
||||
})
|
||||
}
|
||||
/// Read live LSM state, or `None` when the LSM write path is not
|
||||
/// enabled for this table.
|
||||
///
|
||||
/// The default implementation returns `NotSupported`.
|
||||
async fn get_lsm_stats(&self, _include_generation_rows: bool) -> Result<Option<LsmStats>> {
|
||||
Err(Error::NotSupported {
|
||||
message: "get_lsm_stats is not supported on this table type".into(),
|
||||
})
|
||||
}
|
||||
/// Drain and close any cached MemWAL shard writers for this table.
|
||||
///
|
||||
/// The default implementation is a no-op; table types that maintain
|
||||
@@ -1733,7 +1685,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(())
|
||||
@@ -1755,10 +1707,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
|
||||
///
|
||||
@@ -1775,85 +1726,6 @@ impl Table {
|
||||
self.inner.get_lsm_write_spec().await
|
||||
}
|
||||
|
||||
/// Converge this table's LSM write path into its base table.
|
||||
///
|
||||
/// One `flush` to seal every memtable into L0, then compaction triggers
|
||||
/// until every generation that existed at that moment has reached base.
|
||||
/// The loop runs client-side, reading progress from `get_lsm_stats`, so
|
||||
/// there is no held socket and nothing to reconcile if you drop this
|
||||
/// future partway through.
|
||||
///
|
||||
/// **Best-effort.** Generations created *after* the opening flush are
|
||||
/// deliberately not waited on — that is what lets this terminate on a
|
||||
/// table taking writes. Idempotent and safe on a cadence: an
|
||||
/// already-converged table costs two round trips and triggers nothing.
|
||||
///
|
||||
/// **No deadline, and the caller owns that.** It returns when the target
|
||||
/// generations are gone, propagates a terminal server fault, and
|
||||
/// otherwise waits however long the server takes. A slow table and a
|
||||
/// stuck one are the same picture from here: the compactor pool is shared
|
||||
/// across every table on the node, so a checkpoint queued behind
|
||||
/// unrelated work is indistinguishable from one that is merging. Wrap
|
||||
/// this in `tokio::time::timeout` for a wall-clock bound; abandoning it
|
||||
/// partway costs nothing.
|
||||
///
|
||||
/// # Example
|
||||
///
|
||||
/// ```no_run
|
||||
/// # use lancedb::Table;
|
||||
/// # async fn example(table: &Table) -> Result<(), Box<dyn std::error::Error>> {
|
||||
/// let before = table.get_lsm_stats(false).await?;
|
||||
/// table.checkpoint_lsm().await?;
|
||||
/// let after = table.get_lsm_stats(false).await?;
|
||||
/// # Ok(())
|
||||
/// # }
|
||||
/// ```
|
||||
pub async fn checkpoint_lsm(&self) -> Result<()> {
|
||||
checkpoint::checkpoint_lsm(self).await
|
||||
}
|
||||
|
||||
/// Seal every bucket's active memtable into L0 without touching the
|
||||
/// base table.
|
||||
///
|
||||
/// Independently useful: flushing makes memtable rows readable from L0 at
|
||||
/// a lower per-query cost. On a node that has not claimed this table it
|
||||
/// claims it and replays the WAL log first — reporting "nothing to flush"
|
||||
/// without replaying would lie about durable data.
|
||||
pub async fn flush_lsm(&self) -> Result<()> {
|
||||
self.inner.flush_lsm().await
|
||||
}
|
||||
|
||||
/// Run one bounded L0 → base compaction pass per bucket, reporting what
|
||||
/// it merged and what is left.
|
||||
///
|
||||
/// One pass, not convergence: that bounds each request's cost and gives a
|
||||
/// caller driving its own cadence a progress signal per round trip.
|
||||
pub async fn compact_lsm(&self) -> Result<()> {
|
||||
self.inner.compact_lsm().await
|
||||
}
|
||||
|
||||
/// Read live per-bucket LSM state.
|
||||
///
|
||||
/// Answers "how far behind is my fresh tier", "which bucket is hot", and
|
||||
/// "why is my fresh-tier vector search brute-force". Mutates no table
|
||||
/// state, though on a node that has not claimed this table it claims it,
|
||||
/// exactly as a read would.
|
||||
///
|
||||
/// `include_generation_rows` reports a row count per L0 generation. Off by
|
||||
/// default: each count opens an uncached Lance dataset, and
|
||||
/// `checkpoint_lsm` polls this needing only generation numbers.
|
||||
///
|
||||
/// `Ok(None)` only when the LSM write path is not enabled, matching
|
||||
/// [`Table::get_lsm_write_spec`]. Stats is fresh-tier only, so with the
|
||||
/// WAL off there is no manifest to report and a struct of zeros would
|
||||
/// read as measurements.
|
||||
///
|
||||
/// Do not build a checkpoint's termination on this: the completion
|
||||
/// predicate lives in the `flush` and `compact` responses.
|
||||
pub async fn get_lsm_stats(&self, include_generation_rows: bool) -> Result<Option<LsmStats>> {
|
||||
self.inner.get_lsm_stats(include_generation_rows).await
|
||||
}
|
||||
|
||||
/// Drain and close any cached MemWAL shard writers held for this table.
|
||||
///
|
||||
/// When an [`LsmWriteSpec`] is installed, `merge_insert` opens MemWAL shard
|
||||
@@ -3569,24 +3441,9 @@ impl BaseTable for NativeTable {
|
||||
let num_rows = self.count_rows(None).await?;
|
||||
let num_indices = self.list_indices().await?.len();
|
||||
let ds = self.dataset.get().await?;
|
||||
// Sizes come from the manifest. Summing per-field `bytes_on_disk` instead
|
||||
// would open every data file to read its column metadata, which costs one
|
||||
// IO per fragment and reports 0 for legacy v1 storage.
|
||||
//
|
||||
// The manifest summary covers only the fragments' base data files, so
|
||||
// overlay files (recorded on each fragment) and index files (recorded in
|
||||
// the manifest's index section) are added separately.
|
||||
let mut total_bytes = ds.manifest().summary().total_files_size as usize;
|
||||
for frag in ds.manifest().fragments.iter() {
|
||||
for overlay in &frag.overlays {
|
||||
if let Some(size) = overlay.data_file.file_size_bytes.get() {
|
||||
total_bytes += size.get() as usize;
|
||||
}
|
||||
}
|
||||
}
|
||||
for index in ds.load_indices().await?.iter() {
|
||||
total_bytes += index.total_size_bytes().unwrap_or(0) as usize;
|
||||
}
|
||||
let ds_clone = (*ds).clone();
|
||||
let ds_stats = Arc::new(ds_clone).calculate_data_stats().await?;
|
||||
let total_bytes = ds_stats.fields.iter().map(|f| f.bytes_on_disk).sum::<u64>() as usize;
|
||||
|
||||
let frags = ds.get_fragments();
|
||||
let mut sorted_sizes = join_all(
|
||||
@@ -3658,12 +3515,7 @@ impl BaseTable for NativeTable {
|
||||
#[skip_serializing_none]
|
||||
#[derive(Debug, Deserialize, PartialEq)]
|
||||
pub struct TableStatistics {
|
||||
/// The total size, in bytes, of the table's data files, index files, and
|
||||
/// overlay files
|
||||
///
|
||||
/// Read from the manifest, so this excludes deletion files and manifests,
|
||||
/// and it excludes any file whose size the manifest does not record
|
||||
/// (tables and indices written before writers persisted file sizes).
|
||||
/// The total number of bytes in the table
|
||||
pub total_bytes: usize,
|
||||
|
||||
/// The number of rows in the table
|
||||
@@ -3724,7 +3576,6 @@ mod tests {
|
||||
use super::*;
|
||||
use crate::connect;
|
||||
use crate::connection::ConnectBuilder;
|
||||
use crate::io::object_store::io_tracking::IoTrackingStore;
|
||||
use crate::query::Select;
|
||||
use crate::query::{ExecutableQuery, QueryBase};
|
||||
use crate::test_utils::connection::new_test_connection;
|
||||
@@ -5107,7 +4958,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));
|
||||
@@ -5117,125 +4968,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]
|
||||
@@ -5283,16 +5024,12 @@ mod tests {
|
||||
|
||||
let res = table.stats().await.unwrap();
|
||||
println!("{:#?}", res);
|
||||
// `total_bytes` is the full on-disk size of the 11 data files (this table
|
||||
// has no index or overlay files), so it is well above the 2000 bytes of
|
||||
// column data these 250 int32 pairs hold: each file carries its own footer
|
||||
// and metadata.
|
||||
assert_eq!(
|
||||
res,
|
||||
TableStatistics {
|
||||
num_rows: 250,
|
||||
num_indices: 0,
|
||||
total_bytes: 8925,
|
||||
total_bytes: 2300,
|
||||
fragment_stats: FragmentStatistics {
|
||||
num_fragments: 11,
|
||||
num_small_fragments: 11,
|
||||
@@ -5332,196 +5069,4 @@ mod tests {
|
||||
}
|
||||
)
|
||||
}
|
||||
|
||||
/// `total_bytes` counts more than the base data files: index files and
|
||||
/// overlay files recorded in the manifest are included too.
|
||||
#[tokio::test]
|
||||
pub async fn test_stats_includes_index_and_overlay_files() {
|
||||
use lance::dataset::WriteDestination;
|
||||
use lance::dataset::transaction::{DataOverlayGroup, Operation};
|
||||
use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
|
||||
use lance_file::writer::FileWriterOptions;
|
||||
use lance_io::utils::CachedFileSize;
|
||||
use lance_table::format::DataFile;
|
||||
use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage};
|
||||
use roaring::RoaringBitmap;
|
||||
|
||||
let tmp_dir = tempdir().unwrap();
|
||||
let uri = tmp_dir.path().to_str().unwrap();
|
||||
let conn = ConnectBuilder::new(uri)
|
||||
.read_consistency_interval(Duration::from_secs(0))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let schema = Arc::new(Schema::new(vec![
|
||||
Field::new("id", DataType::Int32, false),
|
||||
Field::new("foo", DataType::Int32, true),
|
||||
]));
|
||||
let batch = RecordBatch::try_new(
|
||||
schema.clone(),
|
||||
vec![
|
||||
Arc::new(Int32Array::from_iter_values(0..100)),
|
||||
Arc::new(Int32Array::from_iter_values(0..100)),
|
||||
],
|
||||
)
|
||||
.unwrap();
|
||||
let table = conn
|
||||
.create_table("test_stats_extra_files", batch)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let data_only = table.stats().await.unwrap().total_bytes;
|
||||
assert!(data_only > 0);
|
||||
|
||||
// A scalar index adds index files whose sizes are recorded in the
|
||||
// manifest's index section.
|
||||
table
|
||||
.create_index(&["id"], Index::Auto)
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let with_index = table.stats().await.unwrap().total_bytes;
|
||||
let dataset = {
|
||||
let native = table.as_native().unwrap();
|
||||
(*native.dataset.get().await.unwrap()).clone()
|
||||
};
|
||||
let index_bytes: usize = dataset
|
||||
.load_indices()
|
||||
.await
|
||||
.unwrap()
|
||||
.iter()
|
||||
.map(|idx| idx.total_size_bytes().unwrap_or(0) as usize)
|
||||
.sum();
|
||||
assert!(index_bytes > 0);
|
||||
assert_eq!(with_index, data_only + index_bytes);
|
||||
|
||||
// Commit an overlay file supplying new `foo` values for the first three
|
||||
// rows of fragment 0. There is no high-level API that writes overlays
|
||||
// yet, so write the overlay's data file and commit the `DataOverlay`
|
||||
// operation by hand.
|
||||
let read_version = dataset.version().version;
|
||||
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 = 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,
|
||||
obj_writer,
|
||||
overlay_schema,
|
||||
FileWriterOptions::default(),
|
||||
)
|
||||
.unwrap();
|
||||
writer
|
||||
.write_column(0, Arc::new(Int32Array::from(vec![1000, 1001, 1002])) as _)
|
||||
.await
|
||||
.unwrap();
|
||||
let summary = writer.finish().await.unwrap();
|
||||
let overlay_bytes = summary.size_bytes as usize;
|
||||
assert!(overlay_bytes > 0);
|
||||
|
||||
let mut data_file = DataFile::new_unstarted(filename, file_version);
|
||||
data_file.fields = writer
|
||||
.field_id_to_column_indices()
|
||||
.iter()
|
||||
.map(|(field_id, _)| *field_id as i32)
|
||||
.collect::<Vec<_>>()
|
||||
.into();
|
||||
data_file.column_indices = writer
|
||||
.field_id_to_column_indices()
|
||||
.iter()
|
||||
.map(|(_, column_index)| *column_index as i32)
|
||||
.collect::<Vec<_>>()
|
||||
.into();
|
||||
data_file.file_size_bytes = CachedFileSize::new(summary.size_bytes);
|
||||
|
||||
let overlay = DataOverlayFile {
|
||||
data_file,
|
||||
coverage: OverlayCoverage::dense(RoaringBitmap::from_iter(0..3)),
|
||||
committed_version: 0,
|
||||
};
|
||||
Dataset::commit(
|
||||
WriteDestination::Dataset(Arc::new(dataset)),
|
||||
Operation::DataOverlay {
|
||||
groups: vec![DataOverlayGroup {
|
||||
fragment_id,
|
||||
overlays: vec![overlay],
|
||||
}],
|
||||
},
|
||||
Some(read_version),
|
||||
None,
|
||||
None,
|
||||
Arc::new(Default::default()),
|
||||
false,
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
table.checkout_latest().await.unwrap();
|
||||
let with_overlay = table.stats().await.unwrap().total_bytes;
|
||||
assert_eq!(with_overlay, with_index + overlay_bytes);
|
||||
}
|
||||
|
||||
/// `stats()` must stay manifest-only. Summing per-field `bytes_on_disk`
|
||||
/// instead opens every data file, so cost would grow with fragment count.
|
||||
#[tokio::test]
|
||||
pub async fn test_stats_does_not_read_data_files() {
|
||||
let tmp_dir = tempdir().unwrap();
|
||||
let uri = tmp_dir.path().to_str().unwrap();
|
||||
|
||||
let conn = ConnectBuilder::new(uri).execute().await.unwrap();
|
||||
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
|
||||
let batch = RecordBatch::try_new(
|
||||
schema.clone(),
|
||||
vec![Arc::new(Int32Array::from_iter_values(0..10))],
|
||||
)
|
||||
.unwrap();
|
||||
|
||||
conn.create_table("test_stats_io", batch.clone())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
let table = conn.open_table("test_stats_io").execute().await.unwrap();
|
||||
const NUM_APPENDS: usize = 20;
|
||||
for _ in 0..NUM_APPENDS {
|
||||
table.add(batch.clone()).execute().await.unwrap();
|
||||
}
|
||||
|
||||
// Reopen through a tracking store so the counters cover `stats()` alone and
|
||||
// not the writes above.
|
||||
let (wrapper, io_stats) = IoTrackingStore::new_wrapper();
|
||||
let table = conn
|
||||
.open_table("test_stats_io")
|
||||
.lance_read_params(ReadParams {
|
||||
store_options: Some(ObjectStoreParams {
|
||||
object_store_wrapper: Some(wrapper),
|
||||
..Default::default()
|
||||
}),
|
||||
..Default::default()
|
||||
})
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
io_stats.lock().unwrap().read_iops = 0;
|
||||
|
||||
let stats = table.stats().await.unwrap();
|
||||
let read_iops = io_stats.lock().unwrap().read_iops;
|
||||
|
||||
assert_eq!(stats.fragment_stats.num_fragments, NUM_APPENDS + 1);
|
||||
assert!(stats.total_bytes > 0);
|
||||
// Reading the fragments' data files would take at least one IOP each.
|
||||
assert!(
|
||||
read_iops < stats.fragment_stats.num_fragments as u64,
|
||||
"stats() issued {} read IOPs across {} fragments",
|
||||
read_iops,
|
||||
stats.fragment_stats.num_fragments
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,315 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! Converging a table's LSM write path into its base table.
|
||||
//!
|
||||
//! `checkpoint_lsm` seals once, then triggers compaction and watches
|
||||
//! generation numbers until the L0 that existed at the start is gone.
|
||||
//!
|
||||
//! The loop runs in the client, not the server: `compact_lsm` dispatches a
|
||||
//! pass and returns, so nothing holds a socket and a client can vanish
|
||||
//! mid-operation with nothing to reconcile. Completion is read from
|
||||
//! generation numbers in the shard manifest — durable state, unlike a count
|
||||
//! in a compact response, which a concurrent write invalidates.
|
||||
//!
|
||||
//! The target set is fixed at the start, so generations created *during* the
|
||||
//! checkpoint are ignored. That is what lets it terminate under write load,
|
||||
//! and what makes it best-effort: it converges the fresh tier as of some
|
||||
//! instant. Idempotent, abandonable at any point, safe on a cadence.
|
||||
//!
|
||||
//! No liveness bound — the caller owns the deadline. The compactor pool is
|
||||
//! shared pod-wide, so a checkpoint queued behind unrelated tables looks
|
||||
//! exactly like one that is merging.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::future::Future;
|
||||
use std::time::Duration;
|
||||
|
||||
use crate::{Error, Result, Table};
|
||||
|
||||
/// The HTTP status a failed request carried, if it carried one.
|
||||
///
|
||||
/// `None` for anything with no retry story: a `TableNotFound` that
|
||||
/// `check_table_response` already translated, or a connection failure that
|
||||
/// never reached the server. Both are terminal.
|
||||
fn status_of(e: &Error) -> Option<u16> {
|
||||
#[cfg(feature = "remote")]
|
||||
{
|
||||
match e {
|
||||
Error::Http {
|
||||
status_code: Some(status),
|
||||
..
|
||||
} => Some(status.as_u16()),
|
||||
_ => None,
|
||||
}
|
||||
}
|
||||
#[cfg(not(feature = "remote"))]
|
||||
{
|
||||
let _ = e;
|
||||
None
|
||||
}
|
||||
}
|
||||
|
||||
/// 429 (latch held, pool saturated, or the pod replaying its WAL) and 503 (a
|
||||
/// draining node, or a proxy between here and it).
|
||||
///
|
||||
/// The status is the whole signal: the server deliberately keeps contention
|
||||
/// off 503, so a latch collision is a 429. A draining node *is* terminal, but
|
||||
/// it is also a 503 that stays a 503, so retrying spends one budget and then
|
||||
/// reports the server's own message — cheaper than parsing the body for the
|
||||
/// namespace code it would take to tell the two apart.
|
||||
fn is_retryable(e: &Error) -> bool {
|
||||
matches!(status_of(e), Some(429 | 503))
|
||||
}
|
||||
|
||||
/// 421: the owning node holds no claim. Only `flush` re-claims and replays,
|
||||
/// so this cannot be retried in place — the caller has to start over.
|
||||
fn is_lost_claim(e: &Error) -> bool {
|
||||
status_of(e) == Some(421)
|
||||
}
|
||||
|
||||
/// Interval between `get_lsm_stats` polls. One interval is roughly one
|
||||
/// compaction pass, the granularity at which the answer can change.
|
||||
///
|
||||
/// Fixed rather than configurable, matching `wait_for_index`. It costs
|
||||
/// nothing on an already-converged table and at most one interval of tail
|
||||
/// latency after the final pass lands.
|
||||
const POLL_INTERVAL: Duration = Duration::from_secs(5);
|
||||
|
||||
/// Cap on re-issues from `flush` after a 421, so a crash-looping node cannot
|
||||
/// turn flush → compact → 421 → flush into a spin.
|
||||
///
|
||||
/// Deliberately not shared with [`MAX_RETRIES`]: a claim that keeps
|
||||
/// evaporating is a broken node, while contention is routine and wants a real
|
||||
/// budget. One shared counter let a merely contended table exhaust this cap
|
||||
/// and then blame a claim it never lost.
|
||||
const MAX_REISSUES: usize = 3;
|
||||
|
||||
/// Retryable faults tolerated on a *single* request, reset on every success —
|
||||
/// scattered contention across a long checkpoint must not accumulate toward a
|
||||
/// cap. Roughly 16s of retrying against the backoff below.
|
||||
const MAX_RETRIES: usize = 8;
|
||||
|
||||
/// Backoff between retries, doubling up to [`RETRY_BACKOFF_MAX`]. Latch
|
||||
/// contention clears in about the time one pass takes, so start small; a
|
||||
/// saturated pool wants the ceiling.
|
||||
const RETRY_BACKOFF_BASE: Duration = Duration::from_millis(100);
|
||||
const RETRY_BACKOFF_MAX: Duration = Duration::from_secs(5);
|
||||
|
||||
/// Sleep before re-issuing a retryable request.
|
||||
async fn backoff(attempt: usize) {
|
||||
let delay = RETRY_BACKOFF_BASE
|
||||
.saturating_mul(1u32 << attempt.min(8) as u32)
|
||||
.min(RETRY_BACKOFF_MAX);
|
||||
tokio::time::sleep(delay).await;
|
||||
}
|
||||
|
||||
/// Whether the drain loop finished or needs the table re-claimed first.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||
enum CheckpointOutcome {
|
||||
Done,
|
||||
ReissueFromFlush,
|
||||
}
|
||||
|
||||
/// What one LSM request produced: its value, or word that the owning node
|
||||
/// holds no claim and only `flush` can get it back.
|
||||
enum Attempt<T> {
|
||||
Ok(T),
|
||||
ReissueFromFlush,
|
||||
}
|
||||
|
||||
/// Issue one LSM request, retrying in place while the fault is retryable.
|
||||
///
|
||||
/// The two recoverable faults have separate budgets: contention clears on its
|
||||
/// own and retries here against [`MAX_RETRIES`], while a 421 needs `flush` to
|
||||
/// re-claim, which only the caller can drive.
|
||||
///
|
||||
/// An exhausted budget propagates the last error *as itself* rather than a
|
||||
/// synthesized one — "429 after nine tries" beats "checkpoint failed", and a
|
||||
/// draining node arrives carrying the server's own message.
|
||||
async fn issue<T, F, Fut>(mut call: F) -> Result<Attempt<T>>
|
||||
where
|
||||
F: FnMut() -> Fut,
|
||||
Fut: Future<Output = Result<T>>,
|
||||
{
|
||||
let mut retries = 0;
|
||||
loop {
|
||||
let e = match call().await {
|
||||
Ok(value) => return Ok(Attempt::Ok(value)),
|
||||
Err(e) => e,
|
||||
};
|
||||
if is_lost_claim(&e) {
|
||||
return Ok(Attempt::ReissueFromFlush);
|
||||
}
|
||||
if !is_retryable(&e) || retries >= MAX_RETRIES {
|
||||
return Err(e);
|
||||
}
|
||||
backoff(retries).await;
|
||||
retries += 1;
|
||||
}
|
||||
}
|
||||
|
||||
/// Drive [`Table::checkpoint_lsm`]: seal once, fix the target watermark
|
||||
/// from the resulting L0, then trigger and poll until it drains.
|
||||
pub(crate) async fn checkpoint_lsm(table: &Table) -> Result<()> {
|
||||
for reissue in 0..=MAX_REISSUES {
|
||||
// The seal turns everything written before this call into a
|
||||
// generation, so the watermark has to be read after it. Idempotent:
|
||||
// sealing an empty memtable is a no-op, so a re-issue does not churn
|
||||
// empty generations.
|
||||
match issue(|| table.flush_lsm()).await? {
|
||||
Attempt::Ok(()) => {}
|
||||
Attempt::ReissueFromFlush => {
|
||||
backoff(reissue).await;
|
||||
continue;
|
||||
}
|
||||
}
|
||||
|
||||
let stats = match issue(|| table.get_lsm_stats(false)).await? {
|
||||
Attempt::Ok(stats) => stats,
|
||||
Attempt::ReissueFromFlush => {
|
||||
backoff(reissue).await;
|
||||
continue;
|
||||
}
|
||||
};
|
||||
let Some(stats) = stats else {
|
||||
// Not WAL-backed; `flush_lsm` would have errored first but for a race.
|
||||
return Ok(());
|
||||
};
|
||||
let targets: HashMap<String, u64> = stats
|
||||
.buckets
|
||||
.iter()
|
||||
.filter_map(|b| Some((b.shard_id.clone(), b.newest_generation()?)))
|
||||
.collect();
|
||||
if targets.is_empty() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
match drain_to_targets(table, &targets).await? {
|
||||
CheckpointOutcome::Done => return Ok(()),
|
||||
CheckpointOutcome::ReissueFromFlush => {
|
||||
backoff(reissue).await;
|
||||
continue;
|
||||
}
|
||||
}
|
||||
}
|
||||
Err(Error::Runtime {
|
||||
message: "checkpoint_lsm: the owning node kept losing its claim; \
|
||||
re-issued from flush the maximum number of times"
|
||||
.into(),
|
||||
})
|
||||
}
|
||||
|
||||
/// Trigger and poll until no bucket holds a generation at or below its
|
||||
/// target.
|
||||
///
|
||||
/// No liveness bound, deliberately. The pod-wide compactor pool (a semaphore
|
||||
/// of 2 by default, shared across every table on the node) is taken *inside*
|
||||
/// the pass, after the bucket latch, so a checkpoint queued behind unrelated
|
||||
/// tables is indistinguishable from one that is merging. An idle-poll counter
|
||||
/// here could only ever have fired on a table that would have finished.
|
||||
async fn drain_to_targets(
|
||||
table: &Table,
|
||||
targets: &HashMap<String, u64>,
|
||||
) -> Result<CheckpointOutcome> {
|
||||
loop {
|
||||
let stats = match issue(|| table.get_lsm_stats(false)).await? {
|
||||
Attempt::Ok(stats) => stats,
|
||||
Attempt::ReissueFromFlush => return Ok(CheckpointOutcome::ReissueFromFlush),
|
||||
};
|
||||
let Some(stats) = stats else {
|
||||
return Ok(CheckpointOutcome::Done);
|
||||
};
|
||||
// `compacting` is the bucket's compaction latch, held from dispatch
|
||||
// until the pass ends — including while it waits on the pod-wide
|
||||
// permit. So it answers one question only: do not pile on. Buckets
|
||||
// with nothing outstanding are skipped, not counted as idle.
|
||||
let mut outstanding = 0;
|
||||
let mut all_compacting = true;
|
||||
for b in &stats.buckets {
|
||||
let Some(target) = targets.get(&b.shard_id) else {
|
||||
continue;
|
||||
};
|
||||
let n = b.outstanding_generations(*target);
|
||||
if n > 0 {
|
||||
outstanding += n;
|
||||
all_compacting &= b.compacting;
|
||||
}
|
||||
}
|
||||
if outstanding == 0 {
|
||||
return Ok(CheckpointOutcome::Done);
|
||||
}
|
||||
|
||||
if !all_compacting {
|
||||
match table.compact_lsm().await {
|
||||
Ok(()) => {}
|
||||
Err(e) if is_lost_claim(&e) => return Ok(CheckpointOutcome::ReissueFromFlush),
|
||||
Err(e) if !is_retryable(&e) => return Err(e),
|
||||
// A 429 here means the server could latch no bucket at all,
|
||||
// which the poll above already handles. Not retried in place:
|
||||
// the latch it would contend for is the one doing the work, so
|
||||
// fall through and re-read — `POLL_INTERVAL` is the backoff.
|
||||
Err(_) => {}
|
||||
}
|
||||
}
|
||||
tokio::time::sleep(POLL_INTERVAL).await;
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(all(test, feature = "remote"))]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn http(status: u16) -> Error {
|
||||
Error::Http {
|
||||
source: "server said no".into(),
|
||||
request_id: "rid".into(),
|
||||
status_code: reqwest::StatusCode::from_u16(status).ok(),
|
||||
}
|
||||
}
|
||||
|
||||
/// Every status the loop acts on. The two predicates are checked together
|
||||
/// because their overlap is what would be wrong: a status must never be
|
||||
/// both, and 421 in particular must not read as retryable — retrying it in
|
||||
/// place re-issues the call that just said the node holds no claim.
|
||||
#[test]
|
||||
fn taxonomy_round_trips() {
|
||||
for status in [429, 503] {
|
||||
assert!(is_retryable(&http(status)), "{status} must retry");
|
||||
assert!(
|
||||
!is_lost_claim(&http(status)),
|
||||
"{status} is not a lost claim"
|
||||
);
|
||||
}
|
||||
assert!(is_lost_claim(&http(421)), "a lost claim must re-claim");
|
||||
assert!(
|
||||
!is_retryable(&http(421)),
|
||||
"retrying a lost claim in place only asks the same node again"
|
||||
);
|
||||
for status in [400, 404, 409, 500] {
|
||||
assert!(!is_retryable(&http(status)), "{status} is terminal");
|
||||
assert!(!is_lost_claim(&http(status)), "{status} is terminal");
|
||||
}
|
||||
}
|
||||
|
||||
/// An error carrying no status has no retry story and must be terminal —
|
||||
/// a connection that never reached the server, or a `TableNotFound` that
|
||||
/// `check_table_response` translated before the loop saw it.
|
||||
#[test]
|
||||
fn errors_without_a_status_are_terminal() {
|
||||
let no_status = Error::Http {
|
||||
source: "connection reset".into(),
|
||||
request_id: "rid".into(),
|
||||
status_code: None,
|
||||
};
|
||||
assert!(!is_retryable(&no_status));
|
||||
assert!(!is_lost_claim(&no_status));
|
||||
|
||||
let translated = Error::TableNotFound {
|
||||
name: "t".into(),
|
||||
source: "gone".into(),
|
||||
};
|
||||
assert!(!is_retryable(&translated));
|
||||
assert!(!is_lost_claim(&translated));
|
||||
}
|
||||
}
|
||||
@@ -1,162 +0,0 @@
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
// SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
//! Live per-bucket LSM state — the shape [`crate::Table::get_lsm_stats`]
|
||||
//! returns and [`super::checkpoint`] polls.
|
||||
//!
|
||||
//! Nothing here is derived: sums and differences (total L0 bytes, WAL lag)
|
||||
//! are the caller's to compute. There is no "WAL is off" shape — that case is
|
||||
//! `None`, because a struct of zeros would read as measurements.
|
||||
|
||||
use serde::Deserialize;
|
||||
|
||||
/// One flushed L0 generation.
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct GenerationStats {
|
||||
pub generation: u64,
|
||||
pub bytes: u64,
|
||||
/// Present only when `include_generation_rows` was requested. Off by
|
||||
/// default because each count opens an uncached Lance dataset, and the
|
||||
/// checkpoint loop polls this route needing only generation numbers.
|
||||
#[serde(default)]
|
||||
pub rows: Option<u64>,
|
||||
}
|
||||
|
||||
/// One in-memory memtable.
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct MemtableStats {
|
||||
pub generation: u64,
|
||||
pub rows: u64,
|
||||
pub bytes: u64,
|
||||
pub batches: u64,
|
||||
/// Names of the indexes this memtable carries. An absent name is the whole
|
||||
/// answer to "why is my fresh-tier search on that column brute-force".
|
||||
pub indexes: Vec<String>,
|
||||
}
|
||||
|
||||
/// Live state of one bucket. A table is N buckets on one node; flattening to
|
||||
/// a single number hides the one hot bucket that is usually why someone
|
||||
/// opened this endpoint.
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct BucketStats {
|
||||
pub shard_id: String,
|
||||
/// `Active` | `Sealed` (drop-table 2PC in flight).
|
||||
pub status: String,
|
||||
pub writer_epoch: u64,
|
||||
pub manifest_version: u64,
|
||||
pub current_generation: u64,
|
||||
pub replay_after_wal_entry_position: u64,
|
||||
pub wal_entry_position_last_seen: u64,
|
||||
pub generations: Vec<GenerationStats>,
|
||||
/// Whether a pass owns this bucket's compaction latch right now. Says *a*
|
||||
/// driver is running, not *whose*, and the latch is held from dispatch —
|
||||
/// including while the pass queues for a pod-wide compactor permit. Read
|
||||
/// it as "do not pile on", never as "mine is progressing".
|
||||
pub compacting: bool,
|
||||
/// Oldest first, active last. Absent for a `Sealed` bucket, whose
|
||||
/// in-memory state is torn down.
|
||||
#[serde(default)]
|
||||
pub memtables: Option<Vec<MemtableStats>>,
|
||||
}
|
||||
|
||||
impl BucketStats {
|
||||
/// The newest flushed generation, or `None` when L0 is empty.
|
||||
pub(crate) fn newest_generation(&self) -> Option<u64> {
|
||||
self.generations.iter().map(|g| g.generation).max()
|
||||
}
|
||||
|
||||
/// How many generations at or below `target` are still in L0.
|
||||
///
|
||||
/// A count, not a boolean: one pass drains a bounded prefix rather than
|
||||
/// the whole target set, so a boolean would read as "no progress" for
|
||||
/// every pass but the last. Compaction drains oldest-first, so this
|
||||
/// decreases monotonically.
|
||||
pub(crate) fn outstanding_generations(&self, target: u64) -> usize {
|
||||
self.generations
|
||||
.iter()
|
||||
.filter(|g| g.generation <= target)
|
||||
.count()
|
||||
}
|
||||
}
|
||||
|
||||
/// Live LSM state, one entry per bucket.
|
||||
#[derive(Debug, Clone, Deserialize)]
|
||||
pub struct LsmStats {
|
||||
pub buckets: Vec<BucketStats>,
|
||||
}
|
||||
|
||||
/// Server-side JSON envelope for `get_lsm_stats`. `lsm_stats` is null when
|
||||
/// the table has no LSM write path.
|
||||
#[derive(Debug, Deserialize)]
|
||||
pub(crate) struct GetLsmStatsResponse {
|
||||
#[serde(default)]
|
||||
pub lsm_stats: Option<LsmStats>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn bucket(shard: &str, generations: &[u64], compacting: bool) -> BucketStats {
|
||||
BucketStats {
|
||||
shard_id: shard.into(),
|
||||
status: "Active".into(),
|
||||
writer_epoch: 1,
|
||||
manifest_version: 1,
|
||||
current_generation: generations.iter().max().copied().unwrap_or(0) + 1,
|
||||
replay_after_wal_entry_position: 0,
|
||||
wal_entry_position_last_seen: 0,
|
||||
generations: generations
|
||||
.iter()
|
||||
.map(|g| GenerationStats {
|
||||
generation: *g,
|
||||
bytes: 1,
|
||||
rows: None,
|
||||
})
|
||||
.collect(),
|
||||
compacting,
|
||||
memtables: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// The target watermark is the newest generation at the start, and a
|
||||
/// generation created after it must not hold the loop open — that is why
|
||||
/// the predicate terminates under write load.
|
||||
#[test]
|
||||
fn newer_generations_do_not_extend_the_target() {
|
||||
let start = bucket("b0", &[7, 8], false);
|
||||
let target = start.newest_generation().expect("L0 is non-empty");
|
||||
assert_eq!(target, 8);
|
||||
|
||||
// Compaction drained 7 and 8; 9 and 10 arrived while it ran.
|
||||
let later = bucket("b0", &[9, 10], false);
|
||||
assert_eq!(
|
||||
later.outstanding_generations(target),
|
||||
0,
|
||||
"generations above the target are somebody else's problem"
|
||||
);
|
||||
|
||||
// Still holding 8 means still outstanding.
|
||||
assert_eq!(
|
||||
bucket("b0", &[8, 9], false).outstanding_generations(target),
|
||||
1
|
||||
);
|
||||
}
|
||||
|
||||
/// The metric counts generations, not buckets: a pass drains a bounded
|
||||
/// prefix, so one bucket going 3 → 2 → 1 → 0 is three steps.
|
||||
#[test]
|
||||
fn progress_is_measured_in_generations() {
|
||||
let target = 3;
|
||||
let counts: Vec<usize> = [&[1u64, 2, 3][..], &[2, 3][..], &[3][..], &[][..]]
|
||||
.iter()
|
||||
.map(|gens| bucket("b0", gens, false).outstanding_generations(target))
|
||||
.collect();
|
||||
assert_eq!(counts, vec![3, 2, 1, 0]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn empty_l0_has_no_target() {
|
||||
assert!(bucket("b0", &[], false).newest_generation().is_none());
|
||||
}
|
||||
}
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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
|
||||
// =============================================================================
|
||||
|
||||
@@ -10,7 +10,7 @@ use arrow_array::{
|
||||
use arrow_schema::{DataType, Field, Fields, Schema};
|
||||
use futures::TryStreamExt;
|
||||
use lance::Dataset;
|
||||
use lance_file::version::LanceFileVersion;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lancedb::{
|
||||
Connection, Error, Result, Table,
|
||||
blob::{BlobRangeRequest, blob},
|
||||
|
||||
Reference in New Issue
Block a user