diff --git a/rust/lance-encoding/src/decoder.rs b/rust/lance-encoding/src/decoder.rs index ddbc3da38e3..492558a056d 100644 --- a/rust/lance-encoding/src/decoder.rs +++ b/rust/lance-encoding/src/decoder.rs @@ -2012,6 +2012,7 @@ fn create_scheduler_decoder( )?; let scheduler_handle = tokio::task::spawn(async move { + let sched_start = std::time::Instant::now(); let mut decode_scheduler = match DecodeBatchScheduler::try_new( target_schema.as_ref(), &column_indices, @@ -2032,7 +2033,9 @@ fn create_scheduler_decoder( return; } }; + let sched_create_us = sched_start.elapsed().as_micros() as u64; + let schedule_start = std::time::Instant::now(); match requested_rows { RequestedRows::Ranges(ranges) => { decode_scheduler.schedule_ranges(&ranges, &filter, tx, config.io) @@ -2041,6 +2044,23 @@ fn create_scheduler_decoder( decode_scheduler.schedule_take(&indices, &filter, tx, config.io) } } + let schedule_us = schedule_start.elapsed().as_micros() as u64; + + use std::sync::atomic::{AtomicU64, Ordering}; + static SCHED_CREATE_TOTAL_US: AtomicU64 = AtomicU64::new(0); + static SCHED_SCHEDULE_TOTAL_US: AtomicU64 = AtomicU64::new(0); + static SCHED_COUNT: AtomicU64 = AtomicU64::new(0); + SCHED_CREATE_TOTAL_US.fetch_add(sched_create_us, Ordering::Relaxed); + SCHED_SCHEDULE_TOTAL_US.fetch_add(schedule_us, Ordering::Relaxed); + let count = SCHED_COUNT.fetch_add(1, Ordering::Relaxed) + 1; + if count % 500 == 0 { + eprintln!( + "[lance-sched-perf] schedulers={} create_total_ms={} schedule_total_ms={}", + count, + SCHED_CREATE_TOTAL_US.load(Ordering::Relaxed) / 1000, + SCHED_SCHEDULE_TOTAL_US.load(Ordering::Relaxed) / 1000, + ); + } }); Ok(check_scheduler_on_drop(decode_stream, scheduler_handle)) diff --git a/rust/lance-io/src/lib.rs b/rust/lance-io/src/lib.rs index e1729db73be..916e830ecf0 100644 --- a/rust/lance-io/src/lib.rs +++ b/rust/lance-io/src/lib.rs @@ -23,6 +23,7 @@ pub mod testing; pub mod traits; pub mod utils; +pub use object_reader::{reset_v1_s3_counters, v1_s3_bytes, v1_s3_calls}; pub use scheduler::{bytes_read_counter, iops_counter}; /// Defines a selection of rows to read from a file/batch diff --git a/rust/lance-io/src/object_reader.rs b/rust/lance-io/src/object_reader.rs index d6b34671aed..8761c7bc44e 100644 --- a/rust/lance-io/src/object_reader.rs +++ b/rust/lance-io/src/object_reader.rs @@ -3,6 +3,7 @@ use std::ops::Range; use std::sync::Arc; +use std::sync::atomic::{AtomicU64, Ordering}; use bytes::Bytes; use deepsize::DeepSizeOf; @@ -17,6 +18,23 @@ use tracing::instrument; use crate::{object_store::DEFAULT_CLOUD_IO_PARALLELISM, traits::Reader}; +// Global counters for V1 (CloudObjectReader) S3 calls +static V1_S3_CALLS: AtomicU64 = AtomicU64::new(0); +static V1_S3_BYTES: AtomicU64 = AtomicU64::new(0); + +pub fn v1_s3_calls() -> u64 { + V1_S3_CALLS.load(Ordering::Acquire) +} + +pub fn v1_s3_bytes() -> u64 { + V1_S3_BYTES.load(Ordering::Acquire) +} + +pub fn reset_v1_s3_counters() { + V1_S3_CALLS.store(0, Ordering::Release); + V1_S3_BYTES.store(0, Ordering::Release); +} + trait StaticGetRange { fn path(&self) -> &Path; fn get_range(&self) -> BoxFuture<'static, OSResult>; @@ -176,6 +194,8 @@ impl Reader for CloudObjectReader { #[instrument(level = "debug", skip(self))] fn get_range(&self, range: Range) -> BoxFuture<'static, OSResult> { + V1_S3_CALLS.fetch_add(1, Ordering::Relaxed); + V1_S3_BYTES.fetch_add((range.end - range.start) as u64, Ordering::Relaxed); let get_request = Arc::new(GetRequest { object_store: self.object_store.clone(), path: self.path.clone(), diff --git a/rust/lance-io/src/scheduler.rs b/rust/lance-io/src/scheduler.rs index e27e9519425..5cd43b4141e 100644 --- a/rust/lance-io/src/scheduler.rs +++ b/rust/lance-io/src/scheduler.rs @@ -32,6 +32,28 @@ const BACKPRESSURE_DEBOUNCE: u64 = 60; static IOPS_COUNTER: AtomicU64 = AtomicU64::new(0); // Global counter of how many bytes were read by the scheduler static BYTES_READ_COUNTER: AtomicU64 = AtomicU64::new(0); +// Global counter of how many S3 requests were submitted (for profiling) +static S3_REQUESTS_COUNTER: AtomicU64 = AtomicU64::new(0); +static S3_IOPS_COUNTER: AtomicU64 = AtomicU64::new(0); +static S3_BYTES_REQUESTED: AtomicU64 = AtomicU64::new(0); + +pub fn s3_requests_counter() -> u64 { + S3_REQUESTS_COUNTER.load(Ordering::Acquire) +} + +pub fn s3_iops_counter() -> u64 { + S3_IOPS_COUNTER.load(Ordering::Acquire) +} + +pub fn s3_bytes_requested() -> u64 { + S3_BYTES_REQUESTED.load(Ordering::Acquire) +} + +pub fn reset_s3_counters() { + S3_REQUESTS_COUNTER.store(0, Ordering::Release); + S3_IOPS_COUNTER.store(0, Ordering::Release); + S3_BYTES_REQUESTED.store(0, Ordering::Release); +} // By default, we limit the number of IOPS across the entire process to 128 // // In theory this is enough for ~10GBps on S3 following the guidelines to issue @@ -844,6 +866,10 @@ impl ScanScheduler { request: Vec>, priority: u128, ) -> impl Future>> + Send + use<> { + S3_REQUESTS_COUNTER.fetch_add(1, Ordering::Relaxed); + S3_IOPS_COUNTER.fetch_add(request.len() as u64, Ordering::Relaxed); + let total_bytes: u64 = request.iter().map(|r| r.end - r.start).sum(); + S3_BYTES_REQUESTED.fetch_add(total_bytes, Ordering::Relaxed); match &self.io_queue { IoQueueType::Standard(io_queue) => futures::future::Either::Left( self.submit_request_standard(reader, request, priority, io_queue), diff --git a/rust/lance/src/dataset.rs b/rust/lance/src/dataset.rs index 6adefc7bf3d..e588971cf22 100644 --- a/rust/lance/src/dataset.rs +++ b/rust/lance/src/dataset.rs @@ -17,6 +17,7 @@ use crate::dataset::metadata::UpdateFieldMetadataBuilder; use crate::dataset::transaction::translate_schema_metadata_updates; use crate::session::caches::{DSMetadataCache, ManifestKey, TransactionKey}; use crate::session::index_caches::DSIndexCache; +use dashmap::DashMap; use itertools::Itertools; use lance_core::ROW_ADDR; use lance_core::datatypes::{OnMissing, OnTypeMismatch, Projectable, Projection}; @@ -28,6 +29,7 @@ use lance_core::utils::tracing::{ }; use lance_datafusion::projection::ProjectionPlan; use lance_file::datatypes::populate_schema_dictionary; +use lance_file::previous::reader::FileReader as PreviousFileReader; use lance_file::reader::FileReaderOptions; use lance_file::version::LanceFileVersion; use lance_index::{DatasetIndexExt, IndexType}; @@ -35,6 +37,7 @@ use lance_io::object_store::{ LanceNamespaceStorageOptionsProvider, ObjectStore, ObjectStoreParams, StorageOptions, StorageOptionsAccessor, StorageOptionsProvider, }; +use lance_io::scheduler::{FileScheduler, ScanScheduler, SchedulerConfig}; use lance_io::utils::{read_last_block, read_message, read_metadata_offset, read_struct}; use lance_namespace::LanceNamespace; use lance_table::format::{ @@ -175,6 +178,20 @@ pub struct Dataset { /// Object store parameters used when opening this dataset. /// These are used when creating object stores for additional base paths. pub(crate) store_params: Option>, + + /// Shared scan scheduler for reading data files. + /// Created once per dataset and reused across all fragment reads to avoid + /// creating new HTTP connection pools per read. + pub scan_scheduler: Arc, + + /// Cache of opened FileScheduler instances, keyed by file path. + /// Avoids re-opening S3 file handles for the same data file across requests. + pub file_scheduler_cache: Arc>, + + /// Cache of opened V1 (legacy) FileReader instances, keyed by file path. + /// Avoids re-creating PreviousFileReader objects (and their S3 metadata reads) + /// for the same data file across requests. + pub v1_reader_cache: Arc>, } impl std::fmt::Debug for Dataset { @@ -706,6 +723,11 @@ impl Dataset { let metadata_cache = Arc::new(session.metadata_cache.for_dataset(&uri)); let index_cache = Arc::new(session.index_cache.for_dataset(&uri)); let fragment_bitmap = Arc::new(manifest.fragments.iter().map(|f| f.id as u32).collect()); + let scan_scheduler = ScanScheduler::new( + object_store.clone(), + SchedulerConfig::max_bandwidth(&object_store), + ); + let file_scheduler_cache = Arc::new(DashMap::new()); Ok(Self { object_store, base: base_path, @@ -720,6 +742,9 @@ impl Dataset { index_cache, file_reader_options, store_params: store_params.map(Box::new), + scan_scheduler, + file_scheduler_cache, + v1_reader_cache: Arc::new(DashMap::new()), }) } @@ -1561,6 +1586,13 @@ impl Dataset { store_params: Option, ) -> Self { let mut cloned = self.clone(); + // Create a new scan scheduler for the new object store + cloned.scan_scheduler = ScanScheduler::new( + object_store.clone(), + SchedulerConfig::max_bandwidth(&object_store), + ); + // Clear the file scheduler cache since it's tied to the old object store + cloned.file_scheduler_cache = Arc::new(DashMap::new()); cloned.object_store = object_store; if let Some(store_params) = store_params { cloned.store_params = Some(Box::new(store_params)); diff --git a/rust/lance/src/dataset/fragment.rs b/rust/lance/src/dataset/fragment.rs index a4db60cf56f..051242b2356 100644 --- a/rust/lance/src/dataset/fragment.rs +++ b/rust/lance/src/dataset/fragment.rs @@ -934,6 +934,14 @@ impl FileFragment { let schema_per_file = Arc::new(projection.intersection_ignore_types(&data_file_schema)?); if data_file.is_legacy_file() { + use std::sync::atomic::{AtomicBool, Ordering as AtomOrd}; + static LOGGED_V1: AtomicBool = AtomicBool::new(false); + if !LOGGED_V1.swap(true, AtomOrd::Relaxed) { + eprintln!( + "[lance-format] data_file uses V1 (legacy) format: version={}.{} path={}", + data_file.file_major_version, data_file.file_minor_version, data_file.path + ); + } let max_field_id = data_file.fields.iter().max().unwrap(); if !schema_per_file.fields.is_empty() { let path = self @@ -941,16 +949,28 @@ impl FileFragment { .data_file_dir(data_file)? .child(data_file.path.as_str()); let field_id_offset = Self::get_field_id_offset(data_file); - let reader = PreviousFileReader::try_new_with_fragment_id( - &self.dataset.object_store, - &path, - self.schema().clone(), - self.id() as u32, - field_id_offset as i32, - *max_field_id, - Some(&self.dataset.metadata_cache.file_metadata_cache(&path)), - ) - .await?; + + // Try to get a cached V1 reader first (avoids re-opening S3 file + // handles and re-reading metadata/page tables on every request) + let reader = if let Some(cached) = self.dataset.v1_reader_cache.get(&path) { + cached.clone() + } else { + let new_reader = PreviousFileReader::try_new_with_fragment_id( + &self.dataset.object_store, + &path, + self.schema().clone(), + self.id() as u32, + field_id_offset as i32, + *max_field_id, + Some(&self.dataset.metadata_cache.file_metadata_cache(&path)), + ) + .await?; + self.dataset + .v1_reader_cache + .insert(path.clone(), new_reader.clone()); + new_reader + }; + let initialized_schema = reader.schema().project_by_schema( schema_per_file.as_ref(), OnMissing::Error, @@ -964,6 +984,16 @@ impl FileFragment { } else if schema_per_file.fields.is_empty() { Ok(None) } else { + { + use std::sync::atomic::{AtomicBool, Ordering as AtomOrd}; + static LOGGED_V2: AtomicBool = AtomicBool::new(false); + if !LOGGED_V2.swap(true, AtomOrd::Relaxed) { + eprintln!( + "[lance-format] data_file uses V2 format: version={}.{} path={}", + data_file.file_major_version, data_file.file_minor_version, data_file.path + ); + } + } let path = self .dataset .data_file_dir(data_file)? @@ -983,17 +1013,38 @@ impl FileFragment { read_config.reader_priority.unwrap_or(0), ) } else { - ( - ScanScheduler::new( - self.dataset.object_store.clone(), - SchedulerConfig::max_bandwidth(&self.dataset.object_store), - ), - 0, - ) + // Reuse the dataset's shared scan scheduler to avoid creating + // new HTTP connection pools per fragment read. + (self.dataset.scan_scheduler.clone(), 0) }; - let file_scheduler = store_scheduler - .open_file_with_priority(&path, reader_priority as u64, &data_file.file_size_bytes) - .await?; + // Check the dataset's file scheduler cache first to avoid + // re-opening S3 file handles for the same data file across requests. + let file_scheduler = + if data_file.base_id.is_none() && read_config.scan_scheduler.is_none() { + if let Some(cached) = self.dataset.file_scheduler_cache.get(&path) { + cached.value().clone() + } else { + let scheduler = store_scheduler + .open_file_with_priority( + &path, + reader_priority as u64, + &data_file.file_size_bytes, + ) + .await?; + self.dataset + .file_scheduler_cache + .insert(path.clone(), scheduler.clone()); + scheduler + } + } else { + store_scheduler + .open_file_with_priority( + &path, + reader_priority as u64, + &data_file.file_size_bytes, + ) + .await? + }; let file_metadata = self.get_file_metadata(&file_scheduler).await?; let path = file_scheduler.reader().path().clone(); let metadata_cache = self.dataset.metadata_cache.file_metadata_cache(&path); diff --git a/rust/lance/src/dataset/write/commit.rs b/rust/lance/src/dataset/write/commit.rs index 0a85d27ce99..32e28ebc84e 100644 --- a/rust/lance/src/dataset/write/commit.rs +++ b/rust/lance/src/dataset/write/commit.rs @@ -4,9 +4,11 @@ use std::collections::HashMap; use std::sync::Arc; +use dashmap::DashMap; use lance_core::utils::mask::RowAddrTreeMap; use lance_file::version::LanceFileVersion; use lance_io::object_store::{ObjectStore, ObjectStoreParams}; +use lance_io::scheduler::{ScanScheduler, SchedulerConfig}; use lance_table::{ format::{DataStorageFormat, is_detached_version}, io::commit::{CommitConfig, CommitHandler, ManifestNamingScheme}, @@ -388,6 +390,12 @@ impl<'a> CommitBuilder<'a> { branch: manifest.branch.clone(), }, ); + let scan_scheduler = ScanScheduler::new( + object_store.clone(), + SchedulerConfig::max_bandwidth(&object_store), + ); + let file_scheduler_cache = Arc::new(DashMap::new()); + let v1_reader_cache = Arc::new(DashMap::new()); Ok(Dataset { object_store, @@ -403,6 +411,9 @@ impl<'a> CommitBuilder<'a> { metadata_cache, file_reader_options: None, store_params: self.store_params.clone().map(Box::new), + scan_scheduler, + file_scheduler_cache, + v1_reader_cache, }) } } diff --git a/rust/lance/src/lib.rs b/rust/lance/src/lib.rs index 934be0e519c..4e06656b231 100644 --- a/rust/lance/src/lib.rs +++ b/rust/lance/src/lib.rs @@ -89,6 +89,14 @@ pub use blob::{BlobArrayBuilder, blob_field}; pub use dataset::Dataset; use lance_index::vector::DIST_COL; +/// Re-export S3 profiling counters from lance-io for external instrumentation +pub mod profiling { + pub use lance_io::scheduler::{ + reset_s3_counters, s3_bytes_requested, s3_iops_counter, s3_requests_counter, + }; + pub use lance_io::{reset_v1_s3_counters, v1_s3_bytes, v1_s3_calls}; +} + /// Creates and loads a [`Dataset`] from the given path. /// Infers the storage backend to use from the scheme in the given table path. ///