mirror of
https://github.com/lancedb/lancedb.git
synced 2026-09-02 11:38:49 +00:00
Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 803dd21ecb |
@@ -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})"
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -4,6 +4,7 @@
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
|
||||
use arrow_array::RecordBatch;
|
||||
use async_trait::async_trait;
|
||||
use http::StatusCode;
|
||||
use lance_io::object_store::StorageOptions;
|
||||
@@ -18,13 +19,14 @@ use lance_namespace::models::{
|
||||
};
|
||||
|
||||
use crate::Error;
|
||||
use crate::data::scannable::Scannable;
|
||||
use crate::database::{
|
||||
CloneTableRequest, CreateTableMode, CreateTableRequest, Database, DatabaseOptions,
|
||||
JobDescription, JobInfo, OpenTableRequest, ReadConsistency, TableNamesRequest,
|
||||
};
|
||||
use crate::error::Result;
|
||||
use crate::remote::util::stream_as_body;
|
||||
use crate::table::BaseTable;
|
||||
use crate::table::{AddDataBuilder, BaseTable};
|
||||
|
||||
use super::ARROW_STREAM_CONTENT_TYPE;
|
||||
use super::client::{
|
||||
@@ -693,7 +695,19 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
}
|
||||
|
||||
async fn create_table(&self, mut request: CreateTableRequest) -> Result<Arc<dyn BaseTable>> {
|
||||
let body = stream_as_body(request.data.scan_as_stream())?;
|
||||
// The create endpoint limits the size of the complete request even though
|
||||
// its body is streamed. Sources without a row-count hint (notably Python
|
||||
// generators / RecordBatchReader) can therefore exceed that limit after
|
||||
// many individually small batches. Create the schema first and feed the
|
||||
// unknown-length source through the multipart insert path instead.
|
||||
let stage_initial_data = request.data.num_rows().is_none();
|
||||
let schema = request.data.schema();
|
||||
let body = if stage_initial_data {
|
||||
let mut empty = RecordBatch::new_empty(schema.clone());
|
||||
stream_as_body(empty.scan_as_stream())?
|
||||
} else {
|
||||
stream_as_body(request.data.scan_as_stream())?
|
||||
};
|
||||
|
||||
let identifier = build_table_identifier(
|
||||
&request.name,
|
||||
@@ -764,6 +778,16 @@ impl<S: HttpSend> Database for RemoteDatabase<S> {
|
||||
table_identifier,
|
||||
version,
|
||||
));
|
||||
table.seed_schema_ref(schema);
|
||||
|
||||
if stage_initial_data {
|
||||
let base_table: Arc<dyn BaseTable> = table.clone();
|
||||
AddDataBuilder::new(base_table, request.data, None)
|
||||
.write_options(request.write_options)
|
||||
.execute()
|
||||
.await?;
|
||||
}
|
||||
|
||||
self.table_cache.insert(cache_key, table.clone()).await;
|
||||
|
||||
Ok(table)
|
||||
@@ -1106,7 +1130,7 @@ mod tests {
|
||||
use std::sync::atomic::{AtomicUsize, Ordering};
|
||||
use std::sync::{Arc, OnceLock};
|
||||
|
||||
use arrow_array::{Int32Array, RecordBatch};
|
||||
use arrow_array::{Int32Array, RecordBatch, RecordBatchIterator};
|
||||
use arrow_schema::{DataType, Field, Schema};
|
||||
use lance_namespace_impls::{DynamicContextProvider, OperationInfo};
|
||||
|
||||
@@ -1371,6 +1395,88 @@ mod tests {
|
||||
assert_eq!(table.name(), "table1");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_create_table_streaming_reader_uses_multipart_insert() {
|
||||
let create_count = Arc::new(AtomicUsize::new(0));
|
||||
let multipart_create_count = Arc::new(AtomicUsize::new(0));
|
||||
let insert_count = Arc::new(AtomicUsize::new(0));
|
||||
let complete_count = Arc::new(AtomicUsize::new(0));
|
||||
|
||||
let create_count_c = create_count.clone();
|
||||
let multipart_create_count_c = multipart_create_count.clone();
|
||||
let insert_count_c = insert_count.clone();
|
||||
let complete_count_c = complete_count.clone();
|
||||
let conn = Connection::new_with_handler_and_config(
|
||||
move |request| {
|
||||
let path = request.url().path();
|
||||
let query = request.url().query().unwrap_or("");
|
||||
match path {
|
||||
"/v1/table/table1/create/" => {
|
||||
create_count_c.fetch_add(1, Ordering::SeqCst);
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.header("phalanx-version", "0.4.0")
|
||||
.body(String::new())
|
||||
.unwrap()
|
||||
}
|
||||
"/v1/table/table1/multipart_write/create" => {
|
||||
multipart_create_count_c.fetch_add(1, Ordering::SeqCst);
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"upload_id":"streaming-create"}"#.to_string())
|
||||
.unwrap()
|
||||
}
|
||||
"/v1/table/table1/insert/" => {
|
||||
assert!(query.contains("upload_id=streaming-create"));
|
||||
assert!(query.contains("upload_part_id="));
|
||||
insert_count_c.fetch_add(1, Ordering::SeqCst);
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(String::new())
|
||||
.unwrap()
|
||||
}
|
||||
"/v1/table/table1/multipart_write/complete" => {
|
||||
assert!(query.contains("upload_id=streaming-create"));
|
||||
complete_count_c.fetch_add(1, Ordering::SeqCst);
|
||||
http::Response::builder()
|
||||
.status(200)
|
||||
.body(r#"{"version":2}"#.to_string())
|
||||
.unwrap()
|
||||
}
|
||||
path => panic!("unexpected path: {path}"),
|
||||
}
|
||||
},
|
||||
ClientConfig {
|
||||
max_bytes_per_request: Some(1),
|
||||
..Default::default()
|
||||
},
|
||||
);
|
||||
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)]));
|
||||
let batches = vec![
|
||||
Ok(RecordBatch::try_new(
|
||||
schema.clone(),
|
||||
vec![Arc::new(Int32Array::from(vec![1, 2, 3]))],
|
||||
)
|
||||
.unwrap()),
|
||||
Ok(RecordBatch::try_new(
|
||||
schema.clone(),
|
||||
vec![Arc::new(Int32Array::from(vec![4, 5, 6]))],
|
||||
)
|
||||
.unwrap()),
|
||||
];
|
||||
let reader: Box<dyn arrow_array::RecordBatchReader + Send> =
|
||||
Box::new(RecordBatchIterator::new(batches, schema));
|
||||
|
||||
let table = conn.create_table("table1", reader).execute().await.unwrap();
|
||||
|
||||
assert_eq!(table.name(), "table1");
|
||||
assert_eq!(create_count.load(Ordering::SeqCst), 1);
|
||||
assert_eq!(multipart_create_count.load(Ordering::SeqCst), 1);
|
||||
assert!(insert_count.load(Ordering::SeqCst) >= 1);
|
||||
assert_eq!(complete_count.load(Ordering::SeqCst), 1);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_create_table_already_exists() {
|
||||
let conn = Connection::new_with_handler(|_| {
|
||||
|
||||
@@ -441,6 +441,11 @@ impl<S: HttpSend> RemoteTable<S> {
|
||||
}
|
||||
}
|
||||
|
||||
/// Seed the schema cache when the caller already has the Arrow schema.
|
||||
pub(crate) fn seed_schema_ref(&self, schema: SchemaRef) {
|
||||
self.schema_cache.seed(schema);
|
||||
}
|
||||
|
||||
/// Return a new handle scoped to `branch`, sharing the client but with fresh
|
||||
/// caches and version/freshness state (the branch tracks its own latest).
|
||||
/// Mirrors `NativeTable`'s handle-per-branch model.
|
||||
@@ -1470,8 +1475,8 @@ impl<S: HttpSend + 'static> RemoteTable<S> {
|
||||
num_partitions: usize,
|
||||
) -> Result<()> {
|
||||
debug_assert!(
|
||||
output.rescannable,
|
||||
"multipart inserts require rescannable input for retry support"
|
||||
output.rescannable || num_partitions == 1,
|
||||
"non-rescannable multipart inserts require a single partition"
|
||||
);
|
||||
|
||||
let plan = Arc::new(
|
||||
@@ -2105,7 +2110,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
|
||||
let table_schema = self.schema().await?;
|
||||
let table_def = TableDefinition::try_from_rich_schema(table_schema.clone())?;
|
||||
|
||||
let num_partitions = if self.server_version.support_multipart_write() {
|
||||
let (num_partitions, use_multipart) = if self.server_version.support_multipart_write() {
|
||||
// Peek at the first batch to estimate write partitions (same as
|
||||
// NativeTable) and, regardless of `write_parallelism`, to detect a
|
||||
// fully empty input. A multipart write creates its upload session
|
||||
@@ -2115,10 +2120,12 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
|
||||
// commit and e.g. `mode=overwrite` would be silently dropped. Route
|
||||
// empty input through the single-request path instead, which always
|
||||
// sends one schema-only request.
|
||||
let unknown_size = add.data.num_rows().is_none();
|
||||
let mut peeked = PeekedScannable::new(add.data);
|
||||
let n = match peeked.peek().await {
|
||||
let first_batch = peeked.peek().await;
|
||||
let n = match first_batch.as_ref() {
|
||||
Some(first_batch) => match add.write_parallelism {
|
||||
Some(parallelism) if parallelism > 1 => parallelism,
|
||||
Some(parallelism) if parallelism > 1 && peeked.rescannable() => parallelism,
|
||||
Some(_) => 1,
|
||||
None => {
|
||||
let max_partitions =
|
||||
@@ -2133,10 +2140,14 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
|
||||
},
|
||||
None => 1,
|
||||
};
|
||||
// Unknown-length readers cannot be sized up-front, so use a
|
||||
// single-partition multipart upload. It remains streaming while
|
||||
// allowing the request body to be split into bounded parts.
|
||||
let use_multipart = first_batch.is_some() && (n > 1 || unknown_size);
|
||||
add.data = Box::new(peeked);
|
||||
n
|
||||
(n, use_multipart)
|
||||
} else {
|
||||
1
|
||||
(1, false)
|
||||
};
|
||||
|
||||
let output = add.into_plan(&table_schema, &table_def)?;
|
||||
@@ -2146,7 +2157,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
|
||||
}
|
||||
let _finish = FinishOnDrop(output.tracker.clone());
|
||||
|
||||
if num_partitions > 1 {
|
||||
if use_multipart {
|
||||
self.add_multipart(output, num_partitions).await
|
||||
} else {
|
||||
self.add_single_partition(output).await
|
||||
|
||||
@@ -304,68 +304,6 @@ mod tests {
|
||||
assert_eq!(all_values, expected);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_parallel_compaction_reserves_fragment_ids_once() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
let schema = Arc::new(Schema::new(vec![Field::new("i", DataType::Int32, false)]));
|
||||
let batch =
|
||||
RecordBatch::try_new(schema, vec![Arc::new(Int32Array::from_iter_values(0..10))])
|
||||
.unwrap();
|
||||
|
||||
let table = conn
|
||||
.create_table("test_parallel_compaction", batch.clone())
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
// Create 64 fragments. With a 20-row target, compaction plans 32 tasks,
|
||||
// which is more than the commit retry limit that used to be exhausted
|
||||
// when each parallel task reserved fragment IDs independently.
|
||||
for _ in 1..64 {
|
||||
table.add(batch.clone()).execute().await.unwrap();
|
||||
}
|
||||
|
||||
// Legacy row IDs require fragment IDs before an index can be remapped.
|
||||
assert!(
|
||||
!table
|
||||
.as_native()
|
||||
.unwrap()
|
||||
.manifest()
|
||||
.await
|
||||
.unwrap()
|
||||
.uses_stable_row_ids()
|
||||
);
|
||||
table
|
||||
.create_index(&["i"], Index::BTree(BTreeIndexBuilder::default()))
|
||||
.execute()
|
||||
.await
|
||||
.unwrap();
|
||||
|
||||
let version_before = table.version().await.unwrap();
|
||||
let stats = table
|
||||
.optimize(OptimizeAction::Compact {
|
||||
options: CompactionOptions {
|
||||
target_rows_per_fragment: 20,
|
||||
num_threads: Some(64),
|
||||
..Default::default()
|
||||
},
|
||||
remap_options: None,
|
||||
})
|
||||
.await
|
||||
.unwrap()
|
||||
.compaction
|
||||
.unwrap();
|
||||
|
||||
assert_eq!(stats.fragments_removed, 64);
|
||||
assert_eq!(stats.fragments_added, 32);
|
||||
assert_eq!(table.count_rows(None).await.unwrap(), 640);
|
||||
assert_eq!(
|
||||
table.version().await.unwrap(),
|
||||
version_before + 2,
|
||||
"parallel compaction should use one fragment reservation commit and one rewrite commit"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn test_optimize_prune_versions() {
|
||||
let conn = connect("memory://").execute().await.unwrap();
|
||||
|
||||
Reference in New Issue
Block a user