From eb129ae50c2f39c2ce48adcdf295b6d8459d17a5 Mon Sep 17 00:00:00 2001 From: Adam Gutglick Date: Tue, 12 Aug 2025 15:23:33 +0100 Subject: [PATCH] [branch-49] Backport #17129 to branch 49 (#17143) * Preserve equivalence properties during projection pushdown (#17129) * Adds parquet data diffs --------- Co-authored-by: Matthew Kim <38759997+friendlymatthew@users.noreply.github.com> --- datafusion/datasource/src/source.rs | 35 +++++++++++++++++- datafusion/sqllogictest/data/1.parquet | Bin 0 -> 1381 bytes datafusion/sqllogictest/data/2.parquet | Bin 0 -> 1403 bytes .../test_files/parquet_filter_pushdown.slt | 32 ++++++++++++++++ 4 files changed, 66 insertions(+), 1 deletion(-) create mode 100644 datafusion/sqllogictest/data/1.parquet create mode 100644 datafusion/sqllogictest/data/2.parquet diff --git a/datafusion/datasource/src/source.rs b/datafusion/datasource/src/source.rs index fde1944ae066a..3a7ff1ef09911 100644 --- a/datafusion/datasource/src/source.rs +++ b/datafusion/datasource/src/source.rs @@ -22,6 +22,7 @@ use std::fmt; use std::fmt::{Debug, Formatter}; use std::sync::Arc; +use datafusion_physical_expr::equivalence::ProjectionMapping; use datafusion_physical_plan::execution_plan::{ Boundedness, EmissionType, SchedulingType, }; @@ -324,7 +325,39 @@ impl ExecutionPlan for DataSourceExec { &self, projection: &ProjectionExec, ) -> Result>> { - self.data_source.try_swapping_with_projection(projection) + match self.data_source.try_swapping_with_projection(projection)? { + Some(new_plan) => { + if let Some(new_data_source_exec) = + new_plan.as_any().downcast_ref::() + { + let projection_mapping = ProjectionMapping::try_new( + projection.expr().iter().cloned(), + &self.schema(), + )?; + + // Project the equivalence properties to the new schema + let projected_eq_properties = self + .cache + .eq_properties + .project(&projection_mapping, new_data_source_exec.schema()); + + let preserved_exec = DataSourceExec { + data_source: Arc::clone(&new_data_source_exec.data_source), + cache: PlanProperties::new( + projected_eq_properties, + new_data_source_exec.cache.partitioning.clone(), + new_data_source_exec.cache.emission_type, + new_data_source_exec.cache.boundedness, + ) + .with_scheduling_type(new_data_source_exec.cache.scheduling_type), + }; + Ok(Some(Arc::new(preserved_exec))) + } else { + Ok(Some(new_plan)) + } + } + None => Ok(None), + } } fn handle_child_pushdown_result( diff --git a/datafusion/sqllogictest/data/1.parquet b/datafusion/sqllogictest/data/1.parquet new file mode 100644 index 0000000000000000000000000000000000000000..a04f669eaeaaeb9edd468f7c28c00ba526627007 GIT binary patch literal 1381 zcmb_c&ubGw6n?v1(;NaO;*7hnhf?TLq_j!3QYDC!);7kpHj6Z=2w|IS+dz^Hn>A67 zMFhc<#~%6zcoe*N_9A!?57Lu(^p8;Rz1jS*p|w3YWOv@Y@6GqV_hvR5!cH-bW!cp{ zQyE+Wn0`O^8%i!fasu#`Qeg56A{fMHFeJ_*EMne(>51gOM@m04B9XuhVLq}{%n@gk z$At(4K>76U>#Hj#h=}$vePErFR72X9?^RDA)yS{Q_b8c>P>r+eI!6do4Gtjb2Fi_L z5r4r_hY`x@xlO*#J}b1J-vPtqM=Go1ip+hDJ(fTeFgTx$Ilk|8%k9dZ+i+L}SZoUP zXy7{)w_K}ELEglDOhf0zcHsCyIjA*Uv>L7Y>x7v`5MkQGt8T0AJ!`nlpzJm~HQ#HJ z-DBXYVH#-*Ocpa1AQHz?`Z-vPtNc*qZ&YjDivGWwW6a=n!6E@)ah%fF3*bK^%;euS zcBlZU(Ryk|i<6>WD*UZt9jVrVY7SdJv^Y!;&SvPvD>0fH_=|DI`L7G?w#?e^!6`kH z$vgZ&vGz6V;~2H%_>~*wPjl4455}>y4-qztrp8q(%D-us>Cp9=r)^>68d3ALqK=Ojf6vSPGg;?yByc@11i7Zm0 z8Cy=FEht2`PWGdk&FyW?Z|*q7u$`bARVQ$Ep0sOTbE4z=a=w0ZHaVL#()dM%KiD4w I*uX#7KgL)qfdBvi literal 0 HcmV?d00001 diff --git a/datafusion/sqllogictest/data/2.parquet b/datafusion/sqllogictest/data/2.parquet new file mode 100644 index 0000000000000000000000000000000000000000..b5e29f81baf152a741dc310da3348bd072532029 GIT binary patch literal 1403 zcmb_c&rcIU6n;Bw+8!EgjWh0=JrEK$jS?tT43fq;1)F#_p-<$9Kn(YR+=BP@u^m?8W z3YP&AA7+0-$pJ`C0KQTROnycNjo37r#At|yjN6cyNdA4Kvgm{4;6N-n1`zQX-w24==Gr0Biy9fIA-`!WOx#yi@=UmM$8ENs$LPj*P`6A z9u*h+z1F^vSuuJ%!#OYDBgR9{dwi+Jb7Bi;E?L18aLxGx0r5aE<7y2bY3Yjysmy|q zUO^=E{IRkeAfD++Ug)n{jmqlNu2YvVtzEm=FDGbv)%NV&C!M<6n&>)hIp4Z7 VlblKFY5cJIKX^F?uz`QOzW^hSGfe;h literal 0 HcmV?d00001 diff --git a/datafusion/sqllogictest/test_files/parquet_filter_pushdown.slt b/datafusion/sqllogictest/test_files/parquet_filter_pushdown.slt index 24e76a570c009..61f4d6fc12a3b 100644 --- a/datafusion/sqllogictest/test_files/parquet_filter_pushdown.slt +++ b/datafusion/sqllogictest/test_files/parquet_filter_pushdown.slt @@ -528,3 +528,35 @@ query TT select val, part from t_pushdown where part = val AND part = 'a'; ---- a a + +statement ok +COPY ( + SELECT + '00000000000000000000000000000001' AS trace_id, + '2023-10-01 00:00:00'::timestamptz AS start_timestamp, + 'prod' as deployment_environment +) +TO 'data/1.parquet'; + +statement ok +COPY ( + SELECT + '00000000000000000000000000000002' AS trace_id, + '2024-10-01 00:00:00'::timestamptz AS start_timestamp, + 'staging' as deployment_environment +) +TO 'data/2.parquet'; + +statement ok +CREATE EXTERNAL TABLE t1 STORED AS PARQUET LOCATION 'data/'; + +statement ok +SET datafusion.execution.parquet.pushdown_filters = true; + +query T +SELECT deployment_environment +FROM t1 +WHERE trace_id = '00000000000000000000000000000002' +ORDER BY start_timestamp, trace_id; +---- +staging