mod sync_reader;
#[cfg(feature = "async")]
mod async_reader;
use arrow::compute::and;
use arrow::compute::kernels::cmp::{gt, lt};
use arrow_array::cast::AsArray;
use arrow_array::types::Int64Type;
use arrow_array::{ArrayRef, BooleanArray, Int64Array, RecordBatch, StringViewArray};
use bytes::Bytes;
#[cfg(feature = "async")]
use futures::FutureExt;
#[cfg(feature = "async")]
use futures::future::BoxFuture;
use parquet::arrow::arrow_reader::{
ArrowPredicateFn, ArrowReaderOptions, ParquetRecordBatchReaderBuilder, RowFilter,
};
#[cfg(feature = "async")]
use parquet::arrow::async_reader::AsyncFileReader;
use parquet::arrow::{ArrowWriter, ProjectionMask};
use parquet::data_type::AsBytes;
use parquet::file::FOOTER_SIZE;
use parquet::file::metadata::PageIndexPolicy;
#[cfg(feature = "async")]
use parquet::file::metadata::ParquetMetaDataReader;
use parquet::file::metadata::{FooterTail, ParquetMetaData, ParquetOffsetIndex};
use parquet::file::page_index::offset_index::PageLocation;
use parquet::file::properties::WriterProperties;
use parquet::schema::types::SchemaDescriptor;
use std::collections::BTreeMap;
use std::fmt::Display;
use std::ops::Range;
use std::sync::{Arc, LazyLock, Mutex};
fn test_file() -> TestParquetFile {
TestParquetFile::new(TEST_FILE_DATA.clone())
}
fn test_options() -> ArrowReaderOptions {
ArrowReaderOptions::default().with_page_index_policy(PageIndexPolicy::from(true))
}
#[cfg(feature = "async")]
#[derive(Clone)]
pub(crate) struct TestReader {
data: Bytes,
metadata: Option<Arc<ParquetMetaData>>,
requests: Arc<Mutex<Vec<Range<usize>>>>,
}
#[cfg(feature = "async")]
impl TestReader {
pub(crate) fn new(data: Bytes) -> Self {
Self {
data,
metadata: Default::default(),
requests: Default::default(),
}
}
pub(crate) fn requests(&self) -> Arc<Mutex<Vec<Range<usize>>>> {
Arc::clone(&self.requests)
}
}
#[cfg(feature = "async")]
impl AsyncFileReader for TestReader {
fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, parquet::errors::Result<Bytes>> {
self.requests
.lock()
.unwrap()
.push(range.start as usize..range.end as usize);
futures::future::ready(Ok(self
.data
.slice(range.start as usize..range.end as usize)))
.boxed()
}
fn get_metadata<'a>(
&'a mut self,
options: Option<&'a ArrowReaderOptions>,
) -> BoxFuture<'a, parquet::errors::Result<Arc<ParquetMetaData>>> {
let mut metadata_reader = ParquetMetaDataReader::new();
if let Some(options) = options {
metadata_reader = metadata_reader
.with_column_index_policy(options.column_index_policy())
.with_offset_index_policy(options.offset_index_policy());
}
self.metadata = Some(Arc::new(
metadata_reader.parse_and_finish(&self.data).unwrap(),
));
futures::future::ready(Ok(self.metadata.clone().unwrap())).boxed()
}
}
fn filter_b_575_625(schema_descr: &SchemaDescriptor) -> RowFilter {
let predicate = ArrowPredicateFn::new(
ProjectionMask::columns(schema_descr, ["b"]),
|batch: RecordBatch| {
let scalar_575 = Int64Array::new_scalar(575);
let scalar_625 = Int64Array::new_scalar(625);
let column = batch.column(0).as_primitive::<Int64Type>();
and(>(column, &scalar_575)?, <(column, &scalar_625)?)
},
);
RowFilter::new(vec![Box::new(predicate)])
}
fn filter_a_175_b_625(schema_descr: &SchemaDescriptor) -> RowFilter {
let predicate_a = ArrowPredicateFn::new(
ProjectionMask::columns(schema_descr, ["a"]),
|batch: RecordBatch| {
let scalar_175 = Int64Array::new_scalar(175);
let column = batch.column(0).as_primitive::<Int64Type>();
gt(column, &scalar_175)
},
);
let predicate_b = ArrowPredicateFn::new(
ProjectionMask::columns(schema_descr, ["b"]),
|batch: RecordBatch| {
let scalar_625 = Int64Array::new_scalar(625);
let column = batch.column(0).as_primitive::<Int64Type>();
lt(column, &scalar_625)
},
);
RowFilter::new(vec![Box::new(predicate_a), Box::new(predicate_b)])
}
fn filter_b_false(schema_descr: &SchemaDescriptor) -> RowFilter {
let predicate = ArrowPredicateFn::new(
ProjectionMask::columns(schema_descr, ["b"]),
|batch: RecordBatch| {
let result =
BooleanArray::from_iter(std::iter::repeat_n(Some(false), batch.num_rows()));
Ok(result)
},
);
RowFilter::new(vec![Box::new(predicate)])
}
static TEST_FILE_DATA: LazyLock<Bytes> = LazyLock::new(|| {
let a: ArrayRef = Arc::new(Int64Array::from_iter_values(0..400));
let b: ArrayRef = Arc::new(Int64Array::from_iter_values(400..800));
let c: ArrayRef = Arc::new(StringViewArray::from_iter_values((0..400).map(|i| {
if i % 2 == 0 {
format!("string_{i}")
} else {
format!("A string larger than 12 bytes and thus not inlined {i}")
}
})));
let input_batch = RecordBatch::try_from_iter(vec![("a", a), ("b", b), ("c", c)]).unwrap();
let mut output = Vec::new();
let writer_options = WriterProperties::builder()
.set_max_row_group_row_count(Some(200))
.set_data_page_row_count_limit(100)
.build();
let mut writer =
ArrowWriter::try_new(&mut output, input_batch.schema(), Some(writer_options)).unwrap();
let mut row_remain = input_batch.num_rows();
while row_remain > 0 {
let chunk_size = row_remain.min(50);
let chunk = input_batch.slice(input_batch.num_rows() - row_remain, chunk_size);
writer.write(&chunk).unwrap();
row_remain -= chunk_size;
}
writer.close().unwrap();
Bytes::from(output)
});
struct TestParquetFile {
bytes: Bytes,
ops: Arc<OperationLog>,
#[cfg(feature = "async")]
parquet_metadata: Arc<ParquetMetaData>,
}
impl TestParquetFile {
fn new(bytes: Bytes) -> Self {
let builder = ParquetRecordBatchReaderBuilder::try_new_with_options(
bytes.clone(),
ArrowReaderOptions::default().with_page_index_policy(PageIndexPolicy::from(true)),
)
.unwrap();
let parquet_metadata = Arc::clone(builder.metadata());
let offset_index = parquet_metadata
.offset_index()
.expect("Parquet metadata should have a page index");
let row_groups = TestRowGroups::new(&parquet_metadata, offset_index);
let footer_location = bytes.len() - FOOTER_SIZE..bytes.len();
let footer = bytes.slice(footer_location.clone());
let footer: &[u8; FOOTER_SIZE] = footer
.as_bytes()
.try_into() .unwrap();
let footer = FooterTail::try_new(footer).unwrap();
let metadata_len = footer.metadata_length();
let metadata_location = footer_location.start - metadata_len..footer_location.start;
let ops = Arc::new(OperationLog::new(
footer_location,
metadata_location,
row_groups,
));
TestParquetFile {
bytes,
ops,
#[cfg(feature = "async")]
parquet_metadata,
}
}
fn bytes(&self) -> &Bytes {
&self.bytes
}
fn ops(&self) -> &Arc<OperationLog> {
&self.ops
}
#[cfg(feature = "async")]
fn parquet_metadata(&self) -> &Arc<ParquetMetaData> {
&self.parquet_metadata
}
}
#[derive(Debug)]
struct TestColumnChunk {
name: String,
location: Range<usize>,
dictionary_page_location: Option<i64>,
page_locations: Vec<PageLocation>,
}
#[derive(Debug)]
struct TestRowGroup {
columns: BTreeMap<String, TestColumnChunk>,
}
#[derive(Debug)]
struct TestRowGroups {
row_groups: Vec<TestRowGroup>,
}
impl TestRowGroups {
fn new(parquet_metadata: &ParquetMetaData, offset_index: &ParquetOffsetIndex) -> Self {
let row_groups = parquet_metadata
.row_groups()
.iter()
.enumerate()
.map(|(rg_index, rg_meta)| {
let columns = rg_meta
.columns()
.iter()
.enumerate()
.map(|(col_idx, col_meta)| {
let column_name = col_meta.column_descr().name().to_string();
let page_locations = offset_index[rg_index][col_idx].page_locations();
let dictionary_page_location = col_meta.dictionary_page_offset();
let (start_offset, length) = col_meta.byte_range();
let start_offset = start_offset as usize;
let end_offset = start_offset + length as usize;
TestColumnChunk {
name: column_name.clone(),
location: start_offset..end_offset,
dictionary_page_location,
page_locations: page_locations.clone(),
}
})
.map(|test_column_chunk| {
(test_column_chunk.name.clone(), test_column_chunk)
})
.collect::<BTreeMap<_, _>>();
TestRowGroup { columns }
})
.collect();
Self { row_groups }
}
fn iter(&self) -> impl Iterator<Item = &TestRowGroup> {
self.row_groups.iter()
}
}
#[derive(Debug, PartialEq)]
enum PageType {
Data {
data_page_index: usize,
},
Dictionary,
Multi {
dictionary_page: bool,
data_page_indices: Vec<usize>,
},
}
impl Display for PageType {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
PageType::Data { data_page_index } => {
write!(f, "DataPage({data_page_index})")
}
PageType::Dictionary => write!(f, "DictionaryPage"),
PageType::Multi {
dictionary_page,
data_page_indices,
} => {
let dictionary_page = if *dictionary_page {
"dictionary_page: true, "
} else {
""
};
write!(
f,
"MultiPage({dictionary_page}data_pages: {data_page_indices:?})",
)
}
}
}
}
#[derive(Debug)]
struct ReadInfo {
row_group_index: usize,
column_name: String,
range: Range<usize>,
read_type: PageType,
num_requests: usize,
}
impl Display for ReadInfo {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let Self {
row_group_index,
column_name,
range,
read_type,
num_requests,
} = self;
let annotation = if (range.len() / num_requests) < 10 {
" [header]"
} else {
" [data]"
};
write!(
f,
"Row Group {row_group_index}, column '{column_name}': {:15} ({:10}, {:8}){annotation}",
format!("{read_type}"),
format!("{} bytes", range.len()),
format!("{num_requests} requests"),
)
}
}
#[derive(Debug)]
enum LogEntry {
ReadFooter(Range<usize>),
ReadMetadata(Range<usize>),
#[allow(dead_code)]
GetProvidedMetadata,
ReadData(ReadInfo),
#[allow(dead_code)]
ReadMultipleData(Vec<LogEntry>),
Unknown(Range<usize>),
Event(String),
}
impl LogEntry {
fn event(event: impl Into<String>) -> Self {
LogEntry::Event(event.into())
}
fn append_string(&self, output: &mut Vec<String>, indent: usize) {
let indent_str = " ".repeat(indent);
match self {
LogEntry::ReadFooter(range) => {
output.push(format!("{indent_str}Footer: {} bytes", range.len()))
}
LogEntry::ReadMetadata(range) => {
output.push(format!("{indent_str}Metadata: {}", range.len()))
}
LogEntry::GetProvidedMetadata => {
output.push(format!("{indent_str}Get Provided Metadata"))
}
LogEntry::ReadData(read_info) => output.push(format!("{indent_str}{read_info}")),
LogEntry::ReadMultipleData(read_infos) => {
output.push(format!("{indent_str}Read Multi:"));
for read_info in read_infos {
let new_indent = indent + 2;
read_info.append_string(output, new_indent);
}
}
LogEntry::Unknown(range) => {
output.push(format!("{indent_str}UNKNOWN: {range:?} (maybe Page Index)"))
}
LogEntry::Event(event) => output.push(format!("Event: {event}")),
}
}
}
#[derive(Debug)]
struct OperationLog {
ops: Mutex<Vec<LogEntry>>,
footer_location: Range<usize>,
metadata_location: Range<usize>,
row_groups: TestRowGroups,
}
impl OperationLog {
fn new(
footer_location: Range<usize>,
metadata_location: Range<usize>,
row_groups: TestRowGroups,
) -> Self {
OperationLog {
ops: Mutex::new(Vec::new()),
metadata_location,
footer_location,
row_groups,
}
}
fn add_entry(&self, entry: LogEntry) {
let mut ops = self.ops.lock().unwrap();
ops.push(entry);
}
fn add_entry_for_range(&self, range: &Range<usize>) {
self.add_entry(self.entry_for_range(range));
}
#[cfg(feature = "async")]
fn add_entry_for_ranges<'a>(&self, ranges: impl IntoIterator<Item = &'a Range<usize>>) {
let entries = ranges
.into_iter()
.map(|range| self.entry_for_range(range))
.collect::<Vec<_>>();
self.add_entry(LogEntry::ReadMultipleData(entries));
}
fn entry_for_range(&self, range: &Range<usize>) -> LogEntry {
let start = range.start as i64;
let end = range.end as i64;
if self.metadata_location.contains(&range.start)
|| self.metadata_location.contains(&(range.end - 1))
{
return LogEntry::ReadMetadata(range.clone());
}
if self.footer_location.contains(&range.start)
|| self.footer_location.contains(&(range.end - 1))
{
return LogEntry::ReadFooter(range.clone());
}
for (row_group_index, row_group) in self.row_groups.iter().enumerate() {
for (column_name, test_column_chunk) in &row_group.columns {
let page_locations = test_column_chunk.page_locations.iter();
let mut data_page_indices = vec![];
for (data_page_index, page_location) in page_locations.enumerate() {
let page_offset = page_location.offset;
let page_end = page_offset + page_location.compressed_page_size as i64;
if start >= page_offset && end <= page_end {
let read_info = ReadInfo {
row_group_index,
column_name: column_name.clone(),
range: range.clone(),
read_type: PageType::Data { data_page_index },
num_requests: 1,
};
return LogEntry::ReadData(read_info);
}
if start < page_end && end > page_offset {
data_page_indices.push(data_page_index);
}
}
let mut dictionary_page = false;
if let Some(dict_page_offset) = test_column_chunk.dictionary_page_location {
let dict_page_end = dict_page_offset + test_column_chunk.location.len() as i64;
if start >= dict_page_offset && end < dict_page_end {
let read_info = ReadInfo {
row_group_index,
column_name: column_name.clone(),
range: range.clone(),
read_type: PageType::Dictionary,
num_requests: 1,
};
return LogEntry::ReadData(read_info);
}
if start < dict_page_end && end > dict_page_offset {
dictionary_page = true;
}
}
let column_byte_range = &test_column_chunk.location;
if column_byte_range.contains(&range.start)
&& column_byte_range.contains(&(range.end - 1))
{
let read_data_entry = ReadInfo {
row_group_index,
column_name: column_name.clone(),
range: range.clone(),
read_type: PageType::Multi {
data_page_indices,
dictionary_page,
},
num_requests: 1,
};
return LogEntry::ReadData(read_data_entry);
}
}
}
LogEntry::Unknown(range.clone())
}
fn coalesce_entries(&self) {
let mut ops = self.ops.lock().unwrap();
let prev_ops = std::mem::take(&mut *ops);
for entry in prev_ops {
let Some(last) = ops.last_mut() else {
ops.push(entry);
continue;
};
let LogEntry::ReadData(ReadInfo {
row_group_index: last_rg_index,
column_name: last_column_name,
range: last_range,
read_type: last_read_type,
num_requests: last_num_reads,
}) = last
else {
ops.push(entry);
continue;
};
let LogEntry::ReadData(ReadInfo {
row_group_index,
column_name,
range,
read_type,
num_requests: num_reads,
}) = &entry
else {
ops.push(entry);
continue;
};
if *row_group_index != *last_rg_index
|| column_name != last_column_name
|| read_type != last_read_type
|| (range.start > last_range.end)
|| (range.end < last_range.start)
|| range.len() > 10
{
ops.push(entry);
continue;
}
*last_range = last_range.start.min(range.start)..last_range.end.max(range.end);
*last_num_reads += num_reads;
}
}
fn snapshot(&self) -> Vec<String> {
self.coalesce_entries();
let ops = self.ops.lock().unwrap();
let mut actual = vec![];
let indent = 0;
ops.iter()
.for_each(|s| s.append_string(&mut actual, indent));
actual
}
}