Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
12 changes: 5 additions & 7 deletions datafusion/datasource-parquet/src/bloom_filter.rs
Original file line number Diff line number Diff line change
Expand Up @@ -250,7 +250,7 @@ mod tests {
use datafusion_physical_expr::planner::logical2physical;
use datafusion_physical_plan::metrics::ExecutionPlanMetricsSet;
use datafusion_pruning::PruningPredicate;
use object_store::ObjectStoreExt;
use object_store::{ObjectStore, ObjectStoreExt};
use parquet::arrow::ArrowWriter;
use parquet::arrow::ParquetRecordBatchStreamBuilder;
use parquet::arrow::async_reader::ParquetObjectReader;
Expand Down Expand Up @@ -644,17 +644,15 @@ mod tests {
let metrics = ExecutionPlanMetricsSet::new();
let file_metrics =
ParquetFileMetrics::new(0, object_meta.location.as_ref(), &metrics);
let store: Arc<dyn ObjectStore> = Arc::new(in_memory);
let inner =
ParquetObjectReader::new(Arc::new(in_memory), object_meta.location.clone())
ParquetObjectReader::new(Arc::clone(&store), object_meta.location.clone())
.with_file_size(object_meta.size);

let partitioned_file = PartitionedFile::new_from_meta(object_meta);

let reader = ParquetFileReader {
inner,
file_metrics: file_metrics.clone(),
partitioned_file,
};
let reader =
ParquetFileReader::new(file_metrics.clone(), store, inner, partitioned_file);

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

ParquetFileReader now has a real new() method

let mut builder = ParquetRecordBatchStreamBuilder::new(reader).await.unwrap();

let access_plan = ParquetAccessPlan::new_all(builder.metadata().num_row_groups());
Expand Down
156 changes: 75 additions & 81 deletions datafusion/datasource-parquet/src/reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -86,62 +86,6 @@ impl DefaultParquetFileReaderFactory {
}
}

/// Implements [`AsyncFileReader`] for a parquet file in object storage.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this was all duplicated with CachedParquetFilerReader and I have unified the two

///
/// This implementation uses the [`ParquetObjectReader`] to read data from the
/// object store on demand, as required, tracking the number of bytes read.
///
/// This implementation does not coalesce I/O operations or cache bytes. Such
/// optimizations can be done either at the object store level or by providing a
/// custom implementation of [`ParquetFileReaderFactory`].
pub struct ParquetFileReader {
pub file_metrics: ParquetFileMetrics,
pub inner: ParquetObjectReader,
pub partitioned_file: PartitionedFile,
}

impl AsyncFileReader for ParquetFileReader {
fn get_bytes(
&mut self,
range: Range<u64>,
) -> BoxFuture<'_, parquet::errors::Result<Bytes>> {
let bytes_scanned = range.end - range.start;
self.file_metrics.bytes_scanned.add(bytes_scanned as usize);
self.inner.get_bytes(range)
}

fn get_byte_ranges(
&mut self,
ranges: Vec<Range<u64>>,
) -> BoxFuture<'_, parquet::errors::Result<Vec<Bytes>>>
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)
}

fn get_metadata<'a>(
&'a mut self,
options: Option<&'a ArrowReaderOptions>,
) -> BoxFuture<'a, parquet::errors::Result<Arc<ParquetMetaData>>> {
self.inner.get_metadata(options)
}
}

impl Drop for ParquetFileReader {
fn drop(&mut self) {
self.file_metrics
.scan_efficiency_ratio
.add_part(self.file_metrics.bytes_scanned.value());
// Multiple ParquetFileReaders may run, so we set_total to avoid adding the total multiple times
self.file_metrics
.scan_efficiency_ratio
.set_total(self.partitioned_file.object_meta.size as usize);
}
}

impl ParquetFileReaderFactory for DefaultParquetFileReaderFactory {
fn create_reader(
&self,
Expand All @@ -166,18 +110,21 @@ impl ParquetFileReaderFactory for DefaultParquetFileReaderFactory {
inner = inner.with_footer_size_hint(hint)
};

Ok(Box::new(ParquetFileReader {
inner,
let reader = ParquetFileReader::new(
file_metrics,
Arc::clone(&self.store),
inner,
partitioned_file,
}))
)
.with_metadata_hint(metadata_size_hint);
Ok(Box::new(reader))
}
}

