From ca5a4ba9ec4f87ee72fff74cfc515686af9a3fe7 Mon Sep 17 00:00:00 2001 From: osipovartem Date: Mon, 7 Sep 2026 23:29:46 +0300 Subject: [PATCH] Preserve projection metadata across data sources --- .../physical_optimizer/projection_pushdown.rs | 33 +++++++++++++++++++ datafusion/datasource/src/source.rs | 8 ++++- 2 files changed, 40 insertions(+), 1 deletion(-) diff --git a/datafusion/core/tests/physical_optimizer/projection_pushdown.rs b/datafusion/core/tests/physical_optimizer/projection_pushdown.rs index 3a8d82f111145..47587e85c6311 100644 --- a/datafusion/core/tests/physical_optimizer/projection_pushdown.rs +++ b/datafusion/core/tests/physical_optimizer/projection_pushdown.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use std::collections::HashMap; use std::sync::Arc; use arrow::compute::SortOptions; @@ -529,6 +530,38 @@ fn test_memory_after_projection() -> Result<()> { Ok(()) } +#[test] +fn test_memory_projection_preserves_field_metadata() -> Result<()> { + let schema = Arc::new(Schema::new(vec![Field::new( + "value", + DataType::Int32, + true, + )])); + let memory = MemorySourceConfig::try_new_exec(&[], Arc::clone(&schema), None)?; + let projected_schema = Schema::new(vec![ + Field::new("value", DataType::Int32, true).with_metadata(HashMap::from([( + "semantic_type".to_string(), + "example".to_string(), + )])), + ]); + let projection: Arc = + Arc::new(ProjectionExec::try_new_with_schema_metadata( + vec![ProjectionExpr::new( + Arc::new(Column::new("value", 0)), + "value", + )], + memory, + &projected_schema, + )?); + + let optimized = + ProjectionPushdown::new().optimize(projection, &ConfigOptions::new())?; + + assert!(optimized.is::()); + assert_eq!(optimized.schema().as_ref(), &projected_schema); + Ok(()) +} + #[test] fn test_streaming_table_after_projection() -> Result<()> { #[derive(Debug)] diff --git a/datafusion/datasource/src/source.rs b/datafusion/datasource/src/source.rs index 741010c595197..1f2913da905c1 100644 --- a/datafusion/datasource/src/source.rs +++ b/datafusion/datasource/src/source.rs @@ -531,7 +531,13 @@ impl ExecutionPlan for DataSourceExec { .try_swapping_with_projection(projection.projection_expr())? { Some(new_data_source) => { - Ok(Some(Arc::new(DataSourceExec::new(new_data_source)))) + let new_exec = Arc::new(DataSourceExec::new(new_data_source)); + // The data source API does not receive the projection's output metadata. + if projection.schema() == new_exec.schema() { + Ok(Some(new_exec)) + } else { + Ok(None) + } } None => Ok(None), }