use crate::io::{
LogEntry, OperationLog, TestParquetFile, filter_a_175_b_625, filter_b_575_625, filter_b_false,
test_file, test_options,
};
use bytes::Bytes;
use futures::future::BoxFuture;
use futures::{FutureExt, StreamExt};
use parquet::arrow::arrow_reader::{ArrowReaderOptions, RowSelection, RowSelector};
use parquet::arrow::async_reader::AsyncFileReader;
use parquet::arrow::{ParquetRecordBatchStreamBuilder, ProjectionMask};
use parquet::errors::Result;
use parquet::file::metadata::PageIndexPolicy;
use parquet::file::metadata::ParquetMetaData;
use std::ops::Range;
use std::sync::Arc;
#[tokio::test]
async fn test_read_entire_file() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options()).await;
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'a': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
" Row Group 0, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
" Row Group 0, column 'c': MultiPage(dictionary_page: true, data_pages: [0, 1]) (7346 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'a': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
" Row Group 1, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
" Row Group 1, column 'c': MultiPage(dictionary_page: true, data_pages: [0, 1]) (7456 bytes, 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_single_group() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options())
.await
.with_row_groups(vec![1]);
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 1, column 'a': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
" Row Group 1, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
" Row Group 1, column 'c': MultiPage(dictionary_page: true, data_pages: [0, 1]) (7456 bytes, 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_single_column() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options()).await;
let schema_descr = builder.metadata().file_metadata().schema_descr_ptr();
let builder = builder.with_projection(ProjectionMask::columns(&schema_descr, ["b"]));
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_row_selection() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options()).await;
let schema_descr = builder.metadata().file_metadata().schema_descr_ptr();
let builder = builder
.with_projection(ProjectionMask::columns(&schema_descr, ["a", "b"]))
.with_row_selection(RowSelection::from(vec![
RowSelector::skip(175),
RowSelector::select(50),
]));
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'a': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 0, column 'a': DataPage(1) (126 bytes , 1 requests) [data]",
" Row Group 0, column 'b': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 0, column 'b': DataPage(1) (126 bytes , 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'a': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 1, column 'a': DataPage(0) (113 bytes , 1 requests) [data]",
" Row Group 1, column 'b': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 1, column 'b': DataPage(0) (113 bytes , 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_limit() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options()).await;
let schema_descr = builder.metadata().file_metadata().schema_descr_ptr();
let builder = builder
.with_projection(ProjectionMask::columns(&schema_descr, ["a"]))
.with_limit(125);
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'a': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 0, column 'a': DataPage(0) (113 bytes , 1 requests) [data]",
" Row Group 0, column 'a': DataPage(1) (126 bytes , 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_single_row_filter() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options()).await;
let schema_descr = builder.metadata().file_metadata().schema_descr_ptr();
let builder = builder
.with_projection(ProjectionMask::columns(&schema_descr, ["a", "b"]))
.with_row_filter(filter_b_575_625(&schema_descr));
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 0, column 'a': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 0, column 'a': DataPage(1) (126 bytes , 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'a': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 1, column 'a': DataPage(0) (113 bytes , 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_single_row_filter_no_page_index() {
let test_file = test_file();
let options = test_options().with_page_index_policy(PageIndexPolicy::from(false));
let builder = async_builder(&test_file, options).await;
let schema_descr = builder.metadata().file_metadata().schema_descr_ptr();
let builder = builder
.with_projection(ProjectionMask::columns(&schema_descr, ["a", "b"]))
.with_row_filter(filter_b_575_625(&schema_descr));
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 0, column 'a': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'a': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_multiple_row_filter() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options()).await;
let schema_descr = builder.metadata().file_metadata().schema_descr_ptr();
let builder = builder
.with_projection(ProjectionMask::columns(&schema_descr, ["c"]))
.with_row_filter(filter_a_175_b_625(&schema_descr));
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'a': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 0, column 'b': DictionaryPage (1617 bytes, 1 requests) [data]",
" Row Group 0, column 'b': DataPage(1) (126 bytes , 1 requests) [data]",
"Read Multi:",
" Row Group 0, column 'c': DictionaryPage (7107 bytes, 1 requests) [data]",
" Row Group 0, column 'c': DataPage(1) (126 bytes , 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'a': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'c': DictionaryPage (7217 bytes, 1 requests) [data]",
" Row Group 1, column 'c': DataPage(0) (113 bytes , 1 requests) [data]",
]
"#);
}
#[tokio::test]
async fn test_read_single_row_filter_all() {
let test_file = test_file();
let builder = async_builder(&test_file, test_options()).await;
let schema_descr = builder.metadata().file_metadata().schema_descr_ptr();
let builder = builder
.with_projection(ProjectionMask::columns(&schema_descr, ["a", "b"]))
.with_row_filter(filter_b_false(&schema_descr));
insta::assert_debug_snapshot!(run(
&test_file,
builder).await, @r#"
[
"Get Provided Metadata",
"Event: Builder Configured",
"Event: Reader Built",
"Read Multi:",
" Row Group 0, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
"Read Multi:",
" Row Group 1, column 'b': MultiPage(dictionary_page: true, data_pages: [0, 1]) (1856 bytes, 1 requests) [data]",
]
"#);
}
async fn async_builder(
test_file: &TestParquetFile,
options: ArrowReaderOptions,
) -> ParquetRecordBatchStreamBuilder<RecordingAsyncFileReader> {
let parquet_meta_data = if options.offset_index_policy() != PageIndexPolicy::Skip
|| options.column_index_policy() != PageIndexPolicy::Skip
{
Arc::clone(test_file.parquet_metadata())
} else {
let metadata = test_file
.parquet_metadata()
.as_ref()
.clone()
.into_builder()
.set_column_index(None)
.set_offset_index(None)
.build();
Arc::new(metadata)
};
let reader = RecordingAsyncFileReader {
bytes: test_file.bytes().clone(),
ops: Arc::clone(test_file.ops()),
parquet_meta_data,
};
ParquetRecordBatchStreamBuilder::new_with_options(reader, options)
.await
.unwrap()
}
async fn run(
test_file: &TestParquetFile,
builder: ParquetRecordBatchStreamBuilder<RecordingAsyncFileReader>,
) -> Vec<String> {
let ops = test_file.ops();
ops.add_entry(LogEntry::event("Builder Configured"));
let mut stream = builder.build().unwrap();
ops.add_entry(LogEntry::event("Reader Built"));
while let Some(batch) = stream.next().await {
match batch {
Ok(_) => {}
Err(e) => panic!("Error reading batch: {e}"),
}
}
ops.snapshot()
}
struct RecordingAsyncFileReader {
bytes: Bytes,
ops: Arc<OperationLog>,
parquet_meta_data: Arc<ParquetMetaData>,
}
impl AsyncFileReader for RecordingAsyncFileReader {
fn get_bytes(&mut self, range: Range<u64>) -> BoxFuture<'_, parquet::errors::Result<Bytes>> {
let ops = Arc::clone(&self.ops);
let data = self
.bytes
.slice(range.start as usize..range.end as usize)
.clone();
let logged_range = Range {
start: range.start as usize,
end: range.end as usize,
};
async move {
ops.add_entry_for_range(&logged_range);
Ok(data)
}
.boxed()
}
fn get_byte_ranges(&mut self, ranges: Vec<Range<u64>>) -> BoxFuture<'_, Result<Vec<Bytes>>> {
let ops = Arc::clone(&self.ops);
let datas = ranges
.iter()
.map(|range| {
self.bytes
.slice(range.start as usize..range.end as usize)
.clone()
})
.collect::<Vec<_>>();
let logged_ranges = ranges
.into_iter()
.map(|r| Range {
start: r.start as usize,
end: r.end as usize,
})
.collect::<Vec<_>>();
async move {
ops.add_entry_for_ranges(&logged_ranges);
Ok(datas)
}
.boxed()
}
fn get_metadata<'a>(
&'a mut self,
_options: Option<&'a ArrowReaderOptions>,
) -> BoxFuture<'a, Result<Arc<ParquetMetaData>>> {
let ops = Arc::clone(&self.ops);
let parquet_meta_data = Arc::clone(&self.parquet_meta_data);
async move {
ops.add_entry(LogEntry::GetProvidedMetadata);
Ok(parquet_meta_data)
}
.boxed()
}
}