Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
14 commits
Select commit Hold shift + click to select a range
3132878
feat: store shared ScanScheduler on Dataset for reuse across fragment…
devin-ai-integration[bot] Mar 23, 2026
40f58f7
style: fix rustfmt formatting for shared scan scheduler tuple
devin-ai-integration[bot] Mar 23, 2026
ac2131b
feat: add FileScheduler cache on Dataset to avoid re-opening S3 files…
devin-ai-integration[bot] Mar 23, 2026
6f8a39b
style: fix import ordering for rustfmt compliance
devin-ai-integration[bot] Mar 23, 2026
521bd53
style: fix rustfmt formatting for FileScheduler cache code
devin-ai-integration[bot] Mar 23, 2026
d243bd7
perf: add S3 call counting and scheduler timing instrumentation
devin-ai-integration[bot] Mar 23, 2026
42f5f1e
perf: add S3 byte counting and V1/V2 format logging for profiling
devin-ai-integration[bot] Mar 23, 2026
42ee910
style: fix rustfmt formatting for V1/V2 format logging
devin-ai-integration[bot] Mar 23, 2026
339918b
perf: add V1 (CloudObjectReader) S3 call counting for profiling
devin-ai-integration[bot] Mar 23, 2026
dbd1244
style: fix import ordering in lib.rs
devin-ai-integration[bot] Mar 23, 2026
7c0a86c
feat: re-export S3 profiling counters through lance crate
devin-ai-integration[bot] Mar 23, 2026
8f1e055
perf: cache V1 PreviousFileReader instances on Dataset to avoid per-r…
devin-ai-integration[bot] Mar 23, 2026
696a900
fix: add missing v1_reader_cache field to Dataset initializer in comm…
devin-ai-integration[bot] Mar 24, 2026
7dc6deb
style: fix rustfmt formatting for v1_reader_cache insert
devin-ai-integration[bot] Mar 24, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 20 additions & 0 deletions rust/lance-encoding/src/decoder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand All @@ -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)
Expand All @@ -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))
Expand Down
1 change: 1 addition & 0 deletions rust/lance-io/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
20 changes: 20 additions & 0 deletions rust/lance-io/src/object_reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@

use std::ops::Range;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};

use bytes::Bytes;
use deepsize::DeepSizeOf;
Expand All @@ -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<GetResult>>;
Expand Down Expand Up @@ -176,6 +194,8 @@ impl Reader for CloudObjectReader {

#[instrument(level = "debug", skip(self))]
fn get_range(&self, range: Range<usize>) -> BoxFuture<'static, OSResult<Bytes>> {
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(),
Expand Down
26 changes: 26 additions & 0 deletions rust/lance-io/src/scheduler.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -844,6 +866,10 @@ impl ScanScheduler {
request: Vec<Range<u64>>,
priority: u128,
) -> impl Future<Output = Result<Vec<Bytes>>> + 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),
Expand Down
32 changes: 32 additions & 0 deletions rust/lance/src/dataset.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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};
Expand All @@ -28,13 +29,15 @@ 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};
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::{
Expand Down Expand Up @@ -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<Box<ObjectStoreParams>>,

/// 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<ScanScheduler>,

/// 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<DashMap<Path, FileScheduler>>,

/// 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<DashMap<Path, PreviousFileReader>>,
}

impl std::fmt::Debug for Dataset {
Expand Down Expand Up @@ -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,
Expand All @@ -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()),
})
}

Expand Down Expand Up @@ -1561,6 +1586,13 @@ impl Dataset {
store_params: Option<ObjectStoreParams>,
) -> 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));
Expand Down
91 changes: 71 additions & 20 deletions rust/lance/src/dataset/fragment.rs
Original file line number Diff line number Diff line change
Expand Up @@ -934,23 +934,43 @@ 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
.dataset
.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,
Expand All @@ -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)?
Expand All @@ -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);
Expand Down
11 changes: 11 additions & 0 deletions rust/lance/src/dataset/write/commit.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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},
Expand Down Expand Up @@ -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,
Expand All @@ -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,
})
}
}
Expand Down
8 changes: 8 additions & 0 deletions rust/lance/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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.
///
Expand Down
Loading