Skip to content

Commit 81c58bf

Browse files
committed
Don't clone the schema in logical2physical
1 parent d0090dd commit 81c58bf

11 files changed

Lines changed: 210 additions & 168 deletions

File tree

datafusion/core/src/datasource/physical_plan/parquet.rs

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -159,7 +159,7 @@ mod tests {
159159
let predicate = self
160160
.predicate
161161
.as_ref()
162-
.map(|p| logical2physical(p, &table_schema));
162+
.map(|p| logical2physical(p, Arc::clone(&table_schema)));
163163

164164
let mut source = ParquetSource::new(table_schema);
165165
if let Some(predicate) = predicate {

datafusion/datasource-parquet/benches/parquet_nested_filter_pushdown.rs

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -65,7 +65,7 @@ fn parquet_nested_filter_pushdown(c: &mut Criterion) {
6565

6666
group.bench_function("no_pushdown", |b| {
6767
let file_schema = setup_reader(&dataset_path);
68-
let predicate = logical2physical(&create_predicate(), &file_schema);
68+
let predicate = logical2physical(&create_predicate(), file_schema);
6969
b.iter(|| {
7070
let matched = scan_with_predicate(&dataset_path, &predicate, false)
7171
.expect("baseline parquet scan with filter succeeded");
@@ -75,7 +75,7 @@ fn parquet_nested_filter_pushdown(c: &mut Criterion) {
7575

7676
group.bench_function("with_pushdown", |b| {
7777
let file_schema = setup_reader(&dataset_path);
78-
let predicate = logical2physical(&create_predicate(), &file_schema);
78+
let predicate = logical2physical(&create_predicate(), file_schema);
7979
b.iter(|| {
8080
let matched = scan_with_predicate(&dataset_path, &predicate, true)
8181
.expect("pushdown parquet scan with filter succeeded");

datafusion/datasource-parquet/src/opener.rs

Lines changed: 20 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -1407,7 +1407,7 @@ mod test {
14071407

14081408
// A filter on "a" should not exclude any rows even if it matches the data
14091409
let expr = col("a").eq(lit(1));
1410-
let predicate = logical2physical(&expr, &schema);
1410+
let predicate = logical2physical(&expr, Arc::clone(&schema));
14111411
let opener = make_opener(predicate);
14121412
let stream = opener.open(file.clone()).unwrap().await.unwrap();
14131413
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1416,7 +1416,7 @@ mod test {
14161416

14171417
// A filter on `b = 5.0` should exclude all rows
14181418
let expr = col("b").eq(lit(ScalarValue::Float32(Some(5.0))));
1419-
let predicate = logical2physical(&expr, &schema);
1419+
let predicate = logical2physical(&expr, Arc::clone(&schema));
14201420
let opener = make_opener(predicate);
14211421
let stream = opener.open(file).unwrap().await.unwrap();
14221422
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1462,7 +1462,8 @@ mod test {
14621462
let expr = col("part").eq(lit(1));
14631463
// Mark the expression as dynamic even if it's not to force partition pruning to happen
14641464
// Otherwise we assume it already happened at the planning stage and won't re-do the work here
1465-
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
1465+
let predicate =
1466+
make_dynamic_expr(logical2physical(&expr, Arc::clone(&table_schema)));
14661467
let opener = make_opener(predicate);
14671468
let stream = opener.open(file.clone()).unwrap().await.unwrap();
14681469
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1473,7 +1474,7 @@ mod test {
14731474
let expr = col("part").eq(lit(2));
14741475
// Mark the expression as dynamic even if it's not to force partition pruning to happen
14751476
// Otherwise we assume it already happened at the planning stage and won't re-do the work here
1476-
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
1477+
let predicate = make_dynamic_expr(logical2physical(&expr, table_schema));
14771478
let opener = make_opener(predicate);
14781479
let stream = opener.open(file).unwrap().await.unwrap();
14791480
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1529,7 +1530,7 @@ mod test {
15291530

15301531
// Filter should match the partition value and file statistics
15311532
let expr = col("part").eq(lit(1)).and(col("b").eq(lit(1.0)));
1532-
let predicate = logical2physical(&expr, &table_schema);
1533+
let predicate = logical2physical(&expr, Arc::clone(&table_schema));
15331534
let opener = make_opener(predicate);
15341535
let stream = opener.open(file.clone()).unwrap().await.unwrap();
15351536
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1538,7 +1539,7 @@ mod test {
15381539

15391540
// Should prune based on partition value but not file statistics
15401541
let expr = col("part").eq(lit(2)).and(col("b").eq(lit(1.0)));
1541-
let predicate = logical2physical(&expr, &table_schema);
1542+
let predicate = logical2physical(&expr, Arc::clone(&table_schema));
15421543
let opener = make_opener(predicate);
15431544
let stream = opener.open(file.clone()).unwrap().await.unwrap();
15441545
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1547,7 +1548,7 @@ mod test {
15471548

15481549
// Should prune based on file statistics but not partition value
15491550
let expr = col("part").eq(lit(1)).and(col("b").eq(lit(7.0)));
1550-
let predicate = logical2physical(&expr, &table_schema);
1551+
let predicate = logical2physical(&expr, Arc::clone(&table_schema));
15511552
let opener = make_opener(predicate);
15521553
let stream = opener.open(file.clone()).unwrap().await.unwrap();
15531554
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1556,7 +1557,7 @@ mod test {
15561557

15571558
// Should prune based on both partition value and file statistics
15581559
let expr = col("part").eq(lit(2)).and(col("b").eq(lit(7.0)));
1559-
let predicate = logical2physical(&expr, &table_schema);
1560+
let predicate = logical2physical(&expr, table_schema);
15601561
let opener = make_opener(predicate);
15611562
let stream = opener.open(file).unwrap().await.unwrap();
15621563
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1602,7 +1603,7 @@ mod test {
16021603

16031604
// Filter should match the partition value and data value
16041605
let expr = col("part").eq(lit(1)).or(col("a").eq(lit(1)));
1605-
let predicate = logical2physical(&expr, &table_schema);
1606+
let predicate = logical2physical(&expr, Arc::clone(&table_schema));
16061607
let opener = make_opener(predicate);
16071608
let stream = opener.open(file.clone()).unwrap().await.unwrap();
16081609
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1611,7 +1612,7 @@ mod test {
16111612

16121613
// Filter should match the partition value but not the data value
16131614
let expr = col("part").eq(lit(1)).or(col("a").eq(lit(3)));
1614-
let predicate = logical2physical(&expr, &table_schema);
1615+
let predicate = logical2physical(&expr, Arc::clone(&table_schema));
16151616
let opener = make_opener(predicate);
16161617
let stream = opener.open(file.clone()).unwrap().await.unwrap();
16171618
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1620,7 +1621,7 @@ mod test {
16201621

16211622
// Filter should not match the partition value but match the data value
16221623
let expr = col("part").eq(lit(2)).or(col("a").eq(lit(1)));
1623-
let predicate = logical2physical(&expr, &table_schema);
1624+
let predicate = logical2physical(&expr, Arc::clone(&table_schema));
16241625
let opener = make_opener(predicate);
16251626
let stream = opener.open(file.clone()).unwrap().await.unwrap();
16261627
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1629,7 +1630,7 @@ mod test {
16291630

16301631
// Filter should not match the partition value or the data value
16311632
let expr = col("part").eq(lit(2)).or(col("a").eq(lit(3)));
1632-
let predicate = logical2physical(&expr, &table_schema);
1633+
let predicate = logical2physical(&expr, table_schema);
16331634
let opener = make_opener(predicate);
16341635
let stream = opener.open(file).unwrap().await.unwrap();
16351636
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1682,7 +1683,7 @@ mod test {
16821683
// This filter could prune based on statistics, but since it's not dynamic it's not applied for pruning
16831684
// (the assumption is this happened already at planning time)
16841685
let expr = col("a").eq(lit(42));
1685-
let predicate = logical2physical(&expr, &table_schema);
1686+
let predicate = logical2physical(&expr, Arc::clone(&table_schema));
16861687
let opener = make_opener(predicate);
16871688
let stream = opener.open(file.clone()).unwrap().await.unwrap();
16881689
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1691,7 +1692,8 @@ mod test {
16911692

16921693
// If we make the filter dynamic, it should prune.
16931694
// This allows dynamic filters to prune partitions/files even if they are populated late into execution.
1694-
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
1695+
let predicate =
1696+
make_dynamic_expr(logical2physical(&expr, Arc::clone(&table_schema)));
16951697
let opener = make_opener(predicate);
16961698
let stream = opener.open(file.clone()).unwrap().await.unwrap();
16971699
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1701,7 +1703,8 @@ mod test {
17011703
// If we have a filter that touches partition columns only and is dynamic, it should prune even if there are no stats.
17021704
file.statistics = Some(Arc::new(Statistics::new_unknown(&file_schema)));
17031705
let expr = col("part").eq(lit(2));
1704-
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
1706+
let predicate =
1707+
make_dynamic_expr(logical2physical(&expr, Arc::clone(&table_schema)));
17051708
let opener = make_opener(predicate);
17061709
let stream = opener.open(file.clone()).unwrap().await.unwrap();
17071710
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1710,7 +1713,8 @@ mod test {
17101713

17111714
// Similarly a filter that combines partition and data columns should prune even if there are no stats.
17121715
let expr = col("part").eq(lit(2)).and(col("a").eq(lit(42)));
1713-
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
1716+
let predicate =
1717+
make_dynamic_expr(logical2physical(&expr, Arc::clone(&table_schema)));
17141718
let opener = make_opener(predicate);
17151719
let stream = opener.open(file.clone()).unwrap().await.unwrap();
17161720
let (num_batches, num_rows) = count_batches_and_rows(stream).await;

datafusion/datasource-parquet/src/row_filter.rs

Lines changed: 21 additions & 24 deletions
Original file line numberDiff line numberDiff line change
@@ -664,15 +664,13 @@ mod test {
664664

665665
let metadata = reader.metadata();
666666

667-
let table_schema =
667+
let table_schema = Arc::new(
668668
parquet_to_arrow_schema(metadata.file_metadata().schema_descr(), None)
669-
.expect("parsing schema");
669+
.expect("parsing schema"),
670+
);
670671

671672
let expr = col("int64_list").is_not_null();
672-
let expr = logical2physical(&expr, &table_schema);
673-
674-
let table_schema = Arc::new(table_schema.clone());
675-
673+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
676674
let list_index = table_schema
677675
.index_of("int64_list")
678676
.expect("list column should exist");
@@ -698,24 +696,24 @@ mod test {
698696

699697
// This is the schema we would like to coerce to,
700698
// which is different from the physical schema of the file.
701-
let table_schema = Schema::new(vec![Field::new(
699+
let table_schema = Arc::new(Schema::new(vec![Field::new(
702700
"timestamp_col",
703701
DataType::Timestamp(Nanosecond, Some(Arc::from("UTC"))),
704702
false,
705-
)]);
703+
)]));
706704

707705
// Test all should fail
708706
let expr = col("timestamp_col").lt(Expr::Literal(
709707
ScalarValue::TimestampNanosecond(Some(1), Some(Arc::from("UTC"))),
710708
None,
711709
));
712-
let expr = logical2physical(&expr, &table_schema);
710+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
713711
let expr = DefaultPhysicalExprAdapterFactory {}
714-
.create(Arc::new(table_schema.clone()), Arc::clone(&file_schema))
712+
.create(Arc::clone(&table_schema), Arc::clone(&file_schema))
715713
.expect("creating expr adapter")
716714
.rewrite(expr)
717715
.expect("rewriting expression");
718-
let candidate = FilterCandidateBuilder::new(expr, file_schema.clone())
716+
let candidate = FilterCandidateBuilder::new(expr, Arc::clone(&file_schema))
719717
.build(&metadata)
720718
.expect("building candidate")
721719
.expect("candidate expected");
@@ -748,10 +746,10 @@ mod test {
748746
ScalarValue::TimestampNanosecond(Some(0), Some(Arc::from("UTC"))),
749747
None,
750748
));
751-
let expr = logical2physical(&expr, &table_schema);
749+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
752750
// Rewrite the expression to add CastExpr for type coercion
753751
let expr = DefaultPhysicalExprAdapterFactory {}
754-
.create(Arc::new(table_schema), Arc::clone(&file_schema))
752+
.create(table_schema, Arc::clone(&file_schema))
755753
.expect("creating expr adapter")
756754
.rewrite(expr)
757755
.expect("rewriting expression");
@@ -784,8 +782,7 @@ mod test {
784782
)]));
785783

786784
let expr = col("struct_col").is_not_null();
787-
let expr = logical2physical(&expr, &table_schema);
788-
785+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
789786
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
790787
}
791788

@@ -811,15 +808,15 @@ mod test {
811808
let expr = col("struct_col")
812809
.is_not_null()
813810
.and(col("int_col").eq(Expr::Literal(ScalarValue::Int32(Some(5)), None)));
814-
let expr = logical2physical(&expr, &table_schema);
811+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
815812

816813
// The entire expression should not be pushed down
817814
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
818815

819816
// However, just the int_col predicate alone should be pushable
820817
let expr_int_only =
821818
col("int_col").eq(Expr::Literal(ScalarValue::Int32(Some(5)), None));
822-
let expr_int_only = logical2physical(&expr_int_only, &table_schema);
819+
let expr_int_only = logical2physical(&expr_int_only, Arc::clone(&table_schema));
823820
assert!(can_expr_be_pushed_down_with_schemas(
824821
&expr_int_only,
825822
&table_schema
@@ -831,7 +828,7 @@ mod test {
831828
let table_schema = Arc::new(get_lists_table_schema());
832829

833830
let expr = col("utf8_list").is_not_null();
834-
let expr = logical2physical(&expr, &table_schema);
831+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
835832
check_expression_can_evaluate_against_schema(&expr, &table_schema);
836833

837834
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
@@ -901,7 +898,7 @@ mod test {
901898
let metadata = parquet_reader_builder.metadata().clone();
902899
let file_schema = parquet_reader_builder.schema().clone();
903900

904-
let expr = logical2physical(&predicate_expr, &file_schema);
901+
let expr = logical2physical(&predicate_expr, Arc::clone(&file_schema));
905902
if expect_list_support {
906903
assert!(supports_list_predicates(&expr));
907904
}
@@ -1015,22 +1012,22 @@ mod test {
10151012

10161013
#[test]
10171014
fn basic_expr_doesnt_prevent_pushdown() {
1018-
let table_schema = get_basic_table_schema();
1015+
let table_schema = Arc::new(get_basic_table_schema());
10191016

10201017
let expr = col("string_col").is_null();
1021-
let expr = logical2physical(&expr, &table_schema);
1018+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
10221019

10231020
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
10241021
}
10251022

10261023
#[test]
10271024
fn complex_expr_doesnt_prevent_pushdown() {
1028-
let table_schema = get_basic_table_schema();
1025+
let table_schema = Arc::new(get_basic_table_schema());
10291026

10301027
let expr = col("string_col")
10311028
.is_not_null()
10321029
.or(col("bigint_col").gt(Expr::Literal(ScalarValue::Int64(Some(5)), None)));
1033-
let expr = logical2physical(&expr, &table_schema);
1030+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
10341031

10351032
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
10361033
}
@@ -1131,7 +1128,7 @@ mod test {
11311128
ScalarValue::Utf8(Some("target".to_string())),
11321129
None,
11331130
));
1134-
let expr = logical2physical(&expr, &file_schema);
1131+
let expr = logical2physical(&expr, Arc::clone(&file_schema));
11351132

11361133
let candidate = FilterCandidateBuilder::new(expr, file_schema)
11371134
.build(&metadata)

0 commit comments

Comments
 (0)