Skip to content
Merged
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
88 changes: 88 additions & 0 deletions docs/source/_static/images/comet-executor-memory.svg
Loading
Sorry, something went wrong. Reload?
Sorry, we cannot display this file.
Sorry, this file is invalid so it cannot be displayed.
29 changes: 19 additions & 10 deletions docs/source/contributor-guide/memory_management.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ under the License.
This page describes how memory is budgeted, accounted, and enforced across the JVM/native
boundary. It is aimed at contributors working on memory pools, operators that reserve memory, or
anyone debugging an out-of-memory report. For user-facing tuning advice, see the
[Tuning Guide](../user-guide/latest/tuning.md).
[Memory Tuning](../user-guide/latest/tuning/memory.md) guide.

This page covers off-heap mode (`spark.memory.offHeap.enabled=true`) only. Comet also has an
on-heap mode, but it exists so that the Spark SQL test suite can run against Comet without changing
Expand Down Expand Up @@ -61,6 +61,15 @@ and no JVM metric measures them, yet they land squarely in container RSS. Comet
its own budget that is meant to shadow the physical one, and declares it to Spark so that the two
compete for a single number. The accuracy of that shadow is the central problem this page is about.

The picture the [Memory Tuning](../user-guide/latest/tuning/memory.md) guide gives users is
deliberately simple:

![Spark and Comet both use the JVM heap and share the off-heap memory pool, and the rest of Comet's native memory has to fit in the executor's memory overhead](../_static/images/comet-executor-memory.svg)

Most of this page is about the line between Comet's share of the off-heap pool and its share of the
memory overhead: which allocations are declared to the pool, and which land in the overhead with
nothing tracking them.

## Who allocates what

