Compare commits

..

1 Commits

Author SHA1 Message Date
Gatefixer 803dd21ecb fix(remote): stream create table readers with multipart writes 2026-08-05 21:15:47 +00:00
14 changed files with 420 additions and 303 deletions
+32 -8
View File
@@ -276,14 +276,38 @@ 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 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: 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: cargo +${{ matrix.msrv }} check
env:
RUSTUP_TOOLCHAIN: ${{ matrix.msrv }}
Generated
+233 -257
View File
File diff suppressed because it is too large Load Diff
+14 -14
View File
@@ -13,20 +13,20 @@ categories = ["database-implementations"]
rust-version = "1.91.0"
[workspace.dependencies]
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" }
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" }
ahash = "0.8"
# Note that this one does not include pyarrow
arrow = { version = "58.0.0", optional = false }
+1 -1
View File
@@ -28,7 +28,7 @@
<properties>
<project.build.sourceEncoding>UTF-8</project.build.sourceEncoding>
<arrow.version>15.0.0</arrow.version>
<lance-core.version>11.0.0-beta.1</lance-core.version>
<lance-core.version>10.1.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>
+2 -2
View File
@@ -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 GooseFS SDK to the version required by Lance's OpenDAL dependency.
goosefs-sdk = { version = "=0.1.9", 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 }
moka = { workspace = true }
pin-project = { workspace = true }
tokio = { version = "1.23", features = ["rt-multi-thread", "sync"] }
+1 -1
View File
@@ -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_file::version::LanceFileVersion;
use lance_encoding::version::LanceFileVersion;
use lance_io::object_store::ObjectStore;
use object_store::path::Path;
+1 -1
View File
@@ -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_file::version::LanceFileVersion;
pub use lance_encoding::version::LanceFileVersion;
#[cfg(feature = "remote")]
use lance_io::object_store::StorageOptions;
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
+1 -1
View File
@@ -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_file::version::LanceFileVersion;
use lance_encoding::version::LanceFileVersion;
use lance_io::object_store::{StorageOptionsAccessor, StorageOptionsProvider};
use lance_table::io::commit::commit_handler_from_url;
use object_store::local::LocalFileSystem;
+2 -2
View File
@@ -201,7 +201,7 @@ impl LanceNamespaceDatabase {
&self,
request: &DbCreateTableRequest,
) -> Result<(
Option<lance_file::version::LanceFileVersion>,
Option<lance_encoding::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_file::version::LanceFileVersion>())
.map(|s| s.parse::<lance_encoding::version::LanceFileVersion>())
.transpose()?;
let v2_manifest_override = storage_options
+109 -3
View File
@@ -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(|_| {
+20 -9
View File
@@ -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
@@ -2939,7 +2950,7 @@ impl<S: HttpSend> BaseTable for RemoteTable<S> {
}
#[derive(Serialize, Clone, Debug)]
pub struct MergeInsertRequest {
pub(crate) struct MergeInsertRequest {
on: String,
when_matched_update_all: bool,
when_matched_update_all_filt: Option<String>,
+1 -1
View File
@@ -90,7 +90,7 @@ struct RemoteBlobState {
/// Seekable Cloud blob handle over HTTP Range.
#[derive(Debug)]
pub struct RemoteBlobFile {
pub(crate) struct RemoteBlobFile {
requester: Arc<dyn BlobRangeRequester>,
state: Mutex<RemoteBlobState>,
closed: AtomicBool,
+2 -2
View File
@@ -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 enum WriteOp {
pub(crate) 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 enum WriteOp {
/// The parsed server response for a completed write, discriminated by the
/// operation that produced it.
#[derive(Debug, Clone)]
pub enum WriteResult {
pub(crate) enum WriteResult {
Add(AddResult),
Merge(MergeResult),
}
+1 -1
View File
@@ -10,7 +10,7 @@ use arrow_array::{
use arrow_schema::{DataType, Field, Fields, Schema};
use futures::TryStreamExt;
use lance::Dataset;
use lance_file::version::LanceFileVersion;
use lance_encoding::version::LanceFileVersion;
use lancedb::{
Connection, Error, Result, Table,
blob::{BlobRangeRequest, blob},