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
11 changes: 11 additions & 0 deletions docs/source/user-guide/latest/tuning/operators.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,6 +59,17 @@ boundaries. A shuffled hash join can still filter probe batches after shuffle, b
its filter back to an earlier scan stage. Compare the [runtime-filter and scan metrics](../metrics.md#hash-joins)
with the setting disabled to distinguish reduced hash-probe work from reader I/O savings.

## Window Functions

`PERCENT_RANK`, `CUME_DIST`, `NTILE`, and aggregates whose frame ends at `UNBOUNDED FOLLOWING`, which includes an
aggregate over `PARTITION BY` with no `ORDER BY` such as `sum(x) OVER (PARTITION BY k)`, need a whole window partition
before they return anything. Comet buffers one window partition at a time for them and reserves it from its native
memory pool, but it does not yet provide spill-to-disk for them. A window partition that does not fit in its share of
the pool, for example because of a heavily skewed key or a window without `PARTITION BY`, fails the task with a memory

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.

This section describes a change in outcome. A window partition that does not fit its share now fails the task, where before it ran on untracked memory. Would it make sense to add an entry for it to the upgrade guide? docs/source/contributor-guide/config_conventions.md (Changing the Behavior of an Existing Config) asks for one, plus a spark.comet.legacy.* key. I'm not sure a key is expected for a memory accounting fix, so a note like the 1.1.0 "Memory Pool Limits" entry may be enough.

error where Spark would spill. Increase `spark.memory.offHeap.size` (see
[Configuring Comet Memory](memory.md#configuring-comet-memory)), or set `spark.comet.exec.window.enabled=false` to run
window functions in Spark.

## Adaptive Partial Aggregation

For high-cardinality grouping, Comet can bypass partial hash aggregation when it is not
Expand Down
2 changes: 2 additions & 0 deletions native/core/src/execution/operators/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,8 @@ mod scan;
mod shuffle_scan;
pub use csv_scan::init_csv_datasource_exec;
pub use shuffle_scan::ShuffleScanExec;
mod window_agg;
pub(crate) use window_agg::CometWindowAggExec;

/// Fixtures for the nested-nullability drift from
/// <https://github.com/apache/datafusion-comet/issues/5137>, shared by the `expand` and
Expand Down
Loading
Loading