use crate::arrow::caching_delete_file_loader::CachingDeleteFileLoader;
use crate::io::FileIO;
use crate::runtime::Runtime;
use crate::util::available_parallelism;
const DEFAULT_RANGE_COALESCE_BYTES: u64 = 1024 * 1024;
const DEFAULT_RANGE_FETCH_CONCURRENCY: usize = 10;
const DEFAULT_METADATA_SIZE_HINT: usize = 512 * 1024;
mod file_reader;
mod options;
mod pipeline;
mod positional_deletes;
mod predicate_visitor;
mod projection;
mod row_filter;
pub use file_reader::ArrowFileReader;
pub(crate) use options::ParquetReadOptions;
use predicate_visitor::{CollectFieldIdVisitor, PredicateConverter};
use projection::{add_fallback_field_ids_to_arrow_schema, apply_name_mapping_to_arrow_schema};
pub struct ArrowReaderBuilder {
batch_size: Option<usize>,
file_io: FileIO,
concurrency_limit_data_files: usize,
row_group_filtering_enabled: bool,
row_selection_enabled: bool,
parquet_read_options: ParquetReadOptions,
runtime: Runtime,
}
impl ArrowReaderBuilder {
pub fn new(file_io: FileIO, runtime: Runtime) -> Self {
let num_cpus = available_parallelism().get();
ArrowReaderBuilder {
batch_size: None,
file_io,
concurrency_limit_data_files: num_cpus,
row_group_filtering_enabled: true,
row_selection_enabled: false,
parquet_read_options: ParquetReadOptions::builder().build(),
runtime,
}
}
pub fn with_data_file_concurrency_limit(mut self, val: usize) -> Self {
self.concurrency_limit_data_files = val;
self
}
pub fn with_batch_size(mut self, batch_size: usize) -> Self {
self.batch_size = Some(batch_size);
self
}
pub fn with_row_group_filtering_enabled(mut self, row_group_filtering_enabled: bool) -> Self {
self.row_group_filtering_enabled = row_group_filtering_enabled;
self
}
pub fn with_row_selection_enabled(mut self, row_selection_enabled: bool) -> Self {
self.row_selection_enabled = row_selection_enabled;
self
}
pub fn with_metadata_size_hint(mut self, metadata_size_hint: usize) -> Self {
self.parquet_read_options.metadata_size_hint = Some(metadata_size_hint);
self
}
pub fn with_range_coalesce_bytes(mut self, range_coalesce_bytes: u64) -> Self {
self.parquet_read_options.range_coalesce_bytes = range_coalesce_bytes;
self
}
pub fn with_range_fetch_concurrency(mut self, range_fetch_concurrency: usize) -> Self {
self.parquet_read_options.range_fetch_concurrency = range_fetch_concurrency;
self
}
pub fn build(self) -> ArrowReader {
ArrowReader {
batch_size: self.batch_size,
file_io: self.file_io.clone(),
delete_file_loader: CachingDeleteFileLoader::new(
self.file_io.clone(),
self.concurrency_limit_data_files,
self.runtime.clone(),
),
concurrency_limit_data_files: self.concurrency_limit_data_files,
row_group_filtering_enabled: self.row_group_filtering_enabled,
row_selection_enabled: self.row_selection_enabled,
parquet_read_options: self.parquet_read_options,
}
}
}
#[derive(Clone)]
pub struct ArrowReader {
batch_size: Option<usize>,
file_io: FileIO,
delete_file_loader: CachingDeleteFileLoader,
concurrency_limit_data_files: usize,
row_group_filtering_enabled: bool,
row_selection_enabled: bool,
parquet_read_options: ParquetReadOptions,
}