From afb7b33f93b6842a60cbe7257a31aaa2229f669c Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Fri, 28 Aug 2026 00:32:57 -0700 Subject: [PATCH 1/9] feat: allow scaling RangePartitioning --- datafusion/ffi/src/execution_plan.rs | 6 + .../ffi/src/physical_expr/partitioning.rs | 59 ++- datafusion/ffi/src/plan_properties.rs | 12 +- datafusion/ffi/src/tests/mod.rs | 29 +- datafusion/ffi/tests/ffi_execution_plan.rs | 26 +- datafusion/physical-expr/src/partitioning.rs | 459 ++++++++++++++++-- datafusion/physical-plan/src/joins/utils.rs | 38 +- .../physical-plan/src/repartition/mod.rs | 19 +- .../proto-models/proto/datafusion.proto | 5 + .../proto-models/src/generated/pbjson.rs | 44 +- .../proto-models/src/generated/prost.rs | 7 + 11 files changed, 627 insertions(+), 77 deletions(-) diff --git a/datafusion/ffi/src/execution_plan.rs b/datafusion/ffi/src/execution_plan.rs index d7ee5dace30cc..75ea170658927 100644 --- a/datafusion/ffi/src/execution_plan.rs +++ b/datafusion/ffi/src/execution_plan.rs @@ -600,6 +600,12 @@ pub mod tests { self } + pub fn with_partitioning(mut self, partitioning: Partitioning) -> Self { + self.props = + Arc::new(self.props.as_ref().clone().with_partitioning(partitioning)); + self + } + pub fn with_expressions( mut self, expressions: Vec>, diff --git a/datafusion/ffi/src/physical_expr/partitioning.rs b/datafusion/ffi/src/physical_expr/partitioning.rs index 2a9a8528c6c3e..05535caefaef3 100644 --- a/datafusion/ffi/src/physical_expr/partitioning.rs +++ b/datafusion/ffi/src/physical_expr/partitioning.rs @@ -33,8 +33,9 @@ use crate::physical_expr::sort::FFI_PhysicalSortExpr; #[repr(C)] #[derive(Debug)] pub struct FFI_RangePartitioning { - split_points: SVec>, + samples: SVec>, ordering: SVec, + partition_count: usize, } /// A stable struct for sharing [`Partitioning`] across FFI boundaries. @@ -62,8 +63,8 @@ impl From<&Partitioning> for FFI_Partitioning { } Partitioning::Range(range) => { // Producer-side conversion should be infallible at ABI boundary - let split_points = range - .split_points() + let samples = range + .samples() .iter() .map(|split_point| { split_point @@ -83,8 +84,9 @@ impl From<&Partitioning> for FFI_Partitioning { .map(FFI_PhysicalSortExpr::from) .collect(); Self::Range(FFI_RangePartitioning { - split_points, + samples, ordering, + partition_count: range.partition_count(), }) } Partitioning::UnknownPartitioning(size) => Self::UnknownPartitioning(*size), @@ -105,8 +107,8 @@ impl TryFrom for Partitioning { Self::Hash(exprs, size) } FFI_Partitioning::Range(range) => { - let split_points = range - .split_points + let samples = range + .samples .into_iter() .map(|split_point| { split_point @@ -126,7 +128,11 @@ impl TryFrom for Partitioning { ) })?; - Self::Range(RangePartitioning::try_new(ordering, split_points)?) + Self::Range(RangePartitioning::try_new_with_samples( + ordering, + samples, + range.partition_count, + )?) } FFI_Partitioning::UnknownPartitioning(size) => { Self::UnknownPartitioning(size) @@ -174,6 +180,20 @@ mod tests { )?)) } + fn sampled_range_partitioning() -> Result { + let ordering = LexOrdering::new([PhysicalSortExpr::new_default(Arc::new( + Column::new("a", 0), + ))]) + .expect("non-empty ordering"); + let samples = [10, 20, 30, 40, 50] + .into_iter() + .map(|value| SplitPoint::new(vec![ScalarValue::Int64(Some(value))])) + .collect(); + Ok(Partitioning::Range( + RangePartitioning::try_new_with_samples(ordering, samples, 3)?, + )) + } + #[test] fn round_trip_ffi_partitioning() -> Result<()> { for partitioning in [ @@ -181,6 +201,7 @@ mod tests { Partitioning::Hash(vec![lit(1)], 10), Partitioning::UnknownPartitioning(10), range_partitioning()?, + sampled_range_partitioning()?, ] { let ffi_partitioning: FFI_Partitioning = (&partitioning).into(); let returned: Partitioning = ffi_partitioning.try_into()?; @@ -210,11 +231,33 @@ mod tests { Ok(()) } + #[test] + fn round_trip_ffi_sampled_range_partitioning() -> Result<()> { + let partitioning = sampled_range_partitioning()?; + + let ffi_partitioning: FFI_Partitioning = (&partitioning).into(); + let returned: Partitioning = ffi_partitioning.try_into()?; + let Partitioning::Range(returned) = returned else { + panic!("expected range partitioning"); + }; + let Partitioning::Range(original) = partitioning else { + panic!("expected range partitioning"); + }; + + assert_eq!(returned, original); + assert_eq!(returned.samples(), original.samples()); + assert_eq!(returned.max_partition_count(), 6); + assert_eq!(returned.partition_count(), 3); + + Ok(()) + } + #[test] fn ffi_range_partitioning_rejects_empty_ordering() { let ffi_partitioning = FFI_Partitioning::Range(FFI_RangePartitioning { - split_points: SVec::new(), + samples: SVec::new(), ordering: SVec::new(), + partition_count: 1, }); let err = Partitioning::try_from(ffi_partitioning).unwrap_err(); diff --git a/datafusion/ffi/src/plan_properties.rs b/datafusion/ffi/src/plan_properties.rs index 09ef26af32349..eccbc9bc6b37a 100644 --- a/datafusion/ffi/src/plan_properties.rs +++ b/datafusion/ffi/src/plan_properties.rs @@ -290,11 +290,14 @@ mod tests { let col = datafusion::physical_plan::expressions::col("a", &schema)?; let ordering = LexOrdering::new([PhysicalSortExpr::new_default(col)]) .expect("non-empty ordering"); - let split_points = vec![ + let samples = vec![ SplitPoint::new(vec![ScalarValue::Int64(Some(10))]), SplitPoint::new(vec![ScalarValue::Int64(Some(20))]), + SplitPoint::new(vec![ScalarValue::Int64(Some(30))]), + SplitPoint::new(vec![ScalarValue::Int64(Some(40))]), + SplitPoint::new(vec![ScalarValue::Int64(Some(50))]), ]; - let range = RangePartitioning::try_new(ordering, split_points)?; + let range = RangePartitioning::try_new_with_samples(ordering, samples, 3)?; Ok(PlanProperties::new( EquivalenceProperties::new(schema), @@ -314,7 +317,6 @@ mod tests { let foreign_props: PlanProperties = local_props_ptr.try_into()?; assert_eq!(format!("{foreign_props:?}"), format!("{original_props:?}")); - Ok(()) } @@ -351,6 +353,10 @@ mod tests { format!("{:?}", original_props.output_partitioning()) ); assert_eq!(format!("{foreign_props:?}"), format!("{original_props:?}")); + let Partitioning::Range(range) = foreign_props.output_partitioning() else { + panic!("expected range partitioning"); + }; + assert_eq!(range.max_partition_count(), 6); Ok(()) } diff --git a/datafusion/ffi/src/tests/mod.rs b/datafusion/ffi/src/tests/mod.rs index fbc3e83ba49fc..3aba5860ca735 100644 --- a/datafusion/ffi/src/tests/mod.rs +++ b/datafusion/ffi/src/tests/mod.rs @@ -27,10 +27,13 @@ use datafusion_catalog::MemTable; use datafusion_catalog::{Session, TableProvider}; use datafusion_common::stats::Precision; use datafusion_common::{ColumnStatistics, Statistics}; -use datafusion_common::{Result, ScalarValue, exec_err}; +use datafusion_common::{Result, ScalarValue, SplitPoint, exec_err}; use datafusion_expr::{Expr, TableType, col, lit}; -use datafusion_physical_expr::PhysicalExpr; -use datafusion_physical_plan::ExecutionPlan; +use datafusion_physical_expr::expressions::Column; +use datafusion_physical_expr::{ + LexOrdering, PhysicalExpr, PhysicalSortExpr, RangePartitioning, +}; +use datafusion_physical_plan::{ExecutionPlan, Partitioning}; use sync_provider::create_sync_table_provider; use udf_udaf_udwf::{ create_ffi_abs_func, create_ffi_first_value_func, create_ffi_random_func, @@ -120,6 +123,8 @@ pub struct ForeignLibraryModule { pub create_exec_with_statistics: extern "C" fn() -> FFI_ExecutionPlan, + pub create_exec_with_range_partitioning: extern "C" fn() -> FFI_ExecutionPlan, + pub create_table_with_statistics: extern "C" fn(codec: FFI_LogicalExtensionCodec) -> FFI_TableProvider, @@ -232,6 +237,23 @@ pub(crate) extern "C" fn create_exec_with_statistics() -> FFI_ExecutionPlan { FFI_ExecutionPlan::new(plan, None) } +pub(crate) extern "C" fn create_exec_with_range_partitioning() -> FFI_ExecutionPlan { + let schema = create_test_schema(); + let ordering = + LexOrdering::new([PhysicalSortExpr::new_default(Arc::new(Column::new("a", 0)))]) + .expect("non-empty ordering"); + let samples = [10, 20, 30, 40, 50] + .into_iter() + .map(|value| SplitPoint::new(vec![ScalarValue::Int32(Some(value))])) + .collect(); + let partitioning = Partitioning::Range( + RangePartitioning::try_new_with_samples(ordering, samples, 3) + .expect("valid sampled range partitioning"), + ); + let plan = Arc::new(EmptyExec::new(schema).with_partitioning(partitioning)); + FFI_ExecutionPlan::new(plan, None) +} + /// Thin wrapper that attaches a fixed [`Statistics`] snapshot to any inner /// [`TableProvider`] without changing its scan behaviour. #[derive(Debug)] @@ -362,6 +384,7 @@ pub extern "C" fn datafusion_ffi_get_module() -> ForeignLibraryModule { create_exec_with_expressions, create_exec_with_dynamic_expressions, create_exec_with_statistics, + create_exec_with_range_partitioning, create_table_with_statistics, create_physical_optimizer_rule: physical_optimizer::create_physical_optimizer_rule, diff --git a/datafusion/ffi/tests/ffi_execution_plan.rs b/datafusion/ffi/tests/ffi_execution_plan.rs index 4067d7eb49b2a..e0a41ffced6e6 100644 --- a/datafusion/ffi/tests/ffi_execution_plan.rs +++ b/datafusion/ffi/tests/ffi_execution_plan.rs @@ -29,7 +29,7 @@ mod tests { use datafusion_ffi::tests::utils::get_module; use datafusion_physical_plan::execution_plan::InvariantLevel; use datafusion_physical_plan::{ - ChildrenPropertiesMode, ExecutionPlan, ReplaceChildrenOptions, + ChildrenPropertiesMode, ExecutionPlan, Partitioning, ReplaceChildrenOptions, }; use std::sync::Arc; @@ -68,6 +68,30 @@ mod tests { Ok(()) } + #[test] + fn test_ffi_range_partitioning_cross_library() -> Result<(), DataFusionError> { + let module = get_module()?; + let plan = (module.create_exec_with_range_partitioning)(); + let plan: Arc = (&plan).try_into()?; + let Partitioning::Range(range) = plan.properties().output_partitioning() else { + panic!("expected range partitioning"); + }; + + assert_eq!(range.partition_count(), 3); + assert_eq!(range.max_partition_count(), 6); + assert_eq!(range.samples().len(), 5); + assert_eq!( + range + .split_points() + .iter() + .map(|point| point.to_string()) + .collect::>(), + vec!["(20)", "(40)"] + ); + + Ok(()) + } + #[test] fn test_ffi_execution_plan_expressions_cross_library() -> Result<(), DataFusionError> { diff --git a/datafusion/physical-expr/src/partitioning.rs b/datafusion/physical-expr/src/partitioning.rs index 70ba92a0a9ecc..6d5cdcb7842d3 100644 --- a/datafusion/physical-expr/src/partitioning.rs +++ b/datafusion/physical-expr/src/partitioning.rs @@ -152,15 +152,21 @@ impl Display for Partitioning { /// Physical range partitioning. /// -/// [`RangePartitioning`] describes an ordered key space with split points. +/// [`RangePartitioning`] describes an ordered key space with sampled split points. /// /// - `ordering` defines the partitioning key and ordering. -/// - `split_points` define the boundaries between adjacent partitions. +/// - `samples` define the maximum-resolution boundaries. +/// - `partition_count` selects how many ranges to derive from those samples. /// /// Comparisons use the lexicographic order defined by `ordering`, including -/// `ASC`/`DESC` and null ordering. Split points must be strictly ordered -/// according to that ordering, and each split point must have one value per -/// ordering expression. See [`SplitPoint`] for the shared boundary convention. +/// `ASC`/`DESC` and null ordering. Samples must be strictly ordered according +/// to that ordering, and each sample must have one value per ordering +/// expression. See [`SplitPoint`] for the shared boundary convention. +/// +/// When `partition_count` is smaller than [`Self::max_partition_count`], the +/// samples are evenly down-sampled to derive the effective split points. This +/// allows planners to reduce or later restore the number of partitions without +/// losing the original distribution sample. /// /// Like other user-specified data properties such as sortedness, if a source /// declares range partitioning, it is responsible for placing each row in the @@ -200,12 +206,16 @@ impl Display for Partitioning { /// NOTE: Optimizer and execution behavior for this partitioning is intentionally /// not implemented and will be introduced incrementally. See /// . -#[derive(Debug, Clone, PartialEq)] +#[derive(Debug, Clone)] pub struct RangePartitioning { /// Ordered partitioning key. ordering: LexOrdering, - /// Boundaries between adjacent partitions. - split_points: Vec, + /// Maximum-resolution boundaries used to derive split points. + samples: Arc<[SplitPoint]>, + /// Effective boundaries for the current partition count. + split_points: Arc<[SplitPoint]>, + /// Number of effective partitions. + partition_count: usize, } impl RangePartitioning { @@ -214,9 +224,13 @@ impl RangePartitioning { /// Use [`Self::try_new`] to validate the contract documented on /// [`RangePartitioning`]. pub fn new(ordering: LexOrdering, split_points: Vec) -> Self { + let partition_count = split_points.len() + 1; + let split_points: Arc<[SplitPoint]> = Arc::from(split_points); Self { ordering, + samples: Arc::clone(&split_points), split_points, + partition_count, } } @@ -233,19 +247,75 @@ impl RangePartitioning { Ok(Self::new(ordering, split_points)) } + /// Creates sample-backed range partitioning and validates the sample shape, + /// ordering, and target partition count. + /// + /// `partition_count` must be at least one and no larger than + /// `samples.len() + 1`. When it is smaller than that maximum, the samples + /// are evenly down-sampled to derive the effective split points. + pub fn try_new_with_samples( + ordering: LexOrdering, + samples: Vec, + partition_count: usize, + ) -> Result { + validate_range_split_points( + &samples, + &ordering + .iter() + .map(|sort_expr| sort_expr.options) + .collect::>(), + )?; + validate_range_partition_count(partition_count, samples.len() + 1)?; + let samples: Arc<[SplitPoint]> = Arc::from(samples); + let split_points = downsample_split_points(&samples, partition_count); + Ok(Self { + ordering, + samples, + split_points, + partition_count, + }) + } + /// Returns the ordering that defines the range key. pub fn ordering(&self) -> &LexOrdering { &self.ordering } - /// Returns the ordered split points between partitions. + /// Returns the maximum-resolution sample points. + pub fn samples(&self) -> &[SplitPoint] { + &self.samples + } + + /// Returns the effective split points between partitions. pub fn split_points(&self) -> &[SplitPoint] { &self.split_points } /// Returns the number of partitions. pub fn partition_count(&self) -> usize { - self.split_points.len() + 1 + self.partition_count + } + + /// Returns the largest partition count supported by the stored samples. + pub fn max_partition_count(&self) -> usize { + self.samples.len() + 1 + } + + /// Returns this range partitioning scaled to `target_partitions`. + /// + /// Scaling retains the original samples, so a range partitioning that was + /// scaled down can later be scaled back up to [`Self::max_partition_count`]. + pub fn scale(&self, target_partitions: usize) -> Result { + validate_range_partition_count(target_partitions, self.max_partition_count())?; + if target_partitions == self.partition_count { + return Ok(self.clone()); + } + Ok(Self { + ordering: self.ordering.clone(), + samples: Arc::clone(&self.samples), + split_points: downsample_split_points(&self.samples, target_partitions), + partition_count: target_partitions, + }) } /// Calculates the range partitioning after applying the given projection. @@ -278,11 +348,19 @@ impl RangePartitioning { Some(Self { ordering, - split_points: self.split_points.clone(), + samples: Arc::clone(&self.samples), + split_points: Arc::clone(&self.split_points), + partition_count: self.partition_count, }) } } +impl PartialEq for RangePartitioning { + fn eq(&self, other: &Self) -> bool { + self.ordering == other.ordering && self.split_points == other.split_points + } +} + impl Display for RangePartitioning { fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result { let split_points = format_range_split_points(&self.split_points); @@ -296,6 +374,44 @@ impl Display for RangePartitioning { } } +fn downsample_split_points( + samples: &Arc<[SplitPoint]>, + partition_count: usize, +) -> Arc<[SplitPoint]> { + if partition_count == samples.len() + 1 { + return Arc::clone(samples); + } + + let sample_count = samples.len(); + (1..partition_count) + .map(|partition| { + // Use a wider intermediate so valid slice lengths cannot overflow + // when calculating the evenly spaced sample index. + let sample_index = ((partition as u128 * sample_count as u128) + / partition_count as u128) as usize; + samples[sample_index].clone() + }) + .collect::>() + .into() +} + +fn validate_range_partition_count( + partition_count: usize, + max_partition_count: usize, +) -> Result<()> { + if partition_count == 0 { + return datafusion_common::plan_err!( + "Range partitioning partition count must be at least 1" + ); + } + if partition_count > max_partition_count { + return datafusion_common::plan_err!( + "Range partitioning partition count {partition_count} exceeds maximum {max_partition_count}" + ); + } + Ok(()) +} + fn format_range_split_points(split_points: &[SplitPoint]) -> String { split_points .iter() @@ -555,22 +671,27 @@ impl Partitioning { } Partitioning::Range(range) => { let sort_expr = sort_exprs_try_to_proto(range.ordering().iter(), ctx)?; - let split_point = range - .split_points() - .iter() - .map(|split_point| { - let value = split_point - .values() - .iter() - .map(|value| value.try_into().map_err(Into::into)) - .collect::>>()?; - Ok(protobuf::PhysicalRangeSplitPoint { value }) - }) - .collect::>>()?; + let encode_split_points = |split_points: &[SplitPoint]| { + split_points + .iter() + .map(|split_point| { + let value = split_point + .values() + .iter() + .map(|value| value.try_into().map_err(Into::into)) + .collect::>>()?; + Ok(protobuf::PhysicalRangeSplitPoint { value }) + }) + .collect::>>() + }; + let split_point = encode_split_points(range.split_points())?; + let sample_point = encode_split_points(range.samples())?; protobuf::partitioning::PartitionMethod::Range( protobuf::PhysicalRangePartitioning { sort_expr, split_point, + sample_point, + partition_count: partition_count(range.partition_count())?, }, ) } @@ -629,19 +750,44 @@ impl Partitioning { "Range partitioning ordering must not contain duplicate expressions" ); } - let split_points = range - .split_point - .iter() - .map(|split_point| { - let values = split_point - .value + let decode_split_points = + |split_points: &[protobuf::PhysicalRangeSplitPoint]| { + split_points .iter() - .map(|value| ScalarValue::try_from(value).map_err(Into::into)) - .collect::>>()?; - Ok(SplitPoint::new(values)) - }) - .collect::>>()?; - Partitioning::Range(RangePartitioning::try_new(ordering, split_points)?) + .map(|split_point| { + let values = split_point + .value + .iter() + .map(|value| { + ScalarValue::try_from(value).map_err(Into::into) + }) + .collect::>>()?; + Ok(SplitPoint::new(values)) + }) + .collect::>>() + }; + let split_points = decode_split_points(&range.split_point)?; + if range.partition_count == 0 { + // Older payloads derive their partition count from the exact + // split points and do not carry this field. + Partitioning::Range(RangePartitioning::try_new( + ordering, + split_points, + )?) + } else { + let samples = decode_split_points(&range.sample_point)?; + let range_partitioning = RangePartitioning::try_new_with_samples( + ordering, + samples, + partition_count(range.partition_count)?, + )?; + if range_partitioning.split_points() != split_points { + return internal_err!( + "Range partitioning effective split points do not match its samples and partition count" + ); + } + Partitioning::Range(range_partitioning) + } } }; Ok(Some(partitioning)) @@ -828,17 +974,6 @@ mod tests { ) -> Partitioning { Partitioning::Range(self.range(indices, split_points)) } - - fn range_partitioning_with_ordering( - &self, - ordering: LexOrdering, - split_points: Vec, - ) -> Partitioning { - Partitioning::Range( - RangePartitioning::try_new(ordering, split_points) - .expect("test range partitioning should be valid"), - ) - } } fn assert_satisfaction( @@ -1107,6 +1242,96 @@ mod tests { Ok(()) } + #[test] + fn test_range_partitioning_scales_from_samples() -> Result<()> { + let fixture = PartitioningTestFixture::int64(&["a"])?; + let samples = (10..=90) + .step_by(10) + .map(|value| int_split_point([value])) + .collect::>(); + let range = RangePartitioning::try_new_with_samples( + fixture.range_ordering([0]), + samples.clone(), + 4, + )?; + + assert_eq!(range.partition_count(), 4); + assert_eq!(range.max_partition_count(), 10); + assert_eq!(range.samples(), samples); + assert_eq!( + range.split_points(), + vec![ + int_split_point([30]), + int_split_point([50]), + int_split_point([70]), + ] + ); + assert_eq!(range.to_string(), "Range([a@0 ASC], [(30), (50), (70)], 4)"); + + let single = range.scale(1)?; + assert_eq!(single.partition_count(), 1); + assert!(single.split_points().is_empty()); + assert_eq!(single.max_partition_count(), 10); + + let restored = single.scale(single.max_partition_count())?; + assert_eq!(restored.split_points(), samples); + assert_eq!(restored.samples(), samples); + + Ok(()) + } + + #[test] + fn test_range_partitioning_rejects_invalid_partition_count() -> Result<()> { + let fixture = PartitioningTestFixture::int64(&["a"])?; + let ordering = fixture.range_ordering([0]); + let samples = vec![int_split_point([10]), int_split_point([20])]; + + let error = + RangePartitioning::try_new_with_samples(ordering.clone(), samples.clone(), 0) + .unwrap_err() + .to_string(); + assert!(error.contains("must be at least 1"), "{error}"); + + let error = + RangePartitioning::try_new_with_samples(ordering.clone(), samples.clone(), 4) + .unwrap_err() + .to_string(); + assert!(error.contains("exceeds maximum 3"), "{error}"); + + let range = RangePartitioning::try_new(ordering, samples)?; + let error = range.scale(4).unwrap_err().to_string(); + assert!(error.contains("exceeds maximum 3"), "{error}"); + + Ok(()) + } + + #[test] + fn test_range_partitioning_equality_uses_effective_split_points() -> Result<()> { + let fixture = PartitioningTestFixture::int64(&["a"])?; + let ordering = fixture.range_ordering([0]); + let sampled = RangePartitioning::try_new_with_samples( + ordering.clone(), + (10..=90) + .step_by(10) + .map(|value| int_split_point([value])) + .collect(), + 4, + )?; + let exact = RangePartitioning::try_new( + ordering, + vec![ + int_split_point([30]), + int_split_point([50]), + int_split_point([70]), + ], + )?; + + assert_eq!(sampled, exact); + assert_eq!(Partitioning::Range(sampled), Partitioning::Range(exact)); + + Ok(()) + } + #[test] fn test_range_partitioning_try_new_validates_split_points() -> Result<()> { let fixture = PartitioningTestFixture::int64(&["a", "b"])?; @@ -1156,18 +1381,29 @@ mod tests { #[test] fn test_range_partitioning_project_preserves_or_degrades() -> Result<()> { let fixture = PartitioningTestFixture::int64(&["a", "b"])?; - let range_partitioning = fixture.range_partitioning_with_ordering( - [fixture.range_sort_expr(1, SortOptions::new(true, false))].into(), - vec![int_split_point([10])], - ); + let range_partitioning = + Partitioning::Range(RangePartitioning::try_new_with_samples( + [fixture.range_sort_expr(1, SortOptions::new(true, false))].into(), + vec![ + int_split_point([30]), + int_split_point([20]), + int_split_point([10]), + ], + 2, + )?); let keep_b_mapping = ProjectionMapping::from_indices(&[1], &fixture.schema)?; let projected = range_partitioning.project(&keep_b_mapping, &fixture.eq_properties); assert_eq!( projected.to_string(), - "Range([b@0 DESC NULLS LAST], [(10)], 2)" + "Range([b@0 DESC NULLS LAST], [(20)], 2)" ); + let Partitioning::Range(projected_range) = &projected else { + panic!("expected range partitioning, got {projected:?}"); + }; + assert_eq!(projected_range.max_partition_count(), 4); + assert_eq!(projected_range.scale(4)?.split_points().len(), 3); let drop_b_mapping = ProjectionMapping::from_indices(&[0], &fixture.schema)?; let projected = @@ -1382,6 +1618,127 @@ mod ordering_proto_tests { } } +#[cfg(all(test, feature = "proto"))] +mod range_partitioning_proto_tests { + use std::sync::Arc; + + use arrow::datatypes::{DataType, Field, Schema}; + use datafusion_common::{Result, ScalarValue, SplitPoint}; + use datafusion_physical_expr_common::physical_expr::proto_decode::PhysicalExprDecodeCtx; + use datafusion_physical_expr_common::physical_expr::proto_encode::PhysicalExprEncodeCtx; + use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr}; + use datafusion_proto_models::protobuf; + + use super::{Partitioning, RangePartitioning}; + use crate::expressions::Column; + use crate::proto_test_util::{StubDecoder, StubEncoder}; + + fn sampled_partitioning() -> Result { + let ordering = LexOrdering::new([PhysicalSortExpr::new_default(Arc::new( + Column::new("a", 0), + ))]) + .expect("non-empty ordering"); + let samples = [10, 20, 30, 40, 50] + .into_iter() + .map(|value| SplitPoint::new(vec![ScalarValue::Int32(Some(value))])) + .collect(); + Ok(Partitioning::Range( + RangePartitioning::try_new_with_samples(ordering, samples, 3)?, + )) + } + + fn decode(partitioning: &protobuf::Partitioning) -> Result { + let schema = Schema::new(vec![Field::new("a", DataType::Int32, false)]); + let decoder = StubDecoder::ok(); + let decode_ctx = PhysicalExprDecodeCtx::new(&schema, &decoder); + Ok(Partitioning::try_from_proto(partitioning, &decode_ctx)? + .expect("partitioning method is present")) + } + + #[test] + fn sampled_range_partitioning_round_trip_preserves_resolution() -> Result<()> { + let partitioning = sampled_partitioning()?; + let encoder = StubEncoder::ok(); + let encode_ctx = PhysicalExprEncodeCtx::new(&encoder); + let encoded = partitioning.try_to_proto(&encode_ctx)?; + let Some(protobuf::partitioning::PartitionMethod::Range(encoded_range)) = + encoded.partition_method.as_ref() + else { + panic!("expected range partitioning"); + }; + + // Field 2 remains the effective boundary list for older readers. + assert_eq!(encoded_range.split_point.len(), 2); + assert_eq!(encoded_range.sample_point.len(), 5); + assert_eq!(encoded_range.partition_count, 3); + + let decoded = decode(&encoded)?; + let Partitioning::Range(decoded) = decoded else { + panic!("expected range partitioning"); + }; + let Partitioning::Range(original) = partitioning else { + panic!("expected range partitioning"); + }; + assert_eq!(decoded.partition_count(), original.partition_count()); + assert_eq!(decoded.split_points(), original.split_points()); + assert_eq!( + decoded.ordering()[0].options, + original.ordering()[0].options + ); + assert_eq!(decoded.samples(), original.samples()); + assert_eq!(decoded.max_partition_count(), 6); + + Ok(()) + } + + #[test] + fn legacy_range_partitioning_payload_remains_exact() -> Result<()> { + let partitioning = sampled_partitioning()?; + let encoder = StubEncoder::ok(); + let encode_ctx = PhysicalExprEncodeCtx::new(&encoder); + let mut encoded = partitioning.try_to_proto(&encode_ctx)?; + let Some(protobuf::partitioning::PartitionMethod::Range(encoded_range)) = + encoded.partition_method.as_mut() + else { + panic!("expected range partitioning"); + }; + encoded_range.sample_point.clear(); + encoded_range.partition_count = 0; + + let decoded = decode(&encoded)?; + let Partitioning::Range(decoded) = decoded else { + panic!("expected range partitioning"); + }; + assert_eq!(decoded.partition_count(), 3); + assert_eq!(decoded.max_partition_count(), 3); + assert_eq!(decoded.split_points().len(), 2); + + Ok(()) + } + + #[test] + fn sampled_range_partitioning_rejects_inconsistent_effective_points() -> Result<()> { + let partitioning = sampled_partitioning()?; + let encoder = StubEncoder::ok(); + let encode_ctx = PhysicalExprEncodeCtx::new(&encoder); + let mut encoded = partitioning.try_to_proto(&encode_ctx)?; + let Some(protobuf::partitioning::PartitionMethod::Range(encoded_range)) = + encoded.partition_method.as_mut() + else { + panic!("expected range partitioning"); + }; + encoded_range.split_point.pop(); + + let error = decode(&encoded).unwrap_err().to_string(); + assert!( + error.contains("effective split points do not match"), + "{error}" + ); + + Ok(()) + } +} + /// Partition counts are `usize` in memory and `u64` on the wire, so every /// counted [`Partitioning`] variant crosses a width boundary in both /// directions. These pin that neither crossing wraps or panics. diff --git a/datafusion/physical-plan/src/joins/utils.rs b/datafusion/physical-plan/src/joins/utils.rs index 5e8f3c5de929b..40cd6c26e4037 100644 --- a/datafusion/physical-plan/src/joins/utils.rs +++ b/datafusion/physical-plan/src/joins/utils.rs @@ -153,10 +153,11 @@ pub fn adjust_right_output_partitioning( "Offsetting range partitioning produced an empty ordering" ) })?; - Partitioning::Range(RangePartitioning::new( + Partitioning::Range(RangePartitioning::try_new_with_samples( ordering, - range.split_points().to_vec(), - )) + range.samples().to_vec(), + range.partition_count(), + )?) } result => result.clone(), }; @@ -4292,8 +4293,20 @@ mod tests { ScalarValue::Int32(Some(20)), ScalarValue::Int32(Some(50)), ]), + SplitPoint::new(vec![ + ScalarValue::Int32(Some(30)), + ScalarValue::Int32(Some(40)), + ]), + SplitPoint::new(vec![ + ScalarValue::Int32(Some(40)), + ScalarValue::Int32(Some(30)), + ]), + SplitPoint::new(vec![ + ScalarValue::Int32(Some(50)), + ScalarValue::Int32(Some(20)), + ]), ]; - let range = RangePartitioning::try_new( + let range = RangePartitioning::try_new_with_samples( LexOrdering::new([ PhysicalSortExpr::new( Arc::new(Column::new("a", 0)), @@ -4306,9 +4319,15 @@ mod tests { ]) .unwrap(), split_points.clone(), + 3, )?; let adjusted = adjust_right_output_partitioning(&Partitioning::Range(range), 3)?; + let Partitioning::Range(adjusted_range) = &adjusted else { + panic!("expected range partitioning"); + }; + assert_eq!(adjusted_range.max_partition_count(), 6); + assert_eq!(adjusted_range.samples(), split_points); let expected = Partitioning::Range(RangePartitioning::new( LexOrdering::new([ PhysicalSortExpr::new( @@ -4321,7 +4340,16 @@ mod tests { ), ]) .unwrap(), - split_points, + vec![ + SplitPoint::new(vec![ + ScalarValue::Int32(Some(20)), + ScalarValue::Int32(Some(50)), + ]), + SplitPoint::new(vec![ + ScalarValue::Int32(Some(40)), + ScalarValue::Int32(Some(30)), + ]), + ], )); assert_eq!(adjusted, expected); diff --git a/datafusion/physical-plan/src/repartition/mod.rs b/datafusion/physical-plan/src/repartition/mod.rs index 8f4b8558a592b..f7824314e2d74 100644 --- a/datafusion/physical-plan/src/repartition/mod.rs +++ b/datafusion/physical-plan/src/repartition/mod.rs @@ -1877,9 +1877,10 @@ impl ExecutionPlan for RepartitionExec { ); }; - Partitioning::Range(RangePartitioning::try_new( + Partitioning::Range(RangePartitioning::try_new_with_samples( ordering, - range_partitioning.split_points().to_vec(), + range_partitioning.samples().to_vec(), + range_partitioning.partition_count(), )?) } others => others.clone(), @@ -3023,9 +3024,18 @@ mod tests { Field::new("region", DataType::Utf8, false), Field::new("payload", DataType::UInt32, false), ])); + let ordering = + LexOrdering::new([PhysicalSortExpr::new_default(col("id", &schema)?)]) + .expect("non-empty ordering"); + let samples = [10, 20, 30, 40, 50] + .into_iter() + .map(|value| SplitPoint::new(vec![ScalarValue::UInt32(Some(value))])) + .collect(); let repartition = Arc::new(RepartitionExec::try_new( Arc::new(EmptyExec::new(Arc::clone(&schema))), - range_partitioning_on_columns(&schema, &["id"], vec![vec![10]])?, + Partitioning::Range(RangePartitioning::try_new_with_samples( + ordering, samples, 2, + )?), )?); let projection = @@ -3041,9 +3051,10 @@ mod tests { assert!(swapped_repartition.input().is::()); let range = expect_range_partitioning(swapped_repartition.partitioning()); assert_eq!(range.ordering()[0].to_string(), "id@1 ASC"); + assert_eq!(range.max_partition_count(), 6); assert_eq!( range.split_points(), - &[SplitPoint::new(vec![ScalarValue::UInt32(Some(10))])] + &[SplitPoint::new(vec![ScalarValue::UInt32(Some(30))])] ); Ok(()) diff --git a/datafusion/proto-models/proto/datafusion.proto b/datafusion/proto-models/proto/datafusion.proto index 4b63631613ae1..176ee90f49f83 100644 --- a/datafusion/proto-models/proto/datafusion.proto +++ b/datafusion/proto-models/proto/datafusion.proto @@ -1566,7 +1566,12 @@ message PhysicalHashRepartition { message PhysicalRangePartitioning { repeated PhysicalSortExprNode sort_expr = 1; + // Effective split points. Kept for compatibility with older readers. repeated PhysicalRangeSplitPoint split_point = 2; + // Maximum-resolution sample points used to derive effective split points. + repeated PhysicalRangeSplitPoint sample_point = 3; + // Zero in legacy payloads means split_point.len() + 1. + uint64 partition_count = 4; } message PhysicalRangeSplitPoint { diff --git a/datafusion/proto-models/src/generated/pbjson.rs b/datafusion/proto-models/src/generated/pbjson.rs index 9811357f1dd5d..529f1f68bb144 100644 --- a/datafusion/proto-models/src/generated/pbjson.rs +++ b/datafusion/proto-models/src/generated/pbjson.rs @@ -9532,7 +9532,7 @@ impl<'de> serde::Deserialize<'de> for HashRepartition { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } @@ -19459,7 +19459,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalHashRepartition { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } @@ -21191,6 +21191,12 @@ impl serde::Serialize for PhysicalRangePartitioning { if !self.split_point.is_empty() { len += 1; } + if !self.sample_point.is_empty() { + len += 1; + } + if self.partition_count != 0 { + len += 1; + } let mut struct_ser = serializer.serialize_struct("datafusion.PhysicalRangePartitioning", len)?; if !self.sort_expr.is_empty() { struct_ser.serialize_field("sortExpr", &self.sort_expr)?; @@ -21198,6 +21204,14 @@ impl serde::Serialize for PhysicalRangePartitioning { if !self.split_point.is_empty() { struct_ser.serialize_field("splitPoint", &self.split_point)?; } + if !self.sample_point.is_empty() { + struct_ser.serialize_field("samplePoint", &self.sample_point)?; + } + if self.partition_count != 0 { + #[allow(clippy::needless_borrow)] + #[allow(clippy::needless_borrows_for_generic_args)] + struct_ser.serialize_field("partitionCount", ToString::to_string(&self.partition_count).as_str())?; + } struct_ser.end() } } @@ -21212,12 +21226,18 @@ impl<'de> serde::Deserialize<'de> for PhysicalRangePartitioning { "sortExpr", "split_point", "splitPoint", + "sample_point", + "samplePoint", + "partition_count", + "partitionCount", ]; #[allow(clippy::enum_variant_names)] enum GeneratedField { SortExpr, SplitPoint, + SamplePoint, + PartitionCount, } impl<'de> serde::Deserialize<'de> for GeneratedField { fn deserialize(deserializer: D) -> std::result::Result @@ -21241,6 +21261,8 @@ impl<'de> serde::Deserialize<'de> for PhysicalRangePartitioning { match value { "sortExpr" | "sort_expr" => Ok(GeneratedField::SortExpr), "splitPoint" | "split_point" => Ok(GeneratedField::SplitPoint), + "samplePoint" | "sample_point" => Ok(GeneratedField::SamplePoint), + "partitionCount" | "partition_count" => Ok(GeneratedField::PartitionCount), _ => Err(serde::de::Error::unknown_field(value, FIELDS)), } } @@ -21262,6 +21284,8 @@ impl<'de> serde::Deserialize<'de> for PhysicalRangePartitioning { { let mut sort_expr__ = None; let mut split_point__ = None; + let mut sample_point__ = None; + let mut partition_count__ = None; while let Some(k) = map_.next_key()? { match k { GeneratedField::SortExpr => { @@ -21276,11 +21300,27 @@ impl<'de> serde::Deserialize<'de> for PhysicalRangePartitioning { } split_point__ = Some(map_.next_value()?); } + GeneratedField::SamplePoint => { + if sample_point__.is_some() { + return Err(serde::de::Error::duplicate_field("samplePoint")); + } + sample_point__ = Some(map_.next_value()?); + } + GeneratedField::PartitionCount => { + if partition_count__.is_some() { + return Err(serde::de::Error::duplicate_field("partitionCount")); + } + partition_count__ = + Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) + ; + } } } Ok(PhysicalRangePartitioning { sort_expr: sort_expr__.unwrap_or_default(), split_point: split_point__.unwrap_or_default(), + sample_point: sample_point__.unwrap_or_default(), + partition_count: partition_count__.unwrap_or_default(), }) } } diff --git a/datafusion/proto-models/src/generated/prost.rs b/datafusion/proto-models/src/generated/prost.rs index a20632860a4c7..98d9af7f21e6b 100644 --- a/datafusion/proto-models/src/generated/prost.rs +++ b/datafusion/proto-models/src/generated/prost.rs @@ -2363,8 +2363,15 @@ pub struct PhysicalHashRepartition { pub struct PhysicalRangePartitioning { #[prost(message, repeated, tag = "1")] pub sort_expr: ::prost::alloc::vec::Vec, + /// Effective split points. Kept for compatibility with older readers. #[prost(message, repeated, tag = "2")] pub split_point: ::prost::alloc::vec::Vec, + /// Maximum-resolution sample points used to derive effective split points. + #[prost(message, repeated, tag = "3")] + pub sample_point: ::prost::alloc::vec::Vec, + /// Zero in legacy payloads means split_point.len() + 1. + #[prost(uint64, tag = "4")] + pub partition_count: u64, } #[derive(Clone, PartialEq, ::prost::Message)] pub struct PhysicalRangeSplitPoint { From 4013966d66a548c0c62cfed66bbc19d6e9c8faba Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Fri, 28 Aug 2026 21:48:36 -0700 Subject: [PATCH 2/9] fix: align generated protobuf and FFI fixtures --- datafusion/ffi/src/tests/mod.rs | 15 +++++---------- datafusion/ffi/tests/ffi_execution_plan.rs | 2 +- datafusion/proto-models/src/generated/pbjson.rs | 6 +++--- 3 files changed, 9 insertions(+), 14 deletions(-) diff --git a/datafusion/ffi/src/tests/mod.rs b/datafusion/ffi/src/tests/mod.rs index 3aba5860ca735..9f166ec0913b1 100644 --- a/datafusion/ffi/src/tests/mod.rs +++ b/datafusion/ffi/src/tests/mod.rs @@ -123,8 +123,6 @@ pub struct ForeignLibraryModule { pub create_exec_with_statistics: extern "C" fn() -> FFI_ExecutionPlan, - pub create_exec_with_range_partitioning: extern "C" fn() -> FFI_ExecutionPlan, - pub create_table_with_statistics: extern "C" fn(codec: FFI_LogicalExtensionCodec) -> FFI_TableProvider, @@ -232,12 +230,6 @@ pub fn make_test_statistics() -> Statistics { } pub(crate) extern "C" fn create_exec_with_statistics() -> FFI_ExecutionPlan { - let schema = create_test_schema(); - let plan = Arc::new(EmptyExec::new(schema).with_statistics(make_test_statistics())); - FFI_ExecutionPlan::new(plan, None) -} - -pub(crate) extern "C" fn create_exec_with_range_partitioning() -> FFI_ExecutionPlan { let schema = create_test_schema(); let ordering = LexOrdering::new([PhysicalSortExpr::new_default(Arc::new(Column::new("a", 0)))]) @@ -250,7 +242,11 @@ pub(crate) extern "C" fn create_exec_with_range_partitioning() -> FFI_ExecutionP RangePartitioning::try_new_with_samples(ordering, samples, 3) .expect("valid sampled range partitioning"), ); - let plan = Arc::new(EmptyExec::new(schema).with_partitioning(partitioning)); + let plan = Arc::new( + EmptyExec::new(schema) + .with_statistics(make_test_statistics()) + .with_partitioning(partitioning), + ); FFI_ExecutionPlan::new(plan, None) } @@ -384,7 +380,6 @@ pub extern "C" fn datafusion_ffi_get_module() -> ForeignLibraryModule { create_exec_with_expressions, create_exec_with_dynamic_expressions, create_exec_with_statistics, - create_exec_with_range_partitioning, create_table_with_statistics, create_physical_optimizer_rule: physical_optimizer::create_physical_optimizer_rule, diff --git a/datafusion/ffi/tests/ffi_execution_plan.rs b/datafusion/ffi/tests/ffi_execution_plan.rs index e0a41ffced6e6..10b33f83bea70 100644 --- a/datafusion/ffi/tests/ffi_execution_plan.rs +++ b/datafusion/ffi/tests/ffi_execution_plan.rs @@ -71,7 +71,7 @@ mod tests { #[test] fn test_ffi_range_partitioning_cross_library() -> Result<(), DataFusionError> { let module = get_module()?; - let plan = (module.create_exec_with_range_partitioning)(); + let plan = (module.create_exec_with_statistics)(); let plan: Arc = (&plan).try_into()?; let Partitioning::Range(range) = plan.properties().output_partitioning() else { panic!("expected range partitioning"); diff --git a/datafusion/proto-models/src/generated/pbjson.rs b/datafusion/proto-models/src/generated/pbjson.rs index 529f1f68bb144..d639a79c3c929 100644 --- a/datafusion/proto-models/src/generated/pbjson.rs +++ b/datafusion/proto-models/src/generated/pbjson.rs @@ -9532,7 +9532,7 @@ impl<'de> serde::Deserialize<'de> for HashRepartition { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } @@ -19459,7 +19459,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalHashRepartition { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } @@ -21310,7 +21310,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalRangePartitioning { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } From c7672b7af9eac5e185df181c287e33997b691ab5 Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Fri, 4 Sep 2026 23:57:04 -0700 Subject: [PATCH 3/9] fix: show range partitioning scale capacity --- datafusion/physical-expr/src/partitioning.rs | 16 ++++++++++++---- 1 file changed, 12 insertions(+), 4 deletions(-) diff --git a/datafusion/physical-expr/src/partitioning.rs b/datafusion/physical-expr/src/partitioning.rs index 6d5cdcb7842d3..f2bfb3b5bd1e1 100644 --- a/datafusion/physical-expr/src/partitioning.rs +++ b/datafusion/physical-expr/src/partitioning.rs @@ -366,11 +366,15 @@ impl Display for RangePartitioning { let split_points = format_range_split_points(&self.split_points); write!( f, - "Range([{}], [{}], {})", + "Range([{}], [{}], {}", self.ordering, split_points, self.partition_count() - ) + )?; + if self.max_partition_count() != self.partition_count() { + write!(f, ", max {}", self.max_partition_count())?; + } + write!(f, ")") } } @@ -1266,12 +1270,16 @@ mod tests { int_split_point([70]), ] ); - assert_eq!(range.to_string(), "Range([a@0 ASC], [(30), (50), (70)], 4)"); + assert_eq!( + range.to_string(), + "Range([a@0 ASC], [(30), (50), (70)], 4, max 10)" + ); let single = range.scale(1)?; assert_eq!(single.partition_count(), 1); assert!(single.split_points().is_empty()); assert_eq!(single.max_partition_count(), 10); + assert_eq!(single.to_string(), "Range([a@0 ASC], [], 1, max 10)"); let restored = single.scale(single.max_partition_count())?; assert_eq!(restored.split_points(), samples); @@ -1397,7 +1405,7 @@ mod tests { range_partitioning.project(&keep_b_mapping, &fixture.eq_properties); assert_eq!( projected.to_string(), - "Range([b@0 DESC NULLS LAST], [(20)], 2)" + "Range([b@0 DESC NULLS LAST], [(20)], 2, max 4)" ); let Partitioning::Range(projected_range) = &projected else { panic!("expected range partitioning, got {projected:?}"); From 2c70f0a1d0535de484b24a0d3d4cab5a4fde0d85 Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Sun, 13 Sep 2026 21:32:55 -0700 Subject: [PATCH 4/9] fix: regenerate range partitioning JSON bindings --- datafusion/proto-models/src/generated/pbjson.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/datafusion/proto-models/src/generated/pbjson.rs b/datafusion/proto-models/src/generated/pbjson.rs index 7f5a17b4b9382..a30c748a1d8f9 100644 --- a/datafusion/proto-models/src/generated/pbjson.rs +++ b/datafusion/proto-models/src/generated/pbjson.rs @@ -21933,7 +21933,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalRangePartitioning { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } From e6b7fb29138336fe0cd23b9f5e129e4aad26ad29 Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Mon, 14 Sep 2026 18:22:23 -0700 Subject: [PATCH 5/9] fix: refine range scaling errors and layout compatibility --- datafusion/physical-expr/src/lib.rs | 1 + datafusion/physical-expr/src/partitioning.rs | 191 +++++++++++++++--- .../enforce_distribution.rs | 140 ++++++++++--- datafusion/proto/src/physical_plan/mod.rs | 12 +- datafusion/proto/tests/cases/plans/misc.rs | 10 +- datafusion/proto/tests/cases/plans/sources.rs | 16 +- .../library-user-guide/upgrading/56.0.0.md | 18 ++ 7 files changed, 313 insertions(+), 75 deletions(-) diff --git a/datafusion/physical-expr/src/lib.rs b/datafusion/physical-expr/src/lib.rs index 80e9f88b510ed..f9b12a5335547 100644 --- a/datafusion/physical-expr/src/lib.rs +++ b/datafusion/physical-expr/src/lib.rs @@ -65,6 +65,7 @@ pub use equivalence::{ pub use expressions::{DynamicFilterTracker, DynamicFilterTracking}; pub use partitioning::{ Distribution, Partitioning, PartitioningSatisfaction, RangePartitioning, + RangePartitioningScaleError, }; pub use physical_expr::{ add_offset_to_expr, add_offset_to_physical_sort_exprs, create_lex_ordering, diff --git a/datafusion/physical-expr/src/partitioning.rs b/datafusion/physical-expr/src/partitioning.rs index 851256e9db758..af50841b9459b 100644 --- a/datafusion/physical-expr/src/partitioning.rs +++ b/datafusion/physical-expr/src/partitioning.rs @@ -23,7 +23,7 @@ use crate::{ }; use arrow::datatypes::Schema; pub use datafusion_common::SplitPoint; -use datafusion_common::{Result, validate_range_split_points}; +use datafusion_common::{DataFusionError, Result, validate_range_split_points}; use datafusion_physical_expr_common::physical_expr::format_physical_expr_list; use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr}; #[cfg(feature = "proto")] @@ -204,9 +204,8 @@ impl Display for Partitioning { /// partition 2: keys at/after (2023, Allston) /// ``` /// -/// NOTE: Optimizer and execution behavior for this partitioning is intentionally -/// not implemented and will be introduced incrementally. See -/// . +/// Equality includes retained samples, since they determine which future scales +/// are possible. Use [`Self::has_same_layout`] to compare only the current layout. #[derive(Debug, Clone, PartialEq)] pub struct RangePartitioning { /// Ordered partitioning key. @@ -215,37 +214,76 @@ pub struct RangePartitioning { samples: Arc<[SplitPoint]>, /// Effective boundaries for the current partition count. split_points: Arc<[SplitPoint]>, - /// Number of effective partitions. - partition_count: usize, +} + +/// Why a [`RangePartitioning`] cannot be scaled to a requested partition count. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum RangePartitioningScaleError { + /// A partitioning must contain at least one partition. + ZeroPartitions, + /// The retained samples cannot support the requested number of partitions. + /// Callers may retain the current layout or choose another partitioning. + InsufficientSamples { + /// Requested number of partitions. + target_partitions: usize, + /// Largest partition count supported by the retained samples. + max_partitions: usize, + }, +} + +impl Display for RangePartitioningScaleError { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + match self { + Self::ZeroPartitions => { + write!(f, "Range partitioning partition count must be at least 1") + } + Self::InsufficientSamples { + target_partitions, + max_partitions, + } => write!( + f, + "Range partitioning partition count {target_partitions} exceeds maximum {max_partitions}" + ), + } + } +} + +impl std::error::Error for RangePartitioningScaleError {} + +impl From for DataFusionError { + fn from(error: RangePartitioningScaleError) -> Self { + Self::Plan(error.to_string()) + } } impl RangePartitioning { /// Creates range partitioning metadata without validating split points. /// - /// Use [`Self::try_new`] to validate the contract documented on - /// [`RangePartitioning`]. + /// Prefer [`Self::try_new_with_samples`] to validate the boundaries and retain + /// additional samples for scaling up. [`Self::try_new`] remains available for + /// validated exact boundaries. + #[deprecated( + since = "56.0.0", + note = "Use RangePartitioning::try_new_with_samples instead" + )] pub fn new(ordering: LexOrdering, split_points: Vec) -> Self { - let partition_count = split_points.len() + 1; let split_points: Arc<[SplitPoint]> = Arc::from(split_points); Self { ordering, samples: Arc::clone(&split_points), split_points, - partition_count, } } /// Creates range partitioning metadata and validates split point shape and /// ordering. + /// + /// The exact boundaries are also the retained samples. This allows scaling + /// down and back up to the original count, but not beyond it. Prefer + /// [`Self::try_new_with_samples`] when additional sample points are available. pub fn try_new(ordering: LexOrdering, split_points: Vec) -> Result { - validate_range_split_points( - &split_points, - &ordering - .iter() - .map(|sort_expr| sort_expr.options) - .collect::>(), - )?; - Ok(Self::new(ordering, split_points)) + let partition_count = split_points.len() + 1; + Self::try_new_with_samples(ordering, split_points, partition_count) } /// Creates sample-backed range partitioning and validates the sample shape, @@ -254,6 +292,12 @@ impl RangePartitioning { /// `partition_count` must be at least one and no larger than /// `samples.len() + 1`. When it is smaller than that maximum, the samples /// are evenly down-sampled to derive the effective split points. + /// + /// Retain at least `maximum_expected_partitions - 1` samples to support that + /// many partitions later. For example, supplying `4 * partition_count` + /// samples provides capacity for up to `4 * partition_count + 1` partitions. + /// Choose the sampling factor for the workload; small inputs may not have + /// enough distinct values, and callers must handle insufficient capacity. pub fn try_new_with_samples( ordering: LexOrdering, samples: Vec, @@ -273,7 +317,6 @@ impl RangePartitioning { ordering, samples, split_points, - partition_count, }) } @@ -294,7 +337,7 @@ impl RangePartitioning { /// Returns the number of partitions. pub fn partition_count(&self) -> usize { - self.partition_count + self.split_points.len() + 1 } /// Returns the largest partition count supported by the stored samples. @@ -302,20 +345,37 @@ impl RangePartitioning { self.samples.len() + 1 } + /// Whether two range partitionings have the same current key ordering and + /// effective boundaries, irrespective of their retained samples. + /// + /// This does not imply equal scaling capacity. In particular, a plan that + /// combines inputs must not use this comparison to inherit one input's samples + /// for all inputs. Structural equality is required for that use case. + pub fn has_same_layout(&self, other: &Self) -> bool { + self.ordering == other.ordering && self.split_points == other.split_points + } + /// Returns this range partitioning scaled to `target_partitions`. /// /// Scaling retains the original samples, so a range partitioning that was /// scaled down can later be scaled back up to [`Self::max_partition_count`]. - pub fn scale(&self, target_partitions: usize) -> Result { + /// Insufficient samples are an expected limitation, reported separately from + /// an invalid zero count by [`RangePartitioningScaleError`]. The caller chooses + /// whether to retain the existing layout or fall back to another partitioning. + /// This method changes metadata only; the caller must ensure that the actual + /// row distribution matches the resulting boundaries. + pub fn scale( + &self, + target_partitions: usize, + ) -> Result { validate_range_partition_count(target_partitions, self.max_partition_count())?; - if target_partitions == self.partition_count { + if target_partitions == self.partition_count() { return Ok(self.clone()); } Ok(Self { ordering: self.ordering.clone(), samples: Arc::clone(&self.samples), split_points: downsample_split_points(&self.samples, target_partitions), - partition_count: target_partitions, }) } @@ -351,7 +411,6 @@ impl RangePartitioning { ordering, samples: Arc::clone(&self.samples), split_points: Arc::clone(&self.split_points), - partition_count: self.partition_count, }) } @@ -402,7 +461,6 @@ impl RangePartitioning { ordering: new_ordering, samples: Arc::clone(&self.samples), split_points: Arc::clone(&self.split_points), - partition_count: self.partition_count, }) } } @@ -448,16 +506,15 @@ fn downsample_split_points( fn validate_range_partition_count( partition_count: usize, max_partition_count: usize, -) -> Result<()> { +) -> Result<(), RangePartitioningScaleError> { if partition_count == 0 { - return datafusion_common::plan_err!( - "Range partitioning partition count must be at least 1" - ); + return Err(RangePartitioningScaleError::ZeroPartitions); } if partition_count > max_partition_count { - return datafusion_common::plan_err!( - "Range partitioning partition count {partition_count} exceeds maximum {max_partition_count}" - ); + return Err(RangePartitioningScaleError::InsufficientSamples { + target_partitions: partition_count, + max_partitions: max_partition_count, + }); } Ok(()) } @@ -1412,6 +1469,8 @@ mod tests { )?; assert_eq!(sampled.split_points(), exact.split_points()); + assert!(sampled.has_same_layout(&exact)); + assert!(exact.has_same_layout(&sampled)); assert!(sampled.scale(10).is_ok()); assert!(exact.scale(10).is_err()); assert_ne!(sampled, exact); @@ -1420,6 +1479,74 @@ mod tests { Ok(()) } + #[test] + fn test_range_partitioning_layout_requires_keys_options_and_boundaries() -> Result<()> + { + let fixture = PartitioningTestFixture::int64(&["a", "b"])?; + let range = RangePartitioning::try_new_with_samples( + fixture.range_ordering([0]), + vec![int_split_point([10])], + 2, + )?; + assert!(range.has_same_layout(&range.clone())); + let different_boundary = RangePartitioning::try_new_with_samples( + fixture.range_ordering([0]), + vec![int_split_point([20])], + 2, + )?; + assert!(!range.has_same_layout(&different_boundary)); + assert!(!range.has_same_layout(&range.scale(1)?)); + let singleton = range.scale(1)?; + for ordering in [ + fixture.range_ordering([1]), + [fixture.range_sort_expr(0, SortOptions::new(true, false))].into(), + [fixture.range_sort_expr(0, SortOptions::new(false, false))].into(), + ] { + let other = RangePartitioning::try_new_with_samples(ordering, vec![], 1)?; + assert!(!singleton.has_same_layout(&other)); + } + Ok(()) + } + + #[test] + fn test_range_partitioning_scaling_errors_are_distinguishable() -> Result<()> { + let fixture = PartitioningTestFixture::int64(&["a"])?; + let range = RangePartitioning::try_new_with_samples( + fixture.range_ordering([0]), + vec![int_split_point([10])], + 1, + )?; + assert_eq!( + range.scale(0), + Err(RangePartitioningScaleError::ZeroPartitions) + ); + for target_partitions in [3, usize::MAX] { + let error = RangePartitioningScaleError::InsufficientSamples { + target_partitions, + max_partitions: 2, + }; + assert_eq!(range.scale(target_partitions), Err(error)); + assert!(matches!( + DataFusionError::from(error), + DataFusionError::Plan(_) + )); + } + assert_eq!(range.scale(1)?, range); + assert_eq!(range.scale(2)?.scale(1)?, range); + assert_eq!(range.scale(2)?.partition_count(), 2); + let empty = RangePartitioning::try_new_with_samples( + fixture.range_ordering([0]), + vec![], + 1, + )?; + assert_eq!(empty.scale(1)?, empty); + assert!(matches!( + empty.scale(2), + Err(RangePartitioningScaleError::InsufficientSamples { .. }) + )); + Ok(()) + } + #[test] fn test_range_partitioning_try_new_validates_split_points() -> Result<()> { let fixture = PartitioningTestFixture::int64(&["a", "b"])?; diff --git a/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs b/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs index 001a251c03319..979f57ad8d9bb 100644 --- a/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs +++ b/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs @@ -48,7 +48,7 @@ use datafusion_physical_expr::expressions::{Column, NoOp}; use datafusion_physical_expr::utils::map_columns_before_projection; use datafusion_physical_expr::{ EquivalenceProperties, OrderingRequirements, PhysicalExpr, PhysicalExprRef, - physical_exprs_equal, + RangePartitioningScaleError, physical_exprs_equal, }; use datafusion_physical_plan::ExecutionPlanProperties; use datafusion_physical_plan::aggregates::{ @@ -1305,7 +1305,9 @@ fn enforce_distribution_relationships( .map(|s| s.is_satisfied()) .unwrap_or(false) } - (Partitioning::Range(r1), Partitioning::Range(r2)) if r1 == r2 => true, + (Partitioning::Range(r1), Partitioning::Range(r2)) => { + r1.has_same_layout(r2) + } _ => false, }; @@ -1594,33 +1596,35 @@ pub fn ensure_distribution( // A satisfying range layout remains useful to the // co-partitioning pass even when its samples cannot // support the preferred degree of parallelism. - if target_partitions <= range.max_partition_count() { - let scaled = Partitioning::Range( - range.scale(target_partitions)?, - ); - // A single partition satisfies any key requirement, - // but scaling it must still use compatible keys. - if scaled - .satisfaction( - &requirement, - child.plan.equivalence_properties(), - false, - ) - .is_satisfied() - { - scaled_native_range = - !child.plan.is::(); - Some(scaled) - } else { - Some( - requirement - .clone() - .create_partitioning(target_partitions), - ) + match range.scale(target_partitions) { + Ok(range) => { + let scaled = Partitioning::Range(range); + // A single partition satisfies any key requirement, + // but scaling it must still use compatible keys. + if scaled + .satisfaction( + &requirement, + child.plan.equivalence_properties(), + false, + ) + .is_satisfied() + { + scaled_native_range = + !child.plan.is::(); + Some(scaled) + } else { + Some( + requirement + .clone() + .create_partitioning(target_partitions), + ) + } + } + Err(RangePartitioningScaleError::InsufficientSamples { .. }) => { + preserved_unscalable_range = true; + None } - } else { - preserved_unscalable_range = true; - None + Err(error) => return Err(error.into()), } } _ => Some( @@ -1853,3 +1857,83 @@ fn update_children(mut dist_context: DistributionContext) -> Result Result<()> { + let schema = + Arc::new(Schema::new(vec![Field::new("key", DataType::Int64, false)])); + let key = Arc::new(Column::new("key", 0)) as Arc; + let ordering = [PhysicalSortExpr::new_default(Arc::clone(&key))].into(); + let samples = (10..=50) + .step_by(10) + .map(|value| SplitPoint::new(vec![ScalarValue::Int64(Some(value))])) + .collect(); + let sampled = RangePartitioning::try_new_with_samples(ordering, samples, 3)?; + let exact = RangePartitioning::try_new_with_samples( + sampled.ordering().clone(), + sampled.split_points().to_vec(), + 3, + )?; + let requirement = Distribution::KeyPartitioned(vec![key]); + let requirements = + InputDistributionRequirements::co_partitioned(vec![requirement.clone(); 3]); + for reverse in [false, true] { + let ranges = if reverse { + [exact.clone(), sampled.clone()] + } else { + [sampled.clone(), exact.clone()] + }; + let mut children = ranges + .into_iter() + .map(Partitioning::Range) + .chain([Partitioning::RoundRobinBatch(3)]) + .enumerate() + .map(|(index, partitioning)| { + let plan = Arc::new(RepartitionExec::try_new( + Arc::new(EmptyExec::new(Arc::clone(&schema))), + partitioning, + )?) as Arc; + Ok(DistributionChildState { + // The first input represents a native layout already scaled + // earlier in EnsureRequirements, and is the reference. + scaled_native_range: index == 0, + preserved_unscalable_range: false, + context: DistributionContext::new_default(plan), + required_input_ordering: None, + maintains_input_order: false, + requirement: requirement.clone(), + }) + }) + .collect::>>()?; + let original = children + .iter() + .map(|c| Arc::clone(&c.context.plan)) + .collect::>(); + enforce_distribution_relationships("test", &requirements, &mut children, 3)?; + assert!(Arc::ptr_eq(&original[0], &children[0].context.plan)); + assert!( + Arc::ptr_eq(&original[1], &children[1].context.plan), + "matching boundaries must not cause another exchange because samples differ" + ); + assert!(!Arc::ptr_eq(&original[2], &children[2].context.plan)); + let plans = children + .iter() + .map(|c| c.context.plan.as_ref()) + .collect::>(); + assert!( + requirements + .unsatisfied_co_partitioned_children("test", &plans)? + .is_empty() + ); + } + Ok(()) + } +} diff --git a/datafusion/proto/src/physical_plan/mod.rs b/datafusion/proto/src/physical_plan/mod.rs index 940c732f81318..2b349035b3292 100644 --- a/datafusion/proto/src/physical_plan/mod.rs +++ b/datafusion/proto/src/physical_plan/mod.rs @@ -301,10 +301,14 @@ mod file_scan_config_serde { Column::new("value", 0), ))]) .expect("single expression ordering"); - Partitioning::Range(RangePartitioning::new( - ordering, - vec![SplitPoint::new(vec![ScalarValue::Int32(Some(10))])], - )) + Partitioning::Range( + RangePartitioning::try_new_with_samples( + ordering, + vec![SplitPoint::new(vec![ScalarValue::Int32(Some(10))])], + 2, + ) + .unwrap(), + ) } fn decode_source(conf: &protobuf::FileScanExecConf) -> Result> { diff --git a/datafusion/proto/tests/cases/plans/misc.rs b/datafusion/proto/tests/cases/plans/misc.rs index a8b5102da3fdc..65c99a3fca7ba 100644 --- a/datafusion/proto/tests/cases/plans/misc.rs +++ b/datafusion/proto/tests/cases/plans/misc.rs @@ -449,10 +449,12 @@ fn roundtrip_repartition_preserve_order() -> Result<()> { fn roundtrip_range_partitioning() -> Result<()> { let schema = Arc::new(Schema::new(vec![Field::new("a", DataType::Int64, false)])); let input = Arc::new(EmptyExec::new(Arc::clone(&schema))); - let range_partitioning = Partitioning::Range(RangePartitioning::new( - [PhysicalSortExpr::new_default(col("a", &schema)?)].into(), - vec![SplitPoint::new(vec![ScalarValue::Int64(Some(10))])], - )); + let range_partitioning = + Partitioning::Range(RangePartitioning::try_new_with_samples( + [PhysicalSortExpr::new_default(col("a", &schema)?)].into(), + vec![SplitPoint::new(vec![ScalarValue::Int64(Some(10))])], + 2, + )?); // RepartitionExec is used only to carry the partitioning through proto. // Executing range repartitioning is intentionally unsupported. let repartition = RepartitionExec::try_new(input, range_partitioning)?; diff --git a/datafusion/proto/tests/cases/plans/sources.rs b/datafusion/proto/tests/cases/plans/sources.rs index 7aff89eb3fdf7..b2562a008ac65 100644 --- a/datafusion/proto/tests/cases/plans/sources.rs +++ b/datafusion/proto/tests/cases/plans/sources.rs @@ -1046,13 +1046,15 @@ fn roundtrip_parquet_exec_range_output_partitioning() -> Result<()> { let file_schema = Arc::new(Schema::new(vec![Field::new("col", DataType::Int32, false)])); let file_source = Arc::new(ParquetSource::new(Arc::clone(&file_schema))); - let output_partitioning = Partitioning::Range(RangePartitioning::new( - LexOrdering::new(vec![PhysicalSortExpr::new_default(Arc::new(Column::new( - "col", 0, - )))]) - .unwrap(), - vec![SplitPoint::new(vec![ScalarValue::Int32(Some(10))])], - )); + let output_partitioning = + Partitioning::Range(RangePartitioning::try_new_with_samples( + LexOrdering::new(vec![PhysicalSortExpr::new_default(Arc::new(Column::new( + "col", 0, + )))]) + .unwrap(), + vec![SplitPoint::new(vec![ScalarValue::Int32(Some(10))])], + 2, + )?); let scan_config = FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), file_source) .with_file_groups(vec![ diff --git a/docs/source/library-user-guide/upgrading/56.0.0.md b/docs/source/library-user-guide/upgrading/56.0.0.md index 8de2abc5a05fc..28e5e40607792 100644 --- a/docs/source/library-user-guide/upgrading/56.0.0.md +++ b/docs/source/library-user-guide/upgrading/56.0.0.md @@ -25,6 +25,24 @@ in this section pertains to features and changes that have already been merged to the main branch and are awaiting release in this version. +### Range partitioning supports retained samples and scaling + +`RangePartitioning::new` is deprecated. Prefer the validated +`RangePartitioning::try_new_with_samples(ordering, samples, partition_count)`. +For exact boundaries, pass `split_points.len() + 1` as the count, or keep using +the validated `try_new` constructor. Retain additional samples when future +scaling above that count is required; scaling does not create new sample values. + +`scale` returns `Result`. +Callers can match `InsufficientSamples` to choose a fallback separately from +the invalid `ZeroPartitions` case. Equality includes retained samples; use +`has_same_layout` to compare the current key ordering and effective boundaries. + +`FFI_RangePartitioning` has a new layout; rebuild FFI consumers against the +compatible release. The generated `PhysicalRangePartitioning` protobuf struct +also gains fields, affecting exhaustive struct literals. Older protobuf payloads +remain readable, and effective split points remain available to older readers. + ### `ForeignSession::create_physical_plan` is unsupported `ForeignSession::create_physical_plan` no longer forwards to the library that From 6f28f9399d8d4f04d01cd20afb2dcdbf33a2144f Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Mon, 14 Sep 2026 18:25:40 -0700 Subject: [PATCH 6/9] docs: avoid upgrade guide conflict with main --- .../library-user-guide/upgrading/56.0.0.md | 34 +++++++++---------- 1 file changed, 17 insertions(+), 17 deletions(-) diff --git a/docs/source/library-user-guide/upgrading/56.0.0.md b/docs/source/library-user-guide/upgrading/56.0.0.md index 28e5e40607792..5934ba9e56766 100644 --- a/docs/source/library-user-guide/upgrading/56.0.0.md +++ b/docs/source/library-user-guide/upgrading/56.0.0.md @@ -25,6 +25,23 @@ in this section pertains to features and changes that have already been merged to the main branch and are awaiting release in this version. +### `ForeignSession::create_physical_plan` is unsupported + +`ForeignSession::create_physical_plan` no longer forwards to the library that +owns the session. It now returns a `NotImplemented` error because forwarding can +re-enter an installed foreign planner, and the execution-plan handle returned by +the old callback cannot restore local Rust type identities for downcasting. +The original `FFI_SessionRef` callback slot remains in place for ABI compatibility +with DataFusion 55 consumers, but calling that callback returns the same error. + +The session-owning library should instead export its original planner as a +`datafusion_ffi::query_planner::FFI_QueryPlanner` before installing a foreign +planner. The foreign planner can retain and invoke that handle to receive a +serialized physical plan reconstructed with local type identities. See the +`datafusion_ffi::query_planner` module documentation for the complete delegation +pattern. `ForeignSession::query_planner`, `optimize`, and `physical_optimizers` +continue to forward to the owning session across the FFI boundary. + ### Range partitioning supports retained samples and scaling `RangePartitioning::new` is deprecated. Prefer the validated @@ -43,23 +60,6 @@ compatible release. The generated `PhysicalRangePartitioning` protobuf struct also gains fields, affecting exhaustive struct literals. Older protobuf payloads remain readable, and effective split points remain available to older readers. -### `ForeignSession::create_physical_plan` is unsupported - -`ForeignSession::create_physical_plan` no longer forwards to the library that -owns the session. It now returns a `NotImplemented` error because forwarding can -re-enter an installed foreign planner, and the execution-plan handle returned by -the old callback cannot restore local Rust type identities for downcasting. -The original `FFI_SessionRef` callback slot remains in place for ABI compatibility -with DataFusion 55 consumers, but calling that callback returns the same error. - -The session-owning library should instead export its original planner as a -`datafusion_ffi::query_planner::FFI_QueryPlanner` before installing a foreign -planner. The foreign planner can retain and invoke that handle to receive a -serialized physical plan reconstructed with local type identities. See the -`datafusion_ffi::query_planner` module documentation for the complete delegation -pattern. `ForeignSession::query_planner`, `optimize`, and `physical_optimizers` -continue to forward to the owning session across the FFI boundary. - ### `GroupColumn` now requires `values_preserving` Custom implementations of the public `GroupColumn` trait must implement From 8bdc21c19f39526eafa96a7cb3f65d32d048a357 Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Sat, 19 Sep 2026 14:11:26 -0700 Subject: [PATCH 7/9] docs: update range scaling migration guidance Signed-off-by: goutamadwant --- docs/source/library-user-guide/upgrading/56.0.0.md | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) diff --git a/docs/source/library-user-guide/upgrading/56.0.0.md b/docs/source/library-user-guide/upgrading/56.0.0.md index 156fb53d97fca..f944daf322ff1 100644 --- a/docs/source/library-user-guide/upgrading/56.0.0.md +++ b/docs/source/library-user-guide/upgrading/56.0.0.md @@ -127,10 +127,11 @@ For exact boundaries, pass `split_points.len() + 1` as the count, or keep using the validated `try_new` constructor. Retain additional samples when future scaling above that count is required; scaling does not create new sample values. -`scale` returns `Result`. -Callers can match `InsufficientSamples` to choose a fallback separately from -the invalid `ZeroPartitions` case. Equality includes retained samples; use -`has_same_layout` to compare the current key ordering and effective boundaries. +`scale` returns `Option`. It returns `None` when the target +partition count is zero or exceeds the retained sample capacity; callers can +then retain the existing layout or choose another partitioning. Equality includes +retained samples; use `has_same_layout` to compare the current key ordering and +effective boundaries. `FFI_RangePartitioning` has a new layout; rebuild FFI consumers against the compatible release. The generated `PhysicalRangePartitioning` protobuf struct From 1d2d59e69955c06d6dc8f4c82f1061e4354d959b Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Sat, 19 Sep 2026 14:45:01 -0700 Subject: [PATCH 8/9] chore: regenerate range partitioning protobuf Signed-off-by: goutamadwant --- datafusion/proto-models/src/generated/pbjson.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/datafusion/proto-models/src/generated/pbjson.rs b/datafusion/proto-models/src/generated/pbjson.rs index a0ea609622530..744429ae0aea3 100644 --- a/datafusion/proto-models/src/generated/pbjson.rs +++ b/datafusion/proto-models/src/generated/pbjson.rs @@ -10091,7 +10091,7 @@ impl<'de> serde::Deserialize<'de> for HashRepartition { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } @@ -20346,7 +20346,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalHashRepartition { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } @@ -22211,7 +22211,7 @@ impl<'de> serde::Deserialize<'de> for PhysicalRangePartitioning { if partition_count__.is_some() { return Err(serde::de::Error::duplicate_field("partitionCount")); } - partition_count__ = + partition_count__ = Some(map_.next_value::<::pbjson::private::NumberDeserialize<_>>()?.0) ; } From a4abf20425bcab4eb727b18a9ecff9a75ea15fd9 Mon Sep 17 00:00:00 2001 From: goutamadwant Date: Mon, 21 Sep 2026 00:37:24 -0700 Subject: [PATCH 9/9] fix: scale range exchanges with fresh execution state Signed-off-by: goutamadwant --- .../enforce_distribution.rs | 45 +++---- .../physical-plan/src/repartition/mod.rs | 126 +++++++++++++++++- 2 files changed, 137 insertions(+), 34 deletions(-) diff --git a/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs b/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs index 01ff9718700df..04045db3fa0a4 100644 --- a/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs +++ b/datafusion/physical-optimizer/src/ensure_requirements/enforce_distribution.rs @@ -1599,9 +1599,6 @@ pub fn ensure_distribution_with_stats( if should_add_repartition { let partitioning = match child.plan.output_partitioning() { Partitioning::Range(range) if partitioning_satisfied => { - // A satisfying range layout remains useful to the - // co-partitioning pass even when its samples cannot - // support the preferred degree of parallelism. match range.scale(target_partitions) { Some(range) => { let scaled = Partitioning::Range(range); @@ -1617,30 +1614,24 @@ pub fn ensure_distribution_with_stats( { scaled_native_range = !child.plan.is::(); - Some(scaled) + scaled } else { - Some( - requirement.clone().create_partitioning( - target_partitions, - ), - ) + requirement + .clone() + .create_partitioning(target_partitions) } } // Insufficient samples are expected. Preserve // main's policy by falling back to key // repartitioning at the requested parallelism. - None => Some( - requirement - .clone() - .create_partitioning(target_partitions), - ), + None => requirement + .clone() + .create_partitioning(target_partitions), } } - _ => Some( - requirement - .clone() - .create_partitioning(target_partitions), - ), + _ => { + requirement.clone().create_partitioning(target_partitions) + } }; // When there is an existing ordering, we preserve ordering during // repartition. This will be rolled back in the future if any of the @@ -1649,15 +1640,13 @@ pub fn ensure_distribution_with_stats( // requirements. // - Usage of order preserving variants is not desirable (per the flag // `config.optimizer.prefer_existing_sort`). - if let Some(partitioning) = partitioning { - let repartition = RepartitionExec::try_new( - Arc::clone(&child.plan), - partitioning, - )? - .with_preserve_order(); - let plan = Arc::new(repartition) as _; - child = DistributionContext::new(plan, true, vec![child]); - } + let repartition = RepartitionExec::try_new( + Arc::clone(&child.plan), + partitioning, + )? + .with_preserve_order(); + let plan = Arc::new(repartition) as _; + child = DistributionContext::new(plan, true, vec![child]); } } Distribution::UnspecifiedDistribution => { diff --git a/datafusion/physical-plan/src/repartition/mod.rs b/datafusion/physical-plan/src/repartition/mod.rs index dffb0e6293062..d2d8ec38dacc4 100644 --- a/datafusion/physical-plan/src/repartition/mod.rs +++ b/datafusion/physical-plan/src/repartition/mod.rs @@ -2053,9 +2053,18 @@ impl ExecutionPlan for RepartitionExec { new_properties.partitioning = match new_properties.partitioning { RoundRobinBatch(_) => RoundRobinBatch(target_partitions), Hash(hash, _) => Hash(hash, target_partitions), - Range(_) => { - // Number of partitions is constrained by the split points and cannot be changed - return Ok(None); + Range(range) => { + let Some(range) = range.scale(target_partitions) else { + return Ok(None); + }; + // A different layout needs its own channels, router, and metrics. + let mut repartition = + Self::try_new(Arc::clone(&self.input), Range(range))?; + if self.preserve_order { + repartition = repartition.with_preserve_order(); + } + repartition.batch_size = self.batch_size; + return Ok(Some(Arc::new(repartition))); } UnknownPartitioning(_) => UnknownPartitioning(target_partitions), }; @@ -3084,6 +3093,112 @@ mod tests { Ok(()) } + #[tokio::test] + async fn range_repartitioned_scales_with_fresh_execution_state() -> Result<()> { + let schema = test_schema(false); + let ordering = + LexOrdering::new([PhysicalSortExpr::new_default(col("c0", &schema)?)]) + .unwrap(); + let samples = [10, 20, 30] + .into_iter() + .map(|value| SplitPoint::new(vec![ScalarValue::UInt32(Some(value))])) + .collect::>(); + let partitions = [vec![5, 15, 25, 35], vec![6, 16, 26, 36]] + .into_iter() + .map(|values| -> Result<_> { + Ok(vec![RecordBatch::try_new( + Arc::clone(&schema), + vec![Arc::new(UInt32Array::from(values))], + )?]) + }) + .collect::>>()?; + + for preserve_order in [false, true] { + let source = TestMemoryExec::try_new(&partitions, Arc::clone(&schema), None)? + .try_with_sort_information(vec![ordering.clone()])?; + let source = Arc::new(TestMemoryExec::update_cache(&Arc::new(source))); + let mut exec = Arc::new( + RepartitionExec::try_new( + source, + Partitioning::Range(RangePartitioning::try_new_with_samples( + ordering.clone(), + samples.clone(), + 2, + )?), + )? + .with_batch_size(2)?, + ); + if preserve_order { + exec = Arc::new(Arc::unwrap_or_clone(exec).with_preserve_order()); + } + let context = Arc::new(TaskContext::default()); + + // Initialize the old exchange before resizing: its channels and router + // must not be reused for a different set of output boundaries. + let initial = crate::collect_partitioned( + Arc::::clone(&exec), + Arc::clone(&context), + ) + .await?; + assert_eq!( + initial + .iter() + .map(|p| partition_row_count(p)) + .sum::(), + 8 + ); + + for target in [4, 1, 4] { + let scaled = exec + .repartitioned(target, &ConfigOptions::default())? + .expect("retained samples support the requested partition count"); + let repartition = scaled.downcast_ref::().unwrap(); + let range = expect_range_partitioning(repartition.partitioning()); + assert_eq!(range.partition_count(), target); + assert_eq!(range.samples(), samples); + assert_eq!(range.ordering(), &ordering); + assert_eq!(repartition.preserve_order, preserve_order); + assert_eq!(repartition.batch_size, Some(2)); + assert_eq!( + repartition.properties().output_ordering(), + exec.properties().output_ordering() + ); + assert!(!Arc::ptr_eq(&repartition.state, &exec.state)); + assert!(Arc::ptr_eq(&repartition.input, &exec.input)); + assert_eq!(repartition.metrics().unwrap().output_rows(), None); + + let output = + crate::collect_partitioned(Arc::clone(&scaled), Arc::clone(&context)) + .await?; + assert_eq!(output.len(), target); + for (index, batches) in output.iter().enumerate() { + let mut values = collect_partition_u32_values(batches); + if !preserve_order { + values.sort_unstable(); + } + let expected = if target == 1 { + vec![5, 6, 15, 16, 25, 26, 35, 36] + } else { + vec![index as u32 * 10 + 5, index as u32 * 10 + 6] + }; + assert_eq!( + values, + expected.into_iter().map(Some).collect::>() + ); + } + assert_eq!(repartition.metrics().unwrap().output_rows(), Some(8)); + exec = Arc::new(repartition.clone()); + } + for unsupported in [0, 5] { + assert!( + exec.repartitioned(unsupported, &ConfigOptions::default())? + .is_none() + ); + } + } + Ok(()) + } + #[tokio::test] async fn range_repartition_routes_rows_desc() -> Result<()> { let schema = test_schema(false); @@ -5014,12 +5129,11 @@ mod test { })?; assert_eq!(expressions, ["c0@0"]); - // Range partition count is fixed by split points, so repartitioned() - // cannot change it to an arbitrary target. + // Scaling cannot exceed the retained sample capacity. let result = exec.repartitioned(10, &Default::default())?; assert!( result.is_none(), - "range repartitioning should not support changing partition count" + "range repartitioning should reject counts above sample capacity" ); Ok(()) }