mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-04 12:38:38 +00:00
Compare commits
2 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| b64284d8ca | |||
| 8347375362 |
@@ -276,38 +276,14 @@ jobs:
|
||||
# unreadable outside their own branch anyway, since GitHub scopes
|
||||
# caches to the creating ref.
|
||||
save-if: ${{ github.ref == 'refs/heads/main' }}
|
||||
- name: Downgrade dependencies
|
||||
# These packages have newer requirements for MSRV
|
||||
run: |
|
||||
cargo update -p aws-sdk-bedrockruntime --precise 1.77.0
|
||||
cargo update -p aws-sdk-dynamodb --precise 1.68.0
|
||||
cargo update -p aws-config --precise 1.6.0
|
||||
cargo update -p aws-sdk-kms --precise 1.63.0
|
||||
cargo update -p aws-sdk-s3 --precise 1.79.0
|
||||
cargo update -p aws-sdk-sso --precise 1.62.0
|
||||
cargo update -p aws-sdk-ssooidc --precise 1.63.0
|
||||
cargo update -p aws-sdk-sts --precise 1.63.0
|
||||
# aws-runtime/sigv4/credential-types/types and the aws-smithy-*
|
||||
# crates bumped their MSRV to 1.91.1 in late 2026; pin to the last
|
||||
# 1.91.0-compatible versions. The order matters — each downgrade
|
||||
# only succeeds once everything that still pins it at a higher
|
||||
# version has itself been downgraded.
|
||||
cargo update -p aws-runtime --precise 1.5.12
|
||||
cargo update -p aws-types --precise 1.3.9
|
||||
cargo update -p aws-sigv4 --precise 1.3.5
|
||||
cargo update -p aws-credential-types --precise 1.2.8
|
||||
cargo update -p aws-smithy-checksums --precise 0.63.9
|
||||
cargo update -p aws-smithy-runtime --precise 1.9.3
|
||||
cargo update -p aws-smithy-http --precise 0.62.4
|
||||
cargo update -p aws-smithy-eventstream --precise 0.60.12
|
||||
cargo update -p aws-smithy-http-client --precise 1.1.3
|
||||
cargo update -p aws-smithy-observability --precise 0.1.4
|
||||
cargo update -p aws-smithy-query --precise 0.60.8
|
||||
cargo update -p aws-smithy-runtime-api --precise 1.9.1
|
||||
cargo update -p aws-smithy-async --precise 1.2.6
|
||||
cargo update -p aws-smithy-types --precise 1.3.5
|
||||
cargo update -p aws-smithy-xml --precise 0.60.11
|
||||
cargo update -p home --precise 0.5.9
|
||||
- name: Downgrade dependencies that exceed our MSRV
|
||||
# Re-resolve the lockfile against `rust-version` instead of hand-pinning
|
||||
# every crate that raises its MSRV. Hand-pinning drifts: the pins keep
|
||||
# ratcheting further back than needed and eventually contradict a real
|
||||
# requirement elsewhere in the graph.
|
||||
env:
|
||||
CARGO_RESOLVER_INCOMPATIBLE_RUST_VERSIONS: fallback
|
||||
run: cargo update
|
||||
- name: cargo +${{ matrix.msrv }} check
|
||||
env:
|
||||
RUSTUP_TOOLCHAIN: ${{ matrix.msrv }}
|
||||
|
||||
Generated
+258
-234
File diff suppressed because it is too large
Load Diff
+14
-14
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
|
||||
rust-version = "1.91.0"
|
||||
|
||||
[workspace.dependencies]
|
||||
lance = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=10.1.0-beta.1", default-features = false, "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=10.1.0-beta.1", "tag" = "v10.1.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance = { "version" = "=11.0.0-beta.1", default-features = false, "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-core = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datagen = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-file = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-io = { "version" = "=11.0.0-beta.1", default-features = false, "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-index = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-linalg = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-namespace-impls = { "version" = "=11.0.0-beta.1", default-features = false, "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-table = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-testing = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-datafusion = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-encoding = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
lance-arrow = { "version" = "=11.0.0-beta.1", "tag" = "v11.0.0-beta.1", "git" = "https://github.com/lance-format/lance.git" }
|
||||
ahash = "0.8"
|
||||
# Note that this one does not include pyarrow
|
||||
arrow = { version = "58.0.0", optional = false }
|
||||
|
||||
@@ -180,21 +180,6 @@ instead of being materialized with the rest of the row.
|
||||
|
||||
::: lancedb.otel.instrument_lancedb_metrics
|
||||
|
||||
## Legacy V2 migration
|
||||
|
||||
Tables created with the experimental V2 format in LanceDB Node 0.5.x can be
|
||||
rewritten with the legacy PyLance reader. In a dedicated environment, install
|
||||
LanceDB normally, then install the legacy reader without its obsolete PyArrow
|
||||
upper bound and run the migration:
|
||||
|
||||
```shell
|
||||
pip install lancedb
|
||||
pip install --no-deps pylance==0.12.1
|
||||
python -m lancedb.legacy_v2 <database-uri>
|
||||
```
|
||||
|
||||
::: lancedb.legacy_v2.migrate_legacy_v2_tables
|
||||
|
||||
## Exceptions
|
||||
|
||||
::: lancedb.exceptions.MissingValueError
|
||||
|
||||
+1
-1
@@ -28,7 +28,7 @@
|
||||
<properties>
|
||||
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
|
||||
<arrow.version>15.0.0</arrow.version>
|
||||
<lance-core.version>10.1.0-beta.1</lance-core.version>
|
||||
<lance-core.version>11.0.0-beta.1</lance-core.version>
|
||||
<spotless.skip>false</spotless.skip>
|
||||
<spotless.version>2.30.0</spotless.version>
|
||||
<spotless.java.googlejavaformat.version>1.7</spotless.java.googlejavaformat.version>
|
||||
|
||||
@@ -707,9 +707,6 @@ class LanceDBConnection(DBConnection):
|
||||
self._namespace_client_properties = namespace_client_properties
|
||||
if _inner is not None:
|
||||
self._conn = _inner
|
||||
# Native-derived wrappers resolve this in their async reconstruction
|
||||
# path so construction never synchronously re-enters LOOP.
|
||||
self._read_consistency_interval = read_consistency_interval
|
||||
self._cached_namespace_client = None
|
||||
return
|
||||
|
||||
@@ -759,14 +756,11 @@ class LanceDBConnection(DBConnection):
|
||||
# storage_options. Also, this class really shouldn't be holding any state
|
||||
# beyond _conn.
|
||||
self._conn = AsyncConnection(LOOP.run(do_connect()))
|
||||
# Keep property access synchronous so debugger introspection cannot wait on
|
||||
# the background loop while that thread is suspended at a breakpoint.
|
||||
self._read_consistency_interval = read_consistency_interval
|
||||
self._cached_namespace_client: Optional[LanceNamespace] = None
|
||||
|
||||
@property
|
||||
def read_consistency_interval(self) -> Optional[timedelta]:
|
||||
return self._read_consistency_interval
|
||||
return LOOP.run(self._conn.get_read_consistency_interval())
|
||||
|
||||
@property
|
||||
def session(self) -> Optional[Session]:
|
||||
@@ -777,16 +771,8 @@ class LanceDBConnection(DBConnection):
|
||||
return self._conn.uri
|
||||
|
||||
@classmethod
|
||||
def from_inner(
|
||||
cls,
|
||||
inner: LanceDbConnection,
|
||||
read_consistency_interval: Optional[timedelta],
|
||||
):
|
||||
return cls(
|
||||
None,
|
||||
read_consistency_interval=read_consistency_interval,
|
||||
_inner=inner,
|
||||
)
|
||||
def from_inner(cls, inner: LanceDbConnection):
|
||||
return cls(None, _inner=inner)
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return f"{self.__class__.__name__}(uri={self._conn.uri!r})"
|
||||
|
||||
@@ -1,263 +0,0 @@
|
||||
# SPDX-License-Identifier: Apache-2.0
|
||||
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
"""Recovery utilities for the experimental V2 format used by old Node releases."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import argparse
|
||||
import asyncio
|
||||
import inspect
|
||||
import warnings
|
||||
from collections.abc import Iterable
|
||||
from typing import Any
|
||||
|
||||
from packaging.version import Version
|
||||
|
||||
import lancedb
|
||||
|
||||
__all__ = ["migrate_legacy_v2_tables"]
|
||||
|
||||
_LEGACY_LANCE_VERSION = Version("0.12.1")
|
||||
_LEGACY_V2_ERROR_MARKERS = (
|
||||
"missing columnencoding encoding description",
|
||||
"missing lance.encodings.columnencoding encoding description",
|
||||
"was missing a columnencoding",
|
||||
"rust future panicked",
|
||||
"panic in async function",
|
||||
)
|
||||
|
||||
|
||||
def _require_legacy_lance() -> Any:
|
||||
try:
|
||||
import lance
|
||||
except ImportError as error:
|
||||
raise RuntimeError(
|
||||
"Legacy V2 migration requires pylance==0.12.1. Install it in a "
|
||||
"dedicated environment with "
|
||||
"`pip install --no-deps pylance==0.12.1`."
|
||||
) from error
|
||||
|
||||
version = Version(lance.__version__)
|
||||
if version != _LEGACY_LANCE_VERSION:
|
||||
raise RuntimeError(
|
||||
"Legacy V2 migration requires pylance==0.12.1, but found "
|
||||
f"pylance=={version}. Reinstall it with "
|
||||
"`pip install --no-deps --force-reinstall pylance==0.12.1`."
|
||||
)
|
||||
return lance
|
||||
|
||||
|
||||
def _exception_messages(error: BaseException) -> Iterable[str]:
|
||||
seen: set[int] = set()
|
||||
current: BaseException | None = error
|
||||
while current is not None and id(current) not in seen:
|
||||
seen.add(id(current))
|
||||
yield str(current).lower()
|
||||
current = current.__cause__ or current.__context__
|
||||
|
||||
|
||||
def _is_legacy_v2_error(error: BaseException) -> bool:
|
||||
return any(
|
||||
marker in message
|
||||
for message in _exception_messages(error)
|
||||
for marker in _LEGACY_V2_ERROR_MARKERS
|
||||
)
|
||||
|
||||
|
||||
def _list_table_names(db: Any) -> list[str]:
|
||||
list_tables = getattr(db, "list_tables", None)
|
||||
if list_tables is not None:
|
||||
# The deprecated table_names() API defaults to only ten results.
|
||||
return list(list_tables(limit=None).tables)
|
||||
|
||||
# Compatibility for LanceDB 0.16, which was used by the original script.
|
||||
with warnings.catch_warnings():
|
||||
warnings.simplefilter("ignore", DeprecationWarning)
|
||||
return list(db.table_names())
|
||||
|
||||
|
||||
async def _needs_migration(db: Any, table_name: str) -> bool:
|
||||
table = await db.open_table(table_name)
|
||||
try:
|
||||
# One row is enough to load and validate the data-file metadata.
|
||||
await table.query().limit(1).to_arrow()
|
||||
except (KeyboardInterrupt, SystemExit, GeneratorExit):
|
||||
raise
|
||||
except BaseException as error:
|
||||
if _is_legacy_v2_error(error):
|
||||
return True
|
||||
raise
|
||||
return False
|
||||
|
||||
|
||||
async def _create_migrated_table(
|
||||
db: Any, table_name: str, reader: Any, storage_format: str
|
||||
) -> Any:
|
||||
parameters = inspect.signature(db.create_table).parameters
|
||||
options: dict[str, Any] = {"mode": "overwrite"}
|
||||
if "data_storage_version" in parameters:
|
||||
# LanceDB 0.16 exposed the format as a direct create_table option.
|
||||
options["data_storage_version"] = storage_format
|
||||
else:
|
||||
options["storage_options"] = {"new_table_data_storage_version": storage_format}
|
||||
return await db.create_table(table_name, reader, **options)
|
||||
|
||||
|
||||
async def _migrate_table(
|
||||
source_db: Any,
|
||||
destination_db: Any,
|
||||
table_name: str,
|
||||
storage_format: str,
|
||||
) -> int:
|
||||
source_table = source_db.open_table(table_name)
|
||||
source_dataset = source_table.to_lance()
|
||||
source_rows = source_dataset.count_rows()
|
||||
reader = source_dataset.scanner().to_reader()
|
||||
|
||||
migrated_table = await _create_migrated_table(
|
||||
destination_db, table_name, reader, storage_format
|
||||
)
|
||||
migrated_rows = await migrated_table.count_rows()
|
||||
if migrated_rows != source_rows:
|
||||
raise RuntimeError(
|
||||
f"Migration of table {table_name!r} wrote {migrated_rows} rows; "
|
||||
f"expected {source_rows}."
|
||||
)
|
||||
|
||||
# Force the current reader to load data-file metadata before reporting success.
|
||||
await migrated_table.query().limit(1).to_arrow()
|
||||
return migrated_rows
|
||||
|
||||
|
||||
async def migrate_legacy_v2_tables(
|
||||
uri: str,
|
||||
*,
|
||||
table_name: str | None = None,
|
||||
destination_uri: str | None = None,
|
||||
storage_format: str = "2.0",
|
||||
show_progress: bool = True,
|
||||
) -> list[str]:
|
||||
"""Migrate tables written with the incompatible experimental V2 format.
|
||||
|
||||
LanceDB Node 0.5.x could enable an experimental data format when an empty
|
||||
table was created and data was added later. Those files panic older modern
|
||||
readers and are rejected by newer readers. This utility streams them through
|
||||
``pylance==0.12.1`` and rewrites them in a supported format.
|
||||
|
||||
Install the legacy reader in a dedicated environment before running this
|
||||
function::
|
||||
|
||||
pip install lancedb
|
||||
pip install --no-deps pylance==0.12.1
|
||||
|
||||
``--no-deps`` is required because the legacy wheel declares an obsolete
|
||||
PyArrow upper bound. The migration uses only its dataset scanner and writes
|
||||
through the current LanceDB package.
|
||||
|
||||
This migration is available only for local/OSS databases, including object
|
||||
storage URIs. It is not supported for LanceDB Cloud ``db://`` connections.
|
||||
In-place migration creates a new table version, so old data remains available
|
||||
for recovery until old versions are cleaned up. Table indices are not copied
|
||||
and should be rebuilt after migration.
|
||||
|
||||
Parameters
|
||||
----------
|
||||
uri : str
|
||||
Source LanceDB database URI.
|
||||
table_name : str, optional
|
||||
Migrate only this table. By default, inspect every table.
|
||||
destination_uri : str, optional
|
||||
Write to another database. By default, migrate in place.
|
||||
storage_format : str, default "2.0"
|
||||
Data storage format for the rewritten tables. Use ``"0.1"`` for
|
||||
compatibility with older LanceDB releases.
|
||||
show_progress : bool, default True
|
||||
Display progress bars while inspecting and migrating tables.
|
||||
|
||||
Returns
|
||||
-------
|
||||
list of str
|
||||
Names of the migrated tables.
|
||||
"""
|
||||
if uri.startswith("db://") or (
|
||||
destination_uri is not None and destination_uri.startswith("db://")
|
||||
):
|
||||
raise ValueError("Legacy V2 migration is supported only for local/OSS tables")
|
||||
|
||||
# Import and validate before opening or modifying any table.
|
||||
_require_legacy_lance()
|
||||
|
||||
source_db = lancedb.connect(uri)
|
||||
async_source_db = await lancedb.connect_async(uri)
|
||||
destination_db = (
|
||||
async_source_db
|
||||
if destination_uri is None or destination_uri == uri
|
||||
else await lancedb.connect_async(destination_uri)
|
||||
)
|
||||
|
||||
if table_name is not None:
|
||||
table_names = [table_name]
|
||||
else:
|
||||
table_names = _list_table_names(source_db)
|
||||
|
||||
inspection: Iterable[str] = table_names
|
||||
if show_progress:
|
||||
from tqdm.auto import tqdm
|
||||
|
||||
inspection = tqdm(table_names, desc="Checking tables")
|
||||
|
||||
tables_to_migrate = [
|
||||
name for name in inspection if await _needs_migration(async_source_db, name)
|
||||
]
|
||||
|
||||
migration: Iterable[str] = tables_to_migrate
|
||||
if show_progress:
|
||||
from tqdm.auto import tqdm
|
||||
|
||||
migration = tqdm(tables_to_migrate, desc="Migrating tables")
|
||||
|
||||
migrated = []
|
||||
for name in migration:
|
||||
await _migrate_table(source_db, destination_db, name, storage_format)
|
||||
migrated.append(name)
|
||||
return migrated
|
||||
|
||||
|
||||
def _parser() -> argparse.ArgumentParser:
|
||||
parser = argparse.ArgumentParser(
|
||||
description="Migrate tables written with the old experimental V2 format."
|
||||
)
|
||||
parser.add_argument("uri", help="source LanceDB database URI")
|
||||
parser.add_argument("--table-name", help="migrate only this table")
|
||||
parser.add_argument("--destination-uri", help="write to another database URI")
|
||||
parser.add_argument(
|
||||
"--storage-format",
|
||||
default="2.0",
|
||||
help='destination data format (default: "2.0"; use "0.1" for compatibility)',
|
||||
)
|
||||
parser.add_argument(
|
||||
"--no-progress", action="store_true", help="disable progress bars"
|
||||
)
|
||||
return parser
|
||||
|
||||
|
||||
def main() -> None:
|
||||
args = _parser().parse_args()
|
||||
migrated = asyncio.run(
|
||||
migrate_legacy_v2_tables(
|
||||
args.uri,
|
||||
table_name=args.table_name,
|
||||
destination_uri=args.destination_uri,
|
||||
storage_format=args.storage_format,
|
||||
show_progress=not args.no_progress,
|
||||
)
|
||||
)
|
||||
if migrated:
|
||||
print(f"Migrated {len(migrated)} table(s): {', '.join(migrated)}")
|
||||
else:
|
||||
print("No legacy V2 tables found")
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -226,7 +226,7 @@ class PermutationBuilder:
|
||||
|
||||
async def do_execute():
|
||||
inner_tbl = await self._async.execute()
|
||||
return await LanceTable.from_inner(inner_tbl)
|
||||
return LanceTable.from_inner(inner_tbl)
|
||||
|
||||
return LOOP.run(do_execute())
|
||||
|
||||
|
||||
@@ -2182,15 +2182,11 @@ class LanceTable(Table):
|
||||
return self.name
|
||||
|
||||
@classmethod
|
||||
async def from_inner(cls, tbl: LanceDBTable):
|
||||
from .db import AsyncConnection, LanceDBConnection
|
||||
def from_inner(cls, tbl: LanceDBTable):
|
||||
from .db import LanceDBConnection
|
||||
|
||||
async_tbl = AsyncTable(tbl)
|
||||
inner_conn = tbl.database()
|
||||
read_consistency_interval = await AsyncConnection(
|
||||
inner_conn
|
||||
).get_read_consistency_interval()
|
||||
conn = LanceDBConnection.from_inner(inner_conn, read_consistency_interval)
|
||||
conn = LanceDBConnection.from_inner(tbl.database())
|
||||
return cls(
|
||||
conn,
|
||||
async_tbl.name,
|
||||
|
||||
@@ -77,23 +77,6 @@ def test_sync_repr_does_not_use_background_loop(tmp_path, monkeypatch):
|
||||
assert repr(table) == f"LanceTable(name='test', _conn={db!r})"
|
||||
|
||||
|
||||
def test_read_consistency_interval_does_not_use_background_loop(tmp_path, monkeypatch):
|
||||
from lancedb.background_loop import LOOP
|
||||
from lancedb.db import LanceDBConnection
|
||||
|
||||
consistency_interval = timedelta(seconds=5)
|
||||
db = lancedb.connect(tmp_path, read_consistency_interval=consistency_interval)
|
||||
db_from_inner = LanceDBConnection.from_inner(db._inner, consistency_interval)
|
||||
|
||||
def fail_run(*args, **kwargs):
|
||||
raise AssertionError("properties should not use the Python background loop")
|
||||
|
||||
monkeypatch.setattr(LOOP, "run", fail_run)
|
||||
|
||||
assert db.read_consistency_interval == consistency_interval
|
||||
assert db_from_inner.read_consistency_interval == consistency_interval
|
||||
|
||||
|
||||
def test_ingest_pd(tmp_path):
|
||||
db = lancedb.connect(tmp_path)
|
||||
|
||||
|
||||
@@ -1,165 +0,0 @@
|
||||
from unittest.mock import AsyncMock, Mock
|
||||
|
||||
import pytest
|
||||
|
||||
from lancedb import legacy_v2
|
||||
|
||||
|
||||
class FakeQuery:
|
||||
def __init__(self, error=None):
|
||||
self.error = error
|
||||
|
||||
def limit(self, _limit):
|
||||
return self
|
||||
|
||||
async def to_arrow(self):
|
||||
if self.error is not None:
|
||||
raise self.error
|
||||
return None
|
||||
|
||||
|
||||
class FakeAsyncTable:
|
||||
def __init__(self, rows=2, error=None):
|
||||
self.rows = rows
|
||||
self.error = error
|
||||
|
||||
def query(self):
|
||||
return FakeQuery(self.error)
|
||||
|
||||
async def count_rows(self):
|
||||
return self.rows
|
||||
|
||||
|
||||
class FakeDataset:
|
||||
def __init__(self, reader, rows=2):
|
||||
self.reader = reader
|
||||
self.rows = rows
|
||||
|
||||
def count_rows(self):
|
||||
return self.rows
|
||||
|
||||
def scanner(self):
|
||||
scanner = Mock()
|
||||
scanner.to_reader.return_value = self.reader
|
||||
return scanner
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize(
|
||||
"message",
|
||||
[
|
||||
"rust future panicked: unknown error",
|
||||
"Panic in async function",
|
||||
"Missing ColumnEncoding encoding description",
|
||||
"Missing lance.encodings.ColumnEncoding encoding description",
|
||||
"the column at index 0 was missing a ColumnEncoding",
|
||||
],
|
||||
)
|
||||
async def test_needs_migration_recognizes_legacy_reader_errors(message):
|
||||
db = Mock()
|
||||
db.open_table = AsyncMock(return_value=FakeAsyncTable(error=RuntimeError(message)))
|
||||
|
||||
assert await legacy_v2._needs_migration(db, "legacy")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_needs_migration_propagates_unrelated_errors():
|
||||
db = Mock()
|
||||
db.open_table = AsyncMock(
|
||||
return_value=FakeAsyncTable(error=RuntimeError("permission denied"))
|
||||
)
|
||||
|
||||
with pytest.raises(RuntimeError, match="permission denied"):
|
||||
await legacy_v2._needs_migration(db, "legacy")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_migration_streams_and_verifies_rows(monkeypatch):
|
||||
reader = object()
|
||||
source_dataset = FakeDataset(reader)
|
||||
source_table = Mock()
|
||||
source_table.to_lance.return_value = source_dataset
|
||||
source_db = Mock()
|
||||
source_db.list_tables.return_value.tables = ["healthy", "legacy"]
|
||||
source_db.open_table.return_value = source_table
|
||||
|
||||
legacy_error = RuntimeError(
|
||||
"Missing lance.encodings.ColumnEncoding encoding description"
|
||||
)
|
||||
async_source_db = Mock()
|
||||
|
||||
async def open_table(name):
|
||||
if name == "legacy":
|
||||
return FakeAsyncTable(error=legacy_error)
|
||||
return FakeAsyncTable()
|
||||
|
||||
async_source_db.open_table = open_table
|
||||
create_calls = []
|
||||
|
||||
async def create_table(name, data, *, mode, storage_options):
|
||||
create_calls.append((name, data, mode, storage_options))
|
||||
return FakeAsyncTable()
|
||||
|
||||
async_source_db.create_table = create_table
|
||||
|
||||
monkeypatch.setattr(legacy_v2, "_require_legacy_lance", Mock())
|
||||
monkeypatch.setattr(legacy_v2.lancedb, "connect", Mock(return_value=source_db))
|
||||
|
||||
async def connect_async(_uri):
|
||||
return async_source_db
|
||||
|
||||
monkeypatch.setattr(legacy_v2.lancedb, "connect_async", connect_async)
|
||||
|
||||
migrated = await legacy_v2.migrate_legacy_v2_tables("/data/db", show_progress=False)
|
||||
|
||||
assert migrated == ["legacy"]
|
||||
source_db.list_tables.assert_called_once_with(limit=None)
|
||||
assert create_calls == [
|
||||
(
|
||||
"legacy",
|
||||
reader,
|
||||
"overwrite",
|
||||
{"new_table_data_storage_version": "2.0"},
|
||||
)
|
||||
]
|
||||
|
||||
|
||||
def test_list_table_names_supports_legacy_connection():
|
||||
db = Mock(spec=["table_names"])
|
||||
db.table_names.return_value = [f"table_{index}" for index in range(12)]
|
||||
|
||||
assert len(legacy_v2._list_table_names(db)) == 12
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_create_table_uses_legacy_storage_parameter():
|
||||
calls = []
|
||||
|
||||
class LegacyConnection:
|
||||
async def create_table(self, name, data, *, mode, data_storage_version=None):
|
||||
calls.append((name, data, mode, data_storage_version))
|
||||
return FakeAsyncTable()
|
||||
|
||||
reader = object()
|
||||
await legacy_v2._create_migrated_table(LegacyConnection(), "legacy", reader, "0.1")
|
||||
|
||||
assert calls == [("legacy", reader, "overwrite", "0.1")]
|
||||
|
||||
|
||||
def test_requires_exact_legacy_lance_version(monkeypatch):
|
||||
fake_lance = Mock(__version__="9.0.0")
|
||||
monkeypatch.setitem(__import__("sys").modules, "lance", fake_lance)
|
||||
|
||||
with pytest.raises(RuntimeError, match="requires pylance==0.12.1"):
|
||||
legacy_v2._require_legacy_lance()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_cloud_migration_is_rejected_before_dependency_check(monkeypatch):
|
||||
require_lance = Mock()
|
||||
monkeypatch.setattr(legacy_v2, "_require_legacy_lance", require_lance)
|
||||
|
||||
with pytest.raises(ValueError, match="only for local/OSS"):
|
||||
await legacy_v2.migrate_legacy_v2_tables("db://example")
|
||||
|
||||
require_lance.assert_not_called()
|
||||
@@ -6,7 +6,6 @@ import math
|
||||
import pytest
|
||||
|
||||
from lancedb import DBConnection, Table, connect
|
||||
from lancedb.background_loop import LOOP
|
||||
from lancedb.permutation import Permutation, Permutations, permutation_builder
|
||||
|
||||
|
||||
@@ -32,25 +31,6 @@ def test_split_random_ratios(mem_db):
|
||||
assert 65 <= split_1_count <= 75 # ~70% ± tolerance
|
||||
|
||||
|
||||
def test_execute_does_not_reenter_background_loop(tmp_path, monkeypatch):
|
||||
import threading
|
||||
|
||||
db = connect(tmp_path)
|
||||
tbl = db.create_table("test_table", pa.table({"x": range(10)}))
|
||||
original_run = LOOP.run
|
||||
|
||||
def fail_on_reentry(future):
|
||||
assert threading.current_thread() is not LOOP.thread
|
||||
return original_run(future)
|
||||
|
||||
monkeypatch.setattr(LOOP, "run", fail_on_reentry)
|
||||
|
||||
permutation_tbl = permutation_builder(tbl).execute()
|
||||
|
||||
assert permutation_tbl.count_rows() == 10
|
||||
assert permutation_tbl._conn.read_consistency_interval is None
|
||||
|
||||
|
||||
def test_split_random_counts(mem_db):
|
||||
"""Test random splitting with absolute counts."""
|
||||
tbl = mem_db.create_table(
|
||||
|
||||
@@ -6,7 +6,6 @@ import os
|
||||
import sys
|
||||
import threading
|
||||
import warnings
|
||||
from concurrent.futures import ThreadPoolExecutor
|
||||
from datetime import date, datetime, timedelta
|
||||
from time import sleep
|
||||
from typing import List
|
||||
@@ -2125,27 +2124,6 @@ def test_delete(mem_db: DBConnection):
|
||||
assert table.to_arrow()["id"].to_pylist() == [1]
|
||||
|
||||
|
||||
def test_concurrent_deletes_are_thread_safe(mem_db: DBConnection):
|
||||
num_workers = 8
|
||||
table = mem_db.create_table(
|
||||
"my_table", data=[{"id": row_id} for row_id in range(num_workers)]
|
||||
)
|
||||
barrier = threading.Barrier(num_workers)
|
||||
|
||||
def delete(row_id: int):
|
||||
barrier.wait()
|
||||
return table.delete(f"id = {row_id}")
|
||||
|
||||
with ThreadPoolExecutor(max_workers=num_workers) as pool:
|
||||
results = list(pool.map(delete, range(num_workers)))
|
||||
|
||||
assert all(result.num_deleted_rows == 1 for result in results)
|
||||
assert sorted(result.version for result in results) == list(
|
||||
range(2, num_workers + 2)
|
||||
)
|
||||
assert table.count_rows() == 0
|
||||
|
||||
|
||||
def test_delete_expr(mem_db: DBConnection):
|
||||
table = mem_db.create_table(
|
||||
"my_table",
|
||||
|
||||
@@ -745,9 +745,6 @@ impl Table {
|
||||
|
||||
#[allow(private_interfaces)]
|
||||
pub fn delete(self_: PyRef<'_, Self>, condition: PredicateArg) -> PyResult<Bound<'_, PyAny>> {
|
||||
// Do not hold the Python borrow across the await. The cloned Rust table
|
||||
// handle is thread-safe and allows deletes on the same Python table to
|
||||
// run concurrently without PyO3 reporting "Already borrowed".
|
||||
let inner = self_.inner_ref()?.clone();
|
||||
future_into_py(self_.py(), async move {
|
||||
let result = match &condition {
|
||||
|
||||
@@ -49,8 +49,8 @@ lance-namespace = { workspace = true }
|
||||
lance-namespace-impls = { workspace = true }
|
||||
metrics = { workspace = true, optional = true }
|
||||
metrics-util = { workspace = true, optional = true }
|
||||
# Pin the transitive GooseFS SDK until the 0.1.6 compile break is fixed upstream.
|
||||
goosefs-sdk = { version = "=0.1.5", optional = true }
|
||||
# Pin the GooseFS SDK to the version required by Lance's OpenDAL dependency.
|
||||
goosefs-sdk = { version = "=0.1.9", optional = true }
|
||||
moka = { workspace = true }
|
||||
pin-project = { workspace = true }
|
||||
tokio = { version = "1.23", features = ["rt-multi-thread", "sync"] }
|
||||
|
||||
@@ -17,7 +17,7 @@ use arrow_array::builder::LargeBinaryBuilder;
|
||||
use arrow_schema::{DataType, Field, Schema};
|
||||
use lance::dataset::{BlobRangeRequest as LanceBlobRangeRequest, Dataset, WriteParams};
|
||||
use lance_arrow::FieldExt;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lance_file::version::LanceFileVersion;
|
||||
use lance_io::object_store::ObjectStore;
|
||||
use object_store::path::Path;
|
||||
|
||||
|
||||
@@ -34,7 +34,7 @@ use crate::remote::{
|
||||
db::{OPT_REMOTE_API_KEY, OPT_REMOTE_HOST_OVERRIDE, OPT_REMOTE_REGION},
|
||||
};
|
||||
use lance::io::ObjectStoreParams;
|
||||
pub use lance_encoding::version::LanceFileVersion;
|
||||
pub use lance_file::version::LanceFileVersion;
|
||||
#[cfg(feature = "remote")]
|
||||
use lance_io::object_store::StorageOptions;
|
||||
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
|
||||
|
||||
@@ -12,7 +12,7 @@ use lance::dataset::refs::Ref;
|
||||
use lance::dataset::{ReadParams, WriteMode, builder::DatasetBuilder};
|
||||
use lance::io::{ObjectStore, ObjectStoreParams, WrappingObjectStore};
|
||||
use lance_datafusion::utils::StreamingWriteSource;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lance_file::version::LanceFileVersion;
|
||||
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
|
||||
use lance_table::io::commit::commit_handler_from_url;
|
||||
use object_store::local::LocalFileSystem;
|
||||
|
||||
@@ -201,7 +201,7 @@ impl LanceNamespaceDatabase {
|
||||
&self,
|
||||
request: &DbCreateTableRequest,
|
||||
) -> Result<(
|
||||
Option<lance_encoding::version::LanceFileVersion>,
|
||||
Option<lance_file::version::LanceFileVersion>,
|
||||
Option<bool>,
|
||||
Option<bool>,
|
||||
)> {
|
||||
@@ -214,7 +214,7 @@ impl LanceNamespaceDatabase {
|
||||
|
||||
let storage_version_override = storage_options
|
||||
.and_then(|opts| opts.get(OPT_NEW_TABLE_STORAGE_VERSION))
|
||||
.map(|s| s.parse::<lance_encoding::version::LanceFileVersion>())
|
||||
.map(|s| s.parse::<lance_file::version::LanceFileVersion>())
|
||||
.transpose()?;
|
||||
|
||||
let v2_manifest_override = storage_options
|
||||
|
||||
@@ -2939,7 +2939,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
|
||||
}
|
||||
|
||||
#[derive(Serialize, Clone, Debug)]
|
||||
pub(crate) struct MergeInsertRequest {
|
||||
pub struct MergeInsertRequest {
|
||||
on: String,
|
||||
when_matched_update_all: bool,
|
||||
when_matched_update_all_filt: Option<String>,
|
||||
|
||||
@@ -90,7 +90,7 @@ struct RemoteBlobState {
|
||||
|
||||
/// Seekable Cloud blob handle over HTTP Range.
|
||||
#[derive(Debug)]
|
||||
pub(crate) struct RemoteBlobFile {
|
||||
pub struct RemoteBlobFile {
|
||||
requester: Arc<dyn BlobRangeRequester>,
|
||||
state: Mutex<RemoteBlobState>,
|
||||
closed: AtomicBool,
|
||||
|
||||
@@ -33,7 +33,7 @@ use crate::table::{AddResult, MergeResult};
|
||||
/// same Arrow-IPC streaming body and error side-channel; only the target
|
||||
/// endpoint, query parameters, and parsed result type differ.
|
||||
#[derive(Debug, Clone)]
|
||||
pub(crate) enum WriteOp {
|
||||
pub enum WriteOp {
|
||||
/// `add`: stream to `/v1/table/{id}/insert/`, optionally overwriting.
|
||||
Insert { overwrite: bool },
|
||||
/// `merge_insert`: stream to `/v1/table/{id}/merge_insert/` with the merge
|
||||
@@ -49,7 +49,7 @@ pub(crate) enum WriteOp {
|
||||
/// The parsed server response for a completed write, discriminated by the
|
||||
/// operation that produced it.
|
||||
#[derive(Debug, Clone)]
|
||||
pub(crate) enum WriteResult {
|
||||
pub enum WriteResult {
|
||||
Add(AddResult),
|
||||
Merge(MergeResult),
|
||||
}
|
||||
|
||||
@@ -10,7 +10,7 @@ use arrow_array::{
|
||||
use arrow_schema::{DataType, Field, Fields, Schema};
|
||||
use futures::TryStreamExt;
|
||||
use lance::Dataset;
|
||||
use lance_encoding::version::LanceFileVersion;
|
||||
use lance_file::version::LanceFileVersion;
|
||||
use lancedb::{
|
||||
Connection, Error, Result, Table,
|
||||
blob::{BlobRangeRequest, blob},
|
||||
|
||||
Reference in New Issue
Block a user