/// Implementation of [`ParquetFileReaderFactory`] supporting the caching of footer and page
/// metadata. Reads and updates the [`FileMetadataCache`] with the [`ParquetMetaData`] data.
///
/// [`CachedParquetFileReader::get_metadata`] forwards the [`parquet::file::metadata::PageIndexPolicy`] from
/// [`ParquetFileReader::get_metadata`] forwards the [`parquet::file::metadata::PageIndexPolicy`] from
/// [`ArrowReaderOptions`] to [`DFParquetMetadata::fetch_metadata`], so callers such as the
/// parquet opener can skip page-index I/O during the initial metadata load.
#[derive(Debug)]
Expand Down Expand Up @@ -223,50 +170,97 @@ impl ParquetFileReaderFactory for CachedParquetFileReaderFactory {
inner = inner.with_footer_size_hint(hint)
};

Ok(Box::new(CachedParquetFileReader::new(
let reader = ParquetFileReader::new(
file_metrics,
Arc::clone(&self.store),
inner,
partitioned_file,
Arc::clone(&self.metadata_cache),
metadata_size_hint,
)))
)
.with_metadata_hint(metadata_size_hint)
.with_metadata_cache(Some(Arc::clone(&self.metadata_cache)));

Ok(Box::new(reader))
}
}

