You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
{{ message }}
Repository navigation
Decide scan parallelism from the estimated row count, not from statistics precision #26087
Is your feature request related to a problem or challenge?
EnforceDistribution decides whether a child is worth parallelizing (round-robin repartition, and repartition_file_scans splitting a scan into target_partitions file groups) from roundrobin_beneficial_stats in get_repartition_requirement_status (physical-optimizer/src/ensure_requirements/enforce_distribution.rs):
let roundrobin_beneficial_stats = match stats.num_rows{Precision::Exact(n_rows) => n_rows > batch_size,Precision::Inexact(n_rows) => !should_use_estimates || (n_rows > batch_size),Precision::Absent => true,};
With the default datafusion.execution.use_row_number_estimates_to_optimize_partitioning = false, an Inexact row count is treated as "always worth parallelizing" regardless of its value. Precision is used as a proxy for trust, so the decision depends on whether anything touched the statistics rather than on how many rows there are:
FileScanConfig::statistics() downgrades num_rows to Inexact as soon as a predicate is pushed into the source. The same 100-row scan is left as one group without a pushed filter (Exact(100), not worth it) and split into 4 groups with one (Inexact(100), "worth it").
Because of that, the plan for small tables depends on whether FilterPushdown runs before or after enforcement. We noticed this while working on feat: run physical requirement enforcement first, as a PhysicalAnalyzer phase #25688: moving FilterPushdown after enforcement changed repartition_scan.slt from 4 file groups to 1 for tables far below batch_size.
Sources already protect themselves against over-splitting through repartition_file_min_size inside ExecutionPlan::repartitioned, so the precision gate mostly affects small inputs, where splitting is pure overhead.
Describe the solution you'd like
Decide from the estimate regardless of precision: n_rows > batch_size for both Exact and Inexact; keep "assume beneficial" only for Absent.
That is what use_row_number_estimates_to_optimize_partitioning = true already does, and the config docs say "We plan to make this the default in the future". Concretely:
Flip the default of use_row_number_estimates_to_optimize_partitioning to true.
Update the affected golden plans (repartition_scan.slt, the filter_pushdown snapshots and similar): sources below batch_size rows that carry a pushed-down predicate stop being split or round-robined.
Deprecate the flag and remove it after the usual period, so the decision has one code path.
Describe alternatives you've considered
Move the parallelism decision into ExecutionPlan::repartitioned entirely (the source knows its file sizes), and keep enforcement to inserting the exchanges that declared requirements need. Larger change, could be a follow-up.
Additional context
Related to the ordering discussion in #25688. I can take this.
Is your feature request related to a problem or challenge?
EnforceDistributiondecides whether a child is worth parallelizing (round-robin repartition, andrepartition_file_scanssplitting a scan intotarget_partitionsfile groups) fromroundrobin_beneficial_statsinget_repartition_requirement_status(physical-optimizer/src/ensure_requirements/enforce_distribution.rs):With the default
datafusion.execution.use_row_number_estimates_to_optimize_partitioning = false, anInexactrow count is treated as "always worth parallelizing" regardless of its value. Precision is used as a proxy for trust, so the decision depends on whether anything touched the statistics rather than on how many rows there are:FileScanConfig::statistics()downgradesnum_rowstoInexactas soon as a predicate is pushed into the source. The same 100-row scan is left as one group without a pushed filter (Exact(100), not worth it) and split into 4 groups with one (Inexact(100), "worth it").FilterPushdownruns before or after enforcement. We noticed this while working on feat: run physical requirement enforcement first, as a PhysicalAnalyzer phase #25688: movingFilterPushdownafter enforcement changedrepartition_scan.sltfrom 4 file groups to 1 for tables far belowbatch_size.Sources already protect themselves against over-splitting through
repartition_file_min_sizeinsideExecutionPlan::repartitioned, so the precision gate mostly affects small inputs, where splitting is pure overhead.Describe the solution you'd like
Decide from the estimate regardless of precision:
n_rows > batch_sizefor bothExactandInexact; keep "assume beneficial" only forAbsent.That is what
use_row_number_estimates_to_optimize_partitioning = truealready does, and the config docs say "We plan to make this the default in the future". Concretely:use_row_number_estimates_to_optimize_partitioningtotrue.repartition_scan.slt, thefilter_pushdownsnapshots and similar): sources belowbatch_sizerows that carry a pushed-down predicate stop being split or round-robined.Describe alternatives you've considered
Move the parallelism decision into
ExecutionPlan::repartitionedentirely (the source knows its file sizes), and keep enforcement to inserting the exchanges that declared requirements need. Larger change, could be a follow-up.Additional context
Related to the ordering discussion in #25688. I can take this.