Compare commits

...

5 Commits

Author SHA1 Message Date
Gatefixer da86d804ba test(node): cover tables across database connections 2026-08-08 12:48:40 +00:00
Dan Tasse 77a93fee76 fix: get table size from metadata, not files (#3790)
Some issues:
- file_size_bytes is optional in the manifest, so if it's not there (old
writer I guess) it'll under-report the table size.
- it changes results a little bit from the old way by including per-file
footers and metadata (probably not a big difference at real scale)

---------

Co-authored-by: Will Jones <willjones127@gmail.com>
2026-08-07 17:41:41 -04:00
Lance Release 7bb501839a Bump version: 0.37.1-beta.0 → 0.37.1-beta.1 2026-08-07 21:16:07 +00:00
Andrew Chen 5b347afd99 fix: avoid AttributeError in JinaEmbeddings image input for str/Path (#3670)
## What

`JinaEmbeddings._generate_image_input_dict()` crashes with
`AttributeError: 'function' object has no attribute 'urlparse'` on any
image given as a URL string, local path string, or `pathlib.Path` — i.e.
every documented `jina-clip-v1` image-embedding use case except raw
`bytes`.

## Why

```python
from urllib.parse import urlparse
...
parsed = urlparse.urlparse(image)
```

`urlparse` is imported as a function, then called as if it were the
`urllib.parse` module (`urlparse.urlparse(...)`). The module-level
`is_valid_url()` a few lines above does it correctly (`urlparse(text)`),
which is why this reads as a typo rather than intentional. Fixed to
`urlparse(str(image))` — `str()` is needed because `urlparse()` only
accepts `str`/`bytes` and raises a different `AttributeError` on a raw
`Path`.

## Testing

Added `test_jina_generate_image_input_dict_local_path`, which fails with
the original `AttributeError` before the fix and passes after, covering
both a `str` path and a `pathlib.Path`. Verified locally (built the Rust
extension, ran red→green, then the full `test_embeddings.py` file: 15
passed / 8 skipped, no regressions) and with `ruff check`/`ruff format`.

---
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-08-07 14:05:45 -07:00
Dan Rammer 706a9c327f feat: infer maintained indexes when an LsmWriteSpec omits them (#3748)
## What

`LsmWriteSpec::maintained_indexes` becomes `Option<Vec<String>>`:

| value | meaning |
|---|---|
| `None` (new default) | every index the MemWAL supports, resolved when
the spec is installed |
| `Some([])` | maintain nothing — a scan/filter-only WAL table |
| `Some([..])` | exactly these, taken verbatim |

`with_maintained_indexes` keeps its signature;
`with_no_maintained_indexes()` is new. Surfaced through the remote path
(null on the wire), Python, and Node.

## Why

Callers had to state the maintained set by hand every time, which is
both tedious and easy to get wrong — the common case is "maintain what I
already built."

Resolution filters on `IndexConfig::is_memwal_maintainable`, delegating
to lance's `is_maintainable_index_type`. This is load-bearing rather
than cosmetic: lance does **not** skip an index type its memtable cannot
build, it errors when the shard writer opens, so sweeping up a bitmap
index would fail every memtable claim and leave the table unwritable.
The inferred set excludes those, and an explicit list naming one is now
rejected at spec time instead of at claim time.

## Behavior change

A freshly constructed spec used to maintain **nothing**; it now
maintains **everything supported**. This flipped because napi collapses
`undefined` and `null` to `None`, so TypeScript cannot express "absent
means nothing, null means all" — any other choice makes the bindings
disagree with the wire. The error direction also favors it: an unwanted
maintained index costs memory, while a silently unmaintained one
degrades FTS to an unscored scan.

Three existing tests encoded the old default and are updated rather than
worked around.

## Caveat

The resolved set is a snapshot, not a subscription. An index created
after the spec is installed is not maintained until the spec is unset
and set again. `get_lsm_write_spec` therefore always reports a concrete
list — `None` never round-trips.

## Dependency

Needs a lance release carrying `is_maintainable_index_type`
(lance-format/lance#8095) before this builds against the pinned tag.
Draft until then.

## Testing

38 Rust LSM tests and 10 Python tests pass against a local lance build,
including new coverage that a bitmap index is excluded from inference
and rejected when named, and that `[]` stays distinguishable from null
on the wire.

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

---------

Co-authored-by: Claude Opus 5 (1M context) <noreply@anthropic.com>
2026-08-07 14:50:22 -05:00
36 changed files with 684 additions and 127 deletions
+1 -1
View File
@@ -1,5 +1,5 @@
[tool.bumpversion]
current_version = "0.37.1-beta.0"
current_version = "0.37.1-beta.1"
parse = """(?x)
(?P<major>0|[1-9]\\d*)\\.
(?P<minor>0|[1-9]\\d*)\\.
Generated
+4 -3
View File
@@ -5411,7 +5411,7 @@ dependencies = [
[[package]]
name = "lancedb"
version = "0.37.1-beta.0"
version = "0.37.1-beta.1"
dependencies = [
"ahash",
"anyhow",
@@ -5480,6 +5480,7 @@ dependencies = [
"random_word",
"regex",
"reqwest 0.12.28",
"roaring",
"rstest",
"semver",
"serde",
@@ -5499,7 +5500,7 @@ dependencies = [
[[package]]
name = "lancedb-nodejs"
version = "0.37.1-beta.0"
version = "0.37.1-beta.1"
dependencies = [
"arrow-array",
"arrow-buffer",
@@ -5524,7 +5525,7 @@ dependencies = [
[[package]]
name = "lancedb-python"
version = "0.37.1-beta.0"
version = "0.37.1-beta.1"
dependencies = [
"arrow",
"async-trait",
+1 -1
View File
@@ -14,7 +14,7 @@ Add the following dependency to your `pom.xml`:
<dependency>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-core</artifactId>
<version>0.37.1-beta.0</version>
<version>0.37.1-beta.1</version>
</dependency>
```
+9 -3
View File
@@ -431,9 +431,10 @@ Read the [LsmWriteSpec](../interfaces/LsmWriteSpec.md) currently installed on th
Resolves to `undefined` when the MemWAL LSM write path is not enabled (no
spec has been set, or it was removed with [Table#unsetLsmWriteSpec](Table.md#unsetlsmwritespec)).
The returned spec — including its `maintainedIndexes` and
`writerConfigDefaults` — mirrors what was passed to
[Table#setLsmWriteSpec](Table.md#setlsmwritespec).
The returned spec mirrors what was passed to
[Table#setLsmWriteSpec](Table.md#setlsmwritespec), except that `maintainedIndexes` always
reports the concrete list resolved when the spec was set — `undefined`
never round-trips.
#### Returns
@@ -806,6 +807,11 @@ All variants require the table to have an unenforced primary key
([Table#setUnenforcedPrimaryKey](Table.md#setunenforcedprimarykey)); bucket sharding additionally
requires it to be the single column being bucketed.
Omitting `maintainedIndexes` maintains every index on the table, resolved
here, failing if one cannot be maintained — name them to install anyway.
Naming them pins an exact set, and a still-building index is rejected
rather than quietly omitted.
#### Parameters
* **spec**: [`LsmWriteSpec`](../interfaces/LsmWriteSpec.md)
+3 -1
View File
@@ -34,7 +34,9 @@ Bucket and identity variants: the sharding column.
optional maintainedIndexes: string[];
```
Names of indexes the MemWAL should keep up to date during writes.
Indexes the MemWAL keeps up to date. Omit to maintain every supported
index, resolved on install — a snapshot, so indexes created later are not
maintained. Pass `[]` for none.
***
+4 -1
View File
@@ -44,4 +44,7 @@ The number of rows in the table
totalBytes: number;
```
The total number of bytes in the table
The total size, in bytes, of the table's data files, index files, and
overlay files
Read from the manifest, so this excludes deletion files and manifests.
+1 -1
View File
@@ -8,7 +8,7 @@
<parent>
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.37.1-beta.0</version>
<version>0.37.1-beta.1</version>
<relativePath>../pom.xml</relativePath>
</parent>
+1 -1
View File
@@ -6,7 +6,7 @@
<groupId>com.lancedb</groupId>
<artifactId>lancedb-parent</artifactId>
<version>0.37.1-beta.0</version>
<version>0.37.1-beta.1</version>
<packaging>pom</packaging>
<name>${project.artifactId}</name>
<description>LanceDB Java SDK Parent POM</description>
+1 -1
View File
@@ -1,7 +1,7 @@
[package]
name = "lancedb-nodejs"
edition.workspace = true
version = "0.37.1-beta.0"
version = "0.37.1-beta.1"
publish = false
license.workspace = true
description.workspace = true
+27
View File
@@ -69,6 +69,33 @@ describe("given a connection", () => {
await expect(tbl.countRows()).resolves.toBe(1);
});
it("should isolate object-form table creation across databases", async () => {
const otherTmpDir = tmp.dirSync({ unsafeCleanup: true });
const otherDb = await connect(otherTmpDir.name);
try {
const firstTable = await db.createTable({
name: "defaultTable",
data: [{ rowId: "id1", vector: Array(384).fill(0) }],
});
const secondTable = await otherDb.createTable({
name: "defaultTable",
data: [{ rowId: "id2", vector: Array(384).fill(0) }],
});
await expect(db.tableNames()).resolves.toEqual(["defaultTable"]);
await expect(otherDb.tableNames()).resolves.toEqual(["defaultTable"]);
const firstRows = await firstTable.query().select(["rowId"]).toArray();
const secondRows = await secondTable.query().select(["rowId"]).toArray();
expect(firstRows.map((row) => row.rowId)).toEqual(["id1"]);
expect(secondRows.map((row) => row.rowId)).toEqual(["id2"]);
} finally {
otherDb.close();
otherTmpDir.removeCallback();
}
});
it("should be able to drop tables`", async () => {
await db.createTable("test", [{ id: 1 }, { id: 2 }]);
await db.createTable("test2", [{ id: 1 }, { id: 2 }]);
+9 -1
View File
@@ -277,8 +277,16 @@ describe.each([arrow15, arrow16, arrow17, arrow18])(
},
numIndices: 0,
numRows: 3,
totalBytes: 44,
// Full on-disk size of the two data files, footers and metadata included.
totalBytes: 684,
});
// Index files count toward totalBytes too (only deletion files and
// manifests are excluded).
await table.createIndex("id", { config: Index.btree() });
const statsWithIndex = await table.stats();
expect(statsWithIndex.numIndices).toBe(1);
expect(statsWithIndex.totalBytes).toBeGreaterThan(684);
});
it("should overwrite data if asked", async () => {
+14 -4
View File
@@ -197,7 +197,11 @@ export interface LsmWriteSpec {
column?: string;
/** Bucket variant: the number of buckets, in `[1, 1024]`. */
numBuckets?: number;
/** Names of indexes the MemWAL should keep up to date during writes. */
/**
* Indexes the MemWAL keeps up to date. Omit to maintain every supported
* index, resolved on install — a snapshot, so indexes created later are not
* maintained. Pass `[]` for none.
*/
maintainedIndexes?: string[];
/** Default `ShardWriter` configuration recorded in the MemWAL index. */
writerConfigDefaults?: Record<string, string>;
@@ -595,6 +599,11 @@ export abstract class Table {
* All variants require the table to have an unenforced primary key
* ({@link Table#setUnenforcedPrimaryKey}); bucket sharding additionally
* requires it to be the single column being bucketed.
*
* Omitting `maintainedIndexes` maintains every index on the table, resolved
* here, failing if one cannot be maintained — name them to install anyway.
* Naming them pins an exact set, and a still-building index is rejected
* rather than quietly omitted.
* @param {LsmWriteSpec} spec The sharding spec to install.
* @returns {Promise<void>}
* @example
@@ -622,9 +631,10 @@ export abstract class Table {
*
* Resolves to `undefined` when the MemWAL LSM write path is not enabled (no
* spec has been set, or it was removed with {@link Table#unsetLsmWriteSpec}).
* The returned spec — including its `maintainedIndexes` and
* `writerConfigDefaults` — mirrors what was passed to
* {@link Table#setLsmWriteSpec}.
* The returned spec mirrors what was passed to
* {@link Table#setLsmWriteSpec}, except that `maintainedIndexes` always
* reports the concrete list resolved when the spec was set — `undefined`
* never round-trips.
* @returns {Promise<LsmWriteSpec | undefined>}
*/
abstract getLsmWriteSpec(): Promise<LsmWriteSpec | undefined>;
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-darwin-arm64",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"os": ["darwin"],
"cpu": ["arm64"],
"main": "lancedb.darwin-arm64.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-gnu",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-arm64-musl",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"os": ["linux"],
"cpu": ["arm64"],
"main": "lancedb.linux-arm64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-gnu",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-gnu.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-linux-x64-musl",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"os": ["linux"],
"cpu": ["x64"],
"main": "lancedb.linux-x64-musl.node",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-arm64-msvc",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"os": [
"win32"
],
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "@lancedb/lancedb-win32-x64-msvc",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"os": ["win32"],
"cpu": ["x64"],
"main": "lancedb.win32-x64-msvc.node",
+2 -2
View File
@@ -1,12 +1,12 @@
{
"name": "@lancedb/lancedb",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "@lancedb/lancedb",
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"cpu": [
"x64",
"arm64"
+1 -1
View File
@@ -11,7 +11,7 @@
"ann"
],
"private": false,
"version": "0.37.1-beta.0",
"version": "0.37.1-beta.1",
"main": "dist/index.js",
"exports": {
".": "./dist/index.js",
+10 -7
View File
@@ -772,7 +772,8 @@ pub struct LsmWriteSpec {
pub column: Option<String>,
/// Bucket variant: the number of buckets, in `[1, 1024]`.
pub num_buckets: Option<u32>,
/// Names of indexes the MemWAL should keep up to date during writes.
/// Indexes the MemWAL keeps up to date. Omitted resolves every
/// maintainable index on install; an empty array means none.
pub maintained_indexes: Option<Vec<String>>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
pub writer_config_defaults: Option<HashMap<String, String>>,
@@ -782,7 +783,6 @@ impl TryFrom<LsmWriteSpec> for lancedb::table::LsmWriteSpec {
type Error = napi::Error;
fn try_from(value: LsmWriteSpec) -> napi::Result<Self> {
let maintained = value.maintained_indexes.unwrap_or_default();
let writer_config_defaults = value.writer_config_defaults.unwrap_or_default();
let spec = match value.spec_type.as_str() {
"bucket" => {
@@ -809,7 +809,7 @@ impl TryFrom<LsmWriteSpec> for lancedb::table::LsmWriteSpec {
}
};
Ok(spec
.with_maintained_indexes(maintained)
.with_maintained_indexes(value.maintained_indexes)
.with_writer_config_defaults(writer_config_defaults))
}
}
@@ -827,7 +827,7 @@ impl From<lancedb::table::LsmWriteSpec> for LsmWriteSpec {
spec_type: "bucket".to_string(),
column: Some(column),
num_buckets: Some(num_buckets),
maintained_indexes: Some(maintained_indexes),
maintained_indexes,
writer_config_defaults: Some(writer_config_defaults),
},
Native::Identity {
@@ -838,7 +838,7 @@ impl From<lancedb::table::LsmWriteSpec> for LsmWriteSpec {
spec_type: "identity".to_string(),
column: Some(column),
num_buckets: None,
maintained_indexes: Some(maintained_indexes),
maintained_indexes,
writer_config_defaults: Some(writer_config_defaults),
},
Native::Unsharded {
@@ -848,7 +848,7 @@ impl From<lancedb::table::LsmWriteSpec> for LsmWriteSpec {
spec_type: "unsharded".to_string(),
column: None,
num_buckets: None,
maintained_indexes: Some(maintained_indexes),
maintained_indexes,
writer_config_defaults: Some(writer_config_defaults),
},
}
@@ -1043,7 +1043,10 @@ impl From<lancedb::index::IndexStatistics> for IndexStatistics {
#[napi(object)]
pub struct TableStatistics {
/// The total number of bytes in the table
/// The total size, in bytes, of the table's data files, index files, and
/// overlay files
///
/// Read from the manifest, so this excludes deletion files and manifests.
pub total_bytes: i64,
/// The number of rows in the table
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb-python"
version = "0.37.1-beta.0"
version = "0.37.1-beta.1"
publish = false
edition.workspace = true
description = "Python bindings for LanceDB"
+7 -4
View File
@@ -653,9 +653,10 @@ class LsmWriteSpec:
def identity(column: str) -> "LsmWriteSpec": ...
@staticmethod
def unsharded() -> "LsmWriteSpec": ...
def with_maintained_indexes(self, indexes: List[str]) -> "LsmWriteSpec":
"""Return a copy of this spec asking the MemWAL to keep the named
indexes up to date as rows are appended."""
def with_maintained_indexes(self, indexes: Optional[List[str]]) -> "LsmWriteSpec":
"""Set which indexes the MemWAL keeps up to date. None resolves every
index on the table at install, failing if one cannot be maintained;
a list is verbatim, empty means none."""
...
def with_writer_config_defaults(self, defaults: Dict[str, str]) -> "LsmWriteSpec":
"""Return a copy of this spec recording the given default
@@ -670,7 +671,9 @@ class LsmWriteSpec:
@property
def num_buckets(self) -> Optional[int]: ...
@property
def maintained_indexes(self) -> List[str]: ...
def maintained_indexes(self) -> Optional[List[str]]:
"""Indexes the MemWAL keeps up to date, or None for every supported one."""
...
@property
def writer_config_defaults(self) -> Dict[str, str]: ...
+4 -3
View File
@@ -87,12 +87,13 @@ class JinaEmbeddings(EmbeddingFunction):
if isinstance(image, bytes):
image_dict = {"image": base64.b64encode(image).decode("utf-8")}
elif isinstance(image, (str, Path)):
parsed = urlparse.urlparse(image)
# TODO handle drive letter on windows.
parsed = urlparse(str(image))
PIL_Image = attempt_import_or_raise("PIL.Image", "pillow")
if parsed.scheme == "file":
pil_image = PIL_Image.open(parsed.path)
elif parsed.scheme == "":
elif parsed.scheme == "" or (os.name == "nt" and len(parsed.scheme) == 1):
# A Windows drive letter parses as a one-character scheme
# ("C:\\img.png" -> scheme="c"), so treat it as a local path.
pil_image = PIL_Image.open(image if os.name == "nt" else parsed.path)
elif parsed.scheme.startswith("http"):
pil_image = PIL_Image.open(io.BytesIO(url_retrieve(image)))
+13 -4
View File
@@ -4676,6 +4676,13 @@ class AsyncTable:
via [`set_unenforced_primary_key`]; bucket sharding additionally
requires it to be the single column being bucketed.
By default the MemWAL maintains every index on the table, resolved
here — a snapshot, so an index created afterwards needs the spec unset
and set again. This fails if one cannot be maintained; name the set
with ``with_maintained_indexes`` to install anyway. That pins an exact
set (a still-building index is rejected, not omitted); ``[]`` maintains
none.
Parameters
----------
spec : LsmWriteSpec
@@ -4702,9 +4709,9 @@ class AsyncTable:
Returns ``None`` when the MemWAL LSM write path is not enabled (no
spec has been set, or it was removed with `unset_lsm_write_spec`).
The returned spec — including its ``maintained_indexes`` and
``writer_config_defaults`` — mirrors what was passed to
`set_lsm_write_spec`.
The returned spec mirrors what was passed to `set_lsm_write_spec`,
except that ``maintained_indexes`` always reports the concrete list
resolved when the spec was set — ``None`` never round-trips.
"""
return await self._inner.get_lsm_write_spec()
@@ -6334,7 +6341,9 @@ class TableStatistics:
Attributes
----------
total_bytes: int
The total number of bytes in the table.
The total size, in bytes, of the table's data files, index files, and
overlay files. Read from the manifest, so this excludes deletion files
and manifests.
num_rows: int
The total number of rows in the table.
num_indices: int
+20
View File
@@ -631,3 +631,23 @@ def test_url_retrieve_downloads_image():
image_bytes = url_retrieve(image_url)
img = Image.open(io.BytesIO(image_bytes))
assert img.size[0] > 0 and img.size[1] > 0
def test_jina_generate_image_input_dict_local_path(tmp_path):
"""
JinaEmbeddings._generate_image_input_dict must accept a local image path
(str or Path), not just bytes. Previously it crashed with
`AttributeError: 'function' object has no attribute 'urlparse'` on any
str/Path input because it called `urlparse.urlparse(image)` instead of
`urlparse(image)` (urlparse was imported as a function, not a module).
"""
Image = pytest.importorskip("PIL.Image")
from lancedb.embeddings.jinaai import JinaEmbeddings
image_path = tmp_path / "test.png"
Image.new("RGB", (4, 4), color="red").save(image_path, format="PNG")
for image in (str(image_path), image_path):
image_dict = JinaEmbeddings._generate_image_input_dict(image)
assert "image" in image_dict
assert isinstance(image_dict["image"], str) and len(image_dict["image"]) > 0
+11 -4
View File
@@ -83,7 +83,9 @@ def test_lsm_write_spec_repr():
assert s.spec_type == "bucket"
assert s.column == "id"
assert s.num_buckets == 4
assert s.maintained_indexes == []
# A fresh spec defers its maintained set to install time.
assert s.maintained_indexes is None
assert s.with_maintained_indexes([]).maintained_indexes == []
assert "bucket" in repr(s)
assert "id" in repr(s)
assert "4" in repr(s)
@@ -169,18 +171,23 @@ def test_get_lsm_write_spec(tmp_path):
table.unset_lsm_write_spec()
assert table.get_lsm_write_spec() is None
# Identity round-trips (column recovered from the schema).
# Identity round-trips (column recovered from the schema). Leaving the
# maintained set to be inferred picks up the index on the table, so the
# spec reads back naming it rather than as "infer".
table.set_lsm_write_spec(LsmWriteSpec.identity("id"))
spec = table.get_lsm_write_spec()
assert spec.spec_type == "identity"
assert spec.column == "id"
assert spec.maintained_indexes == [idx_name]
table.unset_lsm_write_spec()
# Unsharded round-trips (no routing column).
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
# Unsharded round-trips (no routing column). Opting out is distinct from
# the inferred default.
table.set_lsm_write_spec(LsmWriteSpec.unsharded().with_maintained_indexes([]))
spec = table.get_lsm_write_spec()
assert spec.spec_type == "unsharded"
assert spec.column is None
assert spec.maintained_indexes == []
@pytest.mark.asyncio
+2 -2
View File
@@ -544,7 +544,7 @@ def test_lsm_read_fts_unmaintained_index_errors(tmp_path):
table.create_index("text", config=FTS())
# No maintained indexes: the active memtable FTS arm cannot serve un-compacted
# docs, so the search would silently omit them — reject instead.
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
table.set_lsm_write_spec(LsmWriteSpec.unsharded().with_maintained_indexes([]))
with pytest.raises(Exception, match="maintained"):
table.search("fox", query_type="fts", fts_columns="text").to_arrow()
@@ -631,7 +631,7 @@ def test_lsm_read_vector_unmaintained_index_errors(tmp_path):
)
# Spec with NO maintained indexes: the base vector index's catch-up is untracked,
# so the scanner rejects rather than risk dropping compacted-but-unindexed rows.
table.set_lsm_write_spec(LsmWriteSpec.unsharded())
table.set_lsm_write_spec(LsmWriteSpec.unsharded().with_maintained_indexes([]))
with pytest.raises(Exception, match="maintained"):
table.search([1.0] * VECTOR_DIM).to_arrow()
+9 -1
View File
@@ -3713,7 +3713,8 @@ def test_stats(mem_db: DBConnection):
stats = table.stats()
print(f"{stats=}")
assert stats == {
"total_bytes": 60,
# Full on-disk size of the data file, footer and metadata included.
"total_bytes": 633,
"num_rows": 2,
"num_indices": 0,
"fragment_stats": {
@@ -3731,6 +3732,13 @@ def test_stats(mem_db: DBConnection):
},
}
# Index files count toward total_bytes too (only deletion files and
# manifests are excluded).
table.create_index("id", config=BTree())
stats_with_index = table.stats()
assert stats_with_index["num_indices"] == 1
assert stats_with_index["total_bytes"] > stats["total_bytes"]
def test_create_table_empty_list_with_schema(mem_db: DBConnection):
"""Test creating table with empty list data and schema
+31 -15
View File
@@ -246,12 +246,22 @@ impl From<lancedb::table::MergeResult> for MergeResult {
}
}
/// Render for `__repr__`, so the default reads as Python's `None` rather than
/// Rust's `Some([..])`.
fn fmt_maintained(maintained: &Option<Vec<String>>) -> String {
match maintained {
Some(names) => format!("{:?}", names),
None => "None".to_string(),
}
}
/// Specification selecting Lance's MemWAL LSM-style write path for
/// `merge_insert`.
///
/// Constructed via the `bucket(...)`, `identity(...)`, or `unsharded()`
/// classmethods, then optionally chain `with_maintained_indexes(...)` and
/// `with_writer_config_defaults(...)`.
/// `with_writer_config_defaults(...)`. A fresh spec maintains every index the
/// MemWAL supports, resolved on install.
#[pyclass(from_py_object)]
#[derive(Clone, Debug)]
pub struct LsmWriteSpec {
@@ -291,11 +301,11 @@ impl LsmWriteSpec {
}
}
/// Replace the list of indexes the MemWAL should keep up to date as
/// rows are appended. Each name must reference an index that
/// already exists on the table at the time `set_lsm_write_spec`
/// is called.
pub fn with_maintained_indexes(&self, indexes: Vec<String>) -> Self {
/// Set which indexes the MemWAL maintains. `None` (the default)
/// resolves every supported index on install; a list is verbatim,
/// and an empty list maintains nothing.
#[pyo3(signature = (indexes))]
pub fn with_maintained_indexes(&self, indexes: Option<Vec<String>>) -> Self {
Self {
inner: self.inner.clone().with_maintained_indexes(indexes),
}
@@ -317,23 +327,29 @@ impl LsmWriteSpec {
maintained_indexes,
writer_config_defaults,
} => format!(
"LsmWriteSpec.bucket(column={:?}, num_buckets={}, maintained_indexes={:?}, writer_config_defaults={:?})",
column, num_buckets, maintained_indexes, writer_config_defaults,
"LsmWriteSpec.bucket(column={:?}, num_buckets={}, maintained_indexes={}, writer_config_defaults={:?})",
column,
num_buckets,
fmt_maintained(maintained_indexes),
writer_config_defaults,
),
lancedb::table::LsmWriteSpec::Identity {
column,
maintained_indexes,
writer_config_defaults,
} => format!(
"LsmWriteSpec.identity(column={:?}, maintained_indexes={:?}, writer_config_defaults={:?})",
column, maintained_indexes, writer_config_defaults,
"LsmWriteSpec.identity(column={:?}, maintained_indexes={}, writer_config_defaults={:?})",
column,
fmt_maintained(maintained_indexes),
writer_config_defaults,
),
lancedb::table::LsmWriteSpec::Unsharded {
maintained_indexes,
writer_config_defaults,
} => format!(
"LsmWriteSpec.unsharded(maintained_indexes={:?}, writer_config_defaults={:?})",
maintained_indexes, writer_config_defaults,
"LsmWriteSpec.unsharded(maintained_indexes={}, writer_config_defaults={:?})",
fmt_maintained(maintained_indexes),
writer_config_defaults,
),
}
}
@@ -368,10 +384,10 @@ impl LsmWriteSpec {
}
}
/// Names of indexes the MemWAL should keep up to date during writes.
/// Indexes the MemWAL keeps up to date, or `None` for every supported one.
#[getter]
pub fn maintained_indexes(&self) -> Vec<String> {
self.inner.maintained_indexes().to_vec()
pub fn maintained_indexes(&self) -> Option<Vec<String>> {
self.inner.maintained_indexes().map(<[String]>::to_vec)
}
/// Default `ShardWriter` configuration recorded by this spec.
+2 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "lancedb"
version = "0.37.1-beta.0"
version = "0.37.1-beta.1"
edition.workspace = true
description = "LanceDB: A serverless, low-latency vector database for AI applications"
license.workspace = true
@@ -100,6 +100,7 @@ anyhow = "1"
lance-testing = { workspace = true }
tempfile = "3.5.0"
random_word = { version = "0.4.3", features = ["en"] }
roaring = "0.11.4"
tokio = { version = "1.23", features = ["io-util", "macros", "net", "rt-multi-thread", "sync", "test-util"] }
uuid = { version = "1.7.0", features = ["v4"] }
walkdir = "2"
+24 -6
View File
@@ -2520,9 +2520,9 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
self.check_mutable().await?;
// Map the spec onto the server's request DTO. `sharding` is internally
// tagged on `mode` to mirror sophon's `Sharding` enum; `maintained_indexes`
// and `writer_config_defaults` are sent verbatim (an empty list means "no
// maintained indexes", not "default to all").
// tagged on `mode` to mirror sophon's `Sharding` enum. A null
// `maintained_indexes` asks the server to resolve every maintainable
// index at HEAD; a list is verbatim, an empty one meaning none.
let sharding = match &spec {
LsmWriteSpec::Bucket {
column,
@@ -6599,7 +6599,7 @@ mod tests {
.unwrap()
});
let spec = crate::table::LsmWriteSpec::unsharded()
.with_maintained_indexes(["id_idx"])
.with_maintained_indexes(vec!["id_idx".to_string()])
.with_writer_config_defaults([("max_memtable_rows", "1000")]);
table.set_lsm_write_spec(spec).await.unwrap();
}
@@ -6618,7 +6618,8 @@ mod tests {
body["sharding"],
serde_json::json!({ "mode": "bucket", "column": "id", "num_buckets": 16 })
);
assert_eq!(body["maintained_indexes"], serde_json::json!([]));
// An unpinned maintained set sends null: resolve server-side.
assert_eq!(body["maintained_indexes"], serde_json::Value::Null);
http::Response::builder().status(200).body("{}").unwrap()
});
table
@@ -6627,6 +6628,23 @@ mod tests {
.unwrap();
}
/// `[]` (none) must stay distinguishable on the wire from null (all).
#[tokio::test]
async fn test_set_lsm_write_spec_no_maintained_indexes() {
let table = Table::new_with_handler("my_table", |request| {
let body = request.body().unwrap().as_bytes().unwrap();
let body: serde_json::Value = serde_json::from_slice(body).unwrap();
assert_eq!(body["maintained_indexes"], serde_json::json!([]));
http::Response::builder().status(200).body("{}").unwrap()
});
table
.set_lsm_write_spec(
crate::table::LsmWriteSpec::bucket("id", 16).with_maintained_indexes(Vec::new()),
)
.await
.unwrap();
}
#[tokio::test]
async fn test_set_lsm_write_spec_identity() {
let table = Table::new_with_handler("my_table", |request| {
@@ -6701,7 +6719,7 @@ mod tests {
} => {
assert_eq!(column, "id");
assert_eq!(num_buckets, 4);
assert_eq!(maintained_indexes, vec!["id_idx".to_string()]);
assert_eq!(maintained_indexes, Some(vec!["id_idx".to_string()]));
assert_eq!(
writer_config_defaults
.get("durable_write")
+389 -41
View File
@@ -95,7 +95,6 @@ pub use delete::DeleteResult;
use futures::future::join_all;
pub use lance::dataset::refs::{BranchContents, Ref, TagContents, Tags as LanceTags};
pub use lance::dataset::scanner::DatasetRecordBatchStream;
use lance::dataset::statistics::DatasetStatisticsExt;
pub use lance_index::optimize::OptimizeOptions;
pub use lsm_stats::{BucketStats, GenerationStats, LsmStats, MemtableStats};
pub use optimize::{CompactionOptions, OptimizeAction, OptimizeStats};
@@ -371,6 +370,8 @@ pub use self::merge::MergeResult;
/// date) and [`LsmWriteSpec::with_writer_config_defaults`] (default
/// `ShardWriter` configuration recorded in the MemWAL index).
///
/// A fresh spec maintains every index on the table, resolved on install.
///
/// Install a spec with [`Table::set_lsm_write_spec`] and remove it with
/// [`Table::unset_lsm_write_spec`]. The actual `merge_insert` dispatch
/// onto the MemWAL writer is a follow-up.
@@ -385,9 +386,12 @@ pub enum LsmWriteSpec {
Bucket {
column: String,
num_buckets: u32,
/// Names of indexes (already created on the table) that the
/// MemWAL should maintain in-memory as rows are appended.
maintained_indexes: Vec<String>,
/// Indexes the MemWAL maintains in-memory as rows are appended.
///
/// `None` means every index it can maintain, resolved on install — a
/// snapshot, so indexes created later need the spec unset and re-set.
/// `Some([])` maintains nothing.
maintained_indexes: Option<Vec<String>>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
writer_config_defaults: HashMap<String, String>,
},
@@ -397,35 +401,41 @@ pub enum LsmWriteSpec {
/// distinct value of `column` becomes its own shard.
Identity {
column: String,
/// Names of indexes (already created on the table) that the
/// MemWAL should maintain in-memory as rows are appended.
maintained_indexes: Vec<String>,
/// Indexes the MemWAL maintains in-memory as rows are appended.
///
/// `None` means every index it can maintain, resolved on install — a
/// snapshot, so indexes created later need the spec unset and re-set.
/// `Some([])` maintains nothing.
maintained_indexes: Option<Vec<String>>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
writer_config_defaults: HashMap<String, String>,
},
/// No sharding — every `merge_insert` call writes to a single MemWAL shard.
Unsharded {
/// Names of indexes (already created on the table) that the
/// MemWAL should maintain in-memory as rows are appended.
maintained_indexes: Vec<String>,
/// Indexes the MemWAL maintains in-memory as rows are appended.
///
/// `None` means every index it can maintain, resolved on install — a
/// snapshot, so indexes created later need the spec unset and re-set.
/// `Some([])` maintains nothing.
maintained_indexes: Option<Vec<String>>,
/// Default `ShardWriter` configuration recorded in the MemWAL index.
writer_config_defaults: HashMap<String, String>,
},
}
impl LsmWriteSpec {
/// Construct a hash-bucket sharding spec with no maintained indexes.
/// Construct a hash-bucket sharding spec maintaining every index on the table.
pub fn bucket(column: impl Into<String>, num_buckets: u32) -> Self {
Self::Bucket {
column: column.into(),
num_buckets,
maintained_indexes: Vec::new(),
maintained_indexes: None,
writer_config_defaults: HashMap::new(),
}
}
/// Construct an identity-sharding spec (shard by the raw value of
/// `column`) with no maintained indexes.
/// `column`) maintaining every index on the table.
///
/// `column` must be a deterministic function of the unenforced primary
/// key: every row with a given primary key must always produce the same
@@ -437,28 +447,37 @@ impl LsmWriteSpec {
pub fn identity(column: impl Into<String>) -> Self {
Self::Identity {
column: column.into(),
maintained_indexes: Vec::new(),
maintained_indexes: None,
writer_config_defaults: HashMap::new(),
}
}
/// Construct an unsharded spec with no maintained indexes.
/// Construct an unsharded spec maintaining every index on the table.
pub fn unsharded() -> Self {
Self::Unsharded {
maintained_indexes: Vec::new(),
maintained_indexes: None,
writer_config_defaults: HashMap::new(),
}
}
/// Replace the list of indexes the MemWAL should keep up to date as
/// rows are appended. Each name must reference an index that already
/// exists on the table at the time `set_lsm_write_spec` is called.
pub fn with_maintained_indexes<I, S>(mut self, indexes: I) -> Self
where
I: IntoIterator<Item = S>,
S: Into<String>,
{
let v: Vec<String> = indexes.into_iter().map(Into::into).collect();
/// Set which indexes the MemWAL maintains.
///
/// `None` (the default) resolves to every index on the table at install,
/// failing if one cannot be maintained — name the set to install anyway. A
/// list is verbatim: each name must already exist and be maintainable, and
/// an empty list maintains nothing.
///
/// ```
/// # use lancedb::table::LsmWriteSpec;
/// // Every index the table has when the spec is installed:
/// LsmWriteSpec::unsharded().with_maintained_indexes(None);
/// // Exactly these:
/// LsmWriteSpec::unsharded().with_maintained_indexes(vec!["id_idx".to_string()]);
/// // None at all:
/// LsmWriteSpec::unsharded().with_maintained_indexes(Vec::new());
/// ```
pub fn with_maintained_indexes(mut self, indexes: impl Into<Option<Vec<String>>>) -> Self {
let indexes = indexes.into();
match &mut self {
Self::Bucket {
maintained_indexes, ..
@@ -468,7 +487,7 @@ impl LsmWriteSpec {
}
| Self::Unsharded {
maintained_indexes, ..
} => *maintained_indexes = v,
} => *maintained_indexes = indexes,
}
self
}
@@ -504,8 +523,9 @@ impl LsmWriteSpec {
self
}
/// Borrow the list of index names this spec asks MemWAL to maintain.
pub fn maintained_indexes(&self) -> &[String] {
/// Borrow the list of index names this spec asks MemWAL to maintain, or
/// `None` when it asks for every index on the table.
pub fn maintained_indexes(&self) -> Option<&[String]> {
match self {
Self::Bucket {
maintained_indexes, ..
@@ -515,7 +535,7 @@ impl LsmWriteSpec {
}
| Self::Unsharded {
maintained_indexes, ..
} => maintained_indexes,
} => maintained_indexes.as_deref(),
}
}
@@ -1713,7 +1733,7 @@ impl Table {
/// # async fn example(table: &Table) -> Result<(), Box<dyn std::error::Error>> {
/// table
/// .set_lsm_write_spec(
/// LsmWriteSpec::bucket("id", 16).with_maintained_indexes(["id_idx"]),
/// LsmWriteSpec::bucket("id", 16).with_maintained_indexes(vec!["id_idx".to_string()]),
/// )
/// .await?;
/// # Ok(())
@@ -1735,9 +1755,10 @@ impl Table {
///
/// Returns `Ok(None)` when the MemWAL LSM write path is not enabled (no
/// spec has been set, or it was removed with [`Table::unset_lsm_write_spec`]).
/// The returned spec — including its [`LsmWriteSpec::maintained_indexes`] and
/// [`LsmWriteSpec::writer_config_defaults`] — mirrors what was passed to
/// [`Table::set_lsm_write_spec`].
/// The returned spec mirrors what was passed to
/// [`Table::set_lsm_write_spec`], except that
/// [`LsmWriteSpec::maintained_indexes`] always reports the concrete list
/// resolved when the spec was set — `None` never round-trips.
///
/// # Example
///
@@ -3548,9 +3569,24 @@ impl BaseTable for NativeTable {
let num_rows = self.count_rows(None).await?;
let num_indices = self.list_indices().await?.len();
let ds = self.dataset.get().await?;
let ds_clone = (*ds).clone();
let ds_stats = Arc::new(ds_clone).calculate_data_stats().await?;
let total_bytes = ds_stats.fields.iter().map(|f| f.bytes_on_disk).sum::<u64>() as usize;
// Sizes come from the manifest. Summing per-field `bytes_on_disk` instead
// would open every data file to read its column metadata, which costs one
// IO per fragment and reports 0 for legacy v1 storage.
//
// The manifest summary covers only the fragments' base data files, so
// overlay files (recorded on each fragment) and index files (recorded in
// the manifest's index section) are added separately.
let mut total_bytes = ds.manifest().summary().total_files_size as usize;
for frag in ds.manifest().fragments.iter() {
for overlay in &frag.overlays {
if let Some(size) = overlay.data_file.file_size_bytes.get() {
total_bytes += size.get() as usize;
}
}
}
for index in ds.load_indices().await?.iter() {
total_bytes += index.total_size_bytes().unwrap_or(0) as usize;
}
let frags = ds.get_fragments();
let mut sorted_sizes = join_all(
@@ -3622,7 +3658,12 @@ impl BaseTable for NativeTable {
#[skip_serializing_none]
#[derive(Debug, Deserialize, PartialEq)]
pub struct TableStatistics {
/// The total number of bytes in the table
/// The total size, in bytes, of the table's data files, index files, and
/// overlay files
///
/// Read from the manifest, so this excludes deletion files and manifests,
/// and it excludes any file whose size the manifest does not record
/// (tables and indices written before writers persisted file sizes).
pub total_bytes: usize,
/// The number of rows in the table
@@ -3683,6 +3724,7 @@ mod tests {
use super::*;
use crate::connect;
use crate::connection::ConnectBuilder;
use crate::io::object_store::io_tracking::IoTrackingStore;
use crate::query::Select;
use crate::query::{ExecutableQuery, QueryBase};
use crate::test_utils::connection::new_test_connection;
@@ -5065,7 +5107,7 @@ mod tests {
// Bucket spec round-trips exactly, including the routing column (recovered
// from its field id), maintained indexes, and writer config defaults.
let spec = LsmWriteSpec::bucket("id", 4)
.with_maintained_indexes([idx_name])
.with_maintained_indexes(vec![idx_name.clone()])
.with_writer_config_defaults([("durable_write", "false")]);
table.set_lsm_write_spec(spec.clone()).await.unwrap();
assert_eq!(table.get_lsm_write_spec().await.unwrap(), Some(spec));
@@ -5075,15 +5117,125 @@ mod tests {
assert_eq!(table.get_lsm_write_spec().await.unwrap(), None);
// Identity sharding round-trips (column recovered from the schema).
// A spec left at its default maintains every index on the table, so it
// reads back naming the one on the table rather than as "infer".
let spec = LsmWriteSpec::identity("region");
table.set_lsm_write_spec(spec.clone()).await.unwrap();
assert_eq!(table.get_lsm_write_spec().await.unwrap(), Some(spec));
assert_eq!(
table.get_lsm_write_spec().await.unwrap(),
Some(spec.with_maintained_indexes(vec![idx_name.clone()]))
);
table.unset_lsm_write_spec().await.unwrap();
// Unsharded round-trips (no routing column).
let spec = LsmWriteSpec::unsharded();
table.set_lsm_write_spec(spec.clone()).await.unwrap();
assert_eq!(table.get_lsm_write_spec().await.unwrap(), Some(spec));
assert_eq!(
table.get_lsm_write_spec().await.unwrap(),
Some(spec.with_maintained_indexes(vec![idx_name]))
);
}
/// The maintained set defaults to every index on the table, resolved at
/// install. An index the memtable cannot build fails the install rather
/// than being dropped: maintaining it would take the table offline for
/// writes, dropping it would hide that from the caller.
#[tokio::test]
async fn test_set_lsm_write_spec_infers_maintained_indexes() {
let tmp_dir = tempdir().unwrap();
let uri = tmp_dir.path().to_str().unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int64, false),
Field::new("tag", DataType::Utf8, true),
]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(arrow_array::Int64Array::from(vec![1, 2, 3])),
Arc::new(StringArray::from(vec!["a", "b", "c"])),
],
)
.unwrap();
let reader: Box<dyn arrow_array::RecordBatchReader + Send> =
Box::new(RecordBatchIterator::new(vec![Ok(batch)], schema.clone()));
let conn = ConnectBuilder::new(uri)
.read_consistency_interval(Duration::from_secs(0))
.execute()
.await
.unwrap();
let table = conn.create_table("t", reader).execute().await.unwrap();
table
.create_index(&["id"], Index::BTree(Default::default()))
.name("id_btree".to_string())
.execute()
.await
.unwrap();
table
.create_index(&["tag"], Index::Bitmap(Default::default()))
.name("tag_bitmap".to_string())
.execute()
.await
.unwrap();
// Explicitly naming the bitmap index fails before anything commits.
let err = table
.set_lsm_write_spec(
LsmWriteSpec::unsharded().with_maintained_indexes(vec!["tag_bitmap".to_string()]),
)
.await
.unwrap_err();
assert!(
matches!(err, Error::InvalidInput { ref message } if message.contains("tag_bitmap")),
"expected the bitmap index to be rejected, got {err:?}"
);
assert_eq!(table.get_lsm_write_spec().await.unwrap(), None);
// The default covers every index, so the bitmap fails it too.
let err = table
.set_lsm_write_spec(LsmWriteSpec::unsharded())
.await
.unwrap_err();
assert!(
matches!(err, Error::InvalidInput { ref message }
if message.contains("tag_bitmap") && message.contains("maintained_indexes")),
"expected the inferred set to be rejected, got {err:?}"
);
assert_eq!(table.get_lsm_write_spec().await.unwrap(), None);
// Naming the maintainable subset installs.
table
.set_lsm_write_spec(
LsmWriteSpec::unsharded().with_maintained_indexes(vec!["id_btree".to_string()]),
)
.await
.unwrap();
assert_eq!(
table
.get_lsm_write_spec()
.await
.unwrap()
.unwrap()
.maintained_indexes(),
Some(["id_btree".to_string()].as_slice())
);
// Opting out entirely is distinct from the default.
table.unset_lsm_write_spec().await.unwrap();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes(Vec::new()))
.await
.unwrap();
assert_eq!(
table
.get_lsm_write_spec()
.await
.unwrap()
.unwrap()
.maintained_indexes(),
Some([].as_slice())
);
}
#[tokio::test]
@@ -5131,12 +5283,16 @@ mod tests {
let res = table.stats().await.unwrap();
println!("{:#?}", res);
// `total_bytes` is the full on-disk size of the 11 data files (this table
// has no index or overlay files), so it is well above the 2000 bytes of
// column data these 250 int32 pairs hold: each file carries its own footer
// and metadata.
assert_eq!(
res,
TableStatistics {
num_rows: 250,
num_indices: 0,
total_bytes: 2300,
total_bytes: 8925,
fragment_stats: FragmentStatistics {
num_fragments: 11,
num_small_fragments: 11,
@@ -5176,4 +5332,196 @@ mod tests {
}
)
}
/// `total_bytes` counts more than the base data files: index files and
/// overlay files recorded in the manifest are included too.
#[tokio::test]
pub async fn test_stats_includes_index_and_overlay_files() {
use lance::dataset::WriteDestination;
use lance::dataset::transaction::{DataOverlayGroup, Operation};
use lance_file::version::{ConcreteFileVersion, LanceFileVersion};
use lance_file::writer::FileWriterOptions;
use lance_io::utils::CachedFileSize;
use lance_table::format::DataFile;
use lance_table::format::overlay::{DataOverlayFile, OverlayCoverage};
use roaring::RoaringBitmap;
let tmp_dir = tempdir().unwrap();
let uri = tmp_dir.path().to_str().unwrap();
let conn = ConnectBuilder::new(uri)
.read_consistency_interval(Duration::from_secs(0))
.execute()
.await
.unwrap();
let schema = Arc::new(Schema::new(vec![
Field::new("id", DataType::Int32, false),
Field::new("foo", DataType::Int32, true),
]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![
Arc::new(Int32Array::from_iter_values(0..100)),
Arc::new(Int32Array::from_iter_values(0..100)),
],
)
.unwrap();
let table = conn
.create_table("test_stats_extra_files", batch)
.execute()
.await
.unwrap();
let data_only = table.stats().await.unwrap().total_bytes;
assert!(data_only > 0);
// A scalar index adds index files whose sizes are recorded in the
// manifest's index section.
table
.create_index(&["id"], Index::Auto)
.execute()
.await
.unwrap();
let with_index = table.stats().await.unwrap().total_bytes;
let dataset = {
let native = table.as_native().unwrap();
(*native.dataset.get().await.unwrap()).clone()
};
let index_bytes: usize = dataset
.load_indices()
.await
.unwrap()
.iter()
.map(|idx| idx.total_size_bytes().unwrap_or(0) as usize)
.sum();
assert!(index_bytes > 0);
assert_eq!(with_index, data_only + index_bytes);
// Commit an overlay file supplying new `foo` values for the first three
// rows of fragment 0. There is no high-level API that writes overlays
// yet, so write the overlay's data file and commit the `DataOverlay`
// operation by hand.
let read_version = dataset.version().version;
let fragment_id = dataset.get_fragments()[0].id() as u64;
let foo_field_id = dataset.schema().field("foo").unwrap().id;
let overlay_schema = dataset.schema().project_by_ids(&[foo_field_id], true);
let file_version = ConcreteFileVersion::from(LanceFileVersion::Stable);
let filename = "overlay.lance".to_string();
let store = dataset.object_store(None).await.unwrap();
let path = dataset.data_dir().child(filename.clone());
let obj_writer = store.create(&path).await.unwrap();
let mut writer = lance_file::versions::create_writer(
file_version,
obj_writer,
overlay_schema,
FileWriterOptions::default(),
)
.unwrap();
writer
.write_column(0, Arc::new(Int32Array::from(vec![1000, 1001, 1002])) as _)
.await
.unwrap();
let summary = writer.finish().await.unwrap();
let overlay_bytes = summary.size_bytes as usize;
assert!(overlay_bytes > 0);
let mut data_file = DataFile::new_unstarted(filename, file_version);
data_file.fields = writer
.field_id_to_column_indices()
.iter()
.map(|(field_id, _)| *field_id as i32)
.collect::<Vec<_>>()
.into();
data_file.column_indices = writer
.field_id_to_column_indices()
.iter()
.map(|(_, column_index)| *column_index as i32)
.collect::<Vec<_>>()
.into();
data_file.file_size_bytes = CachedFileSize::new(summary.size_bytes);
let overlay = DataOverlayFile {
data_file,
coverage: OverlayCoverage::dense(RoaringBitmap::from_iter(0..3)),
committed_version: 0,
};
Dataset::commit(
WriteDestination::Dataset(Arc::new(dataset)),
Operation::DataOverlay {
groups: vec![DataOverlayGroup {
fragment_id,
overlays: vec![overlay],
}],
},
Some(read_version),
None,
None,
Arc::new(Default::default()),
false,
)
.await
.unwrap();
table.checkout_latest().await.unwrap();
let with_overlay = table.stats().await.unwrap().total_bytes;
assert_eq!(with_overlay, with_index + overlay_bytes);
}
/// `stats()` must stay manifest-only. Summing per-field `bytes_on_disk`
/// instead opens every data file, so cost would grow with fragment count.
#[tokio::test]
pub async fn test_stats_does_not_read_data_files() {
let tmp_dir = tempdir().unwrap();
let uri = tmp_dir.path().to_str().unwrap();
let conn = ConnectBuilder::new(uri).execute().await.unwrap();
let schema = Arc::new(Schema::new(vec![Field::new("id", DataType::Int32, false)]));
let batch = RecordBatch::try_new(
schema.clone(),
vec![Arc::new(Int32Array::from_iter_values(0..10))],
)
.unwrap();
conn.create_table("test_stats_io", batch.clone())
.execute()
.await
.unwrap();
let table = conn.open_table("test_stats_io").execute().await.unwrap();
const NUM_APPENDS: usize = 20;
for _ in 0..NUM_APPENDS {
table.add(batch.clone()).execute().await.unwrap();
}
// Reopen through a tracking store so the counters cover `stats()` alone and
// not the writes above.
let (wrapper, io_stats) = IoTrackingStore::new_wrapper();
let table = conn
.open_table("test_stats_io")
.lance_read_params(ReadParams {
store_options: Some(ObjectStoreParams {
object_store_wrapper: Some(wrapper),
..Default::default()
}),
..Default::default()
})
.execute()
.await
.unwrap();
io_stats.lock().unwrap().read_iops = 0;
let stats = table.stats().await.unwrap();
let read_iops = io_stats.lock().unwrap().read_iops;
assert_eq!(stats.fragment_stats.num_fragments, NUM_APPENDS + 1);
assert!(stats.total_bytes > 0);
// Reading the fragments' data files would take at least one IOP each.
assert!(
read_iops < stats.fragment_stats.num_fragments as u64,
"stats() issued {} read IOPs across {} fragments",
read_iops,
stats.fragment_stats.num_fragments
);
}
}
+2 -2
View File
@@ -1161,7 +1161,7 @@ mod lsm_tests {
.unwrap();
let fts_index = table.list_indices().await.unwrap()[0].name.clone();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes([fts_index]))
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes(vec![fts_index]))
.await
.unwrap();
@@ -1254,7 +1254,7 @@ mod lsm_tests {
.unwrap();
let vec_index = table.list_indices().await.unwrap()[0].name.clone();
table
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes([vec_index]))
.set_lsm_write_spec(LsmWriteSpec::unsharded().with_maintained_indexes(vec![vec_index]))
.await
.unwrap();
+74 -8
View File
@@ -29,6 +29,7 @@ use arrow_schema::{DataType, Schema as ArrowSchema, SchemaRef};
use lance::Dataset;
use lance::dataset::mem_wal::{
DatasetMemWalExt, ShardWriter, ShardWriterConfig, evaluate_sharding_spec,
validate_maintained_indexes,
};
use lance::index::DatasetIndexExt;
use lance_core::datatypes::Schema as LanceSchema;
@@ -37,8 +38,9 @@ use tokio::sync::RwLock;
use uuid::Uuid;
use crate::error::{Error, Result};
use crate::index::IndexConfig;
use crate::table::merge::{MergeInsertBuilder, MergeResult};
use crate::table::{LsmWriteSpec, NativeTable};
use crate::table::{BaseTable, LsmWriteSpec, NativeTable};
/// Spec id of the sole sharding spec installed by [`set_lsm_write_spec`].
/// Must match Lance's `InitializeMemWalBuilder` (`SHARDING_SPEC_ID`).
@@ -80,32 +82,44 @@ pub(crate) async fn set_lsm_write_spec(table: &NativeTable, spec: LsmWriteSpec)
}
}
// Before the builder borrows the dataset clone. `list_indices` merges an
// index's segments into one entry, so the result needs no dedup.
let maintained_indexes = {
let dataset = table.dataset.get().await?;
resolve_maintained_indexes(
&dataset,
&table.list_indices().await?,
spec.maintained_indexes(),
)
.await?
};
let mut dataset = (*table.dataset.get().await?).clone();
let mut builder = dataset.initialize_mem_wal();
let (maintained_indexes, writer_config_defaults) = match spec {
let writer_config_defaults = match spec {
LsmWriteSpec::Bucket {
column,
num_buckets,
maintained_indexes,
writer_config_defaults,
..
} => {
builder = builder.bucket_sharding(column, num_buckets);
(maintained_indexes, writer_config_defaults)
writer_config_defaults
}
LsmWriteSpec::Identity {
column,
maintained_indexes,
writer_config_defaults,
..
} => {
builder = builder.identity_sharding(column);
(maintained_indexes, writer_config_defaults)
writer_config_defaults
}
LsmWriteSpec::Unsharded {
maintained_indexes,
writer_config_defaults,
..
} => {
builder = builder.unsharded();
(maintained_indexes, writer_config_defaults)
writer_config_defaults
}
};
builder = builder.maintained_indexes(maintained_indexes);
@@ -117,6 +131,58 @@ pub(crate) async fn set_lsm_write_spec(table: &NativeTable, spec: LsmWriteSpec)
Ok(())
}
/// Resolve a spec's maintained-index selection against `indices`, as reported
/// by [`Table::list_indices`](crate::Table::list_indices).
///
/// `None` means every index on the table, snapshotted now. Lance validates
/// either selection against its shard-writer rules, so a spec that installs is
/// one the MemWAL can open.
///
/// An unmaintainable index fails an inferred set rather than being dropped from
/// it — dropping would leave the caller believing it is maintained.
async fn resolve_maintained_indexes(
dataset: &Dataset,
indices: &[IndexConfig],
requested: Option<&[String]>,
) -> Result<Vec<String>> {
let Some(requested) = requested else {
let all: Vec<String> = indices.iter().map(|index| index.name.clone()).collect();
validate_maintained_indexes(dataset, &all)
.await
.map_err(|source| Error::InvalidInput {
message: format!(
"cannot maintain every index on this table: {source}. Set \
maintained_indexes explicitly to choose from {}",
index_name_list(indices),
),
})?;
return Ok(all);
};
for name in requested {
if !indices.iter().any(|index| &index.name == name) {
return Err(Error::InvalidInput {
message: format!(
"maintained index '{}' does not exist on this table; it has {}",
name,
index_name_list(indices),
),
});
}
}
validate_maintained_indexes(dataset, requested).await?;
Ok(requested.to_vec())
}
/// Index names for an error message.
fn index_name_list(indices: &[IndexConfig]) -> String {
if indices.is_empty() {
return "no indexes".to_string();
}
let mut names: Vec<&str> = indices.iter().map(|index| index.name.as_str()).collect();
names.sort_unstable();
format!("[{}]", names.join(", "))
}
// =============================================================================
// unset_lsm_write_spec
// =============================================================================