Skip to content

Decide scan parallelism from the estimated row count, not from statistics precision #26087

Description

@zhuqi-lucas

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:

  1. Flip the default of use_row_number_estimates_to_optimize_partitioning to true.
  2. 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.
  3. 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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

Labels

enhancementNew feature or request

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions