diff --git a/docs/source/user-guide/latest/metrics.md b/docs/source/user-guide/latest/metrics.md index a8b2d6bdb8..c3fb35d8e6 100644 --- a/docs/source/user-guide/latest/metrics.md +++ b/docs/source/user-guide/latest/metrics.md @@ -79,11 +79,50 @@ Here is a guide to some of the native metrics. | `spilled_bytes` | Actual bytes written to native shuffle spill files on disk. | | `memory_spilled_bytes` | Uncompressed Arrow backing-buffer and partition-index memory spilled. | +### Native Parquet scans + +Native Parquet scans expose these counters in the Spark SQL metric map as well as native +execution metrics. Counters accumulate per scan operator; they do not instrument individual rows. + +| Metric | Description | +| ------------------------------------------ | ---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `scan_io_data_bytes` | Bytes returned to the Parquet reader for projected data-page ranges. | +| `scan_io_metadata_bytes` | Bytes returned for footer prefetches, page indexes, and Bloom filters. A footer prefetch can also contain unused data bytes. | +| `scan_io_footer_reads` | Storage reads of complete serialized footer payloads, counted once per metadata open. Plaintext payloads must decode successfully; encrypted payloads are counted before key retrieval/decryption, even if those later fail. | +| `scan_io_footer_bytes` | Serialized footer payload bytes, excluding the final eight-byte trailer. These bytes are already included in `scan_io_metadata_bytes`. | +| `scan_io_object_store_get_calls` | Nonempty GET operations at the native remote `ObjectStore` API, after range coalescing. Does not count transport retries or HEAD requests. | +| `scan_io_object_store_get_requested_bytes` | Requested coalesced range bytes at that API. A failed or partly consumed request can contribute requested bytes without the same number of response bytes. | +| `scan_io_object_store_response_bytes_read` | Response bytes actually consumed at that API, including bytes fetched between coalesced ranges. Not HTTP wire bytes. | +| `scan_io_metadata_cache_hits` | Successful, cache-eligible metadata opens requiring no storage reads. | +| `scan_io_metadata_cache_misses` | Successful, cache-eligible metadata opens requiring storage reads. Failed opens and encrypted opens, which bypass this shared cache, increment neither cache counter. | + +Reader-level and object-store bytes are two views of the same reads; do not add them together. +Likewise, footer bytes are a subset of metadata bytes, not a third reader-level category. A warm +metadata-cache hit contributes no new metadata or footer I/O. A valid plaintext footer followed by +a page-index failure still contributes footer bytes; an invalid plaintext footer does not. + +Remote counters follow the backend selected during object-store construction, including native +S3 (`s3`/`s3a`), GCS (`gs`), Azure (`az`, `adl`, `azure`, `abfs`, `abfss`), and HTTP(S) stores. +Local files, in-memory stores, and HDFS/custom backends (including cloud-looking schemes selected +through `fs.comet.libhdfs.schemes`) retain reader-level counters but have zero remote counters. +The native cloud wrapper observes default range coalescing; custom `get_ranges` implementations +require an explicit accounting contract before being composed with that wrapper. + +For data-bearing scans, comparing remote response bytes with reader data bytes can reveal +coalescing and metadata overhead. For a metadata-only scan, data bytes are zero: report the +metadata and remote totals instead of dividing by zero. Cancellation can leave late asynchronous +work outside the final metric snapshot; these counters are not a guarantee of complete network +traffic accounting after cancellation. + ## Task-Level Input Metrics on Spark 4.1+ -Comet's native scans set `inputMetrics.bytesRead` to the actual file IO performed by the -DataFusion parquet reader (`bytes_scanned`). This is the truthful number you would see at the -filesystem layer. +Comet's native scans populate `inputMetrics.bytesRead` from the existing `bytes_scanned` +counter. It counts requested data/Bloom-filter ranges through the Parquet reader's byte-read +methods, not all filesystem I/O. Footer and page-index reads through metadata loading bypass +this counter, and range coalescing can fetch more bytes than the logical ranges request. The +additional scan I/O metrics above expose those differences without changing `bytes_scanned`. +The native `scan_efficiency_ratio` still uses `bytes_scanned` as its numerator and has the same +blind spots; it is not the remote read-amplification ratio described above. Spark 4.1 changed its own parquet reader to pre-open the `SeekableInputStream` and read the file footer outside the `FileScanRDD.compute()` thread. Spark's `inputMetrics.bytesRead` is updated @@ -100,4 +139,5 @@ unaffected and remains exactly equal between Comet and Spark. If you compare Comet's `bytesRead` against vanilla Spark's on Spark 4.1+ (via the Spark UI or the REST API), expect Comet's number to be substantially larger for small files, and closer to -Spark's for large files. Comet's value reflects what the storage layer actually delivered. +Spark's for large files in that workload. Neither metric should be interpreted as complete +filesystem or network traffic accounting. diff --git a/native/core/src/execution/operators/parquet_writer.rs b/native/core/src/execution/operators/parquet_writer.rs index d6e81b85e1..94e37d4b85 100644 --- a/native/core/src/execution/operators/parquet_writer.rs +++ b/native/core/src/execution/operators/parquet_writer.rs @@ -311,7 +311,7 @@ impl ParquetWriterExec { #[cfg(feature = "hdfs-opendal")] { // Use prepare_object_store_with_configs to create and register the object store - let (_object_store_url, object_store_path) = prepare_object_store_with_configs( + let (_object_store_url, object_store_path, _) = prepare_object_store_with_configs( _runtime_env, output_file_path.to_string(), object_store_options, diff --git a/native/core/src/execution/planner.rs b/native/core/src/execution/planner.rs index d109627e82..ab9421f4e7 100644 --- a/native/core/src/execution/planner.rs +++ b/native/core/src/execution/planner.rs @@ -1631,11 +1631,12 @@ impl PhysicalPlanner { .iter() .map(|(k, v)| (k.clone(), v.clone())) .collect(); - let (object_store_url, _) = prepare_object_store_with_configs( - self.session_ctx.runtime_env(), - one_file, - &object_store_options, - )?; + let (object_store_url, _, object_store_backend) = + prepare_object_store_with_configs( + self.session_ctx.runtime_env(), + one_file, + &object_store_options, + )?; // Get files for this partition let files = self.get_partitioned_files(partition_files)?; @@ -1646,6 +1647,7 @@ impl PhysicalPlanner { Some(data_schema), Some(partition_schema), object_store_url, + object_store_backend, file_groups, Some(projection_vector), Some(data_filters?), @@ -1683,7 +1685,7 @@ impl PhysicalPlanner { .and_then(|f| f.partitioned_file.first()) .map(|f| f.file_path.clone()) .ok_or(GeneralError("Failed to locate file".to_string()))?; - let (object_store_url, _) = prepare_object_store_with_configs( + let (object_store_url, _, _) = prepare_object_store_with_configs( self.session_ctx.runtime_env(), one_file, &object_store_options, diff --git a/native/core/src/parquet/eager_page_index_reader_factory.rs b/native/core/src/parquet/eager_page_index_reader_factory.rs index 278814c4bf..bdadc72f98 100644 --- a/native/core/src/parquet/eager_page_index_reader_factory.rs +++ b/native/core/src/parquet/eager_page_index_reader_factory.rs @@ -46,6 +46,7 @@ //! Filed upstream as apache/datafusion#23978. Revert this once the opener merges its deferred //! page-index load back into `FileMetadataCache` instead of bypassing it. +use async_trait::async_trait; use bytes::Bytes; use datafusion::common::Result as DFResult; use datafusion::datasource::physical_plan::parquet::metadata::DFParquetMetadata; @@ -53,29 +54,128 @@ use datafusion::datasource::physical_plan::parquet::{ ParquetFileMetrics, ParquetFileReaderFactory, }; use datafusion::execution::cache::cache_manager::FileMetadataCache; -use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet; +use datafusion::physical_plan::metrics::{ + Count, ExecutionPlanMetricsSet, MetricBuilder, MetricCategory, MetricType, +}; use datafusion_datasource::PartitionedFile; use futures::future::BoxFuture; -use futures::FutureExt; -use object_store::ObjectStore; +use futures::{FutureExt, StreamExt, TryStreamExt}; +use object_store::path::Path; +use object_store::{ + coalesce_ranges, CopyOptions, GetOptions, GetRange, GetResult, GetResultPayload, ListResult, + MultipartUpload, ObjectMeta, ObjectStore, ObjectStoreExt, PutMultipartOptions, PutOptions, + PutPayload, PutResult, RenameOptions, Result as ObjectStoreResult, + OBJECT_STORE_COALESCE_DEFAULT, +}; use parquet::arrow::arrow_reader::ArrowReaderOptions; use parquet::arrow::async_reader::{AsyncFileReader, ParquetObjectReader}; -use parquet::file::metadata::{PageIndexPolicy, ParquetMetaData}; -use std::fmt::Debug; +use parquet::file::metadata::{ + FooterTail, PageIndexPolicy, ParquetMetaData, ParquetMetaDataReader, +}; +use std::fmt::{Debug, Display, Formatter}; use std::ops::Range; -use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; +use std::sync::{Arc, OnceLock}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum ScanIoSource { + ObjectStore, + Local, + OtherObjectStore, +} + +#[derive(Debug)] +struct ScanIoMetrics { + data_bytes: Count, + metadata_bytes: Count, + footer_reads: Count, + footer_bytes: Count, + object_store_get_calls: Count, + object_store_get_requested_bytes: Count, + object_store_response_bytes_read: Count, + metadata_cache_hits: Count, + metadata_cache_misses: Count, +} + +impl ScanIoMetrics { + fn new(metrics: &ExecutionPlanMetricsSet) -> Self { + Self { + data_bytes: byte_counter(metrics, "scan_io_data_bytes"), + metadata_bytes: byte_counter(metrics, "scan_io_metadata_bytes"), + footer_reads: count_counter(metrics, "scan_io_footer_reads"), + footer_bytes: byte_counter(metrics, "scan_io_footer_bytes"), + object_store_get_calls: count_counter(metrics, "scan_io_object_store_get_calls"), + object_store_get_requested_bytes: byte_counter( + metrics, + "scan_io_object_store_get_requested_bytes", + ), + object_store_response_bytes_read: byte_counter( + metrics, + "scan_io_object_store_response_bytes_read", + ), + metadata_cache_hits: count_counter(metrics, "scan_io_metadata_cache_hits"), + metadata_cache_misses: count_counter(metrics, "scan_io_metadata_cache_misses"), + } + } + + fn record_metadata_cache_result(&self, storage_reads: usize) { + if storage_reads == 0 { + self.metadata_cache_hits.add(1); + } else { + self.metadata_cache_misses.add(1); + } + } +} + +fn byte_counter(metrics: &ExecutionPlanMetricsSet, name: &'static str) -> Count { + MetricBuilder::new(metrics) + .with_type(MetricType::Summary) + .with_category(MetricCategory::Bytes) + .global_counter(name) +} + +fn count_counter(metrics: &ExecutionPlanMetricsSet, name: &'static str) -> Count { + MetricBuilder::new(metrics) + .with_type(MetricType::Summary) + .global_counter(name) +} + +fn range_bytes(range: &Range) -> usize { + (range.end - range.start) as usize +} + +fn ranges_bytes(ranges: &[Range]) -> usize { + ranges.iter().map(range_bytes).sum() +} #[derive(Debug)] pub struct EagerPageIndexReaderFactory { store: Arc, metadata_cache: Arc, + scan_io_metrics: Arc, } impl EagerPageIndexReaderFactory { - pub fn new(store: Arc, metadata_cache: Arc) -> Self { + pub(crate) fn new( + store: Arc, + metadata_cache: Arc, + source: ScanIoSource, + metrics: &ExecutionPlanMetricsSet, + ) -> Self { + let scan_io_metrics = Arc::new(ScanIoMetrics::new(metrics)); + let store: Arc = if source == ScanIoSource::ObjectStore { + Arc::new(ScanIoObjectStore { + inner: store, + scan_io_metrics: Arc::clone(&scan_io_metrics), + role: ScanIoStoreRole::ObjectStore, + }) + } else { + store + }; Self { store, metadata_cache, + scan_io_metrics, } } } @@ -104,6 +204,7 @@ impl ParquetFileReaderFactory for EagerPageIndexReaderFactory { Ok(Box::new(EagerPageIndexReader { file_metrics, + scan_io_metrics: Arc::clone(&self.scan_io_metrics), store: Arc::clone(&self.store), inner, partitioned_file, @@ -115,6 +216,7 @@ impl ParquetFileReaderFactory for EagerPageIndexReaderFactory { struct EagerPageIndexReader { file_metrics: ParquetFileMetrics, + scan_io_metrics: Arc, store: Arc, inner: ParquetObjectReader, partitioned_file: PartitionedFile, @@ -123,10 +225,23 @@ struct EagerPageIndexReader { } impl AsyncFileReader for EagerPageIndexReader { + // Pinned parquet 58.4 uses this single-range entry point for Bloom filters, and the + // multi-range entry point below for data pages. Footer/page-index loading is instrumented + // separately in get_metadata. The scan tests pin this method contract; revisit it when + // upgrading parquet rather than assuming an offset alone identifies metadata. fn get_bytes(&mut self, range: Range) -> BoxFuture<'_, parquet::errors::Result> { - let bytes_scanned = range.end - range.start; - self.file_metrics.bytes_scanned.add(bytes_scanned as usize); - self.inner.get_bytes(range) + let requested = range_bytes(&range); + // Preserve the existing requested-range metric. Metadata fetched through get_metadata + // bypasses it, and object-store coalescing can fetch more bytes than these ranges. + self.file_metrics.bytes_scanned.add(requested); + let scan_io_metrics = Arc::clone(&self.scan_io_metrics); + let future = self.inner.get_bytes(range); + async move { + let bytes = future.await?; + scan_io_metrics.metadata_bytes.add(bytes.len()); + Ok(bytes) + } + .boxed() } fn get_byte_ranges( @@ -136,9 +251,18 @@ impl AsyncFileReader for EagerPageIndexReader { where Self: Send, { - let total: u64 = ranges.iter().map(|r| r.end - r.start).sum(); - self.file_metrics.bytes_scanned.add(total as usize); - self.inner.get_byte_ranges(ranges) + let requested = ranges_bytes(&ranges); + self.file_metrics.bytes_scanned.add(requested); + let scan_io_metrics = Arc::clone(&self.scan_io_metrics); + let future = self.inner.get_byte_ranges(ranges); + async move { + let bytes = future.await?; + scan_io_metrics + .data_bytes + .add(bytes.iter().map(Bytes::len).sum()); + Ok(bytes) + } + .boxed() } fn get_metadata<'a>( @@ -151,17 +275,33 @@ impl AsyncFileReader for EagerPageIndexReader { let metadata_cache = Arc::clone(&self.metadata_cache); let store = Arc::clone(&self.store); let metadata_size_hint = self.metadata_size_hint; + let scan_io_metrics = Arc::clone(&self.scan_io_metrics); async move { let file_decryption_properties = options .and_then(|o| o.file_decryption_properties()) .map(Arc::clone); + let cache_enabled = file_decryption_properties.is_none(); let page_index_policy = if file_decryption_properties.is_none() { Some(PageIndexPolicy::Optional) } else { options.map(|o| o.column_index_policy()) }; + let metadata_storage_reads = Arc::new(AtomicUsize::new(0)); + let footer_payload_bytes = Arc::new(AtomicUsize::new(0)); + let metadata_store = ScanIoObjectStore { + inner: store, + scan_io_metrics: Arc::clone(&scan_io_metrics), + role: ScanIoStoreRole::Metadata { + storage_reads: Arc::clone(&metadata_storage_reads), + footer_payload_bytes: Arc::clone(&footer_payload_bytes), + footer_payload: OnceLock::new(), + footer_recorded: AtomicBool::new(false), + file_size: object_meta.size, + record_footer_immediately: !cache_enabled, + }, + }; - DFParquetMetadata::new(store.as_ref(), &object_meta) + let metadata = DFParquetMetadata::new(&metadata_store, &object_meta) .with_decryption_properties(file_decryption_properties) .with_file_metadata_cache(Some(metadata_cache)) .with_metadata_size_hint(metadata_size_hint) @@ -173,12 +313,334 @@ impl AsyncFileReader for EagerPageIndexReader { "Failed to fetch metadata for file {}: {e}", object_meta.location, )) - }) + }); + + if metadata.is_ok() { + let footer_bytes = footer_payload_bytes.load(Ordering::Relaxed); + if footer_bytes > 0 && cache_enabled { + metadata_store.record_footer(footer_bytes); + } + if cache_enabled { + scan_io_metrics.record_metadata_cache_result( + metadata_storage_reads.load(Ordering::Relaxed), + ); + } + } else if cache_enabled { + metadata_store.record_valid_footer(); + } + + metadata } .boxed() } } +#[derive(Debug)] +enum ScanIoStoreRole { + ObjectStore, + Metadata { + storage_reads: Arc, + footer_payload_bytes: Arc, + footer_payload: OnceLock, + footer_recorded: AtomicBool, + file_size: u64, + record_footer_immediately: bool, + }, +} + +#[derive(Debug)] +struct ScanIoObjectStore { + inner: Arc, + scan_io_metrics: Arc, + role: ScanIoStoreRole, +} + +impl ScanIoObjectStore { + fn record_valid_footer(&self) { + if let ScanIoStoreRole::Metadata { + footer_payload_bytes, + footer_payload, + footer_recorded, + .. + } = &self.role + { + if !footer_recorded.load(Ordering::Relaxed) + && footer_payload + .get() + .is_some_and(|payload| ParquetMetaDataReader::decode_metadata(payload).is_ok()) + { + self.record_footer(footer_payload_bytes.load(Ordering::Relaxed)); + } + } + } + + fn record_footer(&self, bytes: usize) { + if let ScanIoStoreRole::Metadata { + footer_recorded, .. + } = &self.role + { + if bytes > 0 && !footer_recorded.swap(true, Ordering::Relaxed) { + self.scan_io_metrics.footer_reads.add(1); + self.scan_io_metrics.footer_bytes.add(bytes); + } + } + } + + fn record_request(&self, bytes: usize) { + if bytes == 0 { + return; + } + match &self.role { + ScanIoStoreRole::ObjectStore => { + self.scan_io_metrics.object_store_get_calls.add(1); + self.scan_io_metrics + .object_store_get_requested_bytes + .add(bytes); + } + ScanIoStoreRole::Metadata { storage_reads, .. } => { + storage_reads.fetch_add(1, Ordering::Relaxed); + } + } + } + + // Footer accounting follows the DataFusion 54.1/parquet 58.4 metadata push decoder: a tail + // read supplies the payload length, then one or more reads supply the complete payload. + // Plaintext payloads are retained here and counted only after metadata succeeds, or when + // get_ranges observes the subsequent page-index request wholly below the footer. Thus a + // later index failure does not erase a decoded footer. Encrypted payloads are counted as + // soon as complete, before key retrieval/decryption; their counter measures payload I/O, + // not successful authentication. Recheck this protocol when upgrading the decoder. + fn record_returned(&self, range: Option<&Range>, bytes: &Bytes) { + match &self.role { + ScanIoStoreRole::ObjectStore => self + .scan_io_metrics + .object_store_response_bytes_read + .add(bytes.len()), + ScanIoStoreRole::Metadata { + footer_payload_bytes, + footer_payload, + file_size, + record_footer_immediately, + .. + } => { + self.scan_io_metrics.metadata_bytes.add(bytes.len()); + if range.is_some_and(|range| range.end == *file_size) && bytes.len() >= 8 { + if let Ok(footer) = FooterTail::try_from(&bytes[bytes.len() - 8..]) { + let _ = footer_payload_bytes.compare_exchange( + 0, + footer.metadata_length(), + Ordering::Relaxed, + Ordering::Relaxed, + ); + } + } + + let footer_bytes = footer_payload_bytes.load(Ordering::Relaxed); + let footer_end = file_size.saturating_sub(8); + if footer_bytes > 0 + && footer_end + .checked_sub(footer_bytes as u64) + .is_some_and(|footer_start| { + range.is_some_and(|range| { + range.start <= footer_start + && range.start.saturating_add(bytes.len() as u64) >= footer_end + }) + }) + { + if *record_footer_immediately { + // Encrypted opens bypass the shared cache and retrieve the key only after + // the complete encrypted payload is read. Count completed payload I/O, + // even if key retrieval, authentication, or metadata decoding later fails. + self.record_footer(footer_bytes); + } else if let Some(range) = range { + let payload_start = + (footer_end - footer_bytes as u64 - range.start) as usize; + let _ = footer_payload + .set(bytes.slice(payload_start..payload_start + footer_bytes)); + } + } + } + } + } +} + +impl Display for ScanIoObjectStore { + fn fmt(&self, formatter: &mut Formatter<'_>) -> std::fmt::Result { + write!(formatter, "scan-io({})", self.inner) + } +} + +#[async_trait] +#[deny(clippy::missing_trait_methods)] +impl ObjectStore for ScanIoObjectStore { + async fn put_opts( + &self, + location: &Path, + payload: PutPayload, + options: PutOptions, + ) -> ObjectStoreResult { + self.inner.put_opts(location, payload, options).await + } + + async fn put_multipart_opts( + &self, + location: &Path, + options: PutMultipartOptions, + ) -> ObjectStoreResult> { + self.inner.put_multipart_opts(location, options).await + } + + async fn get_opts(&self, location: &Path, options: GetOptions) -> ObjectStoreResult { + if options.head { + return self.inner.get_opts(location, options).await; + } + + let requested = match options.range.as_ref() { + Some(GetRange::Bounded(range)) => Some(range_bytes(range)), + Some(GetRange::Suffix(bytes)) => Some(*bytes as usize), + Some(GetRange::Offset(_)) | None => None, + }; + if let Some(requested) = requested { + self.record_request(requested); + } + + let result = self.inner.get_opts(location, options).await?; + if requested.is_none() { + self.record_request(range_bytes(&result.range)); + } + + let meta = result.meta.clone(); + let range = result.range.clone(); + let attributes = result.attributes.clone(); + let payload = if matches!(&result.payload, GetResultPayload::File(..)) { + let bytes = result.bytes().await?; + self.record_returned(Some(&range), &bytes); + GetResultPayload::Stream(futures::stream::once(async move { Ok(bytes) }).boxed()) + } else { + let scan_io_metrics = Arc::clone(&self.scan_io_metrics); + let metadata_read = matches!(self.role, ScanIoStoreRole::Metadata { .. }); + GetResultPayload::Stream( + result + .into_stream() + .inspect_ok(move |bytes| { + if metadata_read { + scan_io_metrics.metadata_bytes.add(bytes.len()); + } else { + scan_io_metrics + .object_store_response_bytes_read + .add(bytes.len()); + } + }) + .boxed(), + ) + }; + + Ok(GetResult { + payload, + meta, + range, + attributes, + }) + } + + async fn get_ranges( + &self, + location: &Path, + ranges: &[Range], + ) -> ObjectStoreResult> { + match &self.role { + ScanIoStoreRole::ObjectStore => { + // Supported native cloud stores use object_store's default coalescing. Apply it + // here so get_opts observes the coalesced requests, not just logical ranges. + // This deliberately bypasses an inner get_ranges override: a cache or custom + // backend must define that observation boundary before being composed here. + // Local/custom HDFS backends are not wrapped in this role and keep delegation. + coalesce_ranges( + ranges, + |range| self.get_range(location, range), + OBJECT_STORE_COALESCE_DEFAULT, + ) + .await + } + ScanIoStoreRole::Metadata { + footer_payload_bytes, + file_size, + record_footer_immediately, + .. + } => { + let footer_bytes = footer_payload_bytes.load(Ordering::Relaxed); + if !record_footer_immediately + && footer_bytes > 0 + && !ranges.is_empty() + && file_size + .saturating_sub(8) + .checked_sub(footer_bytes as u64) + .is_some_and(|footer_start| { + ranges.iter().all(|range| range.end <= footer_start) + }) + { + // In parquet 58.4, plaintext page-index requests below the footer begin only + // after its metadata payload has decoded successfully. Record that footer + // before awaiting the indexes, so a later index-read failure does not erase + // a successful footer read. Encrypted opens use the complete-payload path + // above; cached plaintext opens are recorded only after metadata succeeds. + self.record_footer(footer_bytes); + } + self.record_request(ranges_bytes(ranges)); + let bytes = self.inner.get_ranges(location, ranges).await?; + for (range, bytes) in ranges.iter().zip(bytes.iter()) { + self.record_returned(Some(range), bytes); + } + Ok(bytes) + } + } + } + + fn delete_stream( + &self, + locations: futures::stream::BoxStream<'static, ObjectStoreResult>, + ) -> futures::stream::BoxStream<'static, ObjectStoreResult> { + self.inner.delete_stream(locations) + } + + fn list( + &self, + prefix: Option<&Path>, + ) -> futures::stream::BoxStream<'static, ObjectStoreResult> { + self.inner.list(prefix) + } + + fn list_with_offset( + &self, + prefix: Option<&Path>, + offset: &Path, + ) -> futures::stream::BoxStream<'static, ObjectStoreResult> { + self.inner.list_with_offset(prefix, offset) + } + + async fn list_with_delimiter(&self, prefix: Option<&Path>) -> ObjectStoreResult { + self.inner.list_with_delimiter(prefix).await + } + + async fn copy_opts( + &self, + from: &Path, + to: &Path, + options: CopyOptions, + ) -> ObjectStoreResult<()> { + self.inner.copy_opts(from, to, options).await + } + + async fn rename_opts( + &self, + from: &Path, + to: &Path, + options: RenameOptions, + ) -> ObjectStoreResult<()> { + self.inner.rename_opts(from, to, options).await + } +} + impl Drop for EagerPageIndexReader { fn drop(&mut self) { self.file_metrics @@ -191,3 +653,160 @@ impl Drop for EagerPageIndexReader { .set_total(self.partitioned_file.object_meta.size as usize); } } + +#[cfg(test)] +mod tests { + use super::*; + use object_store::memory::InMemory; + + #[derive(Debug)] + struct RecordingRangeStore { + inner: InMemory, + range_calls: AtomicUsize, + get_calls: AtomicUsize, + } + + impl Display for RecordingRangeStore { + fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { + write!(f, "recording-range-store") + } + } + + #[async_trait] + impl ObjectStore for RecordingRangeStore { + async fn put_opts( + &self, + p: &Path, + v: PutPayload, + o: PutOptions, + ) -> ObjectStoreResult { + self.inner.put_opts(p, v, o).await + } + + async fn put_multipart_opts( + &self, + p: &Path, + o: PutMultipartOptions, + ) -> ObjectStoreResult> { + self.inner.put_multipart_opts(p, o).await + } + + async fn get_opts(&self, p: &Path, o: GetOptions) -> ObjectStoreResult { + self.get_calls.fetch_add(1, Ordering::Relaxed); + self.inner.get_opts(p, o).await + } + + async fn get_ranges( + &self, + p: &Path, + ranges: &[Range], + ) -> ObjectStoreResult> { + self.range_calls.fetch_add(1, Ordering::Relaxed); + self.inner.get_ranges(p, ranges).await + } + + fn delete_stream( + &self, + paths: futures::stream::BoxStream<'static, ObjectStoreResult>, + ) -> futures::stream::BoxStream<'static, ObjectStoreResult> { + self.inner.delete_stream(paths) + } + + fn list( + &self, + p: Option<&Path>, + ) -> futures::stream::BoxStream<'static, ObjectStoreResult> { + self.inner.list(p) + } + + async fn list_with_delimiter(&self, p: Option<&Path>) -> ObjectStoreResult { + self.inner.list_with_delimiter(p).await + } + + async fn copy_opts( + &self, + from: &Path, + to: &Path, + options: CopyOptions, + ) -> ObjectStoreResult<()> { + self.inner.copy_opts(from, to, options).await + } + + async fn rename_opts( + &self, + from: &Path, + to: &Path, + options: RenameOptions, + ) -> ObjectStoreResult<()> { + self.inner.rename_opts(from, to, options).await + } + } + + #[tokio::test] + async fn preserves_custom_range_delegation_for_local_and_other_backends() { + assert_range_read_contract(ScanIoSource::Local).await; + // This is also the classification returned for custom libhdfs schemes, including s3. + assert_range_read_contract(ScanIoSource::OtherObjectStore).await; + } + + #[tokio::test] + async fn remote_metrics_observe_default_coalescing_instead_of_inner_override() { + assert_range_read_contract(ScanIoSource::ObjectStore).await; + } + + async fn assert_range_read_contract(source: ScanIoSource) { + let store = Arc::new(RecordingRangeStore { + inner: InMemory::new(), + range_calls: AtomicUsize::new(0), + get_calls: AtomicUsize::new(0), + }); + let location = Path::from("ranges.parquet"); + store + .put(&location, Bytes::from_static(b"0123456789").into()) + .await + .unwrap(); + let runtime = datafusion::execution::runtime_env::RuntimeEnv::default(); + let metrics = ExecutionPlanMetricsSet::new(); + let factory = EagerPageIndexReaderFactory::new( + Arc::clone(&store) as Arc, + runtime.cache_manager.get_file_metadata_cache(), + source, + &metrics, + ); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::new(location.to_string(), 10), + None, + &metrics, + ) + .unwrap(); + let result = reader.get_byte_ranges(vec![0..2, 4..6]).await.unwrap(); + assert_eq!( + result, + vec![Bytes::from_static(b"01"), Bytes::from_static(b"45")] + ); + let remote = source == ScanIoSource::ObjectStore; + assert_eq!( + store.range_calls.load(Ordering::Relaxed), + usize::from(!remote) + ); + assert_eq!(store.get_calls.load(Ordering::Relaxed), usize::from(remote)); + assert_eq!( + metrics + .clone_inner() + .sum_by_name("scan_io_data_bytes") + .unwrap() + .as_usize(), + 4 + ); + assert_eq!( + metrics + .clone_inner() + .sum_by_name("scan_io_object_store_response_bytes_read") + .unwrap() + .as_usize(), + if remote { 6 } else { 0 } + ); + } +} diff --git a/native/core/src/parquet/mod.rs b/native/core/src/parquet/mod.rs index cfa03220c1..6a06a34cd7 100644 --- a/native/core/src/parquet/mod.rs +++ b/native/core/src/parquet/mod.rs @@ -160,11 +160,12 @@ pub unsafe extern "system" fn Java_org_apache_comet_parquet_Native_initRecordBat let path: String = file_path.try_to_string(env).unwrap(); let object_store_config = get_object_store_options(env, object_store_options)?; - let (object_store_url, object_store_path) = prepare_object_store_with_configs( - session_ctx.runtime_env(), - path.clone(), - &object_store_config, - )?; + let (object_store_url, object_store_path, object_store_backend) = + prepare_object_store_with_configs( + session_ctx.runtime_env(), + path.clone(), + &object_store_config, + )?; let required_schema_buffer = env.convert_byte_array(&required_schema)?; let required_schema = Arc::new(deserialize_schema(&required_schema_buffer)?); @@ -213,6 +214,7 @@ pub unsafe extern "system" fn Java_org_apache_comet_parquet_Native_initRecordBat Some(data_schema), None, object_store_url, + object_store_backend, file_groups, None, data_filters, diff --git a/native/core/src/parquet/parquet_exec.rs b/native/core/src/parquet/parquet_exec.rs index 1308ce97fc..d99000ff72 100644 --- a/native/core/src/parquet/parquet_exec.rs +++ b/native/core/src/parquet/parquet_exec.rs @@ -16,8 +16,9 @@ // under the License. use crate::execution::operators::ExecutionError; -use crate::parquet::eager_page_index_reader_factory::EagerPageIndexReaderFactory; +use crate::parquet::eager_page_index_reader_factory::{EagerPageIndexReaderFactory, ScanIoSource}; use crate::parquet::encryption_support::{CometEncryptionConfig, ENCRYPTION_FACTORY_ID}; +use crate::parquet::parquet_support::ObjectStoreBackend; use crate::parquet::parquet_support::SparkParquetOptions; use crate::parquet::schema_adapter::SparkPhysicalExprAdapterFactory; use arrow::datatypes::{Field, SchemaRef}; @@ -62,6 +63,7 @@ pub(crate) fn init_datasource_exec( data_schema: Option, partition_schema: Option, object_store_url: ObjectStoreUrl, + object_store_backend: ObjectStoreBackend, file_groups: Vec>, projection_vector: Option>, data_filters: Option>>, @@ -158,15 +160,21 @@ pub(crate) fn init_datasource_exec( // cached with the footer, at the cost of losing the skip's benefit when it would have // applied. Filed upstream as apache/datafusion#23978; revert this once that's fixed. // - // TODO: metadata I/O is invisible in metrics. `fetch_metadata` reads via `ObjectStore::get_ranges`, - // bypassing the `get_bytes` path where `bytes_scanned` is counted. A byte-counting ObjectStore - // wrapper would surface it. + // Preserve bytes_scanned's existing requested data/Bloom-filter range accounting. Footer + // and page-index reads through get_metadata bypass it, and coalescing may fetch extra bytes. + // The scan I/O counters expose those gaps without changing task inputMetrics.bytesRead or + // scan_efficiency_ratio, which continue to derive from bytes_scanned. let runtime_env = session_ctx.runtime_env(); let store = runtime_env.object_store(&object_store_url)?; let metadata_cache = runtime_env.cache_manager.get_file_metadata_cache(); - parquet_source = parquet_source.with_parquet_file_reader_factory(Arc::new( - EagerPageIndexReaderFactory::new(store, metadata_cache), + let scan_io_source = scan_io_source(object_store_backend); + let reader_factory = Arc::new(EagerPageIndexReaderFactory::new( + store, + metadata_cache, + scan_io_source, + parquet_source.metrics(), )); + parquet_source = parquet_source.with_parquet_file_reader_factory(reader_factory); // Route data filters through `try_pushdown_filters` rather than calling // `with_predicate` directly. This is the contract DataFusion's optimizer @@ -217,6 +225,14 @@ pub(crate) fn init_datasource_exec( Ok(data_source_exec) } +fn scan_io_source(backend: ObjectStoreBackend) -> ScanIoSource { + match backend { + ObjectStoreBackend::Local => ScanIoSource::Local, + ObjectStoreBackend::Remote => ScanIoSource::ObjectStore, + ObjectStoreBackend::Other => ScanIoSource::OtherObjectStore, + } +} + #[allow(clippy::too_many_arguments)] fn get_options( session_timezone: &str, @@ -294,13 +310,125 @@ mod tests { use arrow::array::Int32Array; use arrow::datatypes::{DataType, Field, Schema}; use arrow::record_batch::RecordBatch; - use datafusion::datasource::physical_plan::parquet::metadata::CachedParquetMetaData; + use bytes::Bytes; + use datafusion::datasource::physical_plan::parquet::metadata::{ + CachedParquetMetaData, DFParquetMetadata, + }; + use datafusion::datasource::physical_plan::parquet::ParquetFileReaderFactory; + use datafusion::execution::cache::cache_manager::CachedFileMetadataEntry; + use datafusion::logical_expr::Operator; + use datafusion::physical_expr::expressions::{BinaryExpr, Literal}; + use datafusion::physical_plan::metrics::ExecutionPlanMetricsSet; use datafusion::physical_plan::ExecutionPlan; use datafusion_comet_spark_expr::test_common::file_util::get_temp_filename; use futures::StreamExt; + use object_store::memory::InMemory; + use object_store::throttle::{ThrottleConfig, ThrottledStore}; + use object_store::{ObjectStore, ObjectStoreExt}; + use parquet::arrow::arrow_reader::ArrowReaderOptions; use parquet::arrow::ArrowWriter; + use parquet::encryption::decrypt::{FileDecryptionProperties, KeyRetriever}; + use parquet::encryption::encrypt::FileEncryptionProperties; + use parquet::file::metadata::{KeyValue, PageIndexPolicy, ParquetMetaDataReader}; use parquet::file::properties::{EnabledStatistics, WriterProperties}; use std::fs::File; + use std::time::Duration; + + fn write_scan_io_fixture() -> (String, SchemaRef) { + let schema = Arc::new(Schema::new(vec![ + Field::new("a", DataType::Int32, false), + Field::new("b", DataType::Int32, false), + ])); + let first = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(Int32Array::from((0..500).collect::>())), + Arc::new(Int32Array::from((10_000..10_500).collect::>())), + ], + ) + .unwrap(); + let second = RecordBatch::try_new( + Arc::clone(&schema), + vec![ + Arc::new(Int32Array::from((1_000..1_500).collect::>())), + Arc::new(Int32Array::from((11_000..11_500).collect::>())), + ], + ) + .unwrap(); + + let filename = get_temp_filename() + .as_path() + .as_os_str() + .to_str() + .unwrap() + .to_string(); + let props = WriterProperties::builder() + .set_statistics_enabled(EnabledStatistics::Page) + .set_data_page_row_count_limit(100) + .set_max_row_group_row_count(Some(500)) + .build(); + let file = File::create(&filename).unwrap(); + let mut writer = ArrowWriter::try_new(file, Arc::clone(&schema), Some(props)).unwrap(); + writer.write(&first).unwrap(); + writer.write(&second).unwrap(); + writer.close().unwrap(); + + (filename, schema) + } + + fn init_test_scan( + required_schema: SchemaRef, + data_schema: SchemaRef, + partitioned_file: PartitionedFile, + projection: Option>, + filters: Option>>, + session_ctx: &Arc, + ) -> Arc { + init_datasource_exec( + required_schema, + Some(data_schema), + None, + ObjectStoreUrl::local_filesystem(), + ObjectStoreBackend::Local, + vec![vec![partitioned_file]], + projection, + filters, + None, + "UTC", + true, + false, + false, + false, + session_ctx, + false, + false, + false, + ) + .unwrap() + } + + async fn drain_scan(scan: &Arc, session_ctx: &Arc) { + let mut stream = scan.execute(0, session_ctx.task_ctx()).unwrap(); + while let Some(batch) = stream.next().await { + batch.unwrap(); + } + } + + fn scan_metric(scan: &Arc, name: &str) -> usize { + scan.metrics() + .unwrap() + .sum_by_name(name) + .unwrap_or_else(|| panic!("missing metric {name}")) + .as_usize() + } + + fn reader_metric(metrics: &ExecutionPlanMetricsSet, name: &str) -> usize { + metrics + .clone_inner() + .sum_by_name(name) + .unwrap_or_else(|| panic!("missing metric {name}")) + .as_usize() + } // Regression test for #4990: a fresh `TableParquetOptions::new()` ignored session-level // `datafusion.execution.parquet.*` settings entirely, so `spark.comet.datafusion. @@ -390,6 +518,7 @@ mod tests { None, None, ObjectStoreUrl::local_filesystem(), + ObjectStoreBackend::Local, vec![vec![partitioned_file]], None, None, @@ -432,4 +561,714 @@ mod tests { "cached metadata must include the page index" ); } + + #[tokio::test] + async fn reports_cold_and_warm_metadata_io_without_data_reads() { + let (filename, _schema) = write_scan_io_fixture(); + let file_bytes = std::fs::read(&filename).unwrap(); + let footer_bytes = usize::try_from(u32::from_le_bytes( + file_bytes[file_bytes.len() - 8..file_bytes.len() - 4] + .try_into() + .unwrap(), + )) + .unwrap(); + let partitioned_file = PartitionedFile::from_path(filename).unwrap(); + let session_ctx = Arc::new(SessionContext::new()); + let runtime_env = session_ctx.runtime_env(); + let store = runtime_env + .object_store(ObjectStoreUrl::local_filesystem()) + .unwrap(); + let cold_metrics = ExecutionPlanMetricsSet::new(); + let metadata_cache = runtime_env.cache_manager.get_file_metadata_cache(); + let cold_factory = EagerPageIndexReaderFactory::new( + Arc::clone(&store), + Arc::clone(&metadata_cache), + ScanIoSource::Local, + &cold_metrics, + ); + let mut cold_reader = cold_factory + .create_reader(0, partitioned_file.clone(), Some(512 * 1024), &cold_metrics) + .unwrap(); + cold_reader.get_metadata(None).await.unwrap(); + + let cold_metadata_bytes = reader_metric(&cold_metrics, "scan_io_metadata_bytes"); + assert!(cold_metadata_bytes > footer_bytes); + assert_eq!(reader_metric(&cold_metrics, "scan_io_data_bytes"), 0); + assert_eq!(reader_metric(&cold_metrics, "scan_io_footer_reads"), 1); + assert_eq!( + reader_metric(&cold_metrics, "scan_io_footer_bytes"), + footer_bytes + ); + assert_eq!( + reader_metric(&cold_metrics, "scan_io_object_store_get_calls"), + 0 + ); + assert_eq!( + reader_metric(&cold_metrics, "scan_io_object_store_get_requested_bytes"), + 0 + ); + assert_eq!( + reader_metric(&cold_metrics, "scan_io_object_store_response_bytes_read"), + 0 + ); + assert_eq!( + reader_metric(&cold_metrics, "scan_io_metadata_cache_hits"), + 0 + ); + assert_eq!( + reader_metric(&cold_metrics, "scan_io_metadata_cache_misses"), + 1 + ); + + let warm_metrics = ExecutionPlanMetricsSet::new(); + let warm_factory = EagerPageIndexReaderFactory::new( + store, + metadata_cache, + ScanIoSource::Local, + &warm_metrics, + ); + let mut warm_reader = warm_factory + .create_reader(0, partitioned_file, Some(512 * 1024), &warm_metrics) + .unwrap(); + warm_reader.get_metadata(None).await.unwrap(); + + assert_eq!(reader_metric(&warm_metrics, "scan_io_metadata_bytes"), 0); + assert_eq!(reader_metric(&warm_metrics, "scan_io_footer_reads"), 0); + assert_eq!(reader_metric(&warm_metrics, "scan_io_footer_bytes"), 0); + assert_eq!( + reader_metric(&warm_metrics, "scan_io_metadata_cache_hits"), + 1 + ); + assert_eq!( + reader_metric(&warm_metrics, "scan_io_metadata_cache_misses"), + 0 + ); + } + + #[tokio::test] + async fn upgrades_partial_cached_metadata_with_local_direct_reads() { + let (filename, _schema) = write_scan_io_fixture(); + let partitioned_file = PartitionedFile::from_path(filename).unwrap(); + let session_ctx = Arc::new(SessionContext::new()); + let runtime_env = session_ctx.runtime_env(); + let store = runtime_env + .object_store(ObjectStoreUrl::local_filesystem()) + .unwrap(); + let metadata_cache = runtime_env.cache_manager.get_file_metadata_cache(); + let footer_metadata = DFParquetMetadata::new(store.as_ref(), &partitioned_file.object_meta) + .with_page_index_policy(Some(PageIndexPolicy::Skip)) + .fetch_metadata() + .await + .unwrap(); + assert!(footer_metadata.column_index().is_none()); + assert!(footer_metadata.offset_index().is_none()); + metadata_cache.put( + &partitioned_file.object_meta.location, + CachedFileMetadataEntry::new( + partitioned_file.object_meta.clone(), + Arc::new(CachedParquetMetaData::new(footer_metadata)), + ), + ); + + let metrics = ExecutionPlanMetricsSet::new(); + let factory = + EagerPageIndexReaderFactory::new(store, metadata_cache, ScanIoSource::Local, &metrics); + let mut reader = factory + .create_reader(0, partitioned_file, Some(512 * 1024), &metrics) + .unwrap(); + let metadata = reader.get_metadata(None).await.unwrap(); + + assert!(metadata.column_index().is_some()); + assert!(metadata.offset_index().is_some()); + let metadata_bytes = reader_metric(&metrics, "scan_io_metadata_bytes"); + assert!(metadata_bytes > 0); + assert_eq!(reader_metric(&metrics, "scan_io_data_bytes"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_footer_bytes"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_cache_hits"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_cache_misses"), 1); + + let column = metadata.row_group(0).column(0); + let page_index_offset = u64::try_from(column.column_index_offset().unwrap()).unwrap(); + let page_index_length = u64::try_from(column.column_index_length().unwrap()).unwrap(); + reader + .get_bytes(page_index_offset..page_index_offset + page_index_length) + .await + .unwrap(); + assert_eq!( + reader_metric(&metrics, "scan_io_metadata_bytes"), + metadata_bytes + page_index_length as usize + ); + assert_eq!(reader_metric(&metrics, "scan_io_data_bytes"), 0); + } + + #[test] + fn registers_scan_io_metrics_once_per_execution_plan() { + let (filename, _schema) = write_scan_io_fixture(); + let partitioned_file = PartitionedFile::from_path(filename).unwrap(); + let session_ctx = Arc::new(SessionContext::new()); + let runtime_env = session_ctx.runtime_env(); + let store = runtime_env + .object_store(ObjectStoreUrl::local_filesystem()) + .unwrap(); + let metadata_cache = runtime_env.cache_manager.get_file_metadata_cache(); + let metrics = ExecutionPlanMetricsSet::new(); + let factory = + EagerPageIndexReaderFactory::new(store, metadata_cache, ScanIoSource::Local, &metrics); + + for partition in 0..8 { + factory + .create_reader(partition, partitioned_file.clone(), None, &metrics) + .unwrap(); + } + + let registered_metrics = metrics.clone_inner(); + let scan_io_metrics = registered_metrics + .iter() + .filter(|metric| metric.value().name().starts_with("scan_io_")) + .collect::>(); + assert_eq!(scan_io_metrics.len(), 9); + for metric in scan_io_metrics { + assert!(metric.labels().is_empty()); + assert_eq!(metric.partition(), None); + } + } + + #[tokio::test] + async fn reports_object_store_read_amplification_after_range_coalescing() { + let file_size = 524_416; + let location = object_store::path::Path::from("coalesced.parquet"); + let store: Arc = Arc::new(InMemory::new()); + store + .put(&location, Bytes::from(vec![0; file_size]).into()) + .await + .unwrap(); + + let session_ctx = Arc::new(SessionContext::new()); + let metadata_cache = session_ctx + .runtime_env() + .cache_manager + .get_file_metadata_cache(); + let metrics = ExecutionPlanMetricsSet::new(); + let factory = EagerPageIndexReaderFactory::new( + store, + metadata_cache, + ScanIoSource::ObjectStore, + &metrics, + ); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::new(location.to_string(), file_size as u64), + None, + &metrics, + ) + .unwrap(); + let buffers = reader + .get_byte_ranges(vec![0..64, 524_352..524_416]) + .await + .unwrap(); + + assert_eq!(buffers.iter().map(Bytes::len).sum::(), 128); + assert_eq!(reader_metric(&metrics, "bytes_scanned"), 128); + assert_eq!(reader_metric(&metrics, "scan_io_data_bytes"), 128); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_bytes"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_object_store_get_calls"), 1); + assert_eq!( + reader_metric(&metrics, "scan_io_object_store_get_requested_bytes"), + file_size + ); + assert_eq!( + reader_metric(&metrics, "scan_io_object_store_response_bytes_read"), + file_size + ); + } + + #[tokio::test] + async fn reports_object_store_footer_and_metadata_reads() { + let (filename, _schema) = write_scan_io_fixture(); + let file_bytes = std::fs::read(filename).unwrap(); + let file_size = file_bytes.len(); + let footer_bytes = usize::try_from(u32::from_le_bytes( + file_bytes[file_size - 8..file_size - 4].try_into().unwrap(), + )) + .unwrap(); + let location = object_store::path::Path::from("footer.parquet"); + let store: Arc = Arc::new(InMemory::new()); + store + .put(&location, Bytes::from(file_bytes).into()) + .await + .unwrap(); + + let session_ctx = Arc::new(SessionContext::new()); + let metadata_cache = session_ctx + .runtime_env() + .cache_manager + .get_file_metadata_cache(); + let metrics = ExecutionPlanMetricsSet::new(); + let factory = EagerPageIndexReaderFactory::new( + store, + metadata_cache, + ScanIoSource::ObjectStore, + &metrics, + ); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::new(location.to_string(), file_size as u64), + Some(512 * 1024), + &metrics, + ) + .unwrap(); + reader.get_metadata(None).await.unwrap(); + + assert_eq!(reader_metric(&metrics, "scan_io_data_bytes"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_bytes"), file_size); + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 1); + assert_eq!( + reader_metric(&metrics, "scan_io_footer_bytes"), + footer_bytes + ); + assert_eq!(reader_metric(&metrics, "scan_io_object_store_get_calls"), 1); + assert_eq!( + reader_metric(&metrics, "scan_io_object_store_get_requested_bytes"), + file_size + ); + assert_eq!( + reader_metric(&metrics, "scan_io_object_store_response_bytes_read"), + file_size + ); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_cache_misses"), 1); + } + + #[tokio::test] + async fn reports_plaintext_footer_before_page_index_read() { + let (filename, _schema) = write_scan_io_fixture(); + let file_bytes = std::fs::read(filename).unwrap(); + let file_size = file_bytes.len(); + let footer_bytes = usize::try_from(u32::from_le_bytes( + file_bytes[file_size - 8..file_size - 4].try_into().unwrap(), + )) + .unwrap(); + let location = object_store::path::Path::from("plaintext-footer.parquet"); + let store: Arc = Arc::new(ThrottledStore::new( + InMemory::new(), + ThrottleConfig { + wait_get_per_call: Duration::from_millis(10), + ..Default::default() + }, + )); + store + .put(&location, Bytes::from(file_bytes).into()) + .await + .unwrap(); + + let session_ctx = Arc::new(SessionContext::new()); + let metadata_cache = session_ctx + .runtime_env() + .cache_manager + .get_file_metadata_cache(); + let metrics = ExecutionPlanMetricsSet::new(); + let factory = EagerPageIndexReaderFactory::new( + store, + metadata_cache, + ScanIoSource::ObjectStore, + &metrics, + ); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::new(location.to_string(), file_size as u64), + Some(footer_bytes + 8), + &metrics, + ) + .unwrap(); + let metadata = reader.get_metadata(None); + tokio::pin!(metadata); + + tokio::select! { + result = &mut metadata => { + panic!("metadata load completed before the page-index read: {result:?}"); + } + () = async { + while reader_metric(&metrics, "scan_io_object_store_get_calls") < 2 { + tokio::task::yield_now().await; + } + } => { + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 1); + assert_eq!(reader_metric(&metrics, "scan_io_footer_bytes"), footer_bytes); + } + } + + metadata.await.unwrap(); + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 1); + assert_eq!( + reader_metric(&metrics, "scan_io_footer_bytes"), + footer_bytes + ); + } + + #[tokio::test] + async fn reports_plaintext_footer_when_prefetched_page_index_is_invalid() { + let (filename, _schema) = write_scan_io_fixture(); + let mut file_bytes = std::fs::read(filename).unwrap(); + let file_size = file_bytes.len(); + let footer_bytes = usize::try_from(u32::from_le_bytes( + file_bytes[file_size - 8..file_size - 4].try_into().unwrap(), + )) + .unwrap(); + let footer_start = file_size - 8 - footer_bytes; + let footer = + ParquetMetaDataReader::decode_metadata(&file_bytes[footer_start..file_size - 8]) + .unwrap(); + let column = footer.row_group(0).column(0); + let index_start = usize::try_from(column.column_index_offset().unwrap()).unwrap(); + let index_bytes = usize::try_from(column.column_index_length().unwrap()).unwrap(); + file_bytes[index_start..index_start + index_bytes].fill(0xff); + + let location = object_store::path::Path::from("invalid-prefetched-page-index.parquet"); + let store: Arc = Arc::new(InMemory::new()); + store + .put(&location, Bytes::from(file_bytes).into()) + .await + .unwrap(); + + let session_ctx = Arc::new(SessionContext::new()); + let metadata_cache = session_ctx + .runtime_env() + .cache_manager + .get_file_metadata_cache(); + let metrics = ExecutionPlanMetricsSet::new(); + let factory = EagerPageIndexReaderFactory::new( + store, + metadata_cache, + ScanIoSource::ObjectStore, + &metrics, + ); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::new(location.to_string(), file_size as u64), + Some(512 * 1024), + &metrics, + ) + .unwrap(); + + assert!(reader.get_metadata(None).await.is_err()); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_bytes"), file_size); + assert_eq!(reader_metric(&metrics, "scan_io_object_store_get_calls"), 1); + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 1); + assert_eq!( + reader_metric(&metrics, "scan_io_footer_bytes"), + footer_bytes + ); + } + + #[tokio::test] + async fn reports_encrypted_footer_before_key_retrieval() { + assert_encrypted_footer_before_key_retrieval(0, false).await; + } + + #[tokio::test] + async fn does_not_report_encrypted_footer_before_payload_is_read() { + assert_encrypted_footer_before_key_retrieval(512 * 1024, false).await; + } + + #[tokio::test] + async fn counts_complete_encrypted_footer_even_when_authentication_fails() { + assert_encrypted_footer_before_key_retrieval(0, true).await; + } + + async fn assert_encrypted_footer_before_key_retrieval(footer_padding: usize, corrupt: bool) { + struct ObservingKeyRetriever { + key: Vec, + metrics: Arc, + footer_bytes: usize, + } + + impl KeyRetriever for ObservingKeyRetriever { + fn retrieve_key(&self, _key_metadata: &[u8]) -> parquet::errors::Result> { + assert_eq!(reader_metric(&self.metrics, "scan_io_footer_reads"), 1); + assert_eq!( + reader_metric(&self.metrics, "scan_io_footer_bytes"), + self.footer_bytes + ); + Ok(self.key.clone()) + } + } + + let key = b"0123456789012345".to_vec(); + let encryption = FileEncryptionProperties::builder(key.clone()) + .with_footer_key_metadata(b"footer".to_vec()) + .build() + .unwrap(); + let properties = WriterProperties::builder() + .with_file_encryption_properties(encryption) + .set_key_value_metadata((footer_padding > 0).then(|| { + vec![KeyValue::new( + "padding".to_string(), + "x".repeat(footer_padding), + )] + })) + .build(); + let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])); + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from(vec![1, 2, 3]))], + ) + .unwrap(); + let filename = get_temp_filename(); + let mut writer = + ArrowWriter::try_new(File::create(&filename).unwrap(), schema, Some(properties)) + .unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); + + let mut file_bytes = std::fs::read(&filename).unwrap(); + let file_size = file_bytes.len(); + let footer_bytes = usize::try_from(u32::from_le_bytes( + file_bytes[file_size - 8..file_size - 4].try_into().unwrap(), + )) + .unwrap(); + let location = object_store::path::Path::from("encrypted-footer.parquet"); + if corrupt { + // Corrupt the encrypted payload's authentication tag, preserving its length/trailer + // and key metadata. Complete I/O is counted before decryption is attempted. + file_bytes[file_size - 9] ^= 1; + } + let store: Arc = if footer_padding > 0 { + assert!(footer_bytes > 512 * 1024); + Arc::new(ThrottledStore::new( + InMemory::new(), + ThrottleConfig { + wait_get_per_call: Duration::from_millis(10), + ..Default::default() + }, + )) + } else { + Arc::new(InMemory::new()) + }; + store + .put(&location, Bytes::from(file_bytes).into()) + .await + .unwrap(); + + let session_ctx = Arc::new(SessionContext::new()); + let metadata_cache = session_ctx + .runtime_env() + .cache_manager + .get_file_metadata_cache(); + let metrics = Arc::new(ExecutionPlanMetricsSet::new()); + let factory = EagerPageIndexReaderFactory::new( + store, + metadata_cache, + ScanIoSource::ObjectStore, + &metrics, + ); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::new(location.to_string(), file_size as u64), + Some(512 * 1024), + &metrics, + ) + .unwrap(); + let decryption = + FileDecryptionProperties::with_key_retriever(Arc::new(ObservingKeyRetriever { + key, + metrics: Arc::clone(&metrics), + footer_bytes, + })) + .build() + .unwrap(); + let options = ArrowReaderOptions::new().with_file_decryption_properties(decryption); + + if footer_padding > 0 { + let metadata = reader.get_metadata(Some(&options)); + tokio::pin!(metadata); + + tokio::select! { + result = &mut metadata => { + panic!("metadata load completed before the second footer read: {result:?}"); + } + () = async { + while reader_metric(&metrics, "scan_io_object_store_get_calls") < 2 { + tokio::task::yield_now().await; + } + } => { + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_footer_bytes"), 0); + } + } + + metadata.await.unwrap(); + } else if corrupt { + assert!(reader.get_metadata(Some(&options)).await.is_err()); + } else { + reader.get_metadata(Some(&options)).await.unwrap(); + } + + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 1); + assert_eq!( + reader_metric(&metrics, "scan_io_footer_bytes"), + footer_bytes + ); + } + + #[tokio::test] + async fn does_not_report_footer_metrics_when_metadata_decoding_fails() { + let (filename, _schema) = write_scan_io_fixture(); + let mut file_bytes = std::fs::read(&filename).unwrap(); + let file_size = file_bytes.len(); + let footer_bytes = usize::try_from(u32::from_le_bytes( + file_bytes[file_size - 8..file_size - 4].try_into().unwrap(), + )) + .unwrap(); + file_bytes[file_size - 8 - footer_bytes..file_size - 8].fill(0xff); + std::fs::write(&filename, file_bytes).unwrap(); + + let session_ctx = Arc::new(SessionContext::new()); + let runtime_env = session_ctx.runtime_env(); + let store = runtime_env + .object_store(ObjectStoreUrl::local_filesystem()) + .unwrap(); + let metadata_cache = runtime_env.cache_manager.get_file_metadata_cache(); + let metrics = ExecutionPlanMetricsSet::new(); + let factory = + EagerPageIndexReaderFactory::new(store, metadata_cache, ScanIoSource::Local, &metrics); + let mut reader = factory + .create_reader( + 0, + PartitionedFile::from_path(filename).unwrap(), + Some(512 * 1024), + &metrics, + ) + .unwrap(); + + assert!(reader.get_metadata(None).await.is_err()); + assert!(reader_metric(&metrics, "scan_io_metadata_bytes") > 0); + assert_eq!(reader_metric(&metrics, "scan_io_footer_reads"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_footer_bytes"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_cache_hits"), 0); + assert_eq!(reader_metric(&metrics, "scan_io_metadata_cache_misses"), 0); + } + + #[tokio::test] + async fn classifies_bloom_filter_reads_as_metadata() { + let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])); + let batch = RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(Int32Array::from( + (0..500).map(|value| value * 2).collect::>(), + ))], + ) + .unwrap(); + let filename = get_temp_filename() + .as_path() + .as_os_str() + .to_str() + .unwrap() + .to_string(); + let props = WriterProperties::builder() + .set_statistics_enabled(EnabledStatistics::Page) + .set_bloom_filter_ndv(500) + .set_bloom_filter_fpp(0.0001) + .build(); + let file = File::create(&filename).unwrap(); + let mut writer = ArrowWriter::try_new(file, Arc::clone(&schema), Some(props)).unwrap(); + writer.write(&batch).unwrap(); + writer.close().unwrap(); + + let filter: Arc = Arc::new(BinaryExpr::new( + Arc::new(Column::new("a", 0)), + Operator::Eq, + Arc::new(Literal::new(ScalarValue::Int32(Some(501)))), + )); + let session_ctx = Arc::new(SessionContext::new()); + let scan = init_test_scan( + Arc::clone(&schema), + schema, + PartitionedFile::from_path(filename).unwrap(), + None, + Some(vec![filter]), + &session_ctx, + ); + drain_scan(&scan, &session_ctx).await; + + let bloom_bytes = scan_metric(&scan, "bytes_scanned"); + assert!(bloom_bytes > 0); + assert_eq!(scan_metric(&scan, "scan_io_data_bytes"), 0); + assert!(scan_metric(&scan, "scan_io_metadata_bytes") >= bloom_bytes); + } + + #[tokio::test] + async fn reports_projected_and_predicate_pruned_data_io() { + let (filename, schema) = write_scan_io_fixture(); + let partitioned_file = PartitionedFile::from_path(filename).unwrap(); + let session_ctx = Arc::new(SessionContext::new()); + + let full_scan = init_test_scan( + Arc::clone(&schema), + Arc::clone(&schema), + partitioned_file.clone(), + None, + None, + &session_ctx, + ); + drain_scan(&full_scan, &session_ctx).await; + + let projected_schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int32, false)])); + let projected_scan = init_test_scan( + Arc::clone(&projected_schema), + Arc::clone(&schema), + partitioned_file.clone(), + Some(vec![0]), + None, + &session_ctx, + ); + drain_scan(&projected_scan, &session_ctx).await; + + let filter: Arc = Arc::new(BinaryExpr::new( + Arc::new(Column::new("a", 0)), + Operator::Lt, + Arc::new(Literal::new(ScalarValue::Int32(Some(500)))), + )); + let filtered_scan = init_test_scan( + projected_schema, + schema, + partitioned_file, + Some(vec![0]), + Some(vec![filter]), + &session_ctx, + ); + drain_scan(&filtered_scan, &session_ctx).await; + + let full_data = scan_metric(&full_scan, "scan_io_data_bytes"); + let projected_data = scan_metric(&projected_scan, "scan_io_data_bytes"); + let filtered_data = scan_metric(&filtered_scan, "scan_io_data_bytes"); + assert!( + projected_data < full_data, + "projection should request fewer data bytes: projected={projected_data}, full={full_data}" + ); + assert!( + filtered_data < projected_data, + "row-group pruning should request fewer data bytes: filtered={filtered_data}, projected={projected_data}" + ); + + for scan in [&full_scan, &projected_scan, &filtered_scan] { + assert_eq!( + scan_metric(scan, "bytes_scanned"), + scan_metric(scan, "scan_io_data_bytes") + ); + assert_eq!(scan_metric(scan, "scan_io_object_store_get_calls"), 0); + assert_eq!( + scan_metric(scan, "scan_io_object_store_get_requested_bytes"), + 0 + ); + assert_eq!( + scan_metric(scan, "scan_io_object_store_response_bytes_read"), + 0 + ); + } + } } diff --git a/native/core/src/parquet/parquet_support.rs b/native/core/src/parquet/parquet_support.rs index 5b22afa260..d4f4cce9ab 100644 --- a/native/core/src/parquet/parquet_support.rs +++ b/native/core/src/parquet/parquet_support.rs @@ -37,7 +37,7 @@ use datafusion::physical_plan::ColumnarValue; use datafusion_comet_spark_expr::EvalMode; use log::debug; use object_store::path::Path; -use object_store::{parse_url, ObjectStore}; +use object_store::{parse_url, ObjectStore, ObjectStoreScheme}; use parquet::arrow::PARQUET_FIELD_ID_META_KEY; use std::collections::HashMap; use std::sync::OnceLock; @@ -519,16 +519,43 @@ fn hash_object_store_configs(configs: &HashMap) -> u64 { hasher.finish() } -/// Parses the url, registers the object store with configurations, and returns a tuple of the object store url -/// and object store path +/// The selected backend, independent of the URL used to register it in DataFusion. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub(crate) enum ObjectStoreBackend { + Local, + Remote, + Other, +} + +fn object_store_backend(url: &Url, is_hdfs: bool) -> Result { + if is_hdfs { + // Custom libhdfs schemes may look like cloud URLs but retain their own range-read API. + return Ok(ObjectStoreBackend::Other); + } + let (scheme, _) = + ObjectStoreScheme::parse(url).map_err(|e| ExecutionError::GeneralError(e.to_string()))?; + Ok(match scheme { + ObjectStoreScheme::Local => ObjectStoreBackend::Local, + ObjectStoreScheme::AmazonS3 + | ObjectStoreScheme::GoogleCloudStorage + | ObjectStoreScheme::MicrosoftAzure + | ObjectStoreScheme::Http => ObjectStoreBackend::Remote, + // Memory and future backends are not implicitly classified as remote network traffic. + _ => ObjectStoreBackend::Other, + }) +} + +/// Parses the URL, selects and registers the backend, and returns its registry URL, path, and +/// classification. Callers must use this classification rather than infer it from a scheme alias. pub(crate) fn prepare_object_store_with_configs( runtime_env: Arc, url: String, object_store_configs: &HashMap, -) -> Result<(ObjectStoreUrl, Path), ExecutionError> { +) -> Result<(ObjectStoreUrl, Path, ObjectStoreBackend), ExecutionError> { let mut url = Url::parse(url.as_str()) .map_err(|e| ExecutionError::GeneralError(format!("Error parsing URL {url}: {e}")))?; let is_hdfs_scheme = is_hdfs_scheme(&url, object_store_configs); + let backend = object_store_backend(&url, is_hdfs_scheme)?; let mut scheme = url.scheme(); if !is_hdfs_scheme && scheme == "s3a" { scheme = "s3"; @@ -583,11 +610,63 @@ pub(crate) fn prepare_object_store_with_configs( let object_store_url = ObjectStoreUrl::parse(url_key.clone())?; runtime_env.register_object_store(&url, object_store); - Ok((object_store_url, object_store_path)) + Ok((object_store_url, object_store_path, backend)) } #[cfg(test)] mod tests { + #[test] + fn classifies_the_selected_backend_using_object_store_parser() { + use super::{is_hdfs_scheme, object_store_backend, ObjectStoreBackend}; + let configs = std::collections::HashMap::from([( + "fs.comet.libhdfs.schemes".to_string(), + "s3,abfs".to_string(), + )]); + for address in [ + "s3://bucket/path", + "s3a://bucket/path", + "gs://bucket/path", + "az://container/path", + "adl://container/path", + "azure://container/path", + "abfs://container/path", + "abfss://container/path", + "http://example.com/path", + "https://example.com/path", + "https://account.blob.core.windows.net/container/path", + ] { + let url = url::Url::parse(address).unwrap(); + assert_eq!( + object_store_backend(&url, false).unwrap(), + ObjectStoreBackend::Remote, + "{address}" + ); + if is_hdfs_scheme(&url, &configs) { + assert_eq!( + object_store_backend(&url, true).unwrap(), + ObjectStoreBackend::Other + ); + } + } + assert_eq!( + object_store_backend(&url::Url::parse("file:///tmp/a").unwrap(), false).unwrap(), + ObjectStoreBackend::Local + ); + assert_eq!( + object_store_backend(&url::Url::parse("memory:///a").unwrap(), false).unwrap(), + ObjectStoreBackend::Other + ); + // These spellings are not accepted native backends in pinned object_store 0.13.2. + for scheme in ["gcs", "wasb", "wasbs", "s3n"] { + let url = url::Url::parse(&format!("{scheme}://bucket/path")).unwrap(); + assert!(object_store_backend(&url, false).is_err()); + assert_eq!( + object_store_backend(&url, true).unwrap(), + ObjectStoreBackend::Other + ); + } + } + #[cfg(not(feature = "hdfs-opendal"))] use datafusion::execution::object_store::ObjectStoreUrl; #[cfg(not(feature = "hdfs-opendal"))] @@ -612,6 +691,7 @@ mod tests { ) -> Result<(ObjectStoreUrl, Path), ExecutionError> { use crate::parquet::parquet_support::prepare_object_store_with_configs; prepare_object_store_with_configs(runtime_env, url, &HashMap::new()) + .map(|(url, path, _)| (url, path)) } #[cfg(not(feature = "hdfs-opendal"))] diff --git a/spark/src/main/scala/org/apache/spark/sql/comet/CometMetricNode.scala b/spark/src/main/scala/org/apache/spark/sql/comet/CometMetricNode.scala index a0eb94e83b..53a4cc8f32 100644 --- a/spark/src/main/scala/org/apache/spark/sql/comet/CometMetricNode.scala +++ b/spark/src/main/scala/org/apache/spark/sql/comet/CometMetricNode.scala @@ -298,6 +298,30 @@ object CometMetricNode { SQLMetrics.createMetric(sc, "Number of row groups matched by limit pruning (not pruned)"), "bytes_scanned" -> SQLMetrics.createSizeMetric(sc, "Number of bytes scanned"), + "scan_io_data_bytes" -> + SQLMetrics.createSizeMetric( + sc, + "Projected Parquet data-page bytes returned to the reader"), + "scan_io_metadata_bytes" -> + SQLMetrics.createSizeMetric( + sc, + "Footer-prefetch, page-index, and Bloom-filter bytes returned to the reader"), + "scan_io_footer_reads" -> + SQLMetrics.createMetric(sc, "Number of Parquet footer payloads read from storage"), + "scan_io_footer_bytes" -> + SQLMetrics.createSizeMetric(sc, "Serialized Parquet footer payload bytes read"), + "scan_io_object_store_get_calls" -> + SQLMetrics.createMetric(sc, "ObjectStore GET operations after range coalescing"), + "scan_io_object_store_get_requested_bytes" -> + SQLMetrics.createSizeMetric(sc, "ObjectStore GET range bytes after range coalescing"), + "scan_io_object_store_response_bytes_read" -> + SQLMetrics.createSizeMetric( + sc, + "ObjectStore response bytes consumed after range coalescing"), + "scan_io_metadata_cache_hits" -> + SQLMetrics.createMetric(sc, "Metadata loads served without storage I/O"), + "scan_io_metadata_cache_misses" -> + SQLMetrics.createMetric(sc, "Metadata loads requiring storage I/O"), "pushdown_rows_pruned" -> SQLMetrics.createMetric(sc, "Rows filtered out by predicates pushed into parquet scan"), "pushdown_rows_matched" -> diff --git a/spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala b/spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala index 4eb6d00178..899d3ffc5c 100644 --- a/spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala +++ b/spark/src/test/scala/org/apache/comet/exec/CometExecSuite.scala @@ -2412,14 +2412,35 @@ class CometExecSuite extends CometTestBase { metrics.contains("time_elapsed_scanning_total"), s"Missing time_elapsed_scanning_total. Available: ${metrics.keys}") assert(metrics.contains("bytes_scanned")) + assert(metrics("bytes_scanned").name.contains("Number of bytes scanned")) assert(metrics.contains("output_rows")) assert(metrics.contains("time_elapsed_opening")) assert(metrics.contains("time_elapsed_processing")) assert(metrics.contains("time_elapsed_scanning_until_data")) + Seq( + "scan_io_data_bytes", + "scan_io_metadata_bytes", + "scan_io_footer_reads", + "scan_io_footer_bytes", + "scan_io_object_store_get_calls", + "scan_io_object_store_get_requested_bytes", + "scan_io_object_store_response_bytes_read", + "scan_io_metadata_cache_hits", + "scan_io_metadata_cache_misses").foreach { name => + assert(metrics.contains(name), s"Missing $name. Available: ${metrics.keys}") + } assert( metrics("time_elapsed_scanning_total").value > 0, "time_elapsed_scanning_total should be > 0") assert(metrics("bytes_scanned").value > 0, "bytes_scanned should be > 0") + assert(metrics("scan_io_data_bytes").value > 0) + assert(metrics("scan_io_metadata_bytes").value > 0) + assert(metrics("scan_io_footer_reads").value > 0) + assert(metrics("scan_io_footer_bytes").value > 0) + assert(metrics("scan_io_object_store_get_calls").value == 0) + assert(metrics("scan_io_object_store_get_requested_bytes").value == 0) + assert(metrics("scan_io_object_store_response_bytes_read").value == 0) + assert(metrics("scan_io_metadata_cache_misses").value > 0) assert(metrics("output_rows").value > 0, "output_rows should be > 0") assert(metrics("time_elapsed_opening").value > 0, "time_elapsed_opening should be > 0") assert( diff --git a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometReadBenchmark.scala b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometReadBenchmark.scala index 30090d1d12..1055240cd7 100644 --- a/spark/src/test/scala/org/apache/spark/sql/benchmark/CometReadBenchmark.scala +++ b/spark/src/test/scala/org/apache/spark/sql/benchmark/CometReadBenchmark.scala @@ -34,6 +34,7 @@ import org.apache.spark.TestUtils import org.apache.spark.benchmark.Benchmark import org.apache.spark.sql.{DataFrame, SparkSession} import org.apache.spark.sql.execution.datasources.parquet.VectorizedParquetRecordReader +import org.apache.spark.sql.execution.metric.SQLMetrics import org.apache.spark.sql.types._ import org.apache.spark.sql.vectorized.ColumnVector @@ -47,6 +48,87 @@ import org.apache.comet.{CometConf, WithHdfsCluster} */ class CometReadBaseBenchmark extends CometBenchmarkBase { + /** + * Measure the nine scan I/O accumulators separately from storage work. Run this benchmark with + * `--scan-metric-overhead`; add `--reverse-cases` to reverse the real-job comparison. The first + * two cases isolate driver creation and task-copy/merge costs. The last runs 10,000 Spark tasks + * with zero or nine extra SQL accumulators, including task serialization and scheduler updates. + * This is not an end-to-end scan benchmark: it excludes native counters, JNI traversal, the SQL + * UI's per-node rendering, and storage I/O. + */ + def scanMetricAccumulatorBenchmark(reverseCases: Boolean): Unit = { + val names = Seq( + "scan_io_data_bytes", + "scan_io_metadata_bytes", + "scan_io_footer_reads", + "scan_io_footer_bytes", + "scan_io_object_store_get_calls", + "scan_io_object_store_get_requested_bytes", + "scan_io_object_store_response_bytes_read", + "scan_io_metadata_cache_hits", + "scan_io_metadata_cache_misses") + def createMetrics() = names.map { name => + if (name.endsWith("bytes") || name.endsWith("bytes_read")) { + SQLMetrics.createSizeMetric(spark.sparkContext, name) + } else { + SQLMetrics.createMetric(spark.sparkContext, name) + } + } + + val operators = 1024 + val creation = + new Benchmark("Scan I/O SQL metrics: creation", operators, minNumIters = 5, output = output) + creation.addCase("nine metrics per operator") { _ => + var count = 0 + while (count < operators) { + assert(createMetrics().size == 9) + count += 1 + } + } + creation.run() + + val tasks = 10000 + val driverMetrics = createMetrics() + val updates = new Benchmark( + "Scan I/O SQL metrics: task snapshots", + tasks, + minNumIters = 5, + output = output) + updates.addCase("copy, update and merge nine metrics") { _ => + driverMetrics.foreach(_.reset()) + var task = 0 + while (task < tasks) { + driverMetrics.foreach { driver => + val local = driver.copyAndReset() + local.add(64L) + driver.merge(local) + } + task += 1 + } + assert(driverMetrics.forall(_.value == tasks * 64L)) + } + updates.run() + + val partitions = spark.sparkContext.parallelize(0 until tasks, tasks) + val jobs = new Benchmark( + "Scan I/O SQL metrics: 10000-task Spark job", + tasks, + minNumIters = 3, + output = output) + val cases = if (reverseCases) Seq(9, 0) else Seq(0, 9) + cases.foreach { count => + val metrics = if (count == 0) Seq.empty else createMetrics() + jobs.addCase(s"$count extra SQL accumulators") { _ => + metrics.foreach(_.reset()) + partitions.foreachPartition { _ => + metrics.foreach(_.add(64L)) + } + assert(metrics.forall(_.value == tasks * 64L)) + } + } + jobs.run() + } + def numericScanBenchmark(values: Int, dataType: DataType): Unit = { val sqlBenchmark = new Benchmark(s"SQL Single ${dataType.sql} Column Scan", values, output = output) @@ -316,6 +398,10 @@ class CometReadBaseBenchmark extends CometBenchmarkBase { } override def runCometBenchmark(mainArgs: Array[String]): Unit = { + if (mainArgs.contains("--scan-metric-overhead")) { + scanMetricAccumulatorBenchmark(mainArgs.contains("--reverse-cases")) + return + } runBenchmarkWithTable("Parquet Reader", 1024 * 1024 * 15) { v => Seq( BooleanType,