diff --git a/datafusion/catalog-listing/src/table.rs b/datafusion/catalog-listing/src/table.rs index 6c294fe077db4..15465beda7347 100644 --- a/datafusion/catalog-listing/src/table.rs +++ b/datafusion/catalog-listing/src/table.rs @@ -25,13 +25,13 @@ use async_trait::async_trait; use datafusion_catalog::{ScanArgs, ScanResult, Session, TableProvider}; use datafusion_common::stats::{Precision, is_known_empty}; use datafusion_common::{ - Constraints, DFSchema, SchemaExt, Statistics, internal_datafusion_err, plan_err, - project_schema, + Constraints, DFSchema, SchemaExt, SplitPoint, Statistics, internal_datafusion_err, + plan_err, project_schema, }; use datafusion_datasource::file::FileSource; use datafusion_datasource::file_groups::FileGroup; use datafusion_datasource::file_scan_config::{ - FileScanConfig, FileScanConfigBuilder, output_partitioning_from_partition_fields, + FileScanConfig, FileScanConfigBuilder, range_partitioning_from_partition_fields, }; use datafusion_datasource::file_sink_config::{FileOutputMode, FileSinkConfig}; #[expect(deprecated)] @@ -68,6 +68,8 @@ pub struct ListFilesResult { pub statistics: Statistics, /// Whether files are grouped by partition values. pub grouped_by_partition: bool, + /// Boundaries between file groups, empty unless `grouped_by_partition` is true. + pub partition_split_points: Vec, } /// Built in [`TableProvider`] that reads data from one or more files as a single table. @@ -633,6 +635,7 @@ impl ListingTable { file_groups: mut partitioned_file_lists, statistics, grouped_by_partition: partitioned_by_file_group, + partition_split_points, } = self .list_files_for_scan(state, &partition_filters, statistic_file_limit) .await?; @@ -652,6 +655,7 @@ impl ListingTable { .config_options() .execution .split_file_groups_by_statistics; + let mut regrouped_by_statistics = false; match split_file_groups_by_statistics .then(|| { output_ordering.first().map(|output_ordering| { @@ -669,6 +673,7 @@ impl ListingTable { Some(Ok(new_groups)) => { if new_groups.len() <= file_group_count { partitioned_file_lists = new_groups; + regrouped_by_statistics = true; } else { log::debug!( "attempted to split file groups by statistics, but there were more file groups than target_partitions; falling back to unordered" @@ -711,14 +716,30 @@ impl ListingTable { } Some(output_partitioning) } else if partitioned_by_file_group { - // Files are grouped by partition column values: declare output - // partitioning on those columns so the optimizer can skip - // repartitioning for aggregates and joins on the partition columns. - output_partitioning_from_partition_fields( - &self.table_schema, - &table_partition_cols.clone().into(), - partitioned_file_lists.len(), - ) + // Groups are contiguous key intervals, so `Range` records where each key + // lives and the optimizer can prove co-partitioning instead of assuming it. + if regrouped_by_statistics { + log::debug!( + "file groups were re-cut by statistics, not declaring partition based output partitioning" + ); + None + } else { + debug_assert_eq!( + partition_split_points.len() + 1, + partitioned_file_lists.len() + ); + range_partitioning_from_partition_fields( + &self.table_schema, + &table_partition_cols.clone().into(), + partition_split_points, + ) + .unwrap_or_else(|e| { + log::debug!( + "could not derive range partitioning from partition-grouped file groups: {e}" + ); + None + }) + } } else { None }; @@ -953,6 +974,7 @@ impl ListingTable { file_groups: vec![], statistics: Statistics::new_unknown(&self.file_schema), grouped_by_partition: false, + partition_split_points: vec![], }); }; let (file_group, inexact_stats) = self @@ -966,21 +988,23 @@ impl ListingTable { // skip repartitioning for aggregates and joins on partition columns. let threshold = ctx.config_options().optimizer.preserve_file_partitions; - let (file_groups, grouped_by_partition) = + let (file_groups, grouped_by_partition, partition_split_points) = if threshold > 0 && !self.options.table_partition_cols.is_empty() { - let grouped = file_group.group_by_partition_values(file_group_count); + let (grouped, split_points) = file_group + .group_by_partition_values_with_split_points(file_group_count); if grouped.len() >= threshold { - (grouped, true) + (grouped, true, split_points) } else { let all_files: Vec<_> = grouped.into_iter().flat_map(|g| g.into_inner()).collect(); ( FileGroup::new(all_files).split_files(file_group_count), false, + vec![], ) } } else { - (file_group.split_files(file_group_count), false) + (file_group.split_files(file_group_count), false, vec![]) }; self.list_files_result_from_groups( @@ -988,6 +1012,7 @@ impl ListingTable { file_groups, inexact_stats, grouped_by_partition, + partition_split_points, ) } @@ -1015,6 +1040,7 @@ impl ListingTable { file_groups: vec![], statistics: Statistics::new_unknown(&self.file_schema), grouped_by_partition: false, + partition_split_points: vec![], }); }; let (file_group, inexact_stats) = @@ -1026,7 +1052,7 @@ impl ListingTable { let file_groups = self.filter_declared_file_groups_by_partition_filters(file_groups, filters)?; - self.list_files_result_from_groups(ctx, file_groups, inexact_stats, false) + self.list_files_result_from_groups(ctx, file_groups, inexact_stats, false, vec![]) } fn filter_declared_file_groups_by_partition_filters( @@ -1061,6 +1087,7 @@ impl ListingTable { file_groups: Vec, inexact_stats: bool, grouped_by_partition: bool, + partition_split_points: Vec, ) -> datafusion_common::Result { let (file_groups, stats) = compute_all_files_statistics( file_groups, @@ -1077,6 +1104,7 @@ impl ListingTable { file_groups, statistics: stats, grouped_by_partition, + partition_split_points, }) } diff --git a/datafusion/datasource/src/file_groups.rs b/datafusion/datasource/src/file_groups.rs index 84594be54b504..b7526c29cd1a7 100644 --- a/datafusion/datasource/src/file_groups.rs +++ b/datafusion/datasource/src/file_groups.rs @@ -19,8 +19,8 @@ use crate::{FileRange, PartitionedFile}; use arrow::compute::SortOptions; -use datafusion_common::Statistics; use datafusion_common::utils::compare_rows; +use datafusion_common::{ScalarValue, SplitPoint, Statistics}; use itertools::Itertools; use std::cmp::{Ordering, min}; use std::collections::{BinaryHeap, HashMap}; @@ -488,19 +488,33 @@ impl FileGroup { /// /// Note: May return fewer groups than `max_target_partitions` when the /// number of unique partition values is less than the target. - #[allow(clippy::allow_attributes, clippy::mutable_key_type)] // ScalarValue has interior mutability but is intentionally used as hash key + /// + /// Wrapper around [`Self::group_by_partition_values_with_split_points`] that discards + /// the split points. pub fn group_by_partition_values( self, max_target_partitions: usize, ) -> Vec { + self.group_by_partition_values_with_split_points(max_target_partitions) + .0 + } + + /// Groups files by partition value and returns the [`SplitPoint`]s between groups. + /// + /// Values are sorted with [`SortOptions::default()`] and cut into at most + /// `max_target_partitions` contiguous chunks so routing a file's `partition_values` + /// through the split points yields the group it was placed in. + #[allow(clippy::allow_attributes, clippy::mutable_key_type)] + pub fn group_by_partition_values_with_split_points( + self, + max_target_partitions: usize, + ) -> (Vec, Vec) { if self.is_empty() || max_target_partitions == 0 { - return vec![]; + return (vec![], vec![]); } - let mut partition_groups: HashMap< - Vec, - Vec, - > = HashMap::new(); + let mut partition_groups: HashMap, Vec> = + HashMap::new(); for file in self.files { partition_groups @@ -512,6 +526,8 @@ impl FileGroup { let num_unique_partitions = partition_groups.len(); // Sort for deterministic bucket assignment across query executions. + // Must match the ordering declared by `range_partitioning_from_partition_fields`, + // otherwise the split points would not describe the groups. let mut sorted_partitions: Vec<_> = partition_groups.into_iter().collect(); let sort_options = vec![ @@ -522,24 +538,31 @@ impl FileGroup { compare_rows(&a.0, &b.0, &sort_options).unwrap_or(Ordering::Equal) }); - if num_unique_partitions <= max_target_partitions { - sorted_partitions - .into_iter() - .map(|(_, files)| FileGroup::new(files)) - .collect() - } else { - // Merge into max_target_partitions buckets using round-robin. - // This maintains grouping by partition value as we are merging groups which already - // contain all values for a partition key. - let mut target_groups = vec![vec![]; max_target_partitions]; - - for (idx, (_, files)) in sorted_partitions.into_iter().enumerate() { - let bucket = idx % max_target_partitions; - target_groups[bucket].extend(files); + let bucket_count = min(num_unique_partitions, max_target_partitions); + let mut groups = Vec::with_capacity(bucket_count); + let mut split_points = Vec::with_capacity(bucket_count.saturating_sub(1)); + + // Every chunk is non-empty because `bucket_count <= num_unique_partitions`, so + // its first tuple is a well defined lower bound. + let mut iter = sorted_partitions.into_iter(); + for bucket in 0..bucket_count { + let start = bucket * num_unique_partitions / bucket_count; + let end = (bucket + 1) * num_unique_partitions / bucket_count; + let mut files = Vec::new(); + for offset in start..end { + let (values, group_files) = iter + .next() + .expect("contiguous chunking never runs past the sorted partitions"); + if offset == start && bucket > 0 { + split_points.push(SplitPoint::new(values)); + } + files.extend(group_files); } - - target_groups.into_iter().map(FileGroup::new).collect() + groups.push(FileGroup::new(files)); } + debug_assert_eq!(split_points.len() + 1, groups.len()); + + (groups, split_points) } } @@ -632,7 +655,9 @@ impl DerefMut for CompareByRangeSize { #[cfg(test)] mod test { use super::*; - use datafusion_common::ScalarValue; + use datafusion_physical_expr::RangePartitioning; + use datafusion_physical_expr::expressions::Column; + use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr}; /// Empty file won't get partitioned #[test] @@ -1276,8 +1301,8 @@ mod test { #[test] fn test_group_by_partition_values_more_groups_than_target() { - // Each file has a single partition value. The number of partition values > max_target_partitions, so - // they should be round-robin distributed into groups. + // More values than `max_target_partitions`, so they are cut into contiguous + // chunks whose sizes differ by at most one. let fg = FileGroup::new(vec![ pfile_with_pv("a", "p1"), pfile_with_pv("b", "p2"), @@ -1287,8 +1312,153 @@ mod test { ]); let groups = fg.group_by_partition_values(3); assert_eq!(groups.len(), 3); - assert_eq!(groups[0].len(), 2); + assert_eq!(groups[0].len(), 1); assert_eq!(groups[1].len(), 2); - assert_eq!(groups[2].len(), 1); + assert_eq!(groups[2].len(), 2); + // Groups must be contiguous key intervals in sorted order. + assert_eq!(group_pvs(&groups[0]), vec!["p1"]); + assert_eq!(group_pvs(&groups[1]), vec!["p2", "p3"]); + assert_eq!(group_pvs(&groups[2]), vec!["p4", "p5"]); + } + + fn group_pvs(group: &FileGroup) -> Vec { + group + .iter() + .map(|f| f.partition_values[0].to_string()) + .collect() + } + + /// Mirrors `RangeRouter` semantics: the partition index is the number of split points + /// that are `<=` the key under the given sort options. + fn range_route(key: &[ScalarValue], split_points: &[SplitPoint]) -> usize { + let sort_options = vec![SortOptions::default(); key.len()]; + split_points + .iter() + .take_while(|sp| { + compare_rows(sp.values(), key, &sort_options).unwrap() + != Ordering::Greater + }) + .count() + } + + #[test] + fn test_group_by_partition_values_split_points_describe_groups() { + let fg = FileGroup::new(vec![ + pfile_with_pv("a", "p1"), + pfile_with_pv("b", "p2"), + pfile_with_pv("c", "p3"), + pfile_with_pv("d", "p4"), + pfile_with_pv("e", "p5"), + pfile_with_pv("f", "p5"), + ]); + let (groups, split_points) = fg.group_by_partition_values_with_split_points(3); + assert_eq!(groups.len(), 3); + assert_eq!(split_points.len(), 2); + assert_eq!(split_points[0].values(), &[ScalarValue::from("p2")]); + assert_eq!(split_points[1].values(), &[ScalarValue::from("p4")]); + + // Every file routes (RangeRouter semantics) to the group it was placed in. + for (idx, group) in groups.iter().enumerate() { + for file in group.iter() { + assert_eq!(range_route(&file.partition_values, &split_points), idx); + } + } + } + + #[test] + fn test_group_by_partition_values_split_points_fewer_values_than_target() { + let fg = FileGroup::new(vec![ + pfile_with_pv("a", "p1"), + pfile_with_pv("b", "p1"), + pfile_with_pv("c", "p2"), + ]); + let (groups, split_points) = fg.group_by_partition_values_with_split_points(4); + assert_eq!(groups.len(), 2); + assert_eq!(split_points.len(), 1); + assert_eq!(split_points[0].values(), &[ScalarValue::from("p2")]); + for (idx, group) in groups.iter().enumerate() { + for file in group.iter() { + assert_eq!(range_route(&file.partition_values, &split_points), idx); + } + } + } + + #[test] + fn test_group_by_partition_values_split_points_single_target_partition() { + // One partition holds every value, so there is no boundary to describe. + let fg = FileGroup::new(vec![ + pfile_with_pv("a", "p1"), + pfile_with_pv("b", "p2"), + pfile_with_pv("c", "p3"), + pfile_with_pv("d", "p3"), + ]); + let (groups, split_points) = fg.group_by_partition_values_with_split_points(1); + assert_eq!(groups.len(), 1); + assert_eq!(groups[0].len(), 4); + assert!(split_points.is_empty()); + } + + fn pfile_with_pvs(path: &str, pvs: &[Option<&str>]) -> PartitionedFile { + let mut file = pfile(path, 10); + file.partition_values = pvs + .iter() + .map(|v| ScalarValue::Utf8(v.map(str::to_string))) + .collect(); + file + } + + /// Two partition columns with NULLs in either position. + fn compound_key_file_group() -> FileGroup { + FileGroup::new(vec![ + pfile_with_pvs("a", &[Some("2026"), Some("09")]), + pfile_with_pvs("b", &[Some("2026"), Some("10")]), + pfile_with_pvs("c", &[None, Some("09")]), + pfile_with_pvs("d", &[Some("2025"), None]), + pfile_with_pvs("e", &[Some("2025"), Some("12")]), + ]) + } + + #[test] + fn test_group_by_partition_values_split_points_null_and_compound_keys() { + let fg = compound_key_file_group(); + let (groups, split_points) = fg.group_by_partition_values_with_split_points(2); + assert_eq!(groups.len(), 2); + assert_eq!(split_points.len(), 1); + assert_eq!( + split_points[0].values(), + &[ScalarValue::from("2025"), ScalarValue::from("12")] + ); + for (idx, group) in groups.iter().enumerate() { + for file in group.iter() { + assert_eq!(range_route(&file.partition_values, &split_points), idx); + } + } + } + + #[test] + fn test_group_by_partition_values_split_point_containing_null() { + let fg = compound_key_file_group(); + let (groups, split_points) = fg.group_by_partition_values_with_split_points(3); + assert_eq!(groups.len(), 3); + assert_eq!( + split_points[0].values(), + &[ScalarValue::from("2025"), ScalarValue::Utf8(None)] + ); + assert_eq!( + split_points[1].values(), + &[ScalarValue::from("2026"), ScalarValue::from("09")] + ); + for (idx, group) in groups.iter().enumerate() { + for file in group.iter() { + assert_eq!(range_route(&file.partition_values, &split_points), idx); + } + } + let ordering = LexOrdering::new(vec![ + PhysicalSortExpr::new_default(Arc::new(Column::new("year", 0))), + PhysicalSortExpr::new_default(Arc::new(Column::new("month", 1))), + ]) + .unwrap(); + let range = RangePartitioning::try_new(ordering, split_points).unwrap(); + assert_eq!(range.partition_count(), 3); } } diff --git a/datafusion/datasource/src/file_scan_config/mod.rs b/datafusion/datasource/src/file_scan_config/mod.rs index 4e72b5e83bd23..6e972abe0ec37 100644 --- a/datafusion/datasource/src/file_scan_config/mod.rs +++ b/datafusion/datasource/src/file_scan_config/mod.rs @@ -39,8 +39,8 @@ use arrow::datatypes::{DataType, Schema, SchemaRef}; use datafusion_common::config::ConfigOptions; use datafusion_common::tree_node::TreeNodeRecursion; use datafusion_common::{ - Constraint, Constraints, Result, ScalarValue, Statistics, internal_datafusion_err, - internal_err, + Constraint, Constraints, Result, ScalarValue, SplitPoint, Statistics, + internal_datafusion_err, internal_err, }; use datafusion_execution::{ SendableRecordBatchStream, TaskContext, object_store::ObjectStoreUrl, @@ -52,7 +52,9 @@ use datafusion_common::stats::{Precision, is_known_empty}; use datafusion_physical_expr::expressions::{BinaryExpr, Column}; use datafusion_physical_expr::projection::{ProjectionExprs, ProjectionMapping}; use datafusion_physical_expr::utils::reassign_expr_columns; -use datafusion_physical_expr::{EquivalenceProperties, Partitioning, split_conjunction}; +use datafusion_physical_expr::{ + EquivalenceProperties, Partitioning, RangePartitioning, split_conjunction, +}; use datafusion_physical_expr_adapter::PhysicalExprAdapterFactory; use datafusion_physical_expr_common::physical_expr::{PhysicalExpr, is_volatile}; use datafusion_physical_expr_common::sort_expr::{LexOrdering, PhysicalSortExpr}; @@ -615,15 +617,54 @@ impl From for FileScanConfigBuilder { } } -/// Builds output partitioning over `partition_cols` (resolved to their indices in -/// `schema`) with `partition_count` partitions. Returns `None` when there are no -/// partition columns. Callers use this to declare the output partitioning of a scan -/// whose file groups are organized by partition column values. +/// Builds `Partitioning::Hash` over `partition_cols` (resolved to their indices in +/// `schema`). Returns `None` when there are no partition columns. +/// +/// # Deprecated +/// Use [`range_partitioning_from_partition_fields`] instead. +#[deprecated( + since = "56.0.0", + note = "Hive file groups are value partitioned, not hash partitioned. Use range_partitioning_from_partition_fields" +)] pub fn output_partitioning_from_partition_fields( schema: &Schema, partition_cols: &Fields, partition_count: usize, ) -> Option { + let exprs = partition_column_exprs(schema, partition_cols)?; + Some(Partitioning::Hash(exprs, partition_count)) +} + +/// Builds the [`Partitioning::Range`] that describes a scan whose file groups were +/// produced by [`FileGroup::group_by_partition_values_with_split_points`]. +pub fn range_partitioning_from_partition_fields( + schema: &Schema, + partition_cols: &Fields, + split_points: Vec, +) -> Result> { + let Some(exprs) = partition_column_exprs(schema, partition_cols) else { + return Ok(None); + }; + let sort_exprs = exprs + .into_iter() + .map(PhysicalSortExpr::new_default) + .collect::>(); + let Some(ordering) = LexOrdering::new(sort_exprs) else { + return Ok(None); + }; + if ordering.len() != partition_cols.len() { + return Ok(None); + } + let range = RangePartitioning::try_new(ordering, split_points)?; + Ok(Some(Partitioning::Range(range))) +} + +/// Resolves `partition_cols` to `Column` expressions against `schema`. Returns `None` +/// when there are no partition columns or one is missing. +fn partition_column_exprs( + schema: &Schema, + partition_cols: &Fields, +) -> Option>> { if partition_cols.is_empty() { return None; } @@ -637,8 +678,7 @@ pub fn output_partitioning_from_partition_fields( .position(|field| field.name() == name)?; exprs.push(Arc::new(Column::new(name, idx))); } - - Some(Partitioning::Hash(exprs, partition_count)) + Some(exprs) } fn project_output_partitioning( @@ -3118,7 +3158,6 @@ mod tests { fn test_output_partitioning_with_partition_columns() { let file_schema = aggr_test_schema(); - // Test single partition column let single_partition_col = vec![Field::new( "date", wrap_partition_type_in_dict(DataType::Utf8), @@ -3136,20 +3175,31 @@ mod tests { FileGroup::new(vec![PartitionedFile::new("f2.parquet".to_string(), 1024)]), FileGroup::new(vec![PartitionedFile::new("f3.parquet".to_string(), 1024)]), ]; - config.output_partitioning = output_partitioning_from_partition_fields( + let dict = |v: &str| wrap_partition_value_in_dict(ScalarValue::from(v)); + config.output_partitioning = range_partitioning_from_partition_fields( config.file_source.table_schema().table_schema(), config.table_partition_cols(), - config.file_groups.len(), - ); + vec![ + SplitPoint::new(vec![dict("2026-09")]), + SplitPoint::new(vec![dict("2026-10")]), + ], + ) + .unwrap(); let partitioning = config.output_partitioning(); match partitioning { - Partitioning::Hash(exprs, num_partitions) => { - assert_eq!(num_partitions, 3); - assert_eq!(exprs.len(), 1); - assert_eq!(exprs[0].downcast_ref::().unwrap().name(), "date"); + Partitioning::Range(range) => { + assert_eq!(range.partition_count(), 3); + assert_eq!(range.ordering().len(), 1); + let sort_expr = &range.ordering()[0]; + assert_eq!( + sort_expr.expr.downcast_ref::().unwrap().name(), + "date" + ); + assert_eq!(sort_expr.options, arrow::compute::SortOptions::default()); + assert_eq!(range.split_points().len(), 2); } - _ => panic!("Expected Hash partitioning"), + other => panic!("Expected Range partitioning, got {other}"), } // Test multiple partition columns @@ -3168,27 +3218,63 @@ mod tests { FileGroup::new(vec![PartitionedFile::new("f1.parquet".to_string(), 1024)]), FileGroup::new(vec![PartitionedFile::new("f2.parquet".to_string(), 1024)]), ]; - config.output_partitioning = output_partitioning_from_partition_fields( + config.output_partitioning = range_partitioning_from_partition_fields( config.file_source.table_schema().table_schema(), config.table_partition_cols(), - config.file_groups.len(), - ); + vec![SplitPoint::new(vec![dict("2026"), dict("09")])], + ) + .unwrap(); let partitioning = config.output_partitioning(); match partitioning { - Partitioning::Hash(exprs, num_partitions) => { - assert_eq!(num_partitions, 2); - assert_eq!(exprs.len(), 2); - let col_names: Vec<_> = exprs + Partitioning::Range(range) => { + assert_eq!(range.partition_count(), 2); + let col_names: Vec<_> = range + .ordering() .iter() - .map(|e| e.downcast_ref::().unwrap().name()) + .map(|e| e.expr.downcast_ref::().unwrap().name()) .collect(); assert_eq!(col_names, vec!["year", "month"]); } - _ => panic!("Expected Hash partitioning"), + other => panic!("Expected Range partitioning, got {other}"), } } + #[test] + fn test_range_partitioning_from_partition_fields_rejects_unordered_split_points() { + let file_schema = aggr_test_schema(); + let config = config_for_projection( + Arc::clone(&file_schema), + None, + Statistics::new_unknown(&file_schema), + vec![Field::new( + "date", + wrap_partition_type_in_dict(DataType::Utf8), + false, + )], + ); + let dict = |v: &str| wrap_partition_value_in_dict(ScalarValue::from(v)); + let err = range_partitioning_from_partition_fields( + config.file_source.table_schema().table_schema(), + config.table_partition_cols(), + vec![ + SplitPoint::new(vec![dict("2026-10")]), + SplitPoint::new(vec![dict("2026-09")]), + ], + ) + .unwrap_err(); + assert!(err.to_string().contains("strictly ordered"), "{err}"); + + // No partition columns => no declared partitioning. + let none = range_partitioning_from_partition_fields( + config.file_source.table_schema().table_schema(), + &Fields::empty(), + vec![], + ) + .unwrap(); + assert!(none.is_none()); + } + #[test] fn try_pushdown_sort_reverses_file_groups_only_when_requested_is_reverse() -> Result<()> { diff --git a/datafusion/physical-plan/src/joins/hash_join/exec.rs b/datafusion/physical-plan/src/joins/hash_join/exec.rs index b72e180543f9a..491e3c7e4b57e 100644 --- a/datafusion/physical-plan/src/joins/hash_join/exec.rs +++ b/datafusion/physical-plan/src/joins/hash_join/exec.rs @@ -1002,28 +1002,11 @@ impl HashJoinExec { return false; } - // `preserve_file_partitions` can report Hive-style file groups as Hash - // partitioned even though their partition indexes do not follow the - // hash router used by partitioned dynamic filters. Reject Hash inputs - // because the metadata cannot distinguish those scans from a real hash - // repartition. Compatible Range inputs remain safe because matching - // ordering and split points align each build filter with its probe - // partition. Other unsupported layouts are rejected. + // Hive style file groups declare `Partitioning::Range`, so a `Hash` input here + // follows the hash router that partitioned dynamic filters assume. Range inputs + // go through the co-partitioning check below. // Follow-up work: enable dynamic filtering for preserve_file_partitioned scans (issue #20195). // https://github.com/apache/datafusion/issues/20195 - if config.optimizer.preserve_file_partitions > 0 - && self.mode == PartitionMode::Partitioned - && matches!( - ( - self.left.output_partitioning(), - self.right.output_partitioning() - ), - (Partitioning::Hash(_, _), Partitioning::Hash(_, _)) - ) - { - return false; - } - if self.mode == PartitionMode::Partitioned && !self.has_partitioned_dynamic_filter_routing() { @@ -9151,6 +9134,7 @@ mod tests { .optimizer .preserve_file_partitions = 1; assert!(range_join.allow_join_dynamic_filter_pushdown(session_config.options())); + assert!(hash_join.allow_join_dynamic_filter_pushdown(session_config.options())); Ok(()) } @@ -9158,8 +9142,6 @@ mod tests { #[test] fn test_partitioned_dynamic_filter_pushdown_rejects_unsupported_partitioning() -> Result<()> { - let (range_join, on) = range_partitioned_dynamic_filter_test_join(10, 10)?; - let hash_join = with_hash_partitioned_children(&range_join, &on)?; let (mismatched_range_join, _) = range_partitioned_dynamic_filter_test_join(10, 11)?; let mut session_config = SessionConfig::default(); @@ -9177,7 +9159,10 @@ mod tests { .options_mut() .optimizer .preserve_file_partitions = 1; - assert!(!hash_join.allow_join_dynamic_filter_pushdown(session_config.options())); + assert!( + !mismatched_range_join + .allow_join_dynamic_filter_pushdown(session_config.options()) + ); Ok(()) } diff --git a/datafusion/proto/tests/cases/plans/sources.rs b/datafusion/proto/tests/cases/plans/sources.rs index 7aff89eb3fdf7..87c6586bfc3a5 100644 --- a/datafusion/proto/tests/cases/plans/sources.rs +++ b/datafusion/proto/tests/cases/plans/sources.rs @@ -1075,3 +1075,84 @@ fn roundtrip_parquet_exec_range_output_partitioning() -> Result<()> { Ok(()) } + +#[test] +fn roundtrip_parquet_exec_range_output_partitioning_dict_split_points() -> Result<()> { + let file_schema = Arc::new(Schema::new(vec![Field::new( + "col", + wrap_partition_type_in_dict(DataType::Utf8), + 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![wrap_partition_value_in_dict( + ScalarValue::from("B"), + )])], + )); + let scan_config = + FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), file_source) + .with_file_groups(vec![ + FileGroup::new(vec![PartitionedFile::new( + "/path/to/file-1.parquet".to_string(), + 1024, + )]), + FileGroup::new(vec![PartitionedFile::new( + "/path/to/file-2.parquet".to_string(), + 1024, + )]), + ]) + .with_output_partitioning(Some(output_partitioning.clone())) + .build(); + + assert_eq!( + roundtrip_file_scan_config(scan_config)?.output_partitioning, + Some(output_partitioning) + ); + + Ok(()) +} + +#[test] +fn roundtrip_parquet_exec_range_output_partitioning_utf8view_split_points() -> Result<()> +{ + let file_schema = Arc::new(Schema::new(vec![Field::new( + "col", + DataType::Utf8View, + 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::Utf8View(Some( + "B".to_string(), + ))])], + )); + let scan_config = + FileScanConfigBuilder::new(ObjectStoreUrl::local_filesystem(), file_source) + .with_file_groups(vec![ + FileGroup::new(vec![PartitionedFile::new( + "/path/to/file-1.parquet".to_string(), + 1024, + )]), + FileGroup::new(vec![PartitionedFile::new( + "/path/to/file-2.parquet".to_string(), + 1024, + )]), + ]) + .with_output_partitioning(Some(output_partitioning.clone())) + .build(); + + assert_eq!( + roundtrip_file_scan_config(scan_config)?.output_partitioning, + Some(output_partitioning) + ); + + Ok(()) +} diff --git a/datafusion/sqllogictest/test_files/preserve_file_partitioning.slt b/datafusion/sqllogictest/test_files/preserve_file_partitioning.slt index e2dd22cc82bba..d6de6a47fc925 100644 --- a/datafusion/sqllogictest/test_files/preserve_file_partitioning.slt +++ b/datafusion/sqllogictest/test_files/preserve_file_partitioning.slt @@ -125,7 +125,7 @@ STORED AS PARQUET; 1 # Create high-cardinality fact table (5 partitions > 3 target_partitions) -# For testing partition merging with consistent hashing +# For testing partition merging into contiguous value ranges query I COPY (SELECT column1 as timestamp, column2 as value FROM (VALUES (TIMESTAMP '2023-01-01T09:00:00', 100.0) @@ -258,7 +258,7 @@ logical_plan physical_plan 01)ProjectionExec: expr=[f_dkey@0 as f_dkey, count(Int64(1))@1 as count(*), sum(fact_table.value)@2 as sum(fact_table.value)] 02)--AggregateExec: mode=SinglePartitioned, gby=[f_dkey@1 as f_dkey], aggr=[count(Int64(1)), sum(fact_table.value)] -03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Hash([f_dkey@1], 3), file_type=parquet +03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Range([f_dkey@1 ASC], [(B), (C)], 3), file_type=parquet # Verify results with optimization match results without optimization query TIR rowsort @@ -320,7 +320,7 @@ physical_plan 01)SortPreservingMergeExec: [f_dkey@0 ASC NULLS LAST] 02)--ProjectionExec: expr=[f_dkey@0 as f_dkey, count(Int64(1))@1 as count(*), avg(fact_table_ordered.value)@2 as avg(fact_table_ordered.value)] 03)----AggregateExec: mode=SinglePartitioned, gby=[f_dkey@1 as f_dkey], aggr=[count(Int64(1)), avg(fact_table_ordered.value)], ordering_mode=Sorted -04)------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_ordering=[f_dkey@1 ASC NULLS LAST], output_partitioning=Hash([f_dkey@1], 3), file_type=parquet +04)------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_ordering=[f_dkey@1 ASC NULLS LAST], output_partitioning=Range([f_dkey@1 ASC], [(B), (C)], 3), file_type=parquet query TIR SELECT f_dkey, count(*), avg(value) FROM fact_table_ordered GROUP BY f_dkey ORDER BY f_dkey; @@ -418,7 +418,7 @@ physical_plan 06)----------FilterExec: service@2 = log 07)------------RepartitionExec: partitioning=RoundRobinBatch(3), input_partitions=1 08)--------------DataSourceExec: file_groups={1 group: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension/data.parquet]]}, projection=[d_dkey, env, service], file_type=parquet, predicate=service@2 = log, pruning_predicate=service_null_count@2 != row_count@3 AND service_min@0 <= log AND log <= service_max@1, required_guarantees=[service in (log)] -09)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_ordering=[f_dkey@1 ASC NULLS LAST], output_partitioning=Hash([f_dkey@1], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible +09)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_ordering=[f_dkey@1 ASC NULLS LAST], output_partitioning=Range([f_dkey@1 ASC], [(B), (C)], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible query TTTIR rowsort SELECT f.f_dkey, MAX(d.env), MAX(d.service), count(*), sum(f.value) @@ -493,7 +493,7 @@ logical_plan physical_plan 01)ProjectionExec: expr=[f_dkey@2 as f_dkey, timestamp@0 as timestamp, value@1 as value, row_number() PARTITION BY [fact_table_ordered.f_dkey] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as rn] 02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [fact_table_ordered.f_dkey] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [fact_table_ordered.f_dkey] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted] -03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet +03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet query TPRI rowsort SELECT f_dkey, timestamp, value, @@ -513,8 +513,7 @@ C 2023-01-01T09:00:20 310.2 3 ########## # TEST 9: High-Cardinality Partitions (more partitions than target_partitions) -# Since num_partitions > target_partitions (5 > 3), files are merged using -# round-robin assignment to ensure exactly target_partitions groups are created. +# Sorted partition values are cut into contiguous chunks. ########## # First verify results without optimization @@ -537,7 +536,7 @@ statement ok set datafusion.optimizer.preserve_file_partitions = 1; # Verify the plan uses SinglePartitioned mode with no RepartitionExec -# The 5 partitions are merged into 3 file groups using round-robin assignment +# 5 values merge into 3 contiguous groups: [A], [B, C], [D, E] query TT EXPLAIN SELECT f_dkey, count(*), sum(value) FROM high_cardinality_table GROUP BY f_dkey; ---- @@ -548,7 +547,7 @@ logical_plan physical_plan 01)ProjectionExec: expr=[f_dkey@0 as f_dkey, count(Int64(1))@1 as count(*), sum(high_cardinality_table.value)@2 as sum(high_cardinality_table.value)] 02)--AggregateExec: mode=SinglePartitioned, gby=[f_dkey@1 as f_dkey], aggr=[count(Int64(1)), sum(high_cardinality_table.value)] -03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=A/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=D/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=B/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=E/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Hash([f_dkey@1], 3), file_type=parquet +03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=B/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=C/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=D/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/high_cardinality/f_dkey=E/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Range([f_dkey@1 ASC], [(B), (D)], 3), file_type=parquet # Verify results with optimization match results without optimization query TIR rowsort @@ -658,9 +657,7 @@ C prod 2017.6 ########## # TEST 12: Partitioned Join with Matching Partition Counts - With Optimization # Both tables have 3 partitions matching target_partitions=3 -# No RepartitionExec needed for join - partitions already satisfy the requirement -# Dynamic filter pushdown is disabled in this mode because preserve_file_partitions -# reports Hash partitioning for Hive-style file groups, which are not hash-routed. +# Identical split points prove co-partitioning, so no RepartitionExec is needed. ########## statement ok @@ -685,8 +682,8 @@ physical_plan 02)--RepartitionExec: partitioning=Hash([f_dkey@0, env@1], 3), input_partitions=3 03)----AggregateExec: mode=Partial, gby=[f_dkey@1 as f_dkey, env@2 as env], aggr=[sum(f.value)] 04)------HashJoinExec: mode=Partitioned, join_type=Inner, on=[(d_dkey@1, f_dkey@1)], projection=[value@2, f_dkey@3, env@0] -05)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_partitioned/d_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_partitioned/d_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_partitioned/d_dkey=C/data.parquet]]}, projection=[env, d_dkey], output_partitioning=Hash([d_dkey@1], 3), file_type=parquet -06)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Hash([f_dkey@1], 3), file_type=parquet +05)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_partitioned/d_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_partitioned/d_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_partitioned/d_dkey=C/data.parquet]]}, projection=[env, d_dkey], output_partitioning=Range([d_dkey@1 ASC], [(B), (C)], 3), file_type=parquet +06)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Range([f_dkey@1 ASC], [(B), (C)], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible query TTR rowsort SELECT f.f_dkey, d.env, sum(f.value) @@ -722,7 +719,7 @@ logical_plan physical_plan 01)ProjectionExec: expr=[f_dkey@0 as f_dkey, timestamp@1 as timestamp, count(Int64(1))@2 as count(*), avg(fact_table.value)@3 as avg(fact_table.value)] 02)--AggregateExec: mode=SinglePartitioned, gby=[f_dkey@2 as f_dkey, timestamp@0 as timestamp], aggr=[count(Int64(1)), avg(fact_table.value)] -03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet +03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet query TPIR rowsort SELECT f_dkey, timestamp, @@ -752,6 +749,196 @@ C 2023-01-01T09:00:40 1 300 C 2023-01-01T09:12:40 1 250 C 2023-01-01T09:12:50 1 275.4 +########## +# TEST 14: Partitioned join, same group count, different partition values +# fact has {A, B, C}, dimension has {A, C, D}. Split points differ, so the scans are not +# co-partitioned and one side must be range repartitioned. +########## + +query I +COPY (SELECT 'dev' as env, 'log' as service) +TO 'test_files/scratch/preserve_file_partitioning/dimension_skewed/d_dkey=A/data.parquet' +STORED AS PARQUET; +---- +1 + +query I +COPY (SELECT 'prod' as env, 'log' as service) +TO 'test_files/scratch/preserve_file_partitioning/dimension_skewed/d_dkey=C/data.parquet' +STORED AS PARQUET; +---- +1 + +query I +COPY (SELECT 'prod' as env, 'trace' as service) +TO 'test_files/scratch/preserve_file_partitioning/dimension_skewed/d_dkey=D/data.parquet' +STORED AS PARQUET; +---- +1 + +statement ok +CREATE EXTERNAL TABLE dimension_table_skewed (env STRING, service STRING) +STORED AS PARQUET +PARTITIONED BY (d_dkey STRING) +LOCATION 'test_files/scratch/preserve_file_partitioning/dimension_skewed/'; + +statement ok +set datafusion.execution.target_partitions = 3; + +statement ok +set datafusion.optimizer.preserve_file_partitions = 1; + +statement ok +set datafusion.optimizer.hash_join_single_partition_threshold = 0; + +statement ok +set datafusion.optimizer.hash_join_single_partition_threshold_rows = 0; + +# Exactly one join input gets a RepartitionExec(Range). +query TT +EXPLAIN SELECT f.f_dkey, d.env, sum(f.value) +FROM fact_table f +INNER JOIN dimension_table_skewed d ON f.f_dkey = d.d_dkey +GROUP BY f.f_dkey, d.env; +---- +logical_plan +01)Aggregate: groupBy=[[f.f_dkey, d.env]], aggr=[[sum(f.value)]] +02)--Projection: f.value, f.f_dkey, d.env +03)----Inner Join: f.f_dkey = d.d_dkey +04)------SubqueryAlias: f +05)--------TableScan: fact_table projection=[value, f_dkey] +06)------SubqueryAlias: d +07)--------TableScan: dimension_table_skewed projection=[env, d_dkey] +physical_plan +01)AggregateExec: mode=FinalPartitioned, gby=[f_dkey@0 as f_dkey, env@1 as env], aggr=[sum(f.value)] +02)--RepartitionExec: partitioning=Hash([f_dkey@0, env@1], 3), input_partitions=3 +03)----AggregateExec: mode=Partial, gby=[f_dkey@1 as f_dkey, env@2 as env], aggr=[sum(f.value)] +04)------HashJoinExec: mode=Partitioned, join_type=Inner, on=[(d_dkey@1, f_dkey@1)], projection=[value@2, f_dkey@3, env@0] +05)--------RepartitionExec: partitioning=Range([d_dkey@1 ASC], [(B), (C)], 3), input_partitions=3 +06)----------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_skewed/d_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_skewed/d_dkey=C/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/dimension_skewed/d_dkey=D/data.parquet]]}, projection=[env, d_dkey], output_partitioning=Range([d_dkey@1 ASC], [(C), (D)], 3), file_type=parquet +07)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Range([f_dkey@1 ASC], [(B), (C)], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible + +# Only A and C match. +query TTR rowsort +SELECT f.f_dkey, d.env, sum(f.value) +FROM fact_table f +INNER JOIN dimension_table_skewed d ON f.f_dkey = d.d_dkey +GROUP BY f.f_dkey, d.env; +---- +A dev 772.4 +C prod 2017.6 + +# Same result with the optimization disabled. +statement ok +set datafusion.optimizer.preserve_file_partitions = 0; + +query TTR rowsort +SELECT f.f_dkey, d.env, sum(f.value) +FROM fact_table f +INNER JOIN dimension_table_skewed d ON f.f_dkey = d.d_dkey +GROUP BY f.f_dkey, d.env; +---- +A dev 772.4 +C prod 2017.6 + +########## +# TEST 15: Partitioned join against a non partitioned VALUES input +# The inserted RepartitionExec reuses the scan's split points, so partition i on both +# sides covers the same key interval. +########## + +statement ok +set datafusion.optimizer.preserve_file_partitions = 1; + +query TT +EXPLAIN SELECT f.f_dkey, d.env, sum(f.value) +FROM fact_table f +INNER JOIN (VALUES ('A', 'dev'), ('B', 'prod'), ('C', 'prod'), ('Z', 'none')) AS d(d_dkey, env) + ON f.f_dkey = d.d_dkey +GROUP BY f.f_dkey, d.env; +---- +logical_plan +01)Aggregate: groupBy=[[f.f_dkey, d.env]], aggr=[[sum(f.value)]] +02)--Projection: f.value, f.f_dkey, d.env +03)----Inner Join: f.f_dkey = CAST(d.d_dkey AS Utf8View) +04)------SubqueryAlias: f +05)--------TableScan: fact_table projection=[value, f_dkey] +06)------SubqueryAlias: d +07)--------Projection: column1 AS d_dkey, column2 AS env +08)----------Values: (Utf8("A"), Utf8("dev")), (Utf8("B"), Utf8("prod")), (Utf8("C"), Utf8("prod")), (Utf8("Z"), Utf8("none")) +physical_plan +01)AggregateExec: mode=FinalPartitioned, gby=[f_dkey@0 as f_dkey, env@1 as env], aggr=[sum(f.value)] +02)--RepartitionExec: partitioning=Hash([f_dkey@0, env@1], 3), input_partitions=3 +03)----AggregateExec: mode=Partial, gby=[f_dkey@1 as f_dkey, env@2 as env], aggr=[sum(f.value)] +04)------HashJoinExec: mode=Partitioned, join_type=Inner, on=[(CAST(d.d_dkey AS Utf8View)@2, f_dkey@1)], projection=[value@3, f_dkey@4, env@1] +05)--------RepartitionExec: partitioning=Range([CAST(d.d_dkey AS Utf8View)@2 ASC], [(B), (C)], 3), input_partitions=1 +06)----------ProjectionExec: expr=[column1@0 as d_dkey, column2@1 as env, CAST(column1@0 AS Utf8View) as CAST(d.d_dkey AS Utf8View)] +07)------------DataSourceExec: partitions=1, partition_sizes=[1] +08)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_partitioning=Range([f_dkey@1 ASC], [(B), (C)], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible + +query TTR rowsort +SELECT f.f_dkey, d.env, sum(f.value) +FROM fact_table f +INNER JOIN (VALUES ('A', 'dev'), ('B', 'prod'), ('C', 'prod'), ('Z', 'none')) AS d(d_dkey, env) + ON f.f_dkey = d.d_dkey +GROUP BY f.f_dkey, d.env; +---- +A dev 772.4 +B prod 614.4 +C prod 2017.6 + +########## +# TEST 16: Self join on merged groups keeps every row +# 5 values merge into [A], [B, C], [D, E], and all must stay routable. +########## + +query TII rowsort +SELECT a.f_dkey, count(*), count(b.value) +FROM high_cardinality_table a +INNER JOIN high_cardinality_table b ON a.f_dkey = b.f_dkey +GROUP BY a.f_dkey; +---- +A 1 1 +B 1 1 +C 1 1 +D 1 1 +E 1 1 + +########## +# TEST 17: Statistics regrouping drops the Range claim +# The re-cut groups no longer follow partition values, so no output partitioning is +# declared, even when the group count is unchanged. +########## + +statement ok +set datafusion.optimizer.preserve_file_partitions = 1; + +statement ok +set datafusion.execution.split_file_groups_by_statistics = true; + +query TT +EXPLAIN SELECT f_dkey, count(*), avg(value) FROM fact_table_ordered GROUP BY f_dkey ORDER BY f_dkey; +---- +logical_plan +01)Sort: fact_table_ordered.f_dkey ASC NULLS LAST +02)--Projection: fact_table_ordered.f_dkey, count(Int64(1)) AS count(*), avg(fact_table_ordered.value) +03)----Aggregate: groupBy=[[fact_table_ordered.f_dkey]], aggr=[[count(Int64(1)), avg(fact_table_ordered.value)]] +04)------TableScan: fact_table_ordered projection=[value, f_dkey] +physical_plan +01)SortPreservingMergeExec: [f_dkey@0 ASC NULLS LAST] +02)--ProjectionExec: expr=[f_dkey@0 as f_dkey, count(Int64(1))@1 as count(*), avg(fact_table_ordered.value)@2 as avg(fact_table_ordered.value)] +03)----AggregateExec: mode=FinalPartitioned, gby=[f_dkey@0 as f_dkey], aggr=[count(Int64(1)), avg(fact_table_ordered.value)], ordering_mode=Sorted +04)------RepartitionExec: partitioning=Hash([f_dkey@0], 3), input_partitions=3, preserve_order=true, sort_exprs=f_dkey@0 ASC NULLS LAST +05)--------AggregateExec: mode=Partial, gby=[f_dkey@1 as f_dkey], aggr=[count(Int64(1)), avg(fact_table_ordered.value)], ordering_mode=Sorted +06)----------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/preserve_file_partitioning/fact/f_dkey=C/data.parquet]]}, projection=[value, f_dkey], output_ordering=[f_dkey@1 ASC NULLS LAST], file_type=parquet + +query TIR +SELECT f_dkey, count(*), avg(value) FROM fact_table_ordered GROUP BY f_dkey ORDER BY f_dkey; +---- +A 7 110.342857142857 +B 7 87.771428571429 +C 7 288.228571428571 + ########## # CLEANUP ########## @@ -769,6 +956,9 @@ reset datafusion.optimizer.hash_join_single_partition_threshold; statement ok reset datafusion.optimizer.hash_join_single_partition_threshold_rows; +statement ok +reset datafusion.execution.split_file_groups_by_statistics; + statement ok DROP TABLE fact_table; @@ -781,5 +971,8 @@ DROP TABLE dimension_table; statement ok DROP TABLE dimension_table_partitioned; +statement ok +DROP TABLE dimension_table_skewed; + statement ok DROP TABLE high_cardinality_table; diff --git a/datafusion/sqllogictest/test_files/repartition_subset_satisfaction.slt b/datafusion/sqllogictest/test_files/repartition_subset_satisfaction.slt index 5371ca59beea1..bb1a4b2ad4dc1 100644 --- a/datafusion/sqllogictest/test_files/repartition_subset_satisfaction.slt +++ b/datafusion/sqllogictest/test_files/repartition_subset_satisfaction.slt @@ -164,7 +164,7 @@ physical_plan 03)----AggregateExec: mode=FinalPartitioned, gby=[f_dkey@0 as f_dkey, date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)@1 as date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)], aggr=[count(Int64(1)), avg(fact_table_ordered.value)], ordering_mode=Sorted 04)------RepartitionExec: partitioning=Hash([f_dkey@0, date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)@1], 3), input_partitions=3, preserve_order=true, sort_exprs=f_dkey@0 ASC NULLS LAST, date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)@1 ASC NULLS LAST 05)--------AggregateExec: mode=Partial, gby=[f_dkey@2 as f_dkey, date_bin(IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }, timestamp@0) as date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)], aggr=[count(Int64(1)), avg(fact_table_ordered.value)], ordering_mode=Sorted -06)----------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet +06)----------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet # Verify results without subset satisfaction query TPIR rowsort @@ -204,7 +204,7 @@ physical_plan 01)SortPreservingMergeExec: [f_dkey@0 ASC NULLS LAST, time_bin@1 ASC NULLS LAST] 02)--ProjectionExec: expr=[f_dkey@0 as f_dkey, date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)@1 as time_bin, count(Int64(1))@2 as count(*), avg(fact_table_ordered.value)@3 as avg(fact_table_ordered.value)] 03)----AggregateExec: mode=SinglePartitioned, gby=[f_dkey@2 as f_dkey, date_bin(IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }, timestamp@0) as date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)], aggr=[count(Int64(1)), avg(fact_table_ordered.value)], ordering_mode=Sorted -04)------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet +04)------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet # Verify results match with subset satisfaction query TPIR rowsort @@ -251,7 +251,7 @@ physical_plan 02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [fact_table_ordered.f_dkey, date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [fact_table_ordered.f_dkey, date_bin(IntervalMonthDayNano(\"IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }\"),fact_table_ordered.timestamp)] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted] 03)----SortExec: expr=[f_dkey@2 ASC NULLS LAST, date_bin(IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }, timestamp@0) ASC NULLS LAST, timestamp@0 ASC NULLS LAST], preserve_partitioning=[true] 04)------RepartitionExec: partitioning=Hash([f_dkey@2, date_bin(IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }, timestamp@0)], 3), input_partitions=3 -05)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet +05)--------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet # Verify results without subset satisfaction query TPRI rowsort @@ -292,7 +292,7 @@ logical_plan physical_plan 01)ProjectionExec: expr=[f_dkey@2 as f_dkey, timestamp@0 as timestamp, value@1 as value, row_number() PARTITION BY [fact_table_ordered.f_dkey, date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW@3 as rn] 02)--BoundedWindowAggExec: wdw=[row_number() PARTITION BY [fact_table_ordered.f_dkey, date_bin(IntervalMonthDayNano("IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }"),fact_table_ordered.timestamp)] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW: Field { "row_number() PARTITION BY [fact_table_ordered.f_dkey, date_bin(IntervalMonthDayNano(\"IntervalMonthDayNano { months: 0, days: 0, nanoseconds: 30000000000 }\"),fact_table_ordered.timestamp)] ORDER BY [fact_table_ordered.timestamp ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW": UInt64 }, frame: RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW], mode=[Sorted] -03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet +03)----DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet # Verify results match with subset satisfaction query TPRI rowsort @@ -379,8 +379,8 @@ physical_plan 11)--------------------HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(d_dkey@1, f_dkey@2)], projection=[f_dkey@4, env@0, timestamp@2, value@3] 12)----------------------CoalescePartitionsExec 13)------------------------FilterExec: service@1 = log, projection=[env@0, d_dkey@2] -14)--------------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=A/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=D/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=C/data.parquet]]}, projection=[env, service, d_dkey], output_partitioning=Hash([d_dkey@2], 3), file_type=parquet, predicate=service@1 = log, pruning_predicate=service_null_count@2 != row_count@3 AND service_min@0 <= log AND log <= service_max@1, required_guarantees=[service in (log)] -15)----------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible +14)--------------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=C/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=D/data.parquet]]}, projection=[env, service, d_dkey], output_partitioning=Range([d_dkey@2 ASC], [(B), (C)], 3), file_type=parquet, predicate=service@1 = log, pruning_predicate=service_null_count@2 != row_count@3 AND service_min@0 <= log AND log <= service_max@1, required_guarantees=[service in (log)] +15)----------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible # Verify results without subset satisfaction query TPR rowsort @@ -474,8 +474,8 @@ physical_plan 09)----------------HashJoinExec: mode=CollectLeft, join_type=Inner, on=[(d_dkey@1, f_dkey@2)], projection=[f_dkey@4, env@0, timestamp@2, value@3] 10)------------------CoalescePartitionsExec 11)--------------------FilterExec: service@1 = log, projection=[env@0, d_dkey@2] -12)----------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=A/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=D/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=C/data.parquet]]}, projection=[env, service, d_dkey], output_partitioning=Hash([d_dkey@2], 3), file_type=parquet, predicate=service@1 = log, pruning_predicate=service_null_count@2 != row_count@3 AND service_min@0 <= log AND log <= service_max@1, required_guarantees=[service in (log)] -13)------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Hash([f_dkey@2], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible +12)----------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=C/data.parquet, WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/dimension/d_dkey=D/data.parquet]]}, projection=[env, service, d_dkey], output_partitioning=Range([d_dkey@2 ASC], [(B), (C)], 3), file_type=parquet, predicate=service@1 = log, pruning_predicate=service_null_count@2 != row_count@3 AND service_min@0 <= log AND log <= service_max@1, required_guarantees=[service in (log)] +13)------------------DataSourceExec: file_groups={3 groups: [[WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=A/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=B/data.parquet], [WORKSPACE_ROOT/datafusion/sqllogictest/test_files/scratch/repartition_subset_satisfaction/fact/f_dkey=C/data.parquet]]}, projection=[timestamp, value, f_dkey], output_ordering=[f_dkey@2 ASC NULLS LAST, timestamp@0 ASC NULLS LAST], output_partitioning=Range([f_dkey@2 ASC], [(B), (C)], 3), file_type=parquet, predicate=DynamicFilter [ empty ], dynamic_rg_pruning=eligible # Verify results match with subset satisfaction query TPR rowsort