mirror of
https://github.com/lancedb/lancedb.git
synced 2026-08-29 17:38:31 +00:00
Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| da86d804ba | |||
| 77a93fee76 | |||
| 7bb501839a | |||
| 5b347afd99 | |||
| 706a9c327f |
+1
-1
@@ -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
@@ -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",
|
||||
|
||||
@@ -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>
|
||||
```
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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.
|
||||
|
||||
***
|
||||
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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
@@ -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
@@ -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
|
||||
|
||||
@@ -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 }]);
|
||||
|
||||
@@ -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
@@ -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,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,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,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,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,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,6 +1,6 @@
|
||||
{
|
||||
"name": "@lancedb/lancedb-win32-arm64-msvc",
|
||||
"version": "0.37.1-beta.0",
|
||||
"version": "0.37.1-beta.1",
|
||||
"os": [
|
||||
"win32"
|
||||
],
|
||||
|
||||
@@ -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",
|
||||
|
||||
Generated
+2
-2
@@ -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
@@ -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
@@ -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
@@ -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"
|
||||
|
||||
@@ -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]: ...
|
||||
|
||||
|
||||
@@ -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)))
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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()
|
||||
|
||||
|
||||
@@ -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
@@ -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.
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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
@@ -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
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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();
|
||||
|
||||
|
||||
@@ -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
|
||||
// =============================================================================
|
||||
|
||||
Reference in New Issue
Block a user