Skip to content

Add skewness metric to nodes that execute in partitioned mode - #25977

Merged
gabotechs merged 1 commit into
apache:mainfrom
LiaCastaneda:lia/output-rows-skew-operators
Oct 5, 2026
Merged

gabotechs merged 1 commit into
apache:mainfrom
LiaCastaneda:lia/output-rows-skew-operators

Conversation

@LiaCastaneda

@LiaCastaneda LiaCastaneda commented Oct 2, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Related to #23237 and feat(metric): Add output skewness metric to detect skewed plans easier

  • Closes #.

Rationale for this change

I think it would be useful to add the skewness metric for nodes that execute in partitioned mdoe -- right now the only node that includes that metric in the explain analyze is DataSourceExec for parquet.

This is specially useful to check how many partitions were idle during Aggregations or HashJoins for example.

What changes are included in this PR?

The PR adds a new metric output_rows_skew in the Explain Analyze output for the following nodes:

  • RepartitionExec
  • HashJoinExec
  • AggregateExec
  • BoundedWindowAggExec
  • WindowAggExec

Note that for a plan that has a Parquet DataSource the skewness might not be the same as for the rest of the nodes, for example we might have the following scenario:

DataSourceExec       [100, 100, 100, 100]  → 0% skewness, data is even
  FilterExec         [100,   0,   0,   0]  → 100%   (filter only matches in partition 0)
    AggregateExec    [  3,   0,   0,   0]  → 100%   (Partial: groups, not rows)
      RepartitionExec Hash  [1, 1, 1, 0]   → ~11%   (3 groups re-spread by key)

What is the testing strategy for this PR?

I added a Explain Analyze test

Are there any user-facing changes?

Yes, a new metric will show in the explain analyze plan for the nodes mentioned above.

@github-actions github-actions Bot added physical-expr Changes to the physical-expr crates core Core DataFusion crate sqllogictest SQL Logic Tests (.slt) physical-plan Changes to the physical-plan crate labels Oct 2, 2026
@codecov-commenter

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 90.90909% with 1 line in your changes missing coverage. Please review.
✅ Project coverage is 82.61%. Comparing base (c922f88) to head (ed306cb).
⚠️ Report is 1 commits behind head on main.

Files with missing lines Patch % Lines
...usion/physical-plan/src/windows/window_agg_exec.rs 0.00% 1 Missing ⚠️
Additional details and impacted files
@@            Coverage Diff             @@
##             main   #25977      +/-   ##
==========================================
- Coverage   82.61%   82.61%   -0.01%     
==========================================
  Files        1147     1147              
  Lines      444812   444818       +6     
  Branches   444812   444818       +6     
==========================================
- Hits       367502   367494       -8     
- Misses      54996    55006      +10     
- Partials    22314    22318       +4     

☔ View full report in Codecov by Harness.
📢 Have feedback on the report? Share it here.

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@rgbuilds rgbuilds left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new metric accounting looks sound, and the focused EXPLAIN ANALYZE test and affected SQL logic tests pass locally.

One small documentation suggestion: explain-usage.md currently lists output_rows_skew under Parquet scans; could it also mention the operators added here?

@gabotechs gabotechs left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

👍 Nice! no comments, looks good to me

@gabotechs
gabotechs added this pull request to the merge queue Oct 5, 2026
Merged via the queue into apache:main with commit c2baac5 Oct 5, 2026
44 checks passed
@LiaCastaneda

Copy link
Copy Markdown
Contributor Author

Thanks for the reviews!

One small documentation suggestion: explain-usage.md currently lists output_rows_skew under Parquet scans; could it also mention the operators added here?

I will make a small PR shortly :)

@LiaCastaneda

Copy link
Copy Markdown
Contributor Author

micro PR: #26048

gafiatulin pushed a commit to gafiatulin/datafusion that referenced this pull request Oct 5, 2026
…age (apache#26048)

Follow-up to apache#25977.

## Which issue does this PR close?

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes apache#123` indicates that this PR will close issue apache#123.
-->

- Closes #.

## Rationale for this change

now output_rows_skew is used in other operators than DataSourceExec, but
I forgot to add it to the doc

## What changes are included in this PR?

Small comment in the doc mentioning output_rows_skew is also used for
other operators than DataSourceExec(parquet)

## What is the testing strategy for this PR?

<!--
We typically require tests for all PRs in order to:
1. Prevent the code from being accidentally broken by subsequent changes
2. Serve as another way to document the expected behavior of the code

Briefly describe how this PR is tested, and point to the specific tests
you added. For example: 'This new feature is covered by the
`sqllogictest` cases added in `foo.slt`'.

If this PR does not add tests, explain why. For example, if the change
is already covered by existing tests, please mention it.

You should also check the `codecov` bot reply on this PR to confirm the
changed code is exercised.
-->

## Are there any user-facing changes?

Just doc.
LiaCastaneda added a commit to DataDog/datafusion that referenced this pull request Oct 6, 2026
…#25977) (#188)

<!--
We generally require a GitHub issue to be filed for all bug fixes and
enhancements and this helps us generate change logs for our releases.
You can link an issue to this PR using the GitHub syntax. For example
`Closes #123` indicates that this PR will close issue #123.
-->

Related to apache#23237 and
[feat(metric): Add output skewness metric to detect skewed plans
easier](apache#21211)

- Closes #.

I think it would be useful to add the skewness metric for nodes that
execute in partitioned mdoe -- right now the only node that includes
that metric in the explain analyze is DataSourceExec for parquet.

This is specially useful to check how many partitions were idle during
Aggregations or HashJoins for example.

The PR adds a new metric `output_rows_skew` in the Explain Analyze
output for the following nodes:

- RepartitionExec
- HashJoinExec
- AggregateExec
- BoundedWindowAggExec
- WindowAggExec

Note that for a plan that has a Parquet DataSource the skewness might
not be the same as for the rest of the nodes, for example we might have
the following scenario:

```
DataSourceExec       [100, 100, 100, 100]  → 0% skewness, data is even
  FilterExec         [100,   0,   0,   0]  → 100%   (filter only matches in partition 0)
    AggregateExec    [  3,   0,   0,   0]  → 100%   (Partial: groups, not rows)
      RepartitionExec Hash  [1, 1, 1, 0]   → ~11%   (3 groups re-spread by key)

```

I added a Explain Analyze test

Yes, a new metric will show in the explain analyze plan for the nodes
mentioned above.

(cherry picked from commit c2baac5)
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

core Core DataFusion crate physical-expr Changes to the physical-expr crates physical-plan Changes to the physical-plan crate sqllogictest SQL Logic Tests (.slt) v56.0.0

Projects

None yet

Development

Successfully merging this pull request may close these issues.

4 participants