/// Implements [`AsyncFileReader`] for a Parquet file in object storage. Reads the file metadata
/// from the [`FileMetadataCache`], if available, otherwise reads it directly from the file and then
/// updates the cache.
pub struct CachedParquetFileReader {
pub file_metrics: ParquetFileMetrics,
/// Implements [`AsyncFileReader`] for a parquet file in object storage.
///
/// This implementation uses the [`ParquetObjectReader`] to read data from the
/// object store on demand, as required, tracking the number of bytes read via
/// [`ParquetFileMetrics`].
///
/// When configured via [`Self::with_metadata_cache`], [`Self::get_metadata`]
/// reads footer and page metadata from the cache when available and populates
/// the cache otherwise. Without a cache, metadata is fetched fresh on every call.
///
/// # Notes
///
/// This implementation does not coalesce I/O operations or cache bytes. Such
/// optimizations can be done either at the object store level or by providing
/// a custom implementation of [`ParquetFileReaderFactory`].
pub struct ParquetFileReader {
file_metrics: ParquetFileMetrics,
store: Arc<dyn ObjectStore>,
pub inner: ParquetObjectReader,
inner: ParquetObjectReader,
partitioned_file: PartitionedFile,
metadata_cache: Arc<FileMetadataCache>,
metadata_cache: Option<Arc<FileMetadataCache>>,
metadata_size_hint: Option<usize>,
}

impl CachedParquetFileReader {
pub fn new(
impl ParquetFileReader {
/// Create a new `ParquetFileReader`.
///
/// By default the reader has no [`FileMetadataCache`] and no metadata
/// size hint, so metadata is fetched fresh on every call (as
/// [`DefaultParquetFileReaderFactory`] does). Use
/// [`Self::with_metadata_cache`] to read and populate a cache (as
/// [`CachedParquetFileReaderFactory`] does), and
/// [`Self::with_metadata_hint`] to set the size hint.
pub(crate) fn new(
file_metrics: ParquetFileMetrics,
store: Arc<dyn ObjectStore>,
inner: ParquetObjectReader,
partitioned_file: PartitionedFile,
metadata_cache: Arc<FileMetadataCache>,
metadata_size_hint: Option<usize>,
) -> Self {
Self {
file_metrics,
store,
inner,
partitioned_file,
metadata_cache,
metadata_size_hint,
metadata_cache: None,
metadata_size_hint: None,
}
}

/// Returns the metrics tracked while reading this file.
pub fn file_metrics(&self) -> &ParquetFileMetrics {
&self.file_metrics
}

/// Returns the file this reader is reading.
pub fn partitioned_file(&self) -> &PartitionedFile {
&self.partitioned_file
}

/// Set the [`FileMetadataCache`] for this reader
pub fn with_metadata_cache(
mut self,
metadata_cache: Option<Arc<FileMetadataCache>>,
) -> Self {
self.metadata_cache = metadata_cache;
self
}

/// Set the metadata size hint for this reader.
///
/// See [`DFParquetMetadata::with_metadata_size_hint`] for more details.
pub fn with_metadata_hint(mut self, metadata_size_hint: Option<usize>) -> Self {
self.metadata_size_hint = metadata_size_hint;
self
}
}

impl AsyncFileReader for CachedParquetFileReader {
impl AsyncFileReader for ParquetFileReader {
fn get_bytes(
&mut self,
range: Range<u64>,
Expand All @@ -293,7 +287,7 @@ impl AsyncFileReader for CachedParquetFileReader {
options: Option<&'a ArrowReaderOptions>,
) -> BoxFuture<'a, parquet::errors::Result<Arc<ParquetMetaData>>> {
let object_meta = self.partitioned_file.object_meta.clone();
let metadata_cache = Arc::clone(&self.metadata_cache);
let metadata_cache = self.metadata_cache.clone();

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this is now Option<Arc<..>> so it needs a clone (rather than being able to call Arc::clone directly


async move {
#[cfg(feature = "parquet_encryption")]
Expand All @@ -308,7 +302,7 @@ impl AsyncFileReader for CachedParquetFileReader {

DFParquetMetadata::new(&self.store, &object_meta)
.with_decryption_properties(file_decryption_properties)
.with_file_metadata_cache(Some(Arc::clone(&metadata_cache)))
.with_file_metadata_cache(metadata_cache)
.with_metadata_size_hint(self.metadata_size_hint)
.with_page_index_policy(page_index_policy)
.fetch_metadata()
Expand All @@ -324,7 +318,7 @@ impl AsyncFileReader for CachedParquetFileReader {
}
}

impl Drop for CachedParquetFileReader {
impl Drop for ParquetFileReader {
fn drop(&mut self) {
self.file_metrics
.scan_efficiency_ratio
Expand Down
48 changes: 48 additions & 0 deletions docs/source/library-user-guide/upgrading/55.0.0.md
Original file line number Diff line number Diff line change
Expand Up @@ -948,3 +948,51 @@ See [PR #23827](https://github.com/apache/datafusion/pull/23827) for details.
The Minimum Supported Rust Version (MSRV) has been updated to [`1.94.0`].

[`1.94.0`]: https://releases.rs/docs/1.94.0/

### `CachedParquetFileReader` removed; `ParquetFileReader` fields are now private

`CachedParquetFileReader` duplicated `ParquetFileReader` and has been removed;
`ParquetFileReader`'s fields are also now private, with
`file_metrics()` and `partitioned_file()` accessors added for the two that
were previously public.

**Who is affected:**

- Code that names the `CachedParquetFileReader` type.
- Code that constructs a `ParquetFileReader` directly via a struct literal, or
reads/writes its fields.

**Migration guide:**

`ParquetFileReader::new` is no longer public; build a reader through
`ParquetFileReaderFactory::create_reader` (via `DefaultParquetFileReaderFactory`
or `CachedParquetFileReaderFactory`) instead of constructing one directly:

```rust,ignore
// Before
let inner = ParquetObjectReader::new(Arc::clone(&store), location).with_file_size(size);
let reader = CachedParquetFileReader::new(
file_metrics,
store,
inner,
partitioned_file,
metadata_cache,
metadata_size_hint,
);

// After
let reader = CachedParquetFileReaderFactory::new(store, metadata_cache)
.create_reader(partition_index, partitioned_file, metadata_size_hint, &metrics)?;
```

Replace field access with the new accessor methods:

```rust,ignore
// Before
let bytes_scanned = reader.file_metrics.bytes_scanned.value();
let location = &reader.partitioned_file.object_meta.location;

// After
let bytes_scanned = reader.file_metrics().bytes_scanned.value();
let location = &reader.partitioned_file().object_meta.location;
```
Loading