Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
60 changes: 44 additions & 16 deletions datafusion/catalog-listing/src/table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -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<SplitPoint>,
}

/// Built in [`TableProvider`] that reads data from one or more files as a single table.
Expand Down Expand Up @@ -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?;
Expand All @@ -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| {
Expand All @@ -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"
Expand Down Expand Up @@ -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
};
Expand Down Expand Up @@ -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
Expand All @@ -966,28 +988,31 @@ 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(
ctx,
file_groups,
inexact_stats,
grouped_by_partition,
partition_split_points,
)
}

Expand Down Expand Up @@ -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) =
Expand All @@ -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(
Expand Down Expand Up @@ -1061,6 +1087,7 @@ impl ListingTable {
file_groups: Vec<FileGroup>,
inexact_stats: bool,
grouped_by_partition: bool,
partition_split_points: Vec<SplitPoint>,
) -> datafusion_common::Result<ListFilesResult> {
let (file_groups, stats) = compute_all_files_statistics(
file_groups,
Expand All @@ -1077,6 +1104,7 @@ impl ListingTable {
file_groups,
statistics: stats,
grouped_by_partition,
partition_split_points,
})
}

Expand Down
Loading