Skip to content

Avoid consecutive RepartitionExec #18341

Description

@NGA-TRAN

Reproducer in sqllogictest: #18343

While experimenting with aggregation, I noticed cases where two RepartitionExec operators appear consecutively in the plan. This seems suboptimal and worth improving.

I’m using a dataset stored in both CSV and Parquet formats to test this behavior.

d_dkey,env,service,host
A,dev,log,ma
B,prod,log,ma
C,prod,log,vim
D,prod,trace,vim

The plan of for data in cvs file looks reasonable

EXPLAIN SELECT env, count(*) FROM dimension_csv GROUP BY env;
+---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| plan_type     | plan                                                                                                                                                                                        |
+---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| logical_plan  | Projection: dimension_csv.env, count(Int64(1)) AS count(*)                                                                                                                                  |
|               |   Aggregate: groupBy=[[dimension_csv.env]], aggr=[[count(Int64(1))]]                                                                                                                        |
|               |     TableScan: dimension_csv projection=[env]                                                                                                                                               |
| physical_plan | ProjectionExec: expr=[env@0 as env, count(Int64(1))@1 as count(*)]                                                                                                                          |
|               |   AggregateExec: mode=FinalPartitioned, gby=[env@0 as env], aggr=[count(Int64(1))]                                                                                                          |
|               |     CoalesceBatchesExec: target_batch_size=8192                                                                                                                                             |
|               |       RepartitionExec: partitioning=Hash([env@0], 16), input_partitions=16                                                                                                                  |
|               |         AggregateExec: mode=Partial, gby=[env@0 as env], aggr=[count(Int64(1))]                                                                                                             |
|               |           RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1                                                                                                             |
|               |             DataSourceExec: file_groups={1 group: [[Users/hoabinhnga.tran/datafusion-optimal-plans/testdata/dimension1/dimension_1.csv]]}, projection=[env], file_type=csv, has_header=true |
|               |                                                                                                                                                                                             |
+---------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+

And this is its graphical plan

Image

However, the plan for Parquet data doesn’t appear to push the round-robin repartition down far enough, resulting in two RepartitionExec operators placed back-to-back. This seems quite suboptimal and likely worth improving.

EXPLAIN SELECT env, count(*) FROM dimension_parquet GROUP BY env;
+---------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| plan_type     | plan                                                                                                                                                                               |
+---------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| logical_plan  | Projection: dimension_parquet.env, count(Int64(1)) AS count(*)                                                                                                                     |
|               |   Aggregate: groupBy=[[dimension_parquet.env]], aggr=[[count(Int64(1))]]                                                                                                           |
|               |     TableScan: dimension_parquet projection=[env]                                                                                                                                  |
| physical_plan | ProjectionExec: expr=[env@0 as env, count(Int64(1))@1 as count(*)]                                                                                                                 |
|               |   AggregateExec: mode=FinalPartitioned, gby=[env@0 as env], aggr=[count(Int64(1))]                                                                                                 |
|               |     CoalesceBatchesExec: target_batch_size=8192                                                                                                                                    |
|               |       RepartitionExec: partitioning=Hash([env@0], 16), input_partitions=16              -- Repartition Hash                                                                        |
|               |         RepartitionExec: partitioning=RoundRobinBatch(16), input_partitions=1           -- Repartition Round Robin                                                                 |
|               |           AggregateExec: mode=Partial, gby=[env@0 as env], aggr=[count(Int64(1))]                                                                                                  |
|               |             DataSourceExec: file_groups={1 group: [[Users/hoabinhnga.tran/datafusion-optimal-plans/testdata/dimension1/dimension_1.parquet]]}, projection=[env], file_type=parquet |
|               |                                                                                                                                                                                    |
+---------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
Image

Fix Proposal

I suggest we fix this by either pushing the round-robin repartition further down—similar to the plan for the CSV file above—or eliminating it entirely if it's unnecessary.

Advanced proposal

If we can determine that the input file is small, there's no need to repartition the data. In that case, we should use a single-step aggregate—like in the example below.

-- Option to keep single partition
set datafusion.execution.target_partitions = 1;

