Compare commits

...

23 Commits

Author SHA1 Message Date
Will Jones 011def461c docs(python): fix cross-references that resolved to the wrong page
`mkdocs build --strict` only catches references it cannot resolve. A bare
anchor such as `[limit][]` or `[vector search][search]` is matched by
autorefs against any heading on the site, so six of them silently linked
into the JavaScript reference instead. The relative links in
`permutation.py` and `remote/errors.py` pointed at in-page anchors and
paths that do not exist.

Targets that still exist here or in an imported inventory now use
mkdocstrings references; the guide pages deleted in #2770 use their
lancedb.com URLs.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-30 16:19:38 -07:00
Will Jones ed6be12ad6 docs: clear the mkdocs warning backlog so --strict passes
`mkdocs build` emitted 61 warnings on main, and rendering the previously
undocumented classes in this PR pushed that to 158. That backlog is what
blocks turning on strict mode (#3707), so clear it here rather than leave
it worse than we found it.

Most of it was one systematic false positive: griffe cannot see the
generated `__init__` of a pydantic dataclass, so every documented
parameter looked unknown. `warn_unknown_params` turns that check off.

The rest were real docstring bugs, in 15 docstrings:

* Prose trailing a `Parameters` section is read as parameter names, which
  invented parameters called `The`, `you` and `To`. Moved into `Notes` or
  the summary.
* numpydoc only reads a type when the colon has spaces around it. Where
  the documented name is a pydantic attribute rather than a signature
  parameter, griffe has no signature to fall back on and the type was
  dropped. Affects nine embedding classes.
* `num_partitions, default sqrt(num_rows)` and friends parse as a list of
  names, rendering a bogus `default` parameter.
* One parameter indented five spaces instead of four.

`nodejs/CONTRIBUTING.md` links to the repo-root CONTRIBUTING.md, which
does not resolve once typedoc copies the file into `docs/src/js/_media/`;
an absolute URL works from both places.

`mkdocs build --strict` now exits 0.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-29 14:15:20 -07:00
Will Jones ac2b689cdb docs(python): render index/embeddings/remote/rerankers from __all__
Four packages are now rendered by a single mkdocstrings directive each,
driven by the module's `__all__`, instead of a hand-maintained list of
symbols. These were where most of the drift was: 7 of 12 rerankers and
14 of 17 embedding functions had never been listed.

`lancedb.embeddings` had no `__all__`; without one mkdocstrings renders
no members at all for a re-export package, so one is added.

AGENTS.md gains a section describing how the reference page is wired up
and how to check a docs build locally, plus a step in the "adding a new
method on Table" checklist.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-29 13:59:23 -07:00
Will Jones 4fc8114871 docs(python): add missing public APIs to the Python reference
The Python API reference page had drifted from the public API. Branch
management (`Branches` / `AsyncBranches`, which own `diff` and `merge`),
structured full-text query classes, take queries, blob helpers,
namespace connections, most rerankers and embedding functions, the
PyTorch dataloader, and several other public symbols were never listed,
so they did not appear in the rendered docs.

Also fixes docstring cross-references that pointed at guide pages which
have since moved off this site, and at unresolvable relative targets
(`[Table](Table)`, `[PyArrow Table](pyarrow.Table)`).

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-29 13:46:45 -07:00
Will Jones 03b26d585b fix: deflake test_read_consistency_interval (#3713)
`test_read_consistency_interval` asserted that a table opened with a
100ms `read_consistency_interval` still read stale data immediately
after a concurrent write. The cache timestamp is set when the table is
opened and reads within the interval do not refresh it, so that
assertion only held if the intervening open/count/commit/count sequence
finished within 100ms of real wall-clock time. On a loaded CI runner it
did not: the TTL expired, `count_rows` refreshed synchronously, and the
test failed with `left: 1, right: 0`. This broke the Rust workflow on
`main` at 0bc08160 (a Python-only commit).

This pins the `background_cache` mock clock once `table2` has seeded its
cache, and advances it explicitly in place of `tokio::time::sleep`, so
the test controls when the interval elapses. Same approach as #3547.
With the clock pinned there is no real sleep left to be imprecise, so
the `cfg(not(target_os = "windows"))` guard is dropped and the test now
runs on Windows too.

Verified by inserting a stall before the write: 120ms reproduces the
original failure deterministically, and with this change the test still
passes with a 500ms stall.

Fixes #3712
2026-07-29 13:06:41 -07:00
Yang Cen f7feed48c3 feat(fts): support custom stop-word lists (#3734)
## What

Expose custom FTS stop-word lists in the Python and TypeScript public
APIs, including their standalone tokenize helpers and remote index
creation.

This PR supports concrete string lists only. It does not add file or
LanceDB-table stop-word sources.

## Why

Rust already exposes Lance's custom stop-word list option. The Python
and TypeScript APIs did not pass it through, and local index details did
not retain the full tokenizer parameters needed by index-backed
tokenization after reopening a table.

## How

- Add `custom_stop_words` / `customStopWords` to the Python and
TypeScript FTS and tokenize options.
- Preserve `None` / `undefined`, empty lists, and list contents without
normalization.
- Load the persisted FTS segment parameters when returning local index
details.
- Serialize the concrete list in remote create-index requests.
- Keep Python and TypeScript tests thin; behavior, persistence, query
tokenization, and remote JSON coverage live primarily in Rust.

## Validation

- `cargo check --quiet --features remote --tests --examples`
- `cargo clippy --quiet --features remote --tests --examples`
- `cargo test --quiet --features remote --tests`
- Python extension rebuild with `uv` and `maturin`
- Targeted Python tests: 4 passed
- Python `ruff format --check` and `ruff check`
- TypeScript build, typecheck, Biome lint, generated docs, and targeted
tests

---------

Co-authored-by: Yang Cen <yangcen@Yangs-Mac-mini.local>
2026-07-29 17:40:12 +08:00
Lance Release e5f489818b Bump version: 0.37.0-beta.0 → 0.37.1-beta.0 2026-07-29 07:12:34 +00:00
buduoqiu 98a52267a2 feat(python): configure streaming transform parallelism (#3699)
## Summary

- add a keyword-only `transform_parallelism` option to
`StreamingDataset`
- preserve CPU auto-detection by default and fall back to one worker
when unavailable
- apply the configured limit to both the transform executor and
concurrency semaphore
- document and test explicit, default, fallback, and invalid values

## Testing

- `uv run --extra tests --with torch pytest
python/tests/test_elastic_dataloader.py -q` (`136 passed`)
- `uvx ruff check python/lancedb/streaming.py
python/tests/test_elastic_dataloader.py`
- `uvx ruff format --check python/lancedb/streaming.py
python/tests/test_elastic_dataloader.py`
- `git diff --check origin/main...HEAD`

Closes #3695

Co-authored-by: buduoqiu <yaodong-shen@users.noreply.github.com>
2026-07-28 15:38:34 -07:00
Will Jones ff50e698cf ci: cut Actions cost by moving builds to free runners and fixing caches (#3735)
Standard GitHub-hosted runners are free on public repos, so all Actions
spend here is on the `*-8x-*` / `4x` larger runners. Measured over 30
days at current (post-Jan-2026) larger-runner rates, that is ~$1,400/mo,
and `npm-publish` is ~70% of it.

## Changes

**Fat LTO was forcing builds onto large runners.** `[profile.release]`
in `.cargo/config.toml` sets `lto = "fat"` with `codegen-units = 1`,
which is single-threaded and the peak-memory step. The macOS
`npm-publish` build was 111 of its 113 minutes in one `napi build` step,
making it the critical path of the whole publish pipeline. The ThinLTO
override already applied to Windows now covers macOS too, and both
Windows builds move from `windows-2025-8x-x64` to the free standard
`windows-2025`.

**The npm-publish cargo cache never existed.** There are zero caches
with its key prefix. The key was static, so `actions/cache` (which only
writes on a miss) could never refresh it, and a multi-GB release
`target/` per target could never fit the repo's 10 GB budget anyway. Now
caches only the crate registry, keyed on `Cargo.lock`. The docker builds
also mounted `.cargo/registry/*` while the cache saved `.cargo-cache`,
so containers re-downloaded the registry every run.

**Cache eviction thrash.** Repo cache usage is 10.4 GB against GitHub's
10 GB cap, so every PR run evicted main's warm entries. `rust.yml` and
`nodejs.yml` now restore everywhere but only save from `main`.

**npm-publish moves to nightly + tags** instead of every push to main
(~90/month). The cross-compiled targets do need watching, so
`report-failure` now fires on scheduled runs, and dedupes onto an
existing open issue rather than filing one per night.

**rust.yml aarch64-pc-windows-msvc** cross-compiled its tests and then
skipped them, paying full codegen and link cost for a compile check.
`windows-11-arm` is now GA and free on public repos, so it builds and
tests natively. Its test step also passes `--target` — without it cargo
used `target/ci/` rather than `target/<triple>/ci/` and rebuilt the
entire dependency graph a second time.

**pypi-publish.yml had no concurrency group**, so force-pushes left a
~74 minute Windows job running.

## What is cost vs. wall-clock

| Change | Cost | Wall-clock |
|---|---|---|
| Windows npm-publish → free runners | **−$570/mo** | slower per job
(8→4 cores) |
| npm-publish nightly | **−$125/mo** | — |
| pypi-publish concurrency | small | — |
| macOS ThinLTO | $0 (already free) | **−~50 min** per release |
| rust aarch64 Windows native | $0 (already free) | **−~25 min** |
| rust `--target` on test step | $0 | large, avoids a second full build
|
| rust-cache `save-if` | small | faster via real cache hits |

## Risks

- The two Windows builds now have 4 cores instead of 8 and ~14 GB of
free disk. If they fail, it is most likely disk rather than memory;
fallback is `windows-2025-4x-x64`, which still halves that line.
- `windows-11-arm` has a thinner toolset (choco/vcpkg/protoc under
emulation) and this enables a test step that has never run, so it may
surface real aarch64 failures. That is the point, but it is the change
most likely to need iteration.
- ThinLTO applies to published macOS and Windows binaries, typically
within a few percent of fat LTO. Linux release builds are untouched.

## Follow-ups

- `python.yml` `pydantic1x` (37 min) and `Doctest` (33 min) each rebuild
the extension from source via `pip install -e .` with no Rust cache;
they should consume the wheel the `linux` job already builds. Worth
~$235/mo and ~70 min of compute per run. Separate PR.
- The three `ubuntu-2404-8x-x64` npm-publish builds (~$420/mo at the old
cadence) are the remaining large-runner spend;
`aarch64-unknown-linux-gnu` could run natively on free
`ubuntu-24.04-arm`. Worth doing after this lands so the ThinLTO change
can be validated first.
- The wheel composite actions declare `python-minor-version` as required
but never use it, and every caller omits it (actionlint warns).

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

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-07-28 14:38:21 -07:00
kid b799ebaa69 fix(node): reject non-string Arrow metadata (#3728)
## Summary

- validate Arrow metadata keys and values independently at runtime
- reject malformed foreign schemas before constructing a local Arrow
schema
- cover valid and invalid metadata entries across Arrow 15–18

## Testing

- `node_modules/.bin/jest --runInBand __test__/arrow.test.ts -t "schema
metadata"`
- `node_modules/.bin/jest --runInBand __test__/arrow.test.ts`
- `node node_modules/@biomejs/biome/bin/biome format --write
lancedb/sanitize.ts __test__/arrow.test.ts`
- `pnpm lint`
- `pnpm build`
- `pnpm run docs`

Fixes #3729
2026-07-28 13:36:08 -07:00
kid 72fc660f9e feat(python): expose AsyncTable.to_lance (#3730)
## Summary

- expose the existing async Lance dataset conversion as
`AsyncTable.to_lance`
- preserve table version, branch, and refreshed storage options when
opening the dataset
- route internal async pandas/query paths through the public API
- cover normal tables, checked-out versions, branches, and forwarded
dataset options

## Testing

- `cd python && uv run --no-sync pytest python/tests/test_table.py -q`
- `cd python && uv run --no-sync pytest python/tests/test_query.py -q`
- `cd python && uv run --no-sync pytest --doctest-modules
python/lancedb/table.py -q`
- `uv run --project python --no-sync ruff format --check
python/python/lancedb/table.py python/python/lancedb/query.py
python/python/tests/test_table.py`
- `uv run --project python --no-sync ruff check .`

Fixes #1387
2026-07-28 13:31:06 -07:00
dependabot[bot] 1ebde1f06c chore(deps): bump arrow from 58.3.0 to 58.4.0 (#3722)
Bumps [arrow](https://github.com/apache/arrow-rs) from 58.3.0 to 58.4.0.
<details>
<summary>Release notes</summary>
<p><em>Sourced from <a
href="https://github.com/apache/arrow-rs/releases">arrow's
releases</a>.</em></p>
<blockquote>
<h2>arrow 58.4.0</h2>
<!-- raw HTML omitted -->
<h1>Changelog</h1>
<h2><a href="https://github.com/apache/arrow-rs/tree/58.4.0">58.4.0</a>
(2026-07-17)</h2>
<p><a
href="https://github.com/apache/arrow-rs/compare/58.3.0...58.4.0">Full
Changelog</a></p>
<p><strong>Merged pull requests:</strong></p>
<ul>
<li>[58_maintenance] [parquet] Allow more encryption algorithms (<a
href="https://redirect.github.com/apache/arrow-rs/issues/9203">#9203</a>)
<a
href="https://redirect.github.com/apache/arrow-rs/pull/10351">#10351</a>
[<a
href="https://github.com/apache/arrow-rs/labels/parquet">parquet</a>]
(<a href="https://github.com/mbutrovich">mbutrovich</a>)</li>
<li>[58_maintenance] Backport cargo audit fixes <a
href="https://redirect.github.com/apache/arrow-rs/pull/10369">#10369</a>
(<a href="https://github.com/alamb">alamb</a>)</li>
<li>[58_maintenance] chore: Ignore py03 vulnerabilities until upgrade <a
href="https://redirect.github.com/apache/arrow-rs/pull/10370">#10370</a>
(<a href="https://github.com/alamb">alamb</a>)</li>
<li>[58_maintenance] Add test for `parquet-testing/bad_data/ARROW-<a
href="https://redirect.github.com/apache/arrow-rs/issues/47662">GH-47662</a>.parquet`
(<a
href="https://redirect.github.com/apache/arrow-rs/issues/10077">#10077</a>)
<a
href="https://redirect.github.com/apache/arrow-rs/pull/10371">#10371</a>
[<a
href="https://github.com/apache/arrow-rs/labels/parquet">parquet</a>]
(<a href="https://github.com/alamb">alamb</a>)</li>
</ul>
<p>* <em>This Changelog was automatically generated by <a
href="https://github.com/github-changelog-generator/github-changelog-generator">github_changelog_generator</a></em></p>
</blockquote>
</details>
<details>
<summary>Changelog</summary>
<p><em>Sourced from <a
href="https://github.com/apache/arrow-rs/blob/58.4.0/CHANGELOG.md">arrow's
changelog</a>.</em></p>
<blockquote>
<h2><a href="https://github.com/apache/arrow-rs/tree/58.4.0">58.4.0</a>
(2026-07-17)</h2>
<p><a
href="https://github.com/apache/arrow-rs/compare/58.3.0...58.4.0">Full
Changelog</a></p>
<p><strong>Merged pull requests:</strong></p>
<ul>
<li>[58_maintenance] [parquet] Allow more encryption algorithms (<a
href="https://redirect.github.com/apache/arrow-rs/issues/9203">#9203</a>)
<a
href="https://redirect.github.com/apache/arrow-rs/pull/10351">#10351</a>
[<a
href="https://github.com/apache/arrow-rs/labels/parquet">parquet</a>]
(<a href="https://github.com/mbutrovich">mbutrovich</a>)</li>
<li>[58_maintenance] Backport cargo audit fixes <a
href="https://redirect.github.com/apache/arrow-rs/pull/10369">#10369</a>
(<a href="https://github.com/alamb">alamb</a>)</li>
<li>[58_maintenance] chore: Ignore py03 vulnerabilities until upgrade <a
href="https://redirect.github.com/apache/arrow-rs/pull/10370">#10370</a>
(<a href="https://github.com/alamb">alamb</a>)</li>
<li>[58_maintenance] Add test for `parquet-testing/bad_data/ARROW-<a
href="https://redirect.github.com/apache/arrow-rs/issues/47662">GH-47662</a>.parquet`
(<a
href="https://redirect.github.com/apache/arrow-rs/issues/10077">#10077</a>)
<a
href="https://redirect.github.com/apache/arrow-rs/pull/10371">#10371</a>
[<a
href="https://github.com/apache/arrow-rs/labels/parquet">parquet</a>]
(<a href="https://github.com/alamb">alamb</a>)</li>
</ul>
<p>* <em>This Changelog was automatically generated by <a
href="https://github.com/github-changelog-generator/github-changelog-generator">github_changelog_generator</a></em></p>
</blockquote>
</details>
<details>
<summary>Commits</summary>
<ul>
<li><a
href="https://github.com/apache/arrow-rs/commit/0ff81c1215cc026a1de93ce3d2078df1ecba6f09"><code>0ff81c1</code></a>
[58_maintenance] Update changelog for <a
href="https://redirect.github.com/apache/arrow-rs/issues/10371">#10371</a>
(<a
href="https://redirect.github.com/apache/arrow-rs/issues/10372">#10372</a>)</li>
<li><a
href="https://github.com/apache/arrow-rs/commit/95d7231227e1ce7a1ec049ab2d45a6cffd7a50f9"><code>95d7231</code></a>
[58_maintenance] Add test for `parquet-testing/bad_data/ARROW-<a
href="https://redirect.github.com/apache/arrow-rs/issues/47662">GH-47662</a>.parque...</li>
<li><a
href="https://github.com/apache/arrow-rs/commit/4544deaa434bbf8e7fe930bf9497fd36e8e737d1"><code>4544dea</code></a>
Prepare for <code>58.4.0</code> release (<a
href="https://redirect.github.com/apache/arrow-rs/issues/10367">#10367</a>)</li>
<li><a
href="https://github.com/apache/arrow-rs/commit/32e8c1809642647ddf87703c410c4713df166281"><code>32e8c18</code></a>
chore: Ignore py03 vulnerabilities until upgrade (<a
href="https://redirect.github.com/apache/arrow-rs/issues/10370">#10370</a>)</li>
<li><a
href="https://github.com/apache/arrow-rs/commit/c12030f29639f9ac36fdfedd7d00f3b6b6bbd2c1"><code>c12030f</code></a>
[58_maintenance] Backport cargo audit fixes (<a
href="https://redirect.github.com/apache/arrow-rs/issues/10369">#10369</a>)</li>
<li><a
href="https://github.com/apache/arrow-rs/commit/01046eed275d4fabfd922f6a8924410102ab1802"><code>01046ee</code></a>
[58_maintenance] [parquet] Allow more encryption algorithms (<a
href="https://redirect.github.com/apache/arrow-rs/issues/9203">#9203</a>)
(<a
href="https://redirect.github.com/apache/arrow-rs/issues/10351">#10351</a>)</li>
<li><a
href="https://github.com/apache/arrow-rs/commit/adb77a16adff42fface41664dc2a3cb564f45fcf"><code>adb77a1</code></a>
[58_maintenance] Fix MSRV CI check (pin tonic to 0.14.5, install
cargo-msrv -...</li>
<li>See full diff in <a
href="https://github.com/apache/arrow-rs/compare/58.3.0...58.4.0">compare
view</a></li>
</ul>
</details>
<br />

Signed-off-by: dependabot[bot] <support@github.com>
Co-authored-by: dependabot[bot] <49699333+dependabot[bot]@users.noreply.github.com>
2026-07-28 13:28:54 -07:00
Dan Tasse 29c030f865 fix: accept either timestamp or timestamp_millis for versions (#3733)
`list_versions()` against a remote table on a server that uses
lance-namespace was failing. The server was returning
`timestamp_millis`, while db-catalog deployments were returning
`timestamp`, and the client was only accepting `timestamp`. So, updated
the client to accept both. (assuming we're migrating over time;
eventually we can turn off the `timestamp` code path I suppose.)
2026-07-28 10:44:49 -04:00
Xuanwo ff6ff09998 feat: support batched blob range reads (#3703)
## Summary

Lance can now plan multiple byte ranges for the same blob in one
`read_blob_ranges` operation, but LanceDB users currently cannot expose
a complete set of logical ranges to that planner.

This complements `BlobFile`: file-like consumers such as PyAV can
continue to discover ranges dynamically, while callers that already know
the ranges for a batch can submit them together.

## Motivating example

A training table may store a large video blob together with a small
application-level clip index:

```text
video: blob
clips: [{offset, length}, ...]
```

The caller can select the videos and clips for a batch, obtain their row
IDs from the query, and read all of the selected windows together:

```python
rows = (
    table.search()
    .select(["clips"])
    .with_row_id(True)
    .limit(64)
    .to_arrow()
    .to_pylist()
)

requests = []
for row in rows:
    clip = sample_clip(row["clips"])
    requests.append(
        (row["_rowid"], clip["offset"], clip["length"])
    )

chunks = table.fetch_blob_ranges("video", requests)
```

Here, `_rowid` comes from the LanceDB query, while `offset` and `length`
come from the application's clip index and are relative to that row's
video blob. The caller describes only the logical reads; Lance still
handles validation, source grouping, coalescing, scheduling, and byte
backpressure.

Lance v10.0.0-beta.5 returns one logical result per blob selector or
range request and explicitly distinguishes null blobs from valid empty
values. LanceDB consumes that aligned result contract directly and only
adds a cardinality check for unresolved row IDs.

This PR exposes batched blob-range reads on local Rust and Python
tables. Results preserve request identity, duplicates, null slots, and
valid empty ranges while allowing Lance to execute the physical reads
out of order. Scheduler buffer sizing remains an internal Lance concern,
so the LanceDB API does not expose `io_buffer_size`.

Cloud tables continue to report this operation as unsupported until
there is a corresponding remote API.
2026-07-27 15:08:47 -07:00
Vivek 119b9baf90 fix: preserve row count in MetadataEraserExec for zero-column batches (#3717)
SELECT COUNT(*) FROM t WHERE <predicate> — and any query that plans an
empty-projection scan — panics the executing query task:

InvalidArgumentError("must either specify a row count or at least one
column")

Root cause

MetadataEraserExec wraps every LanceDB table scan to strip schema-level
metadata, rebuilding each batch in execute():

RecordBatch::try_new(schema.clone(), batch.columns().to_vec()).unwrap()

RecordBatch::try_new infers the row count from the columns. COUNT(*)
with a filter is planned with an empty projection, so the scan emits
zero-column batches — there are no columns to infer a length from,
try_new returns Err, and the .unwrap() panics.

(This is specific to the empty-projection case: COUNT(*) with no filter
is answered from statistics and never scans, and COUNT(<col>) projects a
column — both already work.)
2026-07-27 12:25:11 -07:00
LanceDB Robot ba4558a64f chore: update lance dependency to v10.0.0-beta.5 (#3718)
Updates the Rust workspace Lance dependencies and Java lance-core
dependency to v10.0.0-beta.5. No compatibility fixes were required;
full-workspace Clippy passes with warnings denied. Lance tag:
https://github.com/lance-format/lance/releases/tag/v10.0.0-beta.5
2026-07-27 15:27:40 +08:00
Heng Ge f655f62e09 feat(query): add use_lsm to read MemWAL LSM data (#3489)
## What

MemWAL LSM **read** support. When a table has an LSM write spec
(`set_lsm_write_spec`), `merge_insert` upserts live in the MemWAL
active/frozen memtables and flushed SSTables until an external
compaction merges them into the base table, so a normal scan returns
**stale** data. This routes reads through Lance's `LsmScanner` so
queries also surface that in-flight data, deduplicated by primary key
(newest generation wins).

## How

- Adds a **`use_lsm: Option<bool>`** query flag, symmetric with the
`merge_insert` flag:
- **unset** — auto-route through the LSM scanner when the table carries
a write spec
- **`use_lsm(true)`** — force the LSM path; error if there is no spec
    - **`use_lsm(false)`** — read the base table only (the escape hatch)
- Plain scan, single-column full-text search, and single-vector ANN all
run through one `LsmScanner` (assembled from on-disk shard manifests
plus the cached writer's in-memory memtables), so a `where` predicate is
honored as a **prefilter** uniformly — including for vector search.
- **Compaction-aware snapshots:** an SSTable generation is dropped only
once it is both compacted into the base table and covered by the arm's
base-index catch-up (`index_catchup`); plain scans use the compaction
watermark alone.
- Query shapes the scanner cannot honor hard-error with guidance to set
`use_lsm(false)`: hybrid, multi/binary vectors, `with_row_id`,
reranking, `order_by`, dynamic/Substrait projection or filters,
`distance_range`, `use_index(false)`, postfilter, take-by-row-id/offset,
reads from a time-traveled version, and an unmaintained or ambiguous
(multiple) FTS/vector index. Namespace-pushdown queries fall back to
local execution when a spec is present; WAL-only writers are handled.
- Exposed across the Rust core and the Python (`use_lsm`) and TypeScript
(`useLsm`) bindings, including `TakeQuery`.

Rebased from Lance `7.2.0-beta.3` to `10.0.0-beta.3`.
2026-07-25 23:45:27 -07:00
Will Jones bf15655c83 chore: unify SDK versions and release tags on a single line (#3714)
Python was versioned and tagged separately from the Rust, Java, and
Node.js SDKs, and had drifted three minor versions ahead (0.36 vs 0.33).
Users had no way to tell which Python version corresponded to which Rust
or Node release, and the gap had no meaning behind it.

This unifies the two tracks so there is one version and one tag for all
four SDKs.

## Version

The shared version is set to `0.37.0-beta.0`. Python continues its own
sequence (highest published: 0.36 → 0.37) while Rust, Java, and Node.js
jump 0.33 → 0.37 to meet it. Picking Python's next minor means Python
users see no discontinuity at all, and only the other SDKs skip forward.

Note that `main` trails the `release/v0.32` branch on both lines (main
is at 0.32.0-beta.3 / 0.35.0-beta.3; the release branch carries
0.33.0-beta.0 / 0.36.0-beta.0), so 0.37 is chosen to clear the highest
tag on either branch. Every index stays monotonic:

| index | publishes | last published | next |
|---|---|---|---|
| PyPI | stable only | 0.34.0 | 0.37.0 |
| Fury | previews | 0.36.0b0 | 0.37.0-beta.1 |
| npm | both | 0.33.0-beta.0 | 0.37.0-beta.1 |
| crates.io | stable only | 0.31.0 | 0.37.0 |
| Maven | both | 0.33.0-beta.0 | 0.37.0-beta.1 |

A one-time jump for three SDKs, versus explaining the offset
indefinitely.

## Mechanism

* `python/.bumpversion.toml` is removed. `python/Cargo.toml` — the
source of the Python package version, since `pyproject.toml` declares
`dynamic = ["version"]` — becomes a tracked file of the root config. Its
`cargo update -p lancedb-python` pre-commit hook is dropped as
redundant: `ci/update_lockfiles.sh` already refreshes every workspace
member version in `Cargo.lock`.
* `pypi-publish.yml` triggers on `v*` instead of `python-v*`, so one tag
releases all four packages. `ci/bump_version.sh` and
`make-release-commit.yml` lose their now-dead tag-prefix and
per-language plumbing, including the `python` / `other` dispatch inputs.
* The two byte-identical GH release jobs in `npm-publish.yml` and
`pypi-publish.yml` are replaced by a single `gh-release.yml`. One
release per tag, named `LanceDB vX.Y.Z`, instead of separate "Python
LanceDB" and "Node/Rust LanceDB" releases for the same commit.

The trade-off: there is no longer a way to ship a Python-only patch
without also releasing crates.io, Maven, and npm. That is the cost of
making drift structurally impossible.

## Beta releases marked "Latest" (#3666)

Both GH release jobs used:

```yaml
prerelease: ${{ contains('beta', github.ref) }}
```

The arguments are reversed. `contains(search, item)` asks whether
*`search`* contains *`item`*, so this evaluated "does the literal string
`'beta'` contain `refs/tags/python-v0.35.0-beta.2`?" — always `false`.
Every beta was published as a full release, and GitHub awards "Latest"
to the newest non-prerelease.

The new workflow derives the flag from the parsed version rather than
the raw ref, and sets `make_latest` explicitly:

```yaml
prerelease: ${{ steps.extract_version.outputs.prerelease }}
make_latest: ${{ steps.extract_version.outputs.prerelease == 'false' }}
```

npm was never affected (`--tag preview` uses correct bash), and PyPI
already excludes pre-releases from resolution.

This only fixes releases published from here on. Already-published betas
need a one-time backfill:

```shell
gh api --paginate /repos/lancedb/lancedb/releases \
  --jq '.[] | select(.prerelease == false) | select(.tag_name | test("beta")) | .id' \
  | xargs -I{} gh api -X PATCH /repos/lancedb/lancedb/releases/{} -F prerelease=true
```

## Verification

Ran `ci/bump_version.sh` end-to-end against this branch with the release
tooling installed:

* `preview` → tags `v0.37.0-beta.1` (previous tag `v0.33.0-beta.0`
detected, `pre_n` bump)
* `stable` → tags `v0.37.0`
* Both paths update `.bumpversion.toml`, `rust/lancedb/Cargo.toml`,
`nodejs/Cargo.toml`, `python/Cargo.toml`, `nodejs/package.json`, the 7
`nodejs/npm/*/package.json` files, both Java poms, and
`docs/src/java/java.md` together
* `check_breaking_changes.py` resolves the last stable as `v0.31.0`, so
the minor-version gate passes

All five touched workflows parse as valid YAML and the pre-commit hooks
pass.

## Notes for review

* This targets `main` only, so it takes effect at the next
release-branch cut. The in-flight `release/v0.32` branch still carries
`v0.33.0-beta.0` / `python-v0.36.0-beta.0`; if we want the imminent
stable to be 0.37.0, this needs to be applied there too.
* Historical `python-v*` tags are left alone. The changelog builder
scans `^v`, which does not match them, so the first unified release's
notes will compute `fromTag` from the Rust/Node line only — a one-time
gap in the Python-side changelog.
* Pre-existing and not addressed here: `ci/update_lockfiles.sh --amend`
amends the commit that `bump-my-version` has already tagged, so the
lockfile update lands outside the tag on stable releases.

Fixes #3666
2026-07-25 09:22:35 -07:00
Lance Release a00edef0e6 Bump version: 0.32.0-beta.2 → 0.32.0-beta.3 2026-07-24 22:04:02 +00:00
Lance Release 1b2670443e Bump version: 0.35.0-beta.2 → 0.35.0-beta.3 2026-07-24 22:03:30 +00:00
Yang Cen 9dc5ec03aa feat(fts): add block size configuration (#3691)
## What changed

- add `block_size` to Python FTS configuration and the deprecated
local/remote helpers
- add `blockSize` to the TypeScript FTS options and propagate it through
the NAPI binding
- serialize the value as `block_size` for remote index creation
- document the existing Rust builder API and generate the TypeScript API
reference
- add local, remote, metadata, search, and invalid-value regression
coverage

## Why

Lance supports configuring the number of documents per compressed FTS
posting block, but LanceDB's Python and TypeScript APIs did not expose
the setting. This made the experimental FTS V3 layout unavailable
through those clients and allowed the value to be dropped before index
creation.

## How it works

The default remains `128`. Supported values are `128` and `256`;
selecting `256` uses the experimental FTS V3 format. Invalid values are
rejected by the Lance builder and surfaced as Python or JavaScript
errors.

## Validation

- `cargo check --quiet --features remote --tests --examples`
- `cargo +1.94.0 clippy --quiet --features remote --tests --examples --
-D warnings`
- targeted Rust local and remote index tests
- Rust doctests: 34 passed
- Python Ruff checks, doctest, and targeted local/remote tests: 5 passed
- TypeScript build, Biome lint, generated docs, and targeted Jest tests:
9 passed
- `git diff --check`

## Limitations

The Java client remains unchanged because its external remote REST model
does not currently expose `block_size`.

Co-authored-by: Yang Cen <yangcen@Yangs-Mac-mini.local>
2026-07-24 15:02:38 -07:00
Andrew Chen 18760f74cd fix: crash in AnswerdotaiRerankers/ColbertReranker for return_score="all" (#3671)
## What

`AnswerdotaiRerankers(return_score="all").rerank_hybrid(...)` (and
`ColbertReranker`, which subclasses it without overriding
`rerank_hybrid`) raises:

```
pyarrow.lib.ArrowInvalid: Invalid sort key column: No match for FieldRef.Name(_relevance_score) in _rowid: int64 ...
```

## Why

```python
combined_results = self.merge_results(vector_results, fts_results)
combined_results = self._rerank(combined_results, query)
if self.score == "relevance":
    combined_results = self._keep_relevance_score(combined_results)
elif self.score == "all":
    combined_results = self._merge_and_keep_scores(vector_results, fts_results)
```

When `score == "all"`, `combined_results` is unconditionally overwritten
by `_merge_and_keep_scores(vector_results, fts_results)` **after**
`_rerank()` already computed and appended `_relevance_score` —
discarding it. The following `sort_by("_relevance_score", ...)` then has
nothing to sort on.

Every sibling reranker that supports `return_score="all"`
(`cross_encoder`, `openai`, `cohere`, `jinaai`, `voyageai`, `watsonx`)
instead calls `_merge_and_keep_scores()` **before** `_rerank()`. This
file is the one place the ordering got inverted when `"all"` support was
added (#2509) — a copy/paste inconsistency across the six files that PR
touched. Fix mirrors the pattern already used (and tested) by the other
five rerankers.

Also drops the now-stale `"Only 'relevance' is supported for now"`
docstring line on both classes, left over from before `"all"` support
existed.

## Testing

Added `test_answerdotai_reranker_return_all`, mirroring the existing
`test_cross_encoder_reranker_return_all`. Verified locally with the real
built Rust extension: red (reproduces the exact `ArrowInvalid` above) →
green, using the actual `rerank_hybrid`/`_rerank`/`base.py` code path
with the model call mocked out — my local environment's
`rerankers==0.10.0` fails to load the real ColBERT model against the
available `transformers` version (`AttributeError: 'ColBERTModel' object
has no attribute 'all_tied_weights_keys'`), which I confirmed also
breaks the **pre-existing**, unmodified
`test_colbert_reranker`/`test_answerdotai_reranker` baseline tests
identically — an unrelated local dependency-version issue, not a
regression from this change. `ruff check`/`ruff format` clean; full
`test_rerankers.py` run: 9 passed / 8 skipped / 3 failed (the 3 failures
are exactly those two pre-existing tests plus my new one, all failing at
model-loading time for the same unrelated reason before reaching the
changed code).

---
Disclosure: this PR was drafted with AI assistance (Claude); I reviewed,
tested, and take responsibility for the change.

Co-authored-by: Claude Opus 4.8 <noreply@anthropic.com>
2026-07-24 15:02:23 -07:00
LanceDB Robot c9d07ef6fc chore: update lance dependency to v10.0.0-beta.3 (#3710)
Updates the Rust workspace and Java lance-core dependencies to [Lance
v10.0.0-beta.3](https://github.com/lance-format/lance/releases/tag/v10.0.0-beta.3).

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

Some files were not shown because too many files have changed in this diff Show More