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
16 changes: 13 additions & 3 deletions docs/source/contributor-guide/memory_management.md
Original file line number Diff line number Diff line change
Expand Up @@ -283,9 +283,19 @@ reading it as one overstates the memory a multi-operator task can use.

### Task-shared pools and their lifetime

A single Spark task can run more than one native plan concurrently: a shuffle runs the pre-shuffle
operators and the shuffle writer as separate native execution contexts. If each got its own pool,
the per-task limit would be enforced once per plan rather than once per task.
A single Spark task can run more than one native plan at a time. A native shuffle is not one of
these cases, because its writer is planned together with the native operators that feed it. These
operators do split a task's native work into separate plans:

- `CometUnionExec` and `CometCoalesceExec` read their children through the JVM, so the native plan
above them and the native plans below them are separate.
- `CometCollectLimitExec` and `CometTakeOrderedAndProjectExec` apply their limit in a native plan
of their own, over the output of the plan below them.
- A native Parquet or Iceberg write runs its writer as a native plan of its own, over the output of
the plan below it.

If each plan got its own pool, the per-task limit would be enforced once per plan rather than once
per task.

`acquire_task_shared_pool` (`task_shared.rs`) keeps a process-wide
`HashMap<task_attempt_id, Weak<TaskSharedMemoryPool>>`. Plans in the same task upgrade the existing
Expand Down
7 changes: 3 additions & 4 deletions docs/source/user-guide/latest/tuning.md
Original file line number Diff line number Diff line change
Expand Up @@ -127,10 +127,9 @@ The valid pool types are:
- `fair_unified` (default when `spark.memory.offHeap.enabled=true` is set)
- `greedy_unified`

Both pool types are shared across all native execution contexts within the same Spark task. When
Comet executes a shuffle, it runs two native execution contexts concurrently (e.g. one for
pre-shuffle operators and one for the shuffle writer). The shared pool ensures that the combined
memory usage stays within the per-task limit.
Both pool types are shared by all the native plans in the same Spark task. A task can run more than
one native plan at a time, for example the native operators on either side of a union or a
coalesce. The shared pool ensures that their combined memory usage stays within the per-task limit.

The `fair_unified` pool prevents operators from using more than an even fraction of the available memory
(i.e. `pool_size / num_reservations`). This pool works best when you know beforehand
Expand Down