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
63 changes: 42 additions & 21 deletions datafusion/catalog-listing/src/table.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,14 +25,12 @@ 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,
Column, 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,
};
use datafusion_datasource::file_scan_config::{FileScanConfig, FileScanConfigBuilder};
use datafusion_datasource::file_sink_config::{FileOutputMode, FileSinkConfig};
#[expect(deprecated)]
use datafusion_datasource::schema_adapter::SchemaAdapterFactory;
Expand All @@ -46,7 +44,9 @@ use datafusion_expr::dml::InsertOp;
use datafusion_expr::execution_props::ExecutionProps;
use datafusion_expr::physical_planning_context::PhysicalPlanningContext;
use datafusion_expr::{
Expr, Partitioning as LogicalPartitioning, TableProviderFilterPushDown, TableType,
Expr, Partitioning as LogicalPartitioning,
RangePartitioning as LogicalRangePartitioning, TableProviderFilterPushDown,
TableType,
};
use datafusion_physical_expr::{create_lex_ordering, create_physical_partitioning};
use datafusion_physical_expr_adapter::PhysicalExprAdapterFactory;
Expand All @@ -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 All @@ -678,8 +683,26 @@ impl ListingTable {
None => {} // no ordering required
}

// Hive grouped files are contiguous key intervals, so they declare `Range`
// unless the statistics re-cut above changed the groups. The ordering must
// match the `SortOptions::default()` used to cut the groups.
let derived_output_partitioning =
if partitioned_by_file_group && !regrouped_by_statistics {
let ordering = table_partition_cols
.iter()
.map(|field| {
Expr::Column(Column::from_name(field.name())).sort(true, true)
})
.collect();
Some(LogicalPartitioning::Range(
LogicalRangePartitioning::try_new(ordering, partition_split_points)?,
))
} else {
None
};

let output_partitioning = if let Some(output_partitioning) =
declared_output_partitioning
declared_output_partitioning.or(derived_output_partitioning.as_ref())
{
let output_partitioning = match output_partitioning {
LogicalPartitioning::RoundRobinBatch(_) => {
Expand Down Expand Up @@ -710,15 +733,6 @@ 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(),
)
} else {
None
};
Expand Down Expand Up @@ -953,6 +967,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 +981,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 +1033,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 +1045,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 +1080,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 +1097,7 @@ impl ListingTable {
file_groups,
statistics: stats,
grouped_by_partition,
partition_split_points,
})
}

Expand Down
Loading