Skip to content

Commit e11e013

Browse files
committed
build fix and tests
1 parent 01fa26a commit e11e013

3 files changed

Lines changed: 314 additions & 46 deletions

File tree

‎datafusion/datasource-parquet/src/eager_pruning.rs‎

Lines changed: 277 additions & 46 deletions
Original file line numberDiff line numberDiff line change
@@ -746,68 +746,117 @@ mod tests {
746746
use crate::source::ParquetSource;
747747

748748
use arrow::array::{Int64Array, RecordBatch};
749-
use arrow::datatypes::{DataType, Field};
749+
use arrow::datatypes::{DataType, Field, SchemaRef};
750750
use bytes::{BufMut, BytesMut};
751+
use datafusion_datasource::{FileRange, TableSchemaBuilder};
751752
use datafusion_execution::object_store::ObjectStoreUrl;
752753
use datafusion_expr::{Expr, col, lit};
753754
use datafusion_physical_expr::planner::logical2physical;
754755
use object_store::memory::InMemory;
755756
use object_store::path::Path;
756757
use object_store::{ObjectStore, ObjectStoreExt};
757758
use parquet::arrow::ArrowWriter;
759+
use parquet::arrow::arrow_reader::{RowSelection, RowSelector};
758760
use parquet::file::properties::WriterProperties;
759761

760-
/// Prunes a file whose column `a` holds `0..400` in 4 row groups of 100
761-
/// rows, using `predicate`
762-
async fn prune(
763-
level: EagerParquetPruning,
764-
file_limit: usize,
765-
predicate: Expr,
766-
) -> (FileScanConfig, EagerPruningSummary) {
767-
let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int64, false)]));
768-
let batch = RecordBatch::try_new(
769-
Arc::clone(&schema),
770-
vec![Arc::new(Int64Array::from_iter_values(0..400))],
771-
)
772-
.unwrap();
773-
let props = WriterProperties::builder()
774-
.set_max_row_group_row_count(Some(100))
775-
.build();
776-
let mut out = BytesMut::new().writer();
777-
{
778-
let mut writer =
779-
ArrowWriter::try_new(&mut out, Arc::clone(&schema), Some(props)).unwrap();
780-
writer.write(&batch).unwrap();
781-
writer.finish().unwrap();
782-
}
783-
let data = out.into_inner().freeze();
784-
let size = data.len() as u64;
785-
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
786-
store
787-
.put(&Path::from("test.parquet"), data.into())
788-
.await
762+
/// A parquet file whose column `a` holds `0..400` in 4 row groups of 100
763+
/// rows, in an in memory object store
764+
struct TestFile {
765+
store: Arc<dyn ObjectStore>,
766+
schema: SchemaRef,
767+
size: u64,
768+
}
769+
770+
impl TestFile {
771+
async fn new() -> Self {
772+
let schema =
773+
Arc::new(Schema::new(vec![Field::new("a", DataType::Int64, false)]));
774+
let batch = RecordBatch::try_new(
775+
Arc::clone(&schema),
776+
vec![Arc::new(Int64Array::from_iter_values(0..400))],
777+
)
789778
.unwrap();
779+
let props = WriterProperties::builder()
780+
.set_max_row_group_row_count(Some(100))
781+
.build();
782+
let mut out = BytesMut::new().writer();
783+
{
784+
let mut writer =
785+
ArrowWriter::try_new(&mut out, Arc::clone(&schema), Some(props))
786+
.unwrap();
787+
writer.write(&batch).unwrap();
788+
writer.finish().unwrap();
789+
}
790+
let data = out.into_inner().freeze();
791+
let size = data.len() as u64;
792+
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
793+
store
794+
.put(&Path::from("test.parquet"), data.into())
795+
.await
796+
.unwrap();
797+
Self {
798+
store,
799+
schema,
800+
size,
801+
}
802+
}
790803

791-
let table_schema = TableSchema::from(&schema);
792-
let conf = FileScanConfigBuilder::new(
793-
ObjectStoreUrl::local_filesystem(),
794-
Arc::new(ParquetSource::new(table_schema.clone())),
795-
)
796-
.with_file(PartitionedFile::new("test.parquet", size))
797-
.build();
804+
fn file(&self) -> PartitionedFile {
805+
PartitionedFile::new("test.parquet", self.size)
806+
}
807+
808+
fn table_schema(&self) -> TableSchema {
809+
TableSchema::from(&self.schema)
810+
}
798811

812+
/// Prunes `files` as a single file group
813+
async fn prune_files(
814+
&self,
815+
options: &TableParquetOptions,
816+
table_schema: &TableSchema,
817+
predicate: Expr,
818+
files: Vec<PartitionedFile>,
819+
) -> (FileScanConfig, EagerPruningSummary) {
820+
let conf = FileScanConfigBuilder::new(
821+
ObjectStoreUrl::local_filesystem(),
822+
Arc::new(ParquetSource::new(table_schema.clone())),
823+
)
824+
.with_file_group(FileGroup::new(files))
825+
.build();
826+
827+
let filters = [logical2physical(&predicate, &self.schema)];
828+
let pruner = EagerPruner::try_new(
829+
options,
830+
&filters,
831+
Arc::new(DefaultParquetFileReaderFactory::new(Arc::clone(
832+
&self.store,
833+
))),
834+
4,
835+
)
836+
.unwrap();
837+
pruner.prune(table_schema, conf).await.unwrap()
838+
}
839+
}
840+
841+
fn eager_options(level: EagerParquetPruning) -> TableParquetOptions {
799842
let mut options = TableParquetOptions::default();
800843
options.global.eager_pruning = level;
844+
options
845+
}
846+
847+
/// Prunes the test file using `predicate`
848+
async fn prune(
849+
level: EagerParquetPruning,
850+
file_limit: usize,
851+
predicate: Expr,
852+
) -> (FileScanConfig, EagerPruningSummary) {
853+
let test_file = TestFile::new().await;
854+
let mut options = eager_options(level);
801855
options.global.eager_pruning_file_limit = file_limit;
802-
let filters = [logical2physical(&predicate, &schema)];
803-
let pruner = EagerPruner::try_new(
804-
&options,
805-
&filters,
806-
Arc::new(DefaultParquetFileReaderFactory::new(store)),
807-
4,
808-
)
809-
.unwrap();
810-
pruner.prune(&table_schema, conf).await.unwrap()
856+
let files = vec![test_file.file()];
857+
test_file
858+
.prune_files(&options, &test_file.table_schema(), predicate, files)
859+
.await
811860
}
812861

