mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-03 20:18:54 +00:00
Compare commits
3 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 41d9c13fd3 | |||
| 7357d63e87 | |||
| 624a75edf7 |
@@ -162,6 +162,15 @@ def connect(
|
||||
... },
|
||||
... )
|
||||
|
||||
Azure managed identity authentication requires only the storage account name.
|
||||
LanceDB acquires and refreshes managed identity tokens automatically. For a
|
||||
user-assigned identity, also set ``azure_storage_client_id``:
|
||||
|
||||
>>> db = lancedb.connect( # doctest: +SKIP
|
||||
... "az://my-container/lancedb",
|
||||
... storage_options={"account_name": "my-storage-account"},
|
||||
... )
|
||||
|
||||
For tests and temporary data, use an in-memory database:
|
||||
|
||||
>>> db = lancedb.connect("memory://")
|
||||
@@ -455,6 +464,11 @@ async def connect_async(
|
||||
... db = await lancedb.connect_async("s3://my-bucket/lancedb",
|
||||
... storage_options={
|
||||
... "aws_access_key_id": "***"})
|
||||
... # Azure managed identity tokens are acquired and refreshed automatically
|
||||
... db = await lancedb.connect_async(
|
||||
... "az://my-container/lancedb",
|
||||
... storage_options={"account_name": "my-storage-account"},
|
||||
... )
|
||||
... # For tests and temporary data, use an in-memory database
|
||||
... db = await lancedb.connect_async("memory://")
|
||||
... # Connect to LanceDB cloud
|
||||
|
||||
@@ -707,6 +707,9 @@ 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
|
||||
|
||||
@@ -756,11 +759,14 @@ 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 LOOP.run(self._conn.get_read_consistency_interval())
|
||||
return self._read_consistency_interval
|
||||
|
||||
@property
|
||||
def session(self) -> Optional[Session]:
|
||||
@@ -771,8 +777,16 @@ class LanceDBConnection(DBConnection):
|
||||
return self._conn.uri
|
||||
|
||||
@classmethod
|
||||
def from_inner(cls, inner: LanceDbConnection):
|
||||
return cls(None, _inner=inner)
|
||||
def from_inner(
|
||||
cls,
|
||||
inner: LanceDbConnection,
|
||||
read_consistency_interval: Optional[timedelta],
|
||||
):
|
||||
return cls(
|
||||
None,
|
||||
read_consistency_interval=read_consistency_interval,
|
||||
_inner=inner,
|
||||
)
|
||||
|
||||
def __repr__(self) -> str:
|
||||
return f"{self.__class__.__name__}(uri={self._conn.uri!r})"
|
||||
|
||||
@@ -226,7 +226,7 @@ class PermutationBuilder:
|
||||
|
||||
async def do_execute():
|
||||
inner_tbl = await self._async.execute()
|
||||
return LanceTable.from_inner(inner_tbl)
|
||||
return await LanceTable.from_inner(inner_tbl)
|
||||
|
||||
return LOOP.run(do_execute())
|
||||
|
||||
|
||||
@@ -2182,11 +2182,15 @@ class LanceTable(Table):
|
||||
return self.name
|
||||
|
||||
@classmethod
|
||||
def from_inner(cls, tbl: LanceDBTable):
|
||||
from .db import LanceDBConnection
|
||||
async def from_inner(cls, tbl: LanceDBTable):
|
||||
from .db import AsyncConnection, LanceDBConnection
|
||||
|
||||
async_tbl = AsyncTable(tbl)
|
||||
conn = LanceDBConnection.from_inner(tbl.database())
|
||||
inner_conn = tbl.database()
|
||||
read_consistency_interval = await AsyncConnection(
|
||||
inner_conn
|
||||
).get_read_consistency_interval()
|
||||
conn = LanceDBConnection.from_inner(inner_conn, read_consistency_interval)
|
||||
return cls(
|
||||
conn,
|
||||
async_tbl.name,
|
||||
|
||||
@@ -102,6 +102,10 @@ def fs_from_uri(uri: str) -> Tuple[pa_fs.FileSystem, str]:
|
||||
az_blob_fs = adlfs.AzureBlobFileSystem(
|
||||
account_name=os.environ.get("AZURE_STORAGE_ACCOUNT_NAME"),
|
||||
account_key=os.environ.get("AZURE_STORAGE_ACCOUNT_KEY"),
|
||||
# Without an explicit key, authenticate with DefaultAzureCredential
|
||||
# instead of attempting anonymous access. In particular, this enables
|
||||
# managed identity authentication on Azure hosts.
|
||||
anon=False,
|
||||
)
|
||||
|
||||
fs = pa_fs.PyFileSystem(pa_fs.FSSpecHandler(az_blob_fs))
|
||||
|
||||
@@ -77,6 +77,23 @@ 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)
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@
|
||||
# SPDX-FileCopyrightText: Copyright The LanceDB Authors
|
||||
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
|
||||
import lancedb
|
||||
@@ -12,11 +13,19 @@ import pytest
|
||||
# AWS_PROFILE=default TEST_S3_BASE_URL=s3://my_bucket/dataset pytest tests/test_io.py
|
||||
#
|
||||
# Azure:
|
||||
# You need to setup Azure credentials an a base path to run this test. Example
|
||||
# You need to set up Azure credentials and a base path to run this test. Examples:
|
||||
#
|
||||
# Account key:
|
||||
# export AZURE_STORAGE_ACCOUNT_NAME="<account>"
|
||||
# export AZURE_STORAGE_ACCOUNT_KEY="<key>"
|
||||
# export REMOTE_BASE_URL=az://my_blob/dataset
|
||||
# pytest tests/test_io.py
|
||||
#
|
||||
# Managed identity (system-assigned or user-assigned):
|
||||
# export AZURE_STORAGE_ACCOUNT_NAME="<account>"
|
||||
# export AZURE_STORAGE_CLIENT_ID="<client-id>" # user-assigned identity only
|
||||
# export REMOTE_BASE_URL=az://my_blob/dataset
|
||||
# pytest tests/test_io.py
|
||||
|
||||
|
||||
@pytest.fixture(autouse=True, scope="module")
|
||||
@@ -58,3 +67,28 @@ def test_remote_io():
|
||||
assert len(db) == 1
|
||||
|
||||
assert db.open_table("test").name == db["test"].name
|
||||
|
||||
|
||||
@pytest.mark.skipif(
|
||||
(os.environ.get("REMOTE_BASE_URL") is None),
|
||||
reason="please setup remote base url",
|
||||
)
|
||||
def test_remote_io_async():
|
||||
async def run():
|
||||
db = await lancedb.connect_async(os.environ["REMOTE_BASE_URL"])
|
||||
table = await db.create_table(
|
||||
"test_async",
|
||||
data=[
|
||||
{"vector": [3.1, 4.1], "item": "foo"},
|
||||
{"vector": [5.9, 26.5], "item": "bar"},
|
||||
],
|
||||
)
|
||||
|
||||
assert await table.count_rows() == 2
|
||||
assert (await db.open_table("test_async")).name == "test_async"
|
||||
assert "test_async" in await db.table_names()
|
||||
|
||||
await db.drop_table("test_async")
|
||||
assert "test_async" not in await db.table_names()
|
||||
|
||||
asyncio.run(run())
|
||||
|
||||
@@ -6,6 +6,7 @@ import math
|
||||
import pytest
|
||||
|
||||
from lancedb import DBConnection, Table, connect
|
||||
from lancedb.background_loop import LOOP
|
||||
from lancedb.permutation import Permutation, Permutations, permutation_builder
|
||||
|
||||
|
||||
@@ -31,6 +32,25 @@ 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(
|
||||
|
||||
@@ -779,56 +779,6 @@ async def test_distance_range_with_new_rows_async():
|
||||
assert dist >= min_dist
|
||||
|
||||
|
||||
def test_flat_vector_search_large_limit_is_globally_sorted_after_upsert(mem_db):
|
||||
# Regression test for https://github.com/lancedb/lancedb/issues/3669.
|
||||
# Limits above Lance's 8,192-row output batch size used to expose batches in
|
||||
# scheduling order instead of global distance order after plan repartitioning.
|
||||
rng = np.random.default_rng(42)
|
||||
num_rows = 20_000
|
||||
fragment_size = 5_000
|
||||
dimension = 16
|
||||
limit = 12_000
|
||||
vectors = rng.random((num_rows, dimension), dtype=np.float32)
|
||||
|
||||
def make_data(start, end):
|
||||
return pa.table(
|
||||
{
|
||||
"id": pa.array(np.arange(start, end), type=pa.int64()),
|
||||
"vector": pa.FixedSizeListArray.from_arrays(
|
||||
pa.array(vectors[start:end].reshape(-1)), dimension
|
||||
),
|
||||
}
|
||||
)
|
||||
|
||||
table = mem_db.create_table("flat_knn_order", make_data(0, fragment_size))
|
||||
for start in range(fragment_size, num_rows, fragment_size):
|
||||
table.add(make_data(start, start + fragment_size))
|
||||
|
||||
upsert = pa.table(
|
||||
{
|
||||
"id": pa.array([0], type=pa.int64()),
|
||||
"vector": pa.array(
|
||||
[[0.0] * dimension], type=pa.list_(pa.float32(), dimension)
|
||||
),
|
||||
}
|
||||
)
|
||||
(
|
||||
table.merge_insert("id")
|
||||
.when_matched_update_all()
|
||||
.when_not_matched_insert_all()
|
||||
.execute(upsert)
|
||||
)
|
||||
|
||||
query = np.zeros(dimension, dtype=np.float32)
|
||||
for _ in range(10):
|
||||
result = (
|
||||
table.search(query).where("id >= 0").select(["id"]).limit(limit).to_arrow()
|
||||
)
|
||||
distances = result["_distance"].to_numpy()
|
||||
assert len(distances) == limit
|
||||
assert np.all(distances[:-1] <= distances[1:])
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"multivec_table", [pa.float16(), pa.float32(), pa.float64()], indirect=True
|
||||
)
|
||||
|
||||
@@ -6,6 +6,7 @@ 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
|
||||
@@ -2124,6 +2125,27 @@ 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",
|
||||
|
||||
@@ -26,7 +26,14 @@ import pandas as pd
|
||||
import polars as pl
|
||||
import pytest
|
||||
import lancedb
|
||||
from lancedb.util import flatten_columns, get_uri_scheme, join_uri, value_to_sql
|
||||
from lancedb import util
|
||||
from lancedb.util import (
|
||||
flatten_columns,
|
||||
fs_from_uri,
|
||||
get_uri_scheme,
|
||||
join_uri,
|
||||
value_to_sql,
|
||||
)
|
||||
from utils import exception_output
|
||||
|
||||
|
||||
@@ -78,6 +85,33 @@ def test_normalize_uri():
|
||||
assert parsed_scheme == expected_scheme
|
||||
|
||||
|
||||
def test_fs_from_uri_azure_uses_default_credential(monkeypatch):
|
||||
azure_options = {}
|
||||
filesystem = object()
|
||||
|
||||
class MockAdlfs:
|
||||
@staticmethod
|
||||
def AzureBlobFileSystem(**kwargs):
|
||||
azure_options.update(kwargs)
|
||||
return object()
|
||||
|
||||
monkeypatch.setenv("AZURE_STORAGE_ACCOUNT_NAME", "account")
|
||||
monkeypatch.delenv("AZURE_STORAGE_ACCOUNT_KEY", raising=False)
|
||||
monkeypatch.setattr(util, "adlfs", MockAdlfs)
|
||||
monkeypatch.setattr(util.pa_fs, "FSSpecHandler", lambda _: object())
|
||||
monkeypatch.setattr(util.pa_fs, "PyFileSystem", lambda _: filesystem)
|
||||
|
||||
actual_filesystem, path = fs_from_uri("az://container/database")
|
||||
|
||||
assert actual_filesystem is filesystem
|
||||
assert path == "container/database"
|
||||
assert azure_options == {
|
||||
"account_name": "account",
|
||||
"account_key": None,
|
||||
"anon": False,
|
||||
}
|
||||
|
||||
|
||||
def test_join_uri_remote():
|
||||
schemes = ["s3", "az", "gs"]
|
||||
for scheme in schemes:
|
||||
|
||||
@@ -745,6 +745,9 @@ 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 {
|
||||
|
||||
Reference in New Issue
Block a user