Enabling Comet does not add one new memory consumer, it adds several, and they are not all
Expand Down Expand Up @@ -426,7 +435,7 @@ reports the bytes Rust's allocator has handed out next to the pools' reservation
than per interval, enable tracing and compare `native_allocated` against
`comet_memory_reserved_total`; see [Tracing](tracing.md#analyzing-memory-usage).

[memory-usage-log]: ../user-guide/latest/tuning.md#sizing-the-overhead-from-the-memory-usage-log
[memory-usage-log]: ../user-guide/latest/tuning/memory.md#sizing-the-overhead-from-the-memory-usage-log

## What the container sees

Expand Down Expand Up @@ -522,14 +531,14 @@ much they matter:

## Debugging memory issues

| Tool | What it gives you |
| ------------------------------------------------------------------------------------------------ | ------------------------------------------------------------------------------------------- |
| `spark.comet.debug.memory=true` | `LoggingMemoryPool` logs every register/grow/shrink with the consumer name |
| `spark.comet.explain.native.enabled=true` | Native plan with per-operator metrics, including spill counts |
| [Memory usage log](../user-guide/latest/tuning.md#sizing-the-overhead-from-the-memory-usage-log) | Executor-wide native allocation vs pool reservations, logged every 10 seconds by default |
| [Tracing](tracing.md#analyzing-memory-usage) | `native_allocated` vs `comet_memory_reserved_total` per event; the accounting gap over time |
| `TrackConsumersPool` | Names the top 10 consumers in `ResourcesExhausted` messages (always on) |
| [`thresher`](https://github.com/cetra3/thresher) | Third-party crate that dumps a jemalloc heap profile at a threshold |
| Tool | What it gives you |
| ------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------- |
| `spark.comet.debug.memory=true` | `LoggingMemoryPool` logs every register/grow/shrink with the consumer name |
| `spark.comet.explain.native.enabled=true` | Native plan with per-operator metrics, including spill counts |
| [Memory usage log](../user-guide/latest/tuning/memory.md#sizing-the-overhead-from-the-memory-usage-log) | Executor-wide native allocation vs pool reservations, logged every 10 seconds by default |
| [Tracing](tracing.md#analyzing-memory-usage) | `native_allocated` vs `comet_memory_reserved_total` per event; the accounting gap over time |
| `TrackConsumersPool` | Names the top 10 consumers in `ResourcesExhausted` messages (always on) |
| [`thresher`](https://github.com/cetra3/thresher) | Third-party crate that dumps a jemalloc heap profile at a threshold |

A checklist for triaging an executor OOM kill:

Expand Down
2 changes: 1 addition & 1 deletion docs/source/contributor-guide/plugin_overview.md
Original file line number Diff line number Diff line change
Expand Up @@ -52,7 +52,7 @@ and skips the remaining steps. Otherwise it:
`spark.sql.queryExecutionListeners` when `spark.comet.metrics.enabled=true`.
- Logs a warning for settings that are likely to cause problems, such as an unset `spark.executor.memoryOverhead`.

The plugin does not change any executor memory setting. The [Tuning Guide](../user-guide/latest/tuning.md) covers how
The plugin does not change any executor memory setting. The [Memory Tuning](../user-guide/latest/tuning/memory.md) guide covers how
to size them.

When the driver or an executor stops, the plugin shuts down Comet's native tokio runtime in that JVM.
Expand Down
2 changes: 1 addition & 1 deletion docs/source/contributor-guide/tracing.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,7 +31,7 @@ Native memory is traced as `native_allocated`. Comet wraps whichever global allo
selected and counts the bytes it has handed out. It counts only what Rust code allocated, so it can be compared against the
memory pool's reservations without the allocator's own caching in the way. The same figure appears in
the executor's periodic memory usage log, which does not need tracing; see
[Sizing the Overhead from the Memory Usage Log](../user-guide/latest/tuning.md#sizing-the-overhead-from-the-memory-usage-log).
[Sizing the Overhead from the Memory Usage Log](../user-guide/latest/tuning/memory.md#sizing-the-overhead-from-the-memory-usage-log).

Enabling the `jemalloc` feature adds a second measure, `jemalloc_allocated`, which also includes
jemalloc's own metadata and fragmentation:
Expand Down
14 changes: 13 additions & 1 deletion docs/source/user-guide/latest/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -72,11 +72,23 @@ to read more.
:hidden:

Understanding Comet Plans <understanding-comet-plans>
Tuning Guide <tuning>
Metrics Guide <metrics>
In-Memory Cache <in-memory-cache>
PyArrow UDF Acceleration <pyarrow-udfs>

.. toctree::
:maxdepth: 1
:caption: Tuning
:hidden:

Overview <tuning>
Memory <tuning/memory>
Shuffle <tuning/shuffle>
Celeborn Shuffle <tuning/celeborn>
Scans <tuning/scans>
Operators <tuning/operators>
Row/Columnar Transitions <tuning/transitions>

.. toctree::
:maxdepth: 1
:caption: Integrations
Expand Down
2 changes: 1 addition & 1 deletion docs/source/user-guide/latest/installation.md
Original file line number Diff line number Diff line change
Expand Up @@ -260,7 +260,7 @@ Some cluster managers may require additional configuration, see <https://spark.a
### Memory tuning

In addition to Apache Spark memory configuration parameters, Comet introduces additional parameters to configure memory
allocation for native execution. See [Comet Memory Tuning](./tuning.md) for details.
allocation for native execution. See [Comet Memory Tuning](./tuning/memory.md) for details.

### Kryo serialization

Expand Down
16 changes: 8 additions & 8 deletions docs/source/user-guide/latest/metrics.md
Original file line number Diff line number Diff line change
Expand Up @@ -76,21 +76,21 @@ the value would always be 0.

Native aggregates with grouping keys report these additional metrics:

| Metric | Description |
| ------------------------------------ | ---------------------------------------------------------------------------------------------------------------------------------- |
| `rows bypassing partial aggregation` | Input rows passed through without partial aggregation. See [Adaptive Partial Aggregation](tuning.md#adaptive-partial-aggregation). |
| `number of spills` | Number of times the aggregate spilled to disk. |
| `total spilled bytes` | Bytes written to aggregate spill files. |
| `number of spilled rows` | Rows written to aggregate spill files. |
| `peak native aggregate memory` | Peak memory used by the native aggregate. |
| Metric | Description |
| ------------------------------------ | -------------------------------------------------------------------------------------------------------------------------------------------- |
| `rows bypassing partial aggregation` | Input rows passed through without partial aggregation. See [Adaptive Partial Aggregation](tuning/operators.md#adaptive-partial-aggregation). |
| `number of spills` | Number of times the aggregate spilled to disk. |
| `total spilled bytes` | Bytes written to aggregate spill files. |
| `number of spilled rows` | Rows written to aggregate spill files. |
| `peak native aggregate memory` | Peak memory used by the native aggregate. |

Spill bytes from native sorts, aggregates, and sort-merge joins are also added to Spark's task-level
`diskBytesSpilled` metric in every stage, not only in shuffle stages.

### Hash Joins

With `spark.comet.exec.join.dynamicFilter.enabled=true`, native broadcast and shuffled hash joins
report these additional metric keys. See [Join Runtime Filters](tuning.md#join-runtime-filters) for
report these additional metric keys. See [Join Runtime Filters](tuning/operators.md#join-runtime-filters) for
eligibility and reader restrictions.

| Metric | Description |
Expand Down
Loading
Loading