813862
#[tokio::test]
@@ -890,6 +939,188 @@ mod tests {
890939
assert!(!file.extensions.contains::<ParquetAccessPlan>());
891940
}
892941

942+
#[tokio::test]
943+
async fn skips_encrypted_files() {
944+
let test_file = TestFile::new().await;
945+
let mut options = eager_options(EagerParquetPruning::RowGroups);
946+
options.crypto.factory_id = Some("test_factory".to_string());
947+
let files = vec![test_file.file()];
948+
949+
let (conf, summary) = test_file
950+
.prune_files(
951+
&options,
952+
&test_file.table_schema(),
953+
col("a").lt(lit(250i64)),
954+
files,
955+
)
956+
.await;
957+
958+
assert_eq!(
959+
summary,
960+
EagerPruningSummary::Skipped {
961+
level: EagerParquetPruning::RowGroups,
962+
reason: "encrypted files are not supported".to_string(),
963+
}
964+
);
965+
assert!(
966+
!conf.file_groups[0].files()[0]
967+
.extensions
968+
.contains::<ParquetAccessPlan>()
969+
);
970+
}
971+
972+
#[tokio::test]
973+
async fn skips_scans_with_virtual_columns() {
974+
let test_file = TestFile::new().await;
975+
let table_schema = TableSchemaBuilder::new(Arc::clone(&test_file.schema))
976+
.with_virtual_columns(vec![Arc::new(Field::new(
977+
"row_number",
978+
DataType::Int64,
979+
true,
980+
))])
981+
.build();
982+
let files = vec![test_file.file()];
983+
984+
let (_, summary) = test_file
985+
.prune_files(
986+
&eager_options(EagerParquetPruning::RowGroups),
987+
&table_schema,
988+
col("a").lt(lit(250i64)),
989+
files,
990+
)
991+
.await;
992+
993+
assert_eq!(
994+
summary,
995+
EagerPruningSummary::Skipped {
996+
level: EagerParquetPruning::RowGroups,
997+
reason: "virtual columns are not supported".to_string(),
998+
}
999+
);
1000+
}
1001+
1002+
#[tokio::test]
1003+
async fn keeps_files_that_cannot_be_read() {
1004+
let test_file = TestFile::new().await;
1005+
// The first file does not exist, so its metadata cannot be read
1006+
let files = vec![
1007+
PartitionedFile::new("missing.parquet", 1024),
1008+
test_file.file(),
1009+
];
1010+
1011+
let (conf, summary) = test_file
1012+
.prune_files(
1013+
&eager_options(EagerParquetPruning::RowGroups),
1014+
&test_file.table_schema(),
1015+
col("a").lt(lit(250i64)),
1016+
files,
1017+
)
1018+
.await;
1019+
1020+
let EagerPruningSummary::Pruned(stats) = summary else {
1021+
panic!("expected eager pruning");
1022+
};
1023+
assert_eq!(stats.files, 2);
1024+
assert_eq!(stats.files_not_evaluated, 1);
1025+
assert_eq!(stats.rows_pruned, 100);
1026+
1027+
// Both files are kept, the one that could not be read unchanged
1028+
let files = conf.file_groups[0].files();
1029+
assert_eq!(files.len(), 2);
1030+
assert!(!files[0].extensions.contains::<ParquetAccessPlan>());
1031+
assert!(files[0].statistics.is_none());
1032+
assert!(files[1].extensions.contains::<ParquetAccessPlan>());
1033+
// The row count of one file is unknown, so it is unknown for the scan
1034+
assert_eq!(conf.statistics().num_rows, Precision::Absent);
1035+
}
1036+
1037+
#[tokio::test]
1038+
async fn skips_files_with_a_row_selection() {
1039+
let test_file = TestFile::new().await;
1040+
let selection = ParquetRowSelection::new(RowSelection::from(vec![
1041+
RowSelector::select(200),
1042+
RowSelector::skip(200),
1043+
]));
1044+
let files = vec![test_file.file().with_extension(selection)];
1045+
1046+
let (conf, summary) = test_file
1047+
.prune_files(
1048+
&eager_options(EagerParquetPruning::RowGroups),
1049+
&test_file.table_schema(),
1050+
col("a").lt(lit(250i64)),
1051+
files,
1052+
)
1053+
.await;
1054+
1055+
let EagerPruningSummary::Pruned(stats) = summary else {
1056+
panic!("expected eager pruning");
1057+
};
1058+
assert_eq!(stats.files_not_evaluated, 1);
1059+
assert_eq!(stats.rows_pruned, 0);
1060+
let file = &conf.file_groups[0].files()[0];
1061+
assert!(!file.extensions.contains::<ParquetAccessPlan>());
1062+
assert!(file.extensions.contains::<ParquetRowSelection>());
1063+
}
1064+
1065+
#[tokio::test]
1066+
async fn prunes_row_groups_within_a_file_range() {
1067+
let test_file = TestFile::new().await;
1068+
// A range that covers the first part of the file only
1069+
let mut file = test_file.file();
1070+
file.range = Some(FileRange {
1071+
start: 0,
1072+
end: (test_file.size / 2) as i64,
1073+
});
1074+
let files = vec![file];
1075+
1076+
let (conf, summary) = test_file
1077+
.prune_files(
1078+
&eager_options(EagerParquetPruning::RowGroups),
1079+
&test_file.table_schema(),
1080+
col("a").lt(lit(150i64)),
1081+
files,
1082+
)
1083+
.await;
1084+
1085+
let EagerPruningSummary::Pruned(stats) = summary else {
1086+
panic!("expected eager pruning");
1087+
};
1088+
// Row groups outside the range are not part of the evaluated rows
1089+
assert!(stats.rows < 400, "{stats:?}");
1090+
assert_eq!(stats.rows_pruned, stats.rows - 200);
1091+
let access_plan = conf.file_groups[0].files()[0]
1092+
.extensions
1093+
.get::<ParquetAccessPlan>()
1094+
.unwrap();
1095+
assert_eq!(access_plan.row_group_indexes(), vec![0, 1]);
1096+
}
1097+
1098+
#[tokio::test]
1099+
async fn with_int96_coercion_and_file_arrow_schema() {
1100+
let test_file = TestFile::new().await;
1101+
let mut options = eager_options(EagerParquetPruning::RowGroups);
1102+
options.global.coerce_int96 = Some("ms".to_string());
1103+
options.global.coerce_int96_tz = Some("UTC".to_string());
1104+
let mut file = test_file.file();
1105+
file.arrow_schema = Some(Arc::clone(&test_file.schema));
1106+
let files = vec![file];
1107+
1108+
let (conf, summary) = test_file
1109+
.prune_files(
1110+
&options,
1111+
&test_file.table_schema(),
1112+
col("a").lt(lit(250i64)),
1113+
files,
1114+
)
1115+
.await;
1116+
1117+
let EagerPruningSummary::Pruned(stats) = summary else {
1118+
panic!("expected eager pruning");
1119+
};
1120+
assert_eq!(stats.rows_pruned, 100);
1121+
assert_eq!(conf.statistics().num_rows, Precision::Inexact(300));
1122+
}
1123+
8931124
#[test]
8941125
fn disabled_or_without_filters() {
8951126
let reader_factory: Arc<dyn ParquetFileReaderFactory> = Arc::new(

‎datafusion/datasource-parquet/src/source.rs‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1155,6 +1155,8 @@ impl FileSource for ParquetSource {
11551155
encryption_factory: _,
11561156
reverse_row_groups,
11571157
sort_order_for_reorder,
1158+
// Only used to display the outcome of eager pruning in `EXPLAIN`.
1159+
eager_pruning_summary: _,
11581160
} = self;
11591161

11601162
if schema_provider.is_some() {

0 commit comments

Comments
 (0)