Skip to content

Commit 90837e5

Browse files
committed
Don't clone the schema in logical2physical
1 parent c53c79d commit 90837e5

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
@@ -1406,7 +1406,7 @@ mod test {
14061406

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

14161416
// A filter on `b = 5.0` should exclude all rows
14171417
let expr = col("b").eq(lit(ScalarValue::Float32(Some(5.0))));
1418-
let predicate = logical2physical(&expr, &schema);
1418+
let predicate = logical2physical(&expr, Arc::clone(&schema));
14191419
let opener = make_opener(predicate);
14201420
let stream = opener.open(file).unwrap().await.unwrap();
14211421
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1461,7 +1461,8 @@ mod test {
14611461
let expr = col("part").eq(lit(1));
14621462
// Mark the expression as dynamic even if it's not to force partition pruning to happen
14631463
// Otherwise we assume it already happened at the planning stage and won't re-do the work here
1464-
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
1464+
let predicate =
1465+
make_dynamic_expr(logical2physical(&expr, Arc::clone(&table_schema)));
14651466
let opener = make_opener(predicate);
14661467
let stream = opener.open(file.clone()).unwrap().await.unwrap();
14671468
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1472,7 +1473,7 @@ mod test {
14721473
let expr = col("part").eq(lit(2));
14731474
// Mark the expression as dynamic even if it's not to force partition pruning to happen
14741475
// Otherwise we assume it already happened at the planning stage and won't re-do the work here
1475-
let predicate = make_dynamic_expr(logical2physical(&expr, &table_schema));
1476+
let predicate = make_dynamic_expr(logical2physical(&expr, table_schema));
14761477
let opener = make_opener(predicate);
14771478
let stream = opener.open(file).unwrap().await.unwrap();
14781479
let (num_batches, num_rows) = count_batches_and_rows(stream).await;
@@ -1528,7 +1529,7 @@ mod test {
15281529

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

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

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

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

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

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

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

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

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

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

datafusion/datasource-parquet/src/row_filter.rs

Lines changed: 20 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -691,15 +691,13 @@ mod test {
691691

692692
let metadata = reader.metadata();
693693

694-
let table_schema =
694+
let table_schema = Arc::new(
695695
parquet_to_arrow_schema(metadata.file_metadata().schema_descr(), None)
696-
.expect("parsing schema");
696+
.expect("parsing schema"),
697+
);
697698

698699
let expr = col("int64_list").is_not_null();
699-
let expr = logical2physical(&expr, &table_schema);
700-
701-
let table_schema = Arc::new(table_schema.clone());
702-
700+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
703701
let list_index = table_schema
704702
.index_of("int64_list")
705703
.expect("list column should exist");
@@ -725,24 +723,24 @@ mod test {
725723

726724
// This is the schema we would like to coerce to,
727725
// which is different from the physical schema of the file.
728-
let table_schema = Schema::new(vec![Field::new(
726+
let table_schema = Arc::new(Schema::new(vec![Field::new(
729727
"timestamp_col",
730728
DataType::Timestamp(Nanosecond, Some(Arc::from("UTC"))),
731729
false,
732-
)]);
730+
)]));
733731

734732
// Test all should fail
735733
let expr = col("timestamp_col").lt(Expr::Literal(
736734
ScalarValue::TimestampNanosecond(Some(1), Some(Arc::from("UTC"))),
737735
None,
738736
));
739-
let expr = logical2physical(&expr, &table_schema);
737+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
740738
let expr = DefaultPhysicalExprAdapterFactory {}
741-
.create(Arc::new(table_schema.clone()), Arc::clone(&file_schema))
739+
.create(Arc::clone(&table_schema), Arc::clone(&file_schema))
742740
.expect("creating expr adapter")
743741
.rewrite(expr)
744742
.expect("rewriting expression");
745-
let candidate = FilterCandidateBuilder::new(expr, file_schema.clone())
743+
let candidate = FilterCandidateBuilder::new(expr, Arc::clone(&file_schema))
746744
.build(&metadata)
747745
.expect("building candidate")
748746
.expect("candidate expected");
@@ -775,10 +773,10 @@ mod test {
775773
ScalarValue::TimestampNanosecond(Some(0), Some(Arc::from("UTC"))),
776774
None,
777775
));
778-
let expr = logical2physical(&expr, &table_schema);
776+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
779777
// Rewrite the expression to add CastExpr for type coercion
780778
let expr = DefaultPhysicalExprAdapterFactory {}
781-
.create(Arc::new(table_schema), Arc::clone(&file_schema))
779+
.create(table_schema, Arc::clone(&file_schema))
782780
.expect("creating expr adapter")
783781
.rewrite(expr)
784782
.expect("rewriting expression");
@@ -811,8 +809,7 @@ mod test {
811809
)]));
812810

813811
let expr = col("struct_col").is_not_null();
814-
let expr = logical2physical(&expr, &table_schema);
815-
812+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
816813
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
817814
}
818815

@@ -838,15 +835,15 @@ mod test {
838835
let expr = col("struct_col")
839836
.is_not_null()
840837
.and(col("int_col").eq(Expr::Literal(ScalarValue::Int32(Some(5)), None)));
841-
let expr = logical2physical(&expr, &table_schema);
838+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
842839

843840
// The entire expression should not be pushed down
844841
assert!(!can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
845842

846843
// However, just the int_col predicate alone should be pushable
847844
let expr_int_only =
848845
col("int_col").eq(Expr::Literal(ScalarValue::Int32(Some(5)), None));
849-
let expr_int_only = logical2physical(&expr_int_only, &table_schema);
846+
let expr_int_only = logical2physical(&expr_int_only, Arc::clone(&table_schema));
850847
assert!(can_expr_be_pushed_down_with_schemas(
851848
&expr_int_only,
852849
&table_schema
@@ -858,7 +855,7 @@ mod test {
858855
let table_schema = Arc::new(get_lists_table_schema());
859856

860857
let expr = col("utf8_list").is_not_null();
861-
let expr = logical2physical(&expr, &table_schema);
858+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
862859
check_expression_can_evaluate_against_schema(&expr, &table_schema);
863860

864861
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
@@ -928,7 +925,7 @@ mod test {
928925
let metadata = parquet_reader_builder.metadata().clone();
929926
let file_schema = parquet_reader_builder.schema().clone();
930927

931-
let expr = logical2physical(&predicate_expr, &file_schema);
928+
let expr = logical2physical(&predicate_expr, Arc::clone(&file_schema));
932929
if expect_list_support {
933930
assert!(supports_list_predicates(&expr));
934931
}
@@ -1042,22 +1039,22 @@ mod test {
10421039

10431040
#[test]
10441041
fn basic_expr_doesnt_prevent_pushdown() {
1045-
let table_schema = get_basic_table_schema();
1042+
let table_schema = Arc::new(get_basic_table_schema());
10461043

10471044
let expr = col("string_col").is_null();
1048-
let expr = logical2physical(&expr, &table_schema);
1045+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
10491046

10501047
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
10511048
}
10521049

10531050
#[test]
10541051
fn complex_expr_doesnt_prevent_pushdown() {
1055-
let table_schema = get_basic_table_schema();
1052+
let table_schema = Arc::new(get_basic_table_schema());
10561053

10571054
let expr = col("string_col")
10581055
.is_not_null()
10591056
.or(col("bigint_col").gt(Expr::Literal(ScalarValue::Int64(Some(5)), None)));
1060-
let expr = logical2physical(&expr, &table_schema);
1057+
let expr = logical2physical(&expr, Arc::clone(&table_schema));
10611058

10621059
assert!(can_expr_be_pushed_down_with_schemas(&expr, &table_schema));
10631060
}

0 commit comments

Comments
 (0)