EXPLAIN SELECT env, count(*) FROM dimension_parquet GROUP BY env;
+---------------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| plan_type     | plan                                                                                                                                                                       |
+---------------+----------------------------------------------------------------------------------------------------------------------------------------------------------------------------+
| logical_plan  | Projection: dimension_parquet.env, count(Int64(1)) AS count(*)                                                                                                             |
|               |   Aggregate: groupBy=[[dimension_parquet.env]], aggr=[[count(Int64(1))]]                                                                                                   |
|               |     TableScan: dimension_parquet projection=[env]                                                                                                                          |
| physical_plan | ProjectionExec: expr=[env@0 as env, count(Int64(1))@1 as count(*)]                                                                                                         |
|               |   AggregateExec: mode=Single, gby=[env@0 as env], aggr=[count(Int64(1))]                                                                                                   |
|               |     DataSourceExec: file_groups={1 group: [[Users/hoabinhnga.tran/datafusion-optimal-plans/testdata/dimension1/dimension_1.parquet]]}, projection=[env], file_type=parquet |
|               |                                                                                                                                                                            |
Image

Activity

  1. NGA-TRAN commented on Oct 28, 2025

    @NGA-TRAN
    ContributorAuthor

    Some hints: I suggest starting by investigating when the two RepartitionExec nodes are introduced in the plan, and why the behavior differs between the CSV and Parquet cases.

    Running EXPLAIN VERBOSE can help reveal when these nodes are added—during either the logical or physical optimization phases.

    Explain VERBOSE SELECT env, count(*) FROM dimension_parquet GROUP BY env;
  2. gene-bordegaray commented on Oct 28, 2025

    @gene-bordegaray
    Contributor

    take

  3. alamb commented on Oct 28, 2025

    @alamb
    Contributor
  4. 2010YOUY01 commented on Oct 30, 2025

    @2010YOUY01
    Contributor

    I think it might be related to the issue, there is a discord discussion for slow tpch q1 https://discord.com/channels/885562378132000778/1290751484807352412/1432863136612089959

    From the plan, there is only one hash repartition operator, and no round-robin repartition. I believe they're necessary if the data source is not very small (and we have to push them further down as explained in this issue, otherwise they're useless)
    The reason is the push downed filter l_shipdate <= date '1998-09-02' may filter out most of the input parquet scanner partition.

    (Note below is just a imaginary case, I haven't verified that's the actual reason for slow Q1, but I think it's possible, so it's necessary to fix the RepartitionExec issue)

    e.g. Let's say datafusion.execution.target_partitions=4, so there are 4 parallel parquet scanner, also 4 parallel aggregate operator in the upstream.

    partition 1: l_shipdate has range [1990, 2000]
    partition 2: l_shipdate has range [2000, 2005]
    partition 3: l_shipdate has range [2005, 2010]
    partition 4: l_shipdate has range [2010, 2020]

    The consequence is partition 2,3,4 has no output data in the parquet reader, then only one partition is busy, and the available CPUs can't be fully utilized.

    TPCH Q1 plan
    > CREATE EXTERNAL TABLE IF NOT EXISTS lineitem
    STORED AS parquet
    LOCATION '/Users/yongting/Code/datafusion/benchmarks/data/tpch_sf1/lineitem';
    
    > explain select
        l_returnflag,
        l_linestatus,
        sum(l_quantity) as sum_qty,
        sum(l_extendedprice) as sum_base_price,
        sum(l_extendedprice * (1 - l_discount)) as sum_disc_price,
        sum(l_extendedprice * (1 - l_discount) * (1 + l_tax)) as sum_charge,
        avg(l_quantity) as avg_qty,
        avg(l_extendedprice) as avg_price,
        avg(l_discount) as avg_disc,
        count(*) as count_order
    from
        lineitem
    where
            l_shipdate <= date '1998-09-02'
    group by
        l_returnflag,
        l_linestatus
    order by
        l_returnflag,
        l_linestatus;
    +---------------+-------------------------------+
    | plan_type     | plan                          |
    +---------------+-------------------------------+
    | physical_plan | ┌───────────────────────────┐ |
    |               | │  SortPreservingMergeExec  │ |
    |               | │    --------------------   │ |
    |               | │   l_returnflag ASC NULLS  │ |
    |               | │     LAST, l_linestatus    │ |
    |               | │       ASC NULLS LAST      │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │          SortExec         │ |
    |               | │    --------------------   │ |
    |               | │  l_returnflag@0 ASC NULLS │ |
    |               | │    LAST, l_linestatus@1   │ |
    |               | │       ASC NULLS LAST      │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │       ProjectionExec      │ |
    |               | │    --------------------   │ |
    |               | │         avg_disc:         │ |
    |               | │  avg(lineitem.l_discount) │ |
    |               | │                           │ |
    |               | │         avg_price:        │ |
    |               | │        avg(lineitem       │ |
    |               | │        .l_extendedp       │ |
    |               | │           rice)           │ |
    |               | │                           │ |
    |               | │          avg_qty:         │ |
    |               | │  avg(lineitem.l_quantity) │ |
    |               | │                           │ |
    |               | │        count_order:       │ |
    |               | │      count(Int64(1))      │ |
    |               | │                           │ |
    |               | │       l_linestatus:       │ |
    |               | │        l_linestatus       │ |
    |               | │                           │ |
    |               | │       l_returnflag:       │ |
    |               | │        l_returnflag       │ |
    |               | │                           │ |
    |               | │      sum_base_price:      │ |
    |               | │        sum(lineitem       │ |
    |               | │        .l_extendedp       │ |
    |               | │           rice)           │ |
    |               | │                           │ |
    |               | │        sum_charge:        │ |
    |               | │        sum(lineitem       │ |
    |               | │        .l_extendedp       │ |
    |               | │ rice * Int64(1) - lineitem│ |
    |               | │            ...            │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │       AggregateExec       │ |
    |               | │    --------------------   │ |
    |               | │           aggr:           │ |
    |               | │ sum(lineitem.l_quantity), │ |
    |               | │        sum(lineitem       │ |
    |               | │      .l_extendedpric      │ |
    |               | │      e), sum(lineitem     │ |
    |               | │      .l_extendedprice     │ |
    |               | │    * Int64(1) - lineitem  │ |
    |               | │     .l_discount), sum     │ |
    |               | │         (lineitem         │ |
    |               | │       .l_extendedpri      │ |
    |               | │  ce * Int64(1) - lineitem │ |
    |               | │  .l_discount * Int64(1)   │ |
    |               | │   + lineitem.l_tax), avg  │ |
    |               | │   (lineitem.l_quantity),  │ |
    |               | │        avg(lineitem       │ |
    |               | │       .l_extendedpri      │ |
    |               | │     ce), avg(lineitem     │ |
    |               | │       .l_discount),       │ |
    |               | │          count(1)         │ |
    |               | │                           │ |
    |               | │         group_by:         │ |
    |               | │ l_returnflag, l_linestatus│ |
    |               | │                           │ |
    |               | │           mode:           │ |
    |               | │      FinalPartitioned     │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │    CoalesceBatchesExec    │ |
    |               | │    --------------------   │ |
    |               | │     target_batch_size:    │ |
    |               | │            8192           │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │      RepartitionExec      │ |
    |               | │    --------------------   │ |
    |               | │ partition_count(in->out): │ |
    |               | │          14 -> 14         │ |
    |               | │                           │ |
    |               | │    partitioning_scheme:   │ |
    |               | │   Hash([l_returnflag@0,   │ |
    |               | │    l_linestatus@1], 14)   │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │       AggregateExec       │ |
    |               | │    --------------------   │ |
    |               | │           aggr:           │ |
    |               | │ sum(lineitem.l_quantity), │ |
    |               | │        sum(lineitem       │ |
    |               | │      .l_extendedpric      │ |
    |               | │      e), sum(lineitem     │ |
    |               | │      .l_extendedprice     │ |
    |               | │    * Int64(1) - lineitem  │ |
    |               | │     .l_discount), sum     │ |
    |               | │         (lineitem         │ |
    |               | │       .l_extendedpri      │ |
    |               | │  ce * Int64(1) - lineitem │ |
    |               | │  .l_discount * Int64(1)   │ |
    |               | │   + lineitem.l_tax), avg  │ |
    |               | │   (lineitem.l_quantity),  │ |
    |               | │        avg(lineitem       │ |
    |               | │       .l_extendedpri      │ |
    |               | │     ce), avg(lineitem     │ |
    |               | │       .l_discount),       │ |
    |               | │          count(1)         │ |
    |               | │                           │ |
    |               | │         group_by:         │ |
    |               | │ l_returnflag, l_linestatus│ |
    |               | │                           │ |
    |               | │       mode: Partial       │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │       ProjectionExec      │ |
    |               | │    --------------------   │ |
    |               | │      __common_expr_1:     │ |
    |               | │ l_extendedprice * (Some(1)│ |
    |               | │    ,20,0 - l_discount)    │ |
    |               | │                           │ |
    |               | │        l_discount:        │ |
    |               | │         l_discount        │ |
    |               | │                           │ |
    |               | │      l_extendedprice:     │ |
    |               | │      l_extendedprice      │ |
    |               | │                           │ |
    |               | │       l_linestatus:       │ |
    |               | │        l_linestatus       │ |
    |               | │                           │ |
    |               | │        l_quantity:        │ |
    |               | │         l_quantity        │ |
    |               | │                           │ |
    |               | │       l_returnflag:       │ |
    |               | │        l_returnflag       │ |
    |               | │                           │ |
    |               | │        l_tax: l_tax       │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │    CoalesceBatchesExec    │ |
    |               | │    --------------------   │ |
    |               | │     target_batch_size:    │ |
    |               | │            8192           │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │         FilterExec        │ |
    |               | │    --------------------   │ |
    |               | │         predicate:        │ |
    |               | │  l_shipdate <= 1998-09-02 │ |
    |               | └─────────────┬─────────────┘ |
    |               | ┌─────────────┴─────────────┐ |
    |               | │       DataSourceExec      │ |
    |               | │    --------------------   │ |
    |               | │         files: 21         │ |
    |               | │      format: parquet      │ |
    |               | │                           │ |
    |               | │         predicate:        │ |
    |               | │  l_shipdate <= 1998-09-02 │ |
    |               | └───────────────────────────┘ |
    |               |                               |
    +---------------+-------------------------------+
    1 row(s) fetched.
    Elapsed 0.040 seconds.
    
  5. 2010YOUY01 commented on Oct 30, 2025

    @2010YOUY01
    Contributor

    I've checked on tpch sf1 on parquet generated with benchmarks/ script, running q1 under datafusion.execution.target_partitions=14, with explain analyze verbose on the query. And nearly half of the partitions are actually empty.

    I think fixing it can significantly speed up tpch q1.

  6. alamb commented on Nov 1, 2025

    @alamb
    Contributor

    BTW I have filed a separate ticket to work on tpch q1 performance:

  7. gene-bordegaray commented on Nov 7, 2025

    @gene-bordegaray
    Contributor

    i have made a PR for this fix along with a write-up I plan to share with others about my deep dive into the physical optimizer. I have attached the long-version of my write up that goes into depth (meant to help others understand this part of the codebase) and the report just on this issue along with follow up work I the need for:

    Full Report: The Physical Optimizer and Fixing Consecutive Repartitions In the Enforce Distribution Rule.pdf

    Issue Report: Fixing Consecutive Repartitions In the Enforce Distribution Rule.pdf

    If others think my document is useful I would love to share with community to help others understand this part of the codebase

  8. NGA-TRAN commented on Nov 7, 2025

    @NGA-TRAN
    ContributorAuthor

    @gene-bordegaray : Great analysis. I have read both full report that includes your studying how physical rules work and the the issue report that only include the issue

    To reviewers: If you are familiar with DataFusion already, you only need to read the Issue Report

    This is the summary of the fix

    CSV files that lack of statistics:

    We always add Round Robin Repartition and the plan looks like this which is very reasonable to me (because of no stats)

    01)ProjectionExec: expr=[env@0 as env, count(Int64(1))@1 as count(*)]
    02)--AggregateExec: mode=FinalPartitioned, gby=[env], aggr=[count()]
    03)----CoalesceBatchesExec: target_batch_size=8192
    04)------RepartitionExec: partitioning=Hash([env], 4), input_partitions=4
    05)--------AggregateExec: mode=Partial, gby=[env], aggr=[count()]
    06)----------RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1
    07)------------DataSourceExec

    Parquet files with Exact statistics

    There are 2 cases

    1. When the file is small, we do not need to repartition it and the plan looks like
    01)ProjectionExec: expr=[env, count()]
    02)--AggregateExec: mode=FinalPartitioned, gby=[env], aggr=[count()]
    03)----CoalesceBatchesExec: target_batch_size=8192
    04)------RepartitionExec: partitioning=Hash([env], 4), input_partitions=1
    05)--------AggregateExec: mode=Partial, gby=[env], aggr=[count]
    06)----------DataSourceExec
    1. When the file is large, we do add repartition
    01)ProjectionExec: expr=[env@0 as env, count(Int64(1))@1 as count(*)]
    02)--AggregateExec: mode=FinalPartitioned, gby=[env], aggr=[count()]
    03)----CoalesceBatchesExec: target_batch_size=8192
    04)------RepartitionExec: partitioning=Hash([env], 4), input_partitions=4
    05)--------AggregateExec: mode=Partial, gby=[env], aggr=[count()]
    06)----------RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1
    07)------------DataSourceExec

    The fix is exactly what we want.

    @gene-bordegaray : Could you confirm if the summary above is accurate? If it is, this looks like the ideal fix.

  9. NGA-TRAN commented on Nov 7, 2025

    @NGA-TRAN
    ContributorAuthor

    And if you read the @gene-bordegaray 's doc, he also included the benchmark numbers that everything is running faster as expected per summary above

  10. gene-bordegaray commented on Nov 7, 2025

    @gene-bordegaray
    Contributor

    @gene-bordegaray : Great analysis. I have read both full report that includes your studying how physical rules work and the the issue report that only include the issue

    To reviewers: If you are familiar with DataFusion already, you only need to read the Issue Report

    This is the summary of the fix

    CSV files that lack of statistics:

    We always add Round Robin Repartition and the plan looks like this which is very reasonable to me (because of no stats)

    01)ProjectionExec: expr=[env@0 as env, count(Int64(1))@1 as count(*)]
    02)--AggregateExec: mode=FinalPartitioned, gby=[env], aggr=[count()]
    03)----CoalesceBatchesExec: target_batch_size=8192
    04)------RepartitionExec: partitioning=Hash([env], 4), input_partitions=4
    05)--------AggregateExec: mode=Partial, gby=[env], aggr=[count()]
    06)----------RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1
    07)------------DataSourceExec

    Parquet files with Exact statistics

    There are 2 cases

    1. When the file is small, we do not need to repartition it and the plan looks like

    01)ProjectionExec: expr=[env, count()]
    02)--AggregateExec: mode=FinalPartitioned, gby=[env], aggr=[count()]
    03)----CoalesceBatchesExec: target_batch_size=8192
    04)------RepartitionExec: partitioning=Hash([env], 4), input_partitions=1
    05)--------AggregateExec: mode=Partial, gby=[env], aggr=[count]
    06)----------DataSourceExec
    2. When the file is large, we do add repartition

    01)ProjectionExec: expr=[env@0 as env, count(Int64(1))@1 as count(*)]
    02)--AggregateExec: mode=FinalPartitioned, gby=[env], aggr=[count()]
    03)----CoalesceBatchesExec: target_batch_size=8192
    04)------RepartitionExec: partitioning=Hash([env], 4), input_partitions=4
    05)--------AggregateExec: mode=Partial, gby=[env], aggr=[count()]
    06)----------RepartitionExec: partitioning=RoundRobinBatch(4), input_partitions=1
    07)------------DataSourceExec
    The fix is exactly what we want.

    @gene-bordegaray : Could you confirm if the summary above is accurate? If it is, this looks like the ideal fix.

    Yes this is. great summary, thank you Nga

  11. added a commit that references this issue on Nov 11, 2025
    552dbe4
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Labels

No labels
No labels

Type

No type

Projects

No projects

    Milestone

    No milestone

    Relationships

    None yet

    Development

    No branches or pull requests

    Issue actions