use std::{
collections::VecDeque,
fmt, io,
ops::Range,
sync::{Arc, Mutex},
time::{Duration, Instant},
};
use async_trait::async_trait;
use bytes::Bytes;
use futures_util::{StreamExt, stream, stream::BoxStream};
use object_store::{
CopyOptions, GetOptions, GetResult, GetResultPayload, ListResult, MultipartUpload,
OBJECT_STORE_COALESCE_DEFAULT, ObjectMeta, ObjectStore, ObjectStoreScheme, PutMultipartOptions,
PutOptions, PutPayload, PutResult, RenameOptions, Result, path::Path,
};
use tracing::Instrument;
use url::Url;
use super::range_planning::{
ChosenRangePlan, RangePlanDecision, TransportEstimate, bandwidth_delay_bytes,
choose_range_plan, execute_range_plan, merge_ranges, range_bytes, request_waves,
};
use crate::{
DeltaScanMetrics,
reader::{ParquetRangeReadPolicy, options::MAX_CONCURRENT_PARQUET_RANGE_READS},
};
const TRANSPORT_SAMPLE_WINDOW: usize = 9;
const MIN_TRANSPORT_SAMPLES: usize = 3;
const MIN_THROUGHPUT_SAMPLE_BYTES: u128 = 1024 * 1024;
const MIN_THROUGHPUT_SAMPLE_DELIVERY_TIME: Duration = Duration::from_millis(10);
const RANGE_PLANNING_DIAGNOSTIC_TARGET: &str =
"delta_arrow_reader::diagnostics::parquet_range_planning";
pub(crate) struct MeteredParquetObjectStore {
inner: Arc<dyn ObjectStore>,
metrics: DeltaScanMetrics,
multi_range_read_strategy: MultiRangeReadStrategy,
range_read_estimator: Arc<ParquetRangeReadEstimator>,
}
#[derive(Default)]
pub(crate) struct ParquetRangeReadEstimator {
samples: Mutex<TransportSampleWindows>,
}
#[derive(Default)]
struct TransportSampleWindows {
latencies: VecDeque<Duration>,
throughputs: VecDeque<ThroughputSample>,
}
#[derive(Clone, Copy)]
struct ThroughputSample {
bytes_received: u128,
bytes_per_second: u64,
}
struct TransportObservation {
request_latency: Duration,
throughput: Option<ThroughputSample>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
struct CurrentTransportEstimate {
latency_sample_count: usize,
throughput_sample_count: usize,
estimate: Option<TransportEstimate>,
}
struct CompletedRangeRead {
request_latency: Duration,
payload_started: Instant,
payload_finished: Instant,
bytes_received: usize,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum MultiRangeReadStrategy {
UseStoreImplementation,
ReadExactRanges,
MergeRangesWithinOneMegabyte,
ChooseAutomatically,
}
impl MultiRangeReadStrategy {
pub(crate) fn for_policy(policy: ParquetRangeReadPolicy, table_url: &Url) -> Self {
match policy {
ParquetRangeReadPolicy::Automatic => Self::for_table_url(table_url),
ParquetRangeReadPolicy::ExactRanges => Self::ReadExactRanges,
ParquetRangeReadPolicy::MergeRangesWithinOneMegabyte => {
Self::MergeRangesWithinOneMegabyte
}
ParquetRangeReadPolicy::StoreImplementation => Self::UseStoreImplementation,
}
}
pub(crate) fn for_table_url(table_url: &Url) -> Self {
match ObjectStoreScheme::parse(table_url) {
Ok((ObjectStoreScheme::Local | ObjectStoreScheme::Memory, _)) | Err(_) => {
Self::UseStoreImplementation
}
Ok(_) => Self::ChooseAutomatically,
}
}
}
impl MeteredParquetObjectStore {
pub(crate) fn new(
inner: Arc<dyn ObjectStore>,
metrics: DeltaScanMetrics,
multi_range_read_strategy: MultiRangeReadStrategy,
) -> Self {
Self {
inner,
metrics,
multi_range_read_strategy,
range_read_estimator: Arc::new(ParquetRangeReadEstimator::default()),
}
}
pub(crate) fn with_range_read_estimator(
mut self,
range_read_estimator: Arc<ParquetRangeReadEstimator>,
) -> Self {
self.range_read_estimator = range_read_estimator;
self
}
async fn read_range_with_timing(
&self,
location: &Path,
range: Range<u64>,
) -> Result<(Bytes, CompletedRangeRead)> {
let expected_bytes = range.end - range.start;
let request_started = Instant::now();
let result = self
.get_opts(location, GetOptions::new().with_range(Some(range)))
.await?;
let payload_started = Instant::now();
let bytes = result.bytes().await?;
if u64::try_from(bytes.len()).unwrap_or(u64::MAX) != expected_bytes {
return Err(object_store::Error::Generic {
store: "delta-arrow-reader",
source: io::Error::new(
io::ErrorKind::UnexpectedEof,
"object store returned a byte range with an unexpected length",
)
.into(),
});
}
let payload_finished = Instant::now();
let completed_read = CompletedRangeRead {
request_latency: payload_started.saturating_duration_since(request_started),
payload_started,
payload_finished,
bytes_received: bytes.len(),
};
Ok((bytes, completed_read))
}
fn current_transport_estimate(&self) -> CurrentTransportEstimate {
self.range_read_estimator.current_transport_estimate()
}
fn record_completed_range_reads(&self, completed_reads: &[CompletedRangeRead]) {
let Some(observation) = transport_observation(completed_reads) else {
return;
};
self.range_read_estimator.record(observation);
}
fn record_chosen_range_plan_metrics(&self, plan: &ChosenRangePlan) {
self.metrics
.record_parquet_data_file_exact_ranges_requested(
plan.exact_range_count,
plan.exact_bytes,
);
self.metrics.record_parquet_data_file_physical_range_plan(
plan.physical_ranges.len(),
plan.planned_bytes,
);
match plan.decision {
RangePlanDecision::ColdStart => self
.metrics
.record_parquet_data_file_cold_start_range_plan(),
RangePlanDecision::CostBasedExact => self
.metrics
.record_parquet_data_file_cost_based_exact_range_plan(),
RangePlanDecision::CostBasedMerged => self
.metrics
.record_parquet_data_file_cost_based_merged_range_plan(),
}
}
async fn read_physical_ranges(
&self,
location: &Path,
requested_ranges: &[Range<u64>],
physical_ranges: &[Range<u64>],
) -> Result<(Vec<Bytes>, Duration)> {
let plan_started = Instant::now();
let completed_reads = Arc::new(Mutex::new(Vec::with_capacity(physical_ranges.len())));
let results = execute_range_plan(requested_ranges, physical_ranges, |range| {
let completed_reads = Arc::clone(&completed_reads);
async move {
let (bytes, completed_read) = self.read_range_with_timing(location, range).await?;
completed_reads
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(completed_read);
Ok::<Bytes, object_store::Error>(bytes)
}
})
.await?;
let observed_plan_time = plan_started.elapsed();
self.record_completed_range_reads(
&completed_reads
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner),
);
self.metrics
.record_parquet_range_successful_plan_time(plan_started.elapsed());
Ok((results, observed_plan_time))
}
}
impl ParquetRangeReadEstimator {
fn current_transport_estimate(&self) -> CurrentTransportEstimate {
let samples = self
.samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let latency_sample_count = samples.latencies.len();
let throughput_sample_count = samples.throughputs.len();
if latency_sample_count < MIN_TRANSPORT_SAMPLES
|| throughput_sample_count < MIN_TRANSPORT_SAMPLES
{
return CurrentTransportEstimate {
latency_sample_count,
throughput_sample_count,
estimate: None,
};
}
let mut latencies = samples.latencies.iter().copied().collect::<Vec<_>>();
latencies.sort_unstable();
let mut throughputs = samples.throughputs.iter().copied().collect::<Vec<_>>();
throughputs.sort_unstable_by_key(|sample| sample.bytes_per_second);
let middle_byte = throughputs
.iter()
.fold(0_u128, |total, sample| {
total.saturating_add(sample.bytes_received)
})
.div_ceil(2);
let mut accumulated_bytes = 0_u128;
let shared_throughput_bytes_per_second = throughputs
.into_iter()
.find_map(|sample| {
accumulated_bytes = accumulated_bytes.saturating_add(sample.bytes_received);
(accumulated_bytes >= middle_byte).then_some(sample.bytes_per_second)
})
.unwrap_or(0);
CurrentTransportEstimate {
latency_sample_count,
throughput_sample_count,
estimate: Some(TransportEstimate {
request_latency: latencies[latency_sample_count / 2],
shared_throughput_bytes_per_second,
}),
}
}
fn record(&self, observation: TransportObservation) {
let mut samples = self
.samples
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if samples.latencies.len() == TRANSPORT_SAMPLE_WINDOW {
samples.latencies.pop_front();
}
samples.latencies.push_back(observation.request_latency);
if let Some(throughput) = observation.throughput {
if samples.throughputs.len() == TRANSPORT_SAMPLE_WINDOW {
samples.throughputs.pop_front();
}
samples.throughputs.push_back(throughput);
}
}
}
fn trace_completed_automatic_range_plan(
plan: &ChosenRangePlan,
estimate: CurrentTransportEstimate,
observed_plan_time: Duration,
) {
let plan_estimate = plan.transport_estimate;
tracing::debug!(
target: RANGE_PLANNING_DIAGNOSTIC_TARGET,
event = "automatic_parquet_range_plan_completed",
latency_sample_count = u64::try_from(estimate.latency_sample_count).unwrap_or(u64::MAX),
throughput_sample_count =
u64::try_from(estimate.throughput_sample_count).unwrap_or(u64::MAX),
estimated_request_latency_micros =
plan_estimate.map(|value| value.request_latency.as_micros()),
estimated_shared_throughput_bytes_per_second =
plan_estimate.map(|value| value.shared_throughput_bytes_per_second),
estimated_bandwidth_delay_bytes = plan_estimate.map(bandwidth_delay_bytes),
exact_range_count = u64::try_from(plan.exact_range_count).unwrap_or(u64::MAX),
exact_bytes = plan.exact_bytes,
exact_request_waves = u64::try_from(request_waves(
plan.exact_range_count,
MAX_CONCURRENT_PARQUET_RANGE_READS,
))
.unwrap_or(u64::MAX),
baseline_range_count = u64::try_from(plan.baseline_range_count).unwrap_or(u64::MAX),
baseline_bytes = plan.baseline_bytes,
baseline_request_waves = u64::try_from(request_waves(
plan.baseline_range_count,
MAX_CONCURRENT_PARQUET_RANGE_READS,
))
.unwrap_or(u64::MAX),
selected_range_count = u64::try_from(plan.physical_ranges.len()).unwrap_or(u64::MAX),
selected_bytes = plan.planned_bytes,
selected_request_waves = u64::try_from(request_waves(
plan.physical_ranges.len(),
MAX_CONCURRENT_PARQUET_RANGE_READS,
))
.unwrap_or(u64::MAX),
baseline_predicted_cost_bytes = plan.baseline_predicted_cost_bytes,
selected_predicted_cost_bytes = plan.selected_predicted_cost_bytes,
observed_selected_plan_micros = observed_plan_time.as_micros(),
decision = plan.decision.as_str(),
"Automatic Parquet range plan completed"
);
}
fn transport_observation(completed_reads: &[CompletedRangeRead]) -> Option<TransportObservation> {
if completed_reads.is_empty() {
return None;
}
let mut latencies = completed_reads
.iter()
.map(|read| read.request_latency)
.collect::<Vec<_>>();
latencies.sort_unstable();
let request_latency = latencies[latencies.len() / 2];
let mut delivery_intervals = completed_reads
.iter()
.map(|read| (read.payload_started, read.payload_finished))
.collect::<Vec<_>>();
delivery_intervals.sort_unstable();
let (mut interval_start, mut interval_end) = delivery_intervals[0];
let mut delivery_time = Duration::ZERO;
for (next_start, next_end) in delivery_intervals.into_iter().skip(1) {
if next_start <= interval_end {
interval_end = interval_end.max(next_end);
} else {
delivery_time = delivery_time
.saturating_add(interval_end.saturating_duration_since(interval_start));
(interval_start, interval_end) = (next_start, next_end);
}
}
delivery_time =
delivery_time.saturating_add(interval_end.saturating_duration_since(interval_start));
let bytes_received = completed_reads.iter().fold(0_u128, |total, read| {
total.saturating_add(read.bytes_received as u128)
});
let throughput = (bytes_received >= MIN_THROUGHPUT_SAMPLE_BYTES
&& delivery_time >= MIN_THROUGHPUT_SAMPLE_DELIVERY_TIME)
.then(|| ThroughputSample {
bytes_received,
bytes_per_second: u64::try_from(
bytes_received
.saturating_mul(1_000_000_000)
.checked_div(delivery_time.as_nanos())
.unwrap_or(0),
)
.unwrap_or(u64::MAX),
});
Some(TransportObservation {
request_latency,
throughput,
})
}
impl fmt::Debug for MeteredParquetObjectStore {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("MeteredParquetObjectStore")
}
}
impl fmt::Display for MeteredParquetObjectStore {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("MeteredParquetObjectStore")
}
}
#[async_trait]
impl ObjectStore for MeteredParquetObjectStore {
async fn put_opts(
&self,
location: &Path,
payload: PutPayload,
options: PutOptions,
) -> Result<PutResult> {
self.inner.put_opts(location, payload, options).await
}
async fn put_multipart_opts(
&self,
location: &Path,
options: PutMultipartOptions,
) -> Result<Box<dyn MultipartUpload>> {
self.inner.put_multipart_opts(location, options).await
}
async fn get_opts(&self, location: &Path, options: GetOptions) -> Result<GetResult> {
let should_meter_payload = !options.head;
if should_meter_payload {
if options.range.is_some() {
self.metrics.record_parquet_data_file_range_get_operation();
} else {
self.metrics.record_parquet_data_file_full_get_operation();
}
}
let result = self
.inner
.get_opts(location, options)
.instrument(tracing::debug_span!(
target: "delta_arrow_reader::profile",
"Object store transport"
))
.await?;
if should_meter_payload {
Ok(meter_get_result(result, self.metrics.clone()))
} else {
Ok(result)
}
}
async fn get_ranges(&self, location: &Path, ranges: &[Range<u64>]) -> Result<Vec<Bytes>> {
if ranges.is_empty() {
return Ok(Vec::new());
}
match self.multi_range_read_strategy {
MultiRangeReadStrategy::UseStoreImplementation => {
let plan_started = Instant::now();
let exact_ranges = merge_ranges(ranges, 0);
self.metrics
.record_parquet_data_file_exact_ranges_requested(
exact_ranges.len(),
range_bytes(&exact_ranges),
);
self.metrics
.record_parquet_data_file_store_delegated_range_plan();
self.metrics.record_parquet_data_file_range_get_operation();
let results = self
.inner
.get_ranges(location, ranges)
.instrument(tracing::debug_span!(
target: "delta_arrow_reader::profile",
"Object store transport"
))
.await?;
for bytes in &results {
self.metrics
.record_parquet_data_file_bytes_received(bytes.len());
}
self.metrics
.record_parquet_range_successful_plan_time(plan_started.elapsed());
Ok(results)
}
strategy @ (MultiRangeReadStrategy::ReadExactRanges
| MultiRangeReadStrategy::MergeRangesWithinOneMegabyte) => {
let exact_ranges = merge_ranges(ranges, 0);
let max_gap = if strategy == MultiRangeReadStrategy::MergeRangesWithinOneMegabyte {
OBJECT_STORE_COALESCE_DEFAULT
} else {
0
};
let physical_ranges = merge_ranges(ranges, max_gap);
self.metrics
.record_parquet_data_file_exact_ranges_requested(
exact_ranges.len(),
range_bytes(&exact_ranges),
);
self.metrics.record_parquet_data_file_physical_range_plan(
physical_ranges.len(),
range_bytes(&physical_ranges),
);
let (results, _) = self
.read_physical_ranges(location, ranges, &physical_ranges)
.await?;
Ok(results)
}
MultiRangeReadStrategy::ChooseAutomatically => {
let estimate = self.current_transport_estimate();
let plan = choose_range_plan(ranges, estimate.estimate);
self.record_chosen_range_plan_metrics(&plan);
let (results, observed_plan_time) = self
.read_physical_ranges(location, ranges, &plan.physical_ranges)
.await?;
trace_completed_automatic_range_plan(&plan, estimate, observed_plan_time);
Ok(results)
}
}
}
fn delete_stream(
&self,
locations: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
self.inner.delete_stream(locations)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
self.inner.list(prefix)
}
fn list_with_offset(
&self,
prefix: Option<&Path>,
offset: &Path,
) -> BoxStream<'static, Result<ObjectMeta>> {
self.inner.list_with_offset(prefix, offset)
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
self.inner.list_with_delimiter(prefix).await
}
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
self.inner.copy_opts(from, to, options).await
}
async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
self.inner.rename_opts(from, to, options).await
}
}
fn meter_get_result(result: GetResult, metrics: DeltaScanMetrics) -> GetResult {
let GetResult {
payload,
meta,
range,
attributes,
} = result;
let payload = match payload {
GetResultPayload::Stream(payload) => {
let payload = payload
.map(move |result| {
if let Ok(bytes) = &result {
metrics.record_parquet_data_file_bytes_received(bytes.len());
}
result
})
.boxed();
GetResultPayload::Stream(payload)
}
#[cfg(not(target_arch = "wasm32"))]
GetResultPayload::File(file, path) => {
let local_result = GetResult {
payload: GetResultPayload::File(file, path),
meta: meta.clone(),
range: range.clone(),
attributes: attributes.clone(),
};
let payload = stream::once(async move {
let bytes = local_result.bytes().await?;
metrics.record_parquet_data_file_bytes_received(bytes.len());
Ok(bytes)
})
.boxed();
GetResultPayload::Stream(payload)
}
};
GetResult {
payload,
meta,
range,
attributes,
}
}
#[cfg(test)]
mod tests {
use std::{
collections::BTreeMap,
fmt,
fs::File,
io,
ops::Range,
path::PathBuf,
sync::{
Arc, Mutex, Once,
atomic::{AtomicU64, Ordering},
},
time::{Duration, Instant},
};
use async_trait::async_trait;
use bytes::Bytes;
use chrono::{DateTime, Utc};
use futures_util::{StreamExt, stream, stream::BoxStream};
use object_store::{
Attributes, CopyOptions, Error, GetOptions, GetResult, GetResultPayload, ListResult,
MultipartUpload, ObjectMeta, ObjectStore, ObjectStoreExt, PutMultipartOptions, PutOptions,
PutPayload, PutResult, RenameOptions, Result,
memory::InMemory,
path::Path,
throttle::{ThrottleConfig, ThrottledStore},
};
use tracing::{
Event, Level, Metadata, Subscriber,
field::{Field as TracingField, Visit},
span::{Attributes as TracingAttributes, Id, Record},
subscriber::{Interest, with_default},
};
use super::{
ChosenRangePlan, CompletedRangeRead, CurrentTransportEstimate, MIN_THROUGHPUT_SAMPLE_BYTES,
MIN_THROUGHPUT_SAMPLE_DELIVERY_TIME, MeteredParquetObjectStore, MultiRangeReadStrategy,
ParquetRangeReadEstimator, RANGE_PLANNING_DIAGNOSTIC_TARGET, RangePlanDecision,
TransportEstimate,
};
use crate::{
DeltaScanMetrics, ParquetReaderBackend,
reader::{ParquetRangeReadPolicy, metrics::DeltaScanMetricsConfig},
};
static TRACING_TEST_LOCK: Mutex<()> = Mutex::new(());
static TRACING_TEST_GLOBAL_SUBSCRIBER: Once = Once::new();
#[derive(Clone, Default)]
struct EventFields(Arc<Mutex<Vec<BTreeMap<String, String>>>>);
impl Subscriber for EventFields {
fn register_callsite(&self, metadata: &'static Metadata<'static>) -> Interest {
if metadata.target() == RANGE_PLANNING_DIAGNOSTIC_TARGET
&& *metadata.level() == Level::DEBUG
{
Interest::always()
} else {
Interest::sometimes()
}
}
fn enabled(&self, metadata: &Metadata<'_>) -> bool {
metadata.target() == RANGE_PLANNING_DIAGNOSTIC_TARGET
&& *metadata.level() == Level::DEBUG
}
fn new_span(&self, _attributes: &TracingAttributes<'_>) -> Id {
Id::from_u64(1)
}
fn record(&self, _span: &Id, _values: &Record<'_>) {}
fn record_follows_from(&self, _span: &Id, _follows: &Id) {}
fn event(&self, event: &Event<'_>) {
let mut fields = event
.metadata()
.fields()
.iter()
.map(|field| (field.name().to_owned(), "<empty>".to_owned()))
.collect();
event.record(&mut FieldVisitor(&mut fields));
self.0
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.push(fields);
}
fn enter(&self, _span: &Id) {}
fn exit(&self, _span: &Id) {}
}
struct FieldVisitor<'fields>(&'fields mut BTreeMap<String, String>);
impl Visit for FieldVisitor<'_> {
fn record_debug(&mut self, field: &TracingField, value: &dyn fmt::Debug) {
self.0.insert(field.name().to_owned(), format!("{value:?}"));
}
fn record_str(&mut self, field: &TracingField, value: &str) {
self.0.insert(field.name().to_owned(), value.to_owned());
}
fn record_u64(&mut self, field: &TracingField, value: u64) {
self.0.insert(field.name().to_owned(), value.to_string());
}
fn record_u128(&mut self, field: &TracingField, value: u128) {
self.0.insert(field.name().to_owned(), value.to_string());
}
}
fn capture_range_planning_events<T>(
run: impl FnOnce() -> T,
) -> (T, Vec<BTreeMap<String, String>>) {
let _lock = TRACING_TEST_LOCK
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let events = Arc::new(Mutex::new(Vec::new()));
let subscriber = EventFields(Arc::clone(&events));
TRACING_TEST_GLOBAL_SUBSCRIBER.call_once(|| {
let _ = tracing::subscriber::set_global_default(EventFields::default());
});
let result = with_default(subscriber, || {
tracing::callsite::rebuild_interest_cache();
run()
});
tracing::callsite::rebuild_interest_cache();
let captured = events
.lock()
.map(|events| events.clone())
.unwrap_or_default();
(result, captured)
}
fn event_field<'event>(
event: &'event BTreeMap<String, String>,
field: &str,
) -> Option<&'event str> {
event.get(field).map(String::as_str)
}
fn assert_event_fields(event: &BTreeMap<String, String>, expected_fields: &[(&str, &str)]) {
for (field, expected) in expected_fields {
assert_eq!(event_field(event, field), Some(*expected), "{field}");
}
}
fn direct_metrics() -> DeltaScanMetrics {
DeltaScanMetrics::new(DeltaScanMetricsConfig {
snapshot_version: 1,
parquet_backend: ParquetReaderBackend::Direct,
scan_partitions_planned: 1,
files_planned: 1,
add_actions_excluded_during_planning: Some(0),
estimated_input_rows: Some(1),
estimated_input_bytes: Some(1),
})
}
fn completed_read(
payload_started: Instant,
latency_millis: u64,
bytes_received: usize,
) -> CompletedRangeRead {
completed_read_with_delivery(
payload_started,
latency_millis,
Duration::from_secs(1),
bytes_received,
)
}
fn completed_read_with_delivery(
payload_started: Instant,
latency_millis: u64,
delivery_time: Duration,
bytes_received: usize,
) -> CompletedRangeRead {
CompletedRangeRead {
request_latency: Duration::from_millis(latency_millis),
payload_started,
payload_finished: payload_started + delivery_time,
bytes_received,
}
}
fn completed_read_at_throughput(
payload_started: Instant,
latency_millis: u64,
bytes_received: usize,
bytes_per_second: u64,
) -> CompletedRangeRead {
let delivery_nanos = (bytes_received as u128)
.saturating_mul(1_000_000_000)
.checked_div(u128::from(bytes_per_second))
.and_then(|nanos| u64::try_from(nanos).ok())
.expect("test throughput must produce a valid duration");
completed_read_with_delivery(
payload_started,
latency_millis,
Duration::from_nanos(delivery_nanos),
bytes_received,
)
}
async fn memory_store(metrics: DeltaScanMetrics) -> Result<MeteredParquetObjectStore> {
memory_store_with_strategy(metrics, MultiRangeReadStrategy::ChooseAutomatically).await
}
async fn memory_store_with_strategy(
metrics: DeltaScanMetrics,
strategy: MultiRangeReadStrategy,
) -> Result<MeteredParquetObjectStore> {
let inner = Arc::new(InMemory::new());
inner
.put(
&Path::from("data.parquet"),
PutPayload::from_static(b"0123456789abcdef"),
)
.await?;
Ok(MeteredParquetObjectStore::new(inner, metrics, strategy))
}
#[test]
fn cold_start_completion_reports_plan_shape_without_an_estimate()
-> std::result::Result<(), Box<dyn std::error::Error>> {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let (result, events) = capture_range_planning_events(|| {
runtime.block_on(async {
memory_store(direct_metrics())
.await?
.get_ranges(&Path::from("data.parquet"), &[0..4, 8..12])
.await
})
});
result?;
assert_eq!(events.len(), 1);
let event = &events[0];
assert_event_fields(
event,
&[
("event", "automatic_parquet_range_plan_completed"),
("latency_sample_count", "0"),
("throughput_sample_count", "0"),
("estimated_request_latency_micros", "<empty>"),
("estimated_shared_throughput_bytes_per_second", "<empty>"),
("estimated_bandwidth_delay_bytes", "<empty>"),
("exact_range_count", "2"),
("exact_bytes", "8"),
("exact_request_waves", "1"),
("baseline_range_count", "2"),
("baseline_bytes", "8"),
("baseline_request_waves", "1"),
("selected_range_count", "2"),
("selected_bytes", "8"),
("selected_request_waves", "1"),
("baseline_predicted_cost_bytes", "<empty>"),
("selected_predicted_cost_bytes", "<empty>"),
("decision", "cold_start"),
],
);
assert_ne!(
event_field(event, "observed_selected_plan_micros"),
Some("<empty>")
);
Ok(())
}
#[test]
fn warmed_completions_report_the_estimate_and_exact_candidate_scores()
-> std::result::Result<(), Box<dyn std::error::Error>> {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let location =
Path::from("private-bucket/credential-secret/schema-secret/predicate-secret.parquet");
let requested_ranges = (0..=10)
.map(|index| {
let start = index * 1_000;
start..start + 100
})
.collect::<Vec<_>>();
for (bytes_received, throughput, expected_fields) in [
(
1024 * 1024,
1_000,
[
("estimated_shared_throughput_bytes_per_second", "1000"),
("estimated_bandwidth_delay_bytes", "100"),
("selected_range_count", "11"),
("selected_bytes", "1100"),
("selected_request_waves", "2"),
("baseline_predicted_cost_bytes", "1300"),
("selected_predicted_cost_bytes", "1300"),
("decision", "cost_based_exact"),
],
),
(
10_000_000,
1_000_000_000,
[
("estimated_shared_throughput_bytes_per_second", "1000000000"),
("estimated_bandwidth_delay_bytes", "100000000"),
("selected_range_count", "10"),
("selected_bytes", "2000"),
("selected_request_waves", "1"),
("baseline_predicted_cost_bytes", "200001100"),
("selected_predicted_cost_bytes", "100002000"),
("decision", "cost_based_merged"),
],
),
] {
let (result, events) = capture_range_planning_events(|| {
runtime.block_on(async {
let inner = Arc::new(InMemory::new());
inner
.put(&location, PutPayload::from(vec![0_u8; 10_100]))
.await?;
let store = MeteredParquetObjectStore::new(
inner,
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
for _ in 0..3 {
store.record_completed_range_reads(&[completed_read_at_throughput(
started,
100,
bytes_received,
throughput,
)]);
}
store.get_ranges(&location, &requested_ranges).await
})
});
assert!(result?.iter().all(|bytes| bytes.len() == 100));
assert_eq!(events.len(), 1);
let event = &events[0];
assert_event_fields(
event,
&[
("latency_sample_count", "3"),
("throughput_sample_count", "3"),
("estimated_request_latency_micros", "100000"),
("exact_range_count", "11"),
("exact_bytes", "1100"),
("exact_request_waves", "2"),
("baseline_range_count", "11"),
("baseline_bytes", "1100"),
("baseline_request_waves", "2"),
],
);
assert_event_fields(event, &expected_fields);
assert_ne!(
event_field(event, "observed_selected_plan_micros"),
Some("<empty>")
);
let diagnostic = format!("{event:?}");
for private_value in [
"private-bucket",
"credential-secret",
"schema-secret",
"predicate-secret",
] {
assert!(!diagnostic.contains(private_value));
}
}
Ok(())
}
#[test]
fn range_strategy_follows_table_store_kind()
-> std::result::Result<(), Box<dyn std::error::Error>> {
for location in ["file:///tmp/table", "memory:///table", "hdfs://host/table"] {
assert_eq!(
MultiRangeReadStrategy::for_table_url(&url::Url::parse(location)?),
MultiRangeReadStrategy::UseStoreImplementation,
"{location}"
);
}
for location in [
"s3://bucket/table",
"gs://bucket/table",
"https://example.com/table",
] {
assert_eq!(
MultiRangeReadStrategy::for_table_url(&url::Url::parse(location)?),
MultiRangeReadStrategy::ChooseAutomatically,
"{location}"
);
}
Ok(())
}
#[test]
fn diagnostic_policy_names_map_directly_to_range_strategies()
-> std::result::Result<(), Box<dyn std::error::Error>> {
let remote_url = url::Url::parse("https://example.com/table")?;
for (policy, strategy) in [
(
ParquetRangeReadPolicy::Automatic,
MultiRangeReadStrategy::ChooseAutomatically,
),
(
ParquetRangeReadPolicy::ExactRanges,
MultiRangeReadStrategy::ReadExactRanges,
),
(
ParquetRangeReadPolicy::MergeRangesWithinOneMegabyte,
MultiRangeReadStrategy::MergeRangesWithinOneMegabyte,
),
(
ParquetRangeReadPolicy::StoreImplementation,
MultiRangeReadStrategy::UseStoreImplementation,
),
] {
assert_eq!(
MultiRangeReadStrategy::for_policy(policy, &remote_url),
strategy
);
}
Ok(())
}
#[tokio::test]
async fn bounded_and_unbounded_gets_record_exact_operations_and_bytes() -> Result<()> {
let range_metrics = direct_metrics();
let range_store = memory_store(range_metrics.clone()).await?;
let bytes = range_store
.get_range(&Path::from("data.parquet"), 2..7)
.await?;
assert_eq!(bytes.as_ref(), b"23456");
let snapshot = range_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(1));
assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(5));
let full_metrics = direct_metrics();
let full_store = memory_store(full_metrics.clone()).await?;
let bytes = full_store
.get(&Path::from("data.parquet"))
.await?
.bytes()
.await?;
assert_eq!(bytes.as_ref(), b"0123456789abcdef");
let snapshot = full_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(1));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(16));
Ok(())
}
#[tokio::test]
async fn head_and_failed_gets_preserve_attempt_metrics() -> Result<()> {
let head_metrics = direct_metrics();
let head_store = memory_store(head_metrics.clone()).await?;
let options = GetOptions::new().with_range(Some(1_u64..4)).with_head(true);
let _result = head_store
.get_opts(&Path::from("data.parquet"), options)
.await?;
let snapshot = head_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(0));
let failure_metrics = direct_metrics();
let failure_store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
failure_metrics.clone(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let result = failure_store
.get_opts(
&Path::from("missing.parquet"),
GetOptions::new().with_range(Some(0_u64..4)),
)
.await;
assert!(result.is_err());
let snapshot = failure_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(1));
assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(0));
Ok(())
}
#[tokio::test]
async fn multi_range_strategies_preserve_results_and_expose_metrics() -> Result<()> {
let automatic_metrics = direct_metrics();
let automatic_store = memory_store(automatic_metrics.clone()).await?;
let bytes = automatic_store
.get_ranges(&Path::from("data.parquet"), &[0..4, 8..12])
.await?;
assert_eq!(bytes[0].as_ref(), b"0123");
assert_eq!(bytes[1].as_ref(), b"89ab");
let snapshot = automatic_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_exact_ranges_requested, Some(2));
assert_eq!(
snapshot.parquet_data_file_exact_range_bytes_requested,
Some(8)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_requests_planned,
Some(2)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_bytes_planned,
Some(8)
);
assert_eq!(snapshot.parquet_data_file_cold_start_range_plans, Some(1));
assert_eq!(
snapshot.parquet_data_file_cost_based_exact_range_plans,
Some(0)
);
assert_eq!(
snapshot.parquet_data_file_cost_based_merged_range_plans,
Some(0)
);
assert_eq!(
snapshot.parquet_data_file_store_delegated_range_plans,
Some(0)
);
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(2));
assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(8));
let diagnostic = automatic_metrics.parquet_range_planning_diagnostic_snapshot();
assert_eq!(diagnostic.max_concurrent_physical_range_requests, 10);
assert_eq!(diagnostic.physical_range_request_waves_planned, 1);
let delegated_metrics = direct_metrics();
let delegated_store = memory_store_with_strategy(
delegated_metrics.clone(),
MultiRangeReadStrategy::UseStoreImplementation,
)
.await?;
let bytes = delegated_store
.get_ranges(&Path::from("data.parquet"), &[0..4, 8..12])
.await?;
assert_eq!(bytes[0].as_ref(), b"0123");
assert_eq!(bytes[1].as_ref(), b"89ab");
let snapshot = delegated_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_exact_ranges_requested, Some(2));
assert_eq!(
snapshot.parquet_data_file_exact_range_bytes_requested,
Some(8)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_requests_planned,
Some(0)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_bytes_planned,
Some(0)
);
assert_eq!(snapshot.parquet_data_file_cold_start_range_plans, Some(0));
assert_eq!(
snapshot.parquet_data_file_cost_based_exact_range_plans,
Some(0)
);
assert_eq!(
snapshot.parquet_data_file_cost_based_merged_range_plans,
Some(0)
);
assert_eq!(
snapshot.parquet_data_file_store_delegated_range_plans,
Some(1)
);
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(1));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(8));
assert_eq!(
delegated_metrics
.parquet_range_planning_diagnostic_snapshot()
.physical_range_request_waves_planned,
0
);
Ok(())
}
#[tokio::test]
async fn diagnostic_range_strategies_preserve_bytes_and_apply_their_plans() -> Result<()> {
let requested_ranges = [10..14, 0..4, 2..6, 10..14];
let expected = [b"abcd".as_slice(), b"0123", b"2345", b"abcd"];
for (strategy, expected_planned_requests, expected_planned_bytes) in [
(MultiRangeReadStrategy::ReadExactRanges, 2, 10),
(MultiRangeReadStrategy::MergeRangesWithinOneMegabyte, 1, 14),
(MultiRangeReadStrategy::ChooseAutomatically, 2, 10),
(MultiRangeReadStrategy::UseStoreImplementation, 0, 0),
] {
let metrics = direct_metrics();
let store = memory_store_with_strategy(metrics.clone(), strategy).await?;
let actual = store
.get_ranges(&Path::from("data.parquet"), &requested_ranges)
.await?;
assert_eq!(
actual.iter().map(Bytes::as_ref).collect::<Vec<_>>(),
expected,
"{strategy:?}"
);
let snapshot = metrics.snapshot();
assert_eq!(
snapshot.parquet_data_file_physical_range_requests_planned,
Some(expected_planned_requests),
"{strategy:?}"
);
assert_eq!(
snapshot.parquet_data_file_physical_range_bytes_planned,
Some(expected_planned_bytes),
"{strategy:?}"
);
}
Ok(())
}
#[test]
fn chosen_range_plan_metrics_distinguish_decisions() {
let metrics = direct_metrics();
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
metrics.clone(),
MultiRangeReadStrategy::ChooseAutomatically,
);
for (physical_ranges, planned_bytes, decision) in [
(vec![0..2, 4..6, 8..10], 6, RangePlanDecision::ColdStart),
(
vec![0..2, 4..6, 8..10],
6,
RangePlanDecision::CostBasedExact,
),
(vec![0..6, 8..10], 8, RangePlanDecision::CostBasedMerged),
] {
store.record_chosen_range_plan_metrics(&ChosenRangePlan {
exact_range_count: 3,
exact_bytes: 6,
baseline_range_count: 3,
baseline_bytes: 6,
physical_ranges,
planned_bytes,
transport_estimate: None,
baseline_predicted_cost_bytes: None,
selected_predicted_cost_bytes: None,
decision,
});
}
let snapshot = metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_exact_ranges_requested, Some(9));
assert_eq!(
snapshot.parquet_data_file_exact_range_bytes_requested,
Some(18)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_requests_planned,
Some(8)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_bytes_planned,
Some(20)
);
assert_eq!(snapshot.parquet_data_file_cold_start_range_plans, Some(1));
assert_eq!(
snapshot.parquet_data_file_cost_based_exact_range_plans,
Some(1)
);
assert_eq!(
snapshot.parquet_data_file_cost_based_merged_range_plans,
Some(1)
);
}
#[tokio::test]
async fn separate_stores_share_only_an_explicitly_reused_estimator() -> Result<()> {
let inner = Arc::new(InMemory::new());
inner
.put(
&Path::from("data.parquet"),
PutPayload::from_static(b"0123456789abcdef"),
)
.await?;
let estimator = Arc::new(ParquetRangeReadEstimator::default());
let first = MeteredParquetObjectStore::new(
inner.clone(),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
)
.with_range_read_estimator(Arc::clone(&estimator));
let second = MeteredParquetObjectStore::new(
inner.clone(),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
)
.with_range_read_estimator(Arc::clone(&estimator));
let started = Instant::now();
first.record_completed_range_reads(&[completed_read(started, 10, 1024 * 1024)]);
first.record_completed_range_reads(&[completed_read(started, 20, 2 * 1024 * 1024)]);
assert_eq!(
first.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 2,
throughput_sample_count: 2,
estimate: None,
}
);
second.record_completed_range_reads(&[completed_read(started, 30, 3 * 1024 * 1024)]);
assert!(second.current_transport_estimate().estimate.is_some());
let warm_metrics = direct_metrics();
let warm = MeteredParquetObjectStore::new(
inner.clone(),
warm_metrics.clone(),
MultiRangeReadStrategy::ChooseAutomatically,
)
.with_range_read_estimator(estimator);
let ranges = [0..4, 8..12];
let warm_bytes = warm
.get_ranges(&Path::from("data.parquet"), &ranges)
.await?;
let snapshot = warm_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_cold_start_range_plans, Some(0));
assert_eq!(
snapshot
.parquet_data_file_cost_based_exact_range_plans
.unwrap_or_default()
+ snapshot
.parquet_data_file_cost_based_merged_range_plans
.unwrap_or_default(),
1
);
let cold_metrics = direct_metrics();
let cold = MeteredParquetObjectStore::new(
inner,
cold_metrics.clone(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let cold_bytes = cold
.get_ranges(&Path::from("data.parquet"), &ranges)
.await?;
assert_eq!(cold_bytes, warm_bytes);
assert_eq!(
cold_metrics
.snapshot()
.parquet_data_file_cold_start_range_plans,
Some(1)
);
Ok(())
}
#[test]
fn tiny_reads_update_latency_without_displacing_throughput_samples() {
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
for latency_millis in [10, 20, 30] {
store.record_completed_range_reads(&[completed_read(
started,
latency_millis,
2 * 1024 * 1024,
)]);
}
assert_eq!(
store.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 3,
throughput_sample_count: 3,
estimate: Some(TransportEstimate {
request_latency: Duration::from_millis(20),
shared_throughput_bytes_per_second: 2 * 1024 * 1024,
}),
}
);
for latency_millis in 100..=108 {
store.record_completed_range_reads(&[completed_read(started, latency_millis, 15_000)]);
}
let estimate = store.current_transport_estimate();
assert_eq!(estimate.latency_sample_count, 9);
assert_eq!(estimate.throughput_sample_count, 3);
assert_eq!(
estimate.estimate,
Some(TransportEstimate {
request_latency: Duration::from_millis(104),
shared_throughput_bytes_per_second: 2 * 1024 * 1024,
})
);
}
#[test]
fn throughput_estimate_uses_a_byte_weighted_median() {
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
for (bytes_received, bytes_per_second) in [
(1024 * 1024, 1024 * 1024),
(1024 * 1024, 2 * 1024 * 1024),
(8 * 1024 * 1024, 8 * 1024 * 1024),
] {
store.record_completed_range_reads(&[completed_read_at_throughput(
started,
10,
bytes_received,
bytes_per_second,
)]);
}
assert_eq!(
store
.current_transport_estimate()
.estimate
.map(|estimate| estimate.shared_throughput_bytes_per_second),
Some(8 * 1024 * 1024)
);
}
#[test]
fn throughput_requires_enough_bytes_and_delivery_time() {
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
let minimum_bytes = usize::try_from(MIN_THROUGHPUT_SAMPLE_BYTES)
.expect("throughput threshold must fit usize");
store.record_completed_range_reads(&[completed_read(started, 10, minimum_bytes - 1)]);
store.record_completed_range_reads(&[completed_read_with_delivery(
started,
20,
MIN_THROUGHPUT_SAMPLE_DELIVERY_TIME - Duration::from_nanos(1),
minimum_bytes,
)]);
assert_eq!(
store.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 2,
throughput_sample_count: 0,
estimate: None,
}
);
store.record_completed_range_reads(&[completed_read_with_delivery(
started,
30,
MIN_THROUGHPUT_SAMPLE_DELIVERY_TIME,
minimum_bytes,
)]);
assert_eq!(
store.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 3,
throughput_sample_count: 1,
estimate: None,
}
);
}
#[test]
fn interleaved_tiny_reads_do_not_block_new_transport_evidence() {
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
let representative_bytes = 2 * 1024 * 1024;
for _ in 0..6 {
store.record_completed_range_reads(&[completed_read_at_throughput(
started,
20,
representative_bytes,
2 * 1024 * 1024,
)]);
store.record_completed_range_reads(&[completed_read(started, 5, 15_000)]);
}
let estimate = store.current_transport_estimate();
assert_eq!(estimate.latency_sample_count, 9);
assert_eq!(estimate.throughput_sample_count, 6);
assert_eq!(
estimate
.estimate
.map(|value| value.shared_throughput_bytes_per_second),
Some(2 * 1024 * 1024)
);
for _ in 0..9 {
store.record_completed_range_reads(&[completed_read_at_throughput(
started,
20,
representative_bytes,
64 * 1024 * 1024,
)]);
}
let estimate = store.current_transport_estimate();
assert_eq!(estimate.latency_sample_count, 9);
assert_eq!(estimate.throughput_sample_count, 9);
assert_eq!(
estimate
.estimate
.map(|value| value.shared_throughput_bytes_per_second),
Some(64 * 1024 * 1024)
);
}
#[test]
fn concurrent_payloads_produce_one_shared_throughput_sample() {
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
for _ in 0..3 {
store.record_completed_range_reads(&[
completed_read(started, 10, 512 * 1024),
completed_read(started, 10, 512 * 1024),
]);
}
assert_eq!(
store
.current_transport_estimate()
.estimate
.map(|estimate| estimate.shared_throughput_bytes_per_second),
Some(1024 * 1024)
);
}
#[tokio::test]
async fn successful_tiny_plans_update_only_latency_evidence() -> Result<()> {
let inner = InMemory::new();
inner
.put(
&Path::from("data.parquet"),
PutPayload::from_static(b"0123456789abcdef"),
)
.await?;
let throttled = ThrottledStore::new(
inner,
ThrottleConfig {
wait_get_per_call: Duration::from_millis(1),
wait_get_per_byte: Duration::from_micros(10),
..Default::default()
},
);
let store = MeteredParquetObjectStore::new(
Arc::new(throttled),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
for _ in 0..3 {
let results = store
.get_ranges(&Path::from("data.parquet"), &[0..4, 8..12])
.await?;
assert_eq!(results[0].as_ref(), b"0123");
assert_eq!(results[1].as_ref(), b"89ab");
}
assert_eq!(
store.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 3,
throughput_sample_count: 0,
estimate: None,
}
);
let failed_metrics = direct_metrics();
let failed_store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
failed_metrics.clone(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
for _ in 0..2 {
failed_store.record_completed_range_reads(&[completed_read(started, 10, 1_000)]);
}
assert!(
failed_store
.get_ranges(&Path::from("missing.parquet"), &[0..4, 8..12])
.await
.is_err()
);
assert_eq!(
failed_store.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 2,
throughput_sample_count: 0,
estimate: None,
}
);
let snapshot = failed_metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_exact_ranges_requested, Some(2));
assert_eq!(
snapshot.parquet_data_file_exact_range_bytes_requested,
Some(8)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_requests_planned,
Some(2)
);
assert_eq!(
snapshot.parquet_data_file_physical_range_bytes_planned,
Some(8)
);
assert_eq!(snapshot.parquet_data_file_cold_start_range_plans, Some(1));
Ok(())
}
#[test]
fn failed_automatic_plan_emits_no_completion_and_adds_no_sample() {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("test runtime");
let ((result, estimate), events) = capture_range_planning_events(|| {
runtime.block_on(async {
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
direct_metrics(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
for _ in 0..2 {
store.record_completed_range_reads(&[completed_read(started, 10, 1_000)]);
}
let result = store
.get_ranges(&Path::from("private-missing.parquet"), &[0..4, 8..12])
.await;
(result, store.current_transport_estimate())
})
});
assert!(result.is_err());
assert_eq!(
estimate,
CurrentTransportEstimate {
latency_sample_count: 2,
throughput_sample_count: 0,
estimate: None,
}
);
assert!(events.is_empty());
}
#[test]
fn cancelled_partial_automatic_plan_emits_no_completion_and_adds_no_sample()
-> std::result::Result<(), Box<dyn std::error::Error>> {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()?;
let (result, events) = capture_range_planning_events(|| {
runtime.block_on(async {
let first = PutPayload::from_static(b"abc")
.into_iter()
.next()
.ok_or_else(missing_chunk)?;
let payload = stream::iter([Ok(first)])
.chain(stream::pending::<Result<Bytes>>())
.boxed();
let metrics = direct_metrics();
let store = Arc::new(MeteredParquetObjectStore::new(
Arc::new(ScriptedGetStore::new(test_get_result(
GetResultPayload::Stream(payload),
0..8,
))),
metrics.clone(),
MultiRangeReadStrategy::ChooseAutomatically,
));
let started = Instant::now();
for _ in 0..2 {
store.record_completed_range_reads(&[completed_read(started, 10, 1_000)]);
}
let task_store = Arc::clone(&store);
let task = tokio::spawn(async move {
task_store
.get_ranges(&Path::from("data.parquet"), &[0..4, 4..8])
.await
});
tokio::time::timeout(Duration::from_secs(5), async {
while metrics.snapshot().parquet_data_file_bytes_received != Some(3) {
tokio::task::yield_now().await;
}
})
.await?;
task.abort();
assert!(task.await.is_err_and(|error| error.is_cancelled()));
assert_eq!(
store.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 2,
throughput_sample_count: 0,
estimate: None,
}
);
Ok::<_, Box<dyn std::error::Error>>(())
})
});
result?;
assert!(events.is_empty());
Ok(())
}
#[tokio::test]
async fn truncated_automatic_plan_does_not_update_transport_estimates()
-> std::result::Result<(), Box<dyn std::error::Error>> {
let truncated = PutPayload::from_static(b"abc")
.into_iter()
.next()
.ok_or_else(missing_chunk)?;
let metrics = direct_metrics();
let store = MeteredParquetObjectStore::new(
Arc::new(ScriptedGetStore::new(test_get_result(
GetResultPayload::Stream(stream::iter([Ok(truncated)]).boxed()),
0..8,
))),
metrics.clone(),
MultiRangeReadStrategy::ChooseAutomatically,
);
let started = Instant::now();
for _ in 0..2 {
store.record_completed_range_reads(&[completed_read(started, 10, 1_000)]);
}
let error = store
.get_ranges(&Path::from("data.parquet"), &[0..4, 4..8])
.await
.expect_err("truncated range unexpectedly succeeded");
assert!(error.to_string().contains("unexpected length"));
assert_eq!(metrics.snapshot().parquet_data_file_bytes_received, Some(3));
assert_eq!(
store.current_transport_estimate(),
CurrentTransportEstimate {
latency_sample_count: 2,
throughput_sample_count: 0,
estimate: None,
}
);
Ok(())
}
#[tokio::test]
async fn store_multi_range_handles_empty_and_failed_calls() -> Result<()> {
let metrics = direct_metrics();
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
metrics.clone(),
MultiRangeReadStrategy::UseStoreImplementation,
);
assert!(
store
.get_ranges(&Path::from("missing.parquet"), &[])
.await?
.is_empty()
);
assert_eq!(
metrics.snapshot().parquet_data_file_range_get_operations,
Some(0)
);
assert_eq!(
metrics
.snapshot()
.parquet_data_file_store_delegated_range_plans,
Some(0)
);
assert!(
store
.get_ranges(&Path::from("missing.parquet"), &[0..4, 8..12])
.await
.is_err()
);
let snapshot = metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_exact_ranges_requested, Some(2));
assert_eq!(
snapshot.parquet_data_file_exact_range_bytes_requested,
Some(8)
);
assert_eq!(
snapshot.parquet_data_file_store_delegated_range_plans,
Some(1)
);
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(1));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(0));
Ok(())
}
#[tokio::test]
async fn stream_delivery_records_only_successful_consumed_chunks() -> Result<()> {
let first = PutPayload::from_static(b"abc")
.into_iter()
.next()
.ok_or_else(missing_chunk)?;
let second = PutPayload::from_static(b"defgh")
.into_iter()
.next()
.ok_or_else(missing_chunk)?;
let success_metrics = direct_metrics();
let success_store = scripted_store(
test_get_result(
GetResultPayload::Stream(stream::iter(vec![Ok(first), Ok(second)]).boxed()),
0..8,
),
success_metrics.clone(),
);
let result = success_store.get(&Path::from("data.parquet")).await?;
assert_eq!(result.meta.location, Path::from("data.parquet"));
assert_eq!(result.meta.e_tag.as_deref(), Some("opaque-etag"));
assert_eq!(result.range, 0..8);
assert_eq!(result.bytes().await?.as_ref(), b"abcdefgh");
assert_eq!(
success_metrics.snapshot().parquet_data_file_bytes_received,
Some(8)
);
let dropped_metrics = direct_metrics();
let chunks = PutPayload::from_static(b"abc")
.into_iter()
.chain(PutPayload::from_static(b"defgh").into_iter())
.map(Ok)
.collect::<Vec<_>>();
let dropped_store = scripted_store(
test_get_result(GetResultPayload::Stream(stream::iter(chunks).boxed()), 0..8),
dropped_metrics.clone(),
);
let mut payload = dropped_store
.get(&Path::from("data.parquet"))
.await?
.into_stream();
assert_eq!(
payload
.next()
.await
.transpose()?
.ok_or_else(missing_chunk)?
.as_ref(),
b"abc"
);
drop(payload);
assert_eq!(
dropped_metrics.snapshot().parquet_data_file_bytes_received,
Some(3)
);
let error_metrics = direct_metrics();
let error = Error::Generic {
store: "test",
source: io::Error::other("payload failure").into(),
};
let successful = PutPayload::from_static(b"abc")
.into_iter()
.next()
.ok_or_else(missing_chunk)?;
let error_store = scripted_store(
test_get_result(
GetResultPayload::Stream(stream::iter(vec![Ok(successful), Err(error)]).boxed()),
0..8,
),
error_metrics.clone(),
);
let mut payload = error_store
.get(&Path::from("data.parquet"))
.await?
.into_stream();
assert!(payload.next().await.transpose()?.is_some());
assert!(payload.next().await.transpose().is_err());
assert_eq!(
error_metrics.snapshot().parquet_data_file_bytes_received,
Some(3)
);
Ok(())
}
#[cfg(not(target_arch = "wasm32"))]
#[tokio::test]
async fn local_file_payload_stays_lazy_and_yields_one_large_chunk()
-> std::result::Result<(), Box<dyn std::error::Error>> {
let file = TemporaryTestFile::new(&vec![7_u8; 20_000])?;
let range = 100_u64..16_500;
let dropped_metrics = direct_metrics();
let dropped_store =
scripted_store(file.get_result(range.clone())?, dropped_metrics.clone());
let result = dropped_store
.get_opts(
&Path::from("data.parquet"),
GetOptions::new().with_range(Some(range.clone())),
)
.await?;
drop(result);
assert_eq!(
dropped_metrics.snapshot().parquet_data_file_bytes_received,
Some(0)
);
let delivered_metrics = direct_metrics();
let delivered_store =
scripted_store(file.get_result(range.clone())?, delivered_metrics.clone());
let mut payload = delivered_store
.get_opts(
&Path::from("data.parquet"),
GetOptions::new().with_range(Some(range.clone())),
)
.await?
.into_stream();
let bytes = payload
.next()
.await
.transpose()?
.ok_or_else(missing_chunk)?;
assert_eq!(bytes.len(), usize::try_from(range.end - range.start)?);
assert!(payload.next().await.is_none());
assert_eq!(
delivered_metrics
.snapshot()
.parquet_data_file_bytes_received,
Some(range.end - range.start)
);
Ok(())
}
#[tokio::test]
async fn delegated_operations_and_diagnostics_do_not_leak_or_meter() -> Result<()> {
let metrics = direct_metrics();
let store = MeteredParquetObjectStore::new(
Arc::new(InMemory::new()),
metrics.clone(),
MultiRangeReadStrategy::UseStoreImplementation,
);
let first = Path::from("first.parquet");
let second = Path::from("second.parquet");
let third = Path::from("third.parquet");
store.put(&first, PutPayload::from_static(b"data")).await?;
assert!(store.list(None).next().await.transpose()?.is_some());
store.copy(&first, &second).await?;
store.rename(&second, &third).await?;
store.delete(&first).await?;
store.delete(&third).await?;
let snapshot = metrics.snapshot();
assert_eq!(snapshot.parquet_data_file_range_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_full_get_operations, Some(0));
assert_eq!(snapshot.parquet_data_file_bytes_received, Some(0));
assert_eq!(snapshot.estimated_parquet_task_bytes_admitted, Some(0));
assert_eq!(format!("{store:?}"), "MeteredParquetObjectStore");
let redacted_metrics = direct_metrics();
let redacted_store = scripted_store(
test_get_result(GetResultPayload::Stream(stream::empty().boxed()), 0..0),
redacted_metrics.clone(),
);
let location = Path::from("private/user-password-secret-token.parquet");
let mut options = GetOptions::new().with_range(Some(987_654_321_u64..987_654_999_u64));
options.if_match = Some("secret-conditional-header".to_owned());
options.version = Some("secret-object-version".to_owned());
let _result = redacted_store.get_opts(&location, options).await?;
let diagnostics = format!(
"{redacted_store:?} {redacted_store} {:?}",
redacted_metrics.snapshot()
);
for secret in [
"user-password-secret-token",
"secret-conditional-header",
"secret-object-version",
"987654321",
"987654999",
] {
assert!(!diagnostics.contains(secret));
}
Ok(())
}
fn missing_chunk() -> Error {
Error::Generic {
store: "test",
source: io::Error::other("missing test chunk").into(),
}
}
fn scripted_store(result: GetResult, metrics: DeltaScanMetrics) -> MeteredParquetObjectStore {
MeteredParquetObjectStore::new(
Arc::new(ScriptedGetStore::new(result)),
metrics,
MultiRangeReadStrategy::UseStoreImplementation,
)
}
fn test_get_result(payload: GetResultPayload, range: Range<u64>) -> GetResult {
GetResult {
payload,
meta: ObjectMeta {
location: Path::from("data.parquet"),
last_modified: DateTime::<Utc>::UNIX_EPOCH,
size: range.end,
e_tag: Some("opaque-etag".to_owned()),
version: Some("opaque-version".to_owned()),
},
range,
attributes: Attributes::new(),
}
}
struct ScriptedGetStore {
result: Mutex<Option<GetResult>>,
delegate: InMemory,
}
impl ScriptedGetStore {
fn new(result: GetResult) -> Self {
Self {
result: Mutex::new(Some(result)),
delegate: InMemory::new(),
}
}
}
impl fmt::Debug for ScriptedGetStore {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("ScriptedGetStore")
}
}
impl fmt::Display for ScriptedGetStore {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
formatter.write_str("ScriptedGetStore")
}
}
#[async_trait]
impl ObjectStore for ScriptedGetStore {
async fn put_opts(
&self,
location: &Path,
payload: PutPayload,
options: PutOptions,
) -> Result<PutResult> {
self.delegate.put_opts(location, payload, options).await
}
async fn put_multipart_opts(
&self,
location: &Path,
options: PutMultipartOptions,
) -> Result<Box<dyn MultipartUpload>> {
self.delegate.put_multipart_opts(location, options).await
}
async fn get_opts(&self, _location: &Path, _options: GetOptions) -> Result<GetResult> {
self.result
.lock()
.map_err(|_| missing_chunk())?
.take()
.ok_or_else(missing_chunk)
}
fn delete_stream(
&self,
locations: BoxStream<'static, Result<Path>>,
) -> BoxStream<'static, Result<Path>> {
self.delegate.delete_stream(locations)
}
fn list(&self, prefix: Option<&Path>) -> BoxStream<'static, Result<ObjectMeta>> {
self.delegate.list(prefix)
}
fn list_with_offset(
&self,
prefix: Option<&Path>,
offset: &Path,
) -> BoxStream<'static, Result<ObjectMeta>> {
self.delegate.list_with_offset(prefix, offset)
}
async fn list_with_delimiter(&self, prefix: Option<&Path>) -> Result<ListResult> {
self.delegate.list_with_delimiter(prefix).await
}
async fn copy_opts(&self, from: &Path, to: &Path, options: CopyOptions) -> Result<()> {
self.delegate.copy_opts(from, to, options).await
}
async fn rename_opts(&self, from: &Path, to: &Path, options: RenameOptions) -> Result<()> {
self.delegate.rename_opts(from, to, options).await
}
}
#[cfg(not(target_arch = "wasm32"))]
struct TemporaryTestFile {
path: PathBuf,
}
#[cfg(not(target_arch = "wasm32"))]
impl TemporaryTestFile {
fn new(contents: &[u8]) -> io::Result<Self> {
static NEXT_FILE_ID: AtomicU64 = AtomicU64::new(0);
let id = NEXT_FILE_ID.fetch_add(1, Ordering::Relaxed);
let path = std::env::temp_dir().join(format!(
"delta-arrow-reader-metered-object-store-{}-{id}",
std::process::id()
));
std::fs::write(&path, contents)?;
Ok(Self { path })
}
fn get_result(&self, range: Range<u64>) -> io::Result<GetResult> {
Ok(test_get_result(
GetResultPayload::File(File::open(&self.path)?, self.path.clone()),
range,
))
}
}
#[cfg(not(target_arch = "wasm32"))]
impl Drop for TemporaryTestFile {
fn drop(&mut self) {
let _ = std::fs::remove_file(&self.path);
}
}
}