From b3b3afe380153b383e3f4f29cb0d5e3118a5fb95 Mon Sep 17 00:00:00 2001 From: Michael Kleen Date: Tue, 28 Oct 2025 15:11:47 +0100 Subject: [PATCH 1/2] Move generate_series projection logic into LazyMemoryStream --- datafusion/core/tests/execution/coop.rs | 2 +- .../functions-table/src/generate_series.rs | 24 ++++------- datafusion/physical-plan/src/memory.rs | 42 +++++++++++++++---- datafusion/proto/src/physical_plan/mod.rs | 9 ++-- 4 files changed, 49 insertions(+), 28 deletions(-) diff --git a/datafusion/core/tests/execution/coop.rs b/datafusion/core/tests/execution/coop.rs index b6f406e967509..f6f070f2114d6 100644 --- a/datafusion/core/tests/execution/coop.rs +++ b/datafusion/core/tests/execution/coop.rs @@ -148,7 +148,7 @@ fn make_lazy_exec_with_range( let generator: Arc> = Arc::new(RwLock::new(gen)); // Create a LazyMemoryExec with one partition using our generator - let mut exec = LazyMemoryExec::try_new(schema, vec![generator]).unwrap(); + let mut exec = LazyMemoryExec::try_new(schema, None, vec![generator]).unwrap(); exec.add_ordering(vec![PhysicalSortExpr::new( Arc::new(Column::new(column_name, 0)), diff --git a/datafusion/functions-table/src/generate_series.rs b/datafusion/functions-table/src/generate_series.rs index c66e652147eb8..4325705e04824 100644 --- a/datafusion/functions-table/src/generate_series.rs +++ b/datafusion/functions-table/src/generate_series.rs @@ -237,7 +237,6 @@ impl GenerateSeriesTable { pub fn as_generator( &self, batch_size: usize, - projection: Option>, ) -> Result>> { let generator: Arc> = match &self.args { GenSeriesArgs::ContainsNull { name } => Arc::new(RwLock::new(Empty { name })), @@ -256,7 +255,6 @@ impl GenerateSeriesTable { batch_size, include_end: *include_end, name, - projection, })), GenSeriesArgs::TimestampArgs { start, @@ -297,7 +295,6 @@ impl GenerateSeriesTable { batch_size, include_end: *include_end, name, - projection, })) } GenSeriesArgs::DateArgs { @@ -327,7 +324,6 @@ impl GenerateSeriesTable { batch_size, include_end: *include_end, name, - projection, })), }; @@ -345,7 +341,6 @@ pub struct GenericSeriesState { current: T, include_end: bool, name: &'static str, - projection: Option>, } impl GenericSeriesState { @@ -401,11 +396,7 @@ impl LazyBatchGenerator for GenericSeriesState { let array = self.current.create_array(buf)?; let batch = RecordBatch::try_new(Arc::clone(&self.schema), vec![array])?; - let projected = match self.projection.as_ref() { - Some(projection) => batch.project(projection)?, - None => batch, - }; - Ok(Some(projected)) + Ok(Some(batch)) } } @@ -481,14 +472,13 @@ impl TableProvider for GenerateSeriesTable { _limit: Option, ) -> Result> { let batch_size = state.config_options().execution.batch_size; - let schema = match projection { - Some(projection) => Arc::new(self.schema.project(projection)?), - None => self.schema(), - }; - - let generator = self.as_generator(batch_size, projection.cloned())?; + let generator = self.as_generator(batch_size)?; - Ok(Arc::new(LazyMemoryExec::try_new(schema, vec![generator])?)) + Ok(Arc::new(LazyMemoryExec::try_new( + self.schema(), + projection.cloned(), + vec![generator], + )?)) } } diff --git a/datafusion/physical-plan/src/memory.rs b/datafusion/physical-plan/src/memory.rs index 1bf1e04efb53b..44d8b504b6860 100644 --- a/datafusion/physical-plan/src/memory.rs +++ b/datafusion/physical-plan/src/memory.rs @@ -153,6 +153,8 @@ pub trait LazyBatchGenerator: Send + Sync + fmt::Debug + fmt::Display { pub struct LazyMemoryExec { /// Schema representing the data schema: SchemaRef, + /// Optional projection for which columns to load + projection: Option>, /// Functions to generate batches for each partition batch_generators: Vec>>, /// Plan properties cache storing equivalence properties, partitioning, and execution mode @@ -165,6 +167,7 @@ impl LazyMemoryExec { /// Create a new lazy memory execution plan pub fn try_new( schema: SchemaRef, + projection: Option>, generators: Vec>>, ) -> Result { let boundedness = generators @@ -189,6 +192,11 @@ impl LazyMemoryExec { }) .unwrap_or(Boundedness::Bounded); + let schema = match projection.as_ref() { + Some(columns) => Arc::new(schema.project(columns)?), + None => schema, + }; + let cache = PlanProperties::new( EquivalenceProperties::new(Arc::clone(&schema)), Partitioning::RoundRobinBatch(generators.len()), @@ -199,6 +207,7 @@ impl LazyMemoryExec { Ok(Self { schema, + projection, batch_generators: generators, cache, metrics: ExecutionPlanMetricsSet::new(), @@ -320,6 +329,7 @@ impl ExecutionPlan for LazyMemoryExec { let stream = LazyMemoryStream { schema: Arc::clone(&self.schema), + projection: self.projection.clone(), generator: Arc::clone(&self.batch_generators[partition]), baseline_metrics, }; @@ -338,6 +348,8 @@ impl ExecutionPlan for LazyMemoryExec { /// Stream that generates record batches on demand pub struct LazyMemoryStream { schema: SchemaRef, + /// Optional projection for which columns to load + projection: Option>, /// Generator to produce batches /// /// Note: Idiomatically, DataFusion uses plan-time parallelism - each stream @@ -361,7 +373,14 @@ impl Stream for LazyMemoryStream { let batch = self.generator.write().generate_next_batch(); let poll = match batch { - Ok(Some(batch)) => Poll::Ready(Some(Ok(batch))), + Ok(Some(batch)) => { + // return just the columns requested + let batch = match self.projection.as_ref() { + Some(columns) => batch.project(columns)?, + None => batch, + }; + Poll::Ready(Some(Ok(batch))) + } Ok(None) => Poll::Ready(None), Err(e) => Poll::Ready(Some(Err(e))), }; @@ -434,8 +453,11 @@ mod lazy_memory_tests { schema: Arc::clone(&schema), }; - let exec = - LazyMemoryExec::try_new(schema, vec![Arc::new(RwLock::new(generator))])?; + let exec = LazyMemoryExec::try_new( + schema, + None, + vec![Arc::new(RwLock::new(generator))], + )?; // Test schema assert_eq!(exec.schema().fields().len(), 1); @@ -485,8 +507,11 @@ mod lazy_memory_tests { schema: Arc::clone(&schema), }; - let exec = - LazyMemoryExec::try_new(schema, vec![Arc::new(RwLock::new(generator))])?; + let exec = LazyMemoryExec::try_new( + schema, + None, + vec![Arc::new(RwLock::new(generator))], + )?; // Test invalid partition let result = exec.execute(1, Arc::new(TaskContext::default())); @@ -519,8 +544,11 @@ mod lazy_memory_tests { schema: Arc::clone(&schema), }; - let exec = - LazyMemoryExec::try_new(schema, vec![Arc::new(RwLock::new(generator))])?; + let exec = LazyMemoryExec::try_new( + schema, + None, + vec![Arc::new(RwLock::new(generator))], + )?; let task_ctx = Arc::new(TaskContext::default()); let stream = exec.execute(0, task_ctx)?; diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index 0ebbb373f2d10..5046144754b75 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -1940,10 +1940,13 @@ impl protobuf::PhysicalPlanNode { }; let table = GenerateSeriesTable::new(Arc::clone(&schema), args); - let generator = - table.as_generator(generate_series.target_batch_size as usize, None)?; + let generator = table.as_generator(generate_series.target_batch_size as usize)?; - Ok(Arc::new(LazyMemoryExec::try_new(schema, vec![generator])?)) + Ok(Arc::new(LazyMemoryExec::try_new( + schema, + None, + vec![generator], + )?)) } fn try_into_cooperative_physical_plan( From acc8d547870a47fbf42bb5b00d27139104141f28 Mon Sep 17 00:00:00 2001 From: Michael Kleen Date: Mon, 3 Nov 2025 11:21:35 +0100 Subject: [PATCH 2/2] Add with_projection method to avoid api change --- datafusion/core/tests/execution/coop.rs | 2 +- .../functions-table/src/generate_series.rs | 9 ++-- datafusion/physical-plan/src/memory.rs | 44 +++++++++---------- datafusion/proto/src/physical_plan/mod.rs | 6 +-- 4 files changed, 28 insertions(+), 33 deletions(-) diff --git a/datafusion/core/tests/execution/coop.rs b/datafusion/core/tests/execution/coop.rs index f6f070f2114d6..b6f406e967509 100644 --- a/datafusion/core/tests/execution/coop.rs +++ b/datafusion/core/tests/execution/coop.rs @@ -148,7 +148,7 @@ fn make_lazy_exec_with_range( let generator: Arc> = Arc::new(RwLock::new(gen)); // Create a LazyMemoryExec with one partition using our generator - let mut exec = LazyMemoryExec::try_new(schema, None, vec![generator]).unwrap(); + let mut exec = LazyMemoryExec::try_new(schema, vec![generator]).unwrap(); exec.add_ordering(vec![PhysicalSortExpr::new( Arc::new(Column::new(column_name, 0)), diff --git a/datafusion/functions-table/src/generate_series.rs b/datafusion/functions-table/src/generate_series.rs index 4325705e04824..d71c5945aafcc 100644 --- a/datafusion/functions-table/src/generate_series.rs +++ b/datafusion/functions-table/src/generate_series.rs @@ -474,11 +474,10 @@ impl TableProvider for GenerateSeriesTable { let batch_size = state.config_options().execution.batch_size; let generator = self.as_generator(batch_size)?; - Ok(Arc::new(LazyMemoryExec::try_new( - self.schema(), - projection.cloned(), - vec![generator], - )?)) + Ok(Arc::new( + LazyMemoryExec::try_new(self.schema(), vec![generator])? + .with_projection(projection.cloned()), + )) } } diff --git a/datafusion/physical-plan/src/memory.rs b/datafusion/physical-plan/src/memory.rs index 44d8b504b6860..09710ae1e2edb 100644 --- a/datafusion/physical-plan/src/memory.rs +++ b/datafusion/physical-plan/src/memory.rs @@ -167,7 +167,6 @@ impl LazyMemoryExec { /// Create a new lazy memory execution plan pub fn try_new( schema: SchemaRef, - projection: Option>, generators: Vec>>, ) -> Result { let boundedness = generators @@ -192,11 +191,6 @@ impl LazyMemoryExec { }) .unwrap_or(Boundedness::Bounded); - let schema = match projection.as_ref() { - Some(columns) => Arc::new(schema.project(columns)?), - None => schema, - }; - let cache = PlanProperties::new( EquivalenceProperties::new(Arc::clone(&schema)), Partitioning::RoundRobinBatch(generators.len()), @@ -207,13 +201,28 @@ impl LazyMemoryExec { Ok(Self { schema, - projection, + projection: None, batch_generators: generators, cache, metrics: ExecutionPlanMetricsSet::new(), }) } + pub fn with_projection(mut self, projection: Option>) -> Self { + match projection.as_ref() { + Some(columns) => { + let projected = Arc::new(self.schema.project(columns).unwrap()); + self.cache = self.cache.with_eq_properties(EquivalenceProperties::new( + Arc::clone(&projected), + )); + self.schema = projected; + self.projection = projection; + self + } + _ => self, + } + } + pub fn try_set_partitioning(&mut self, partitioning: Partitioning) -> Result<()> { if partitioning.partition_count() != self.batch_generators.len() { internal_err!( @@ -453,11 +462,8 @@ mod lazy_memory_tests { schema: Arc::clone(&schema), }; - let exec = LazyMemoryExec::try_new( - schema, - None, - vec![Arc::new(RwLock::new(generator))], - )?; + let exec = + LazyMemoryExec::try_new(schema, vec![Arc::new(RwLock::new(generator))])?; // Test schema assert_eq!(exec.schema().fields().len(), 1); @@ -507,11 +513,8 @@ mod lazy_memory_tests { schema: Arc::clone(&schema), }; - let exec = LazyMemoryExec::try_new( - schema, - None, - vec![Arc::new(RwLock::new(generator))], - )?; + let exec = + LazyMemoryExec::try_new(schema, vec![Arc::new(RwLock::new(generator))])?; // Test invalid partition let result = exec.execute(1, Arc::new(TaskContext::default())); @@ -544,11 +547,8 @@ mod lazy_memory_tests { schema: Arc::clone(&schema), }; - let exec = LazyMemoryExec::try_new( - schema, - None, - vec![Arc::new(RwLock::new(generator))], - )?; + let exec = + LazyMemoryExec::try_new(schema, vec![Arc::new(RwLock::new(generator))])?; let task_ctx = Arc::new(TaskContext::default()); let stream = exec.execute(0, task_ctx)?; diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index 5046144754b75..e5f4a1f7d0267 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -1942,11 +1942,7 @@ impl protobuf::PhysicalPlanNode { let table = GenerateSeriesTable::new(Arc::clone(&schema), args); let generator = table.as_generator(generate_series.target_batch_size as usize)?; - Ok(Arc::new(LazyMemoryExec::try_new( - schema, - None, - vec![generator], - )?)) + Ok(Arc::new(LazyMemoryExec::try_new(schema, vec![generator])?)) } fn try_into_cooperative_physical_plan(