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
4 changes: 3 additions & 1 deletion docs/source/contributor-guide/memory_management.md
Original file line number Diff line number Diff line change
Expand Up @@ -432,7 +432,9 @@ diverge for several structural reasons:
(see [The unified pools](#the-unified-pools)). `reserved()` includes it, but Spark's memory
manager does not, so until it is repaid Spark can hand the same bytes to another consumer or task.
The `overcommit` figure in the pool's `Display` output and `try_grow` errors shows how much is
outstanding.
outstanding. The executor's memory usage log leaves it out of the `reserved` figure it reports,
so that the log counts it as untracked: like an undeclared allocation, it needs room beyond what
Spark has handed out. Tracing's `comet_memory_reserved_total` still includes it.

The practical consequence is that `reserved()` is a lower bound on Comet's real footprint, and the
gap is workload-dependent. The margin that covers it has to come from
Expand Down
15 changes: 9 additions & 6 deletions docs/source/user-guide/latest/tuning/memory.md
Original file line number Diff line number Diff line change
Expand Up @@ -141,8 +141,11 @@ Comet native memory usage: allocated 5412.3 MiB, reserved 3890.0 MiB (16 native

- `allocated` is the memory that Comet's native code has allocated and not yet freed, whether or not
a pool tracks it.
- `reserved` is the part that Comet's memory pools track. It is charged against
`spark.memory.offHeap.size`, so the container already has room for it.
- `reserved` is the part that Comet's memory pools have reserved from Spark's off-heap memory. It is
charged against `spark.memory.offHeap.size`, so the container already has room for it. A pool
sometimes has to track memory that Spark could not grant, such as a spilled batch read back from
disk while the off-heap memory is full. `reserved` leaves that memory out, since nothing charges it
against `spark.memory.offHeap.size`.
- `JVM Arrow allocated` is the Arrow memory Comet holds on the JVM side, such as batches read from
Comet's in-memory cache, broadcast data, and batches exchanged with native code or Python workers.
The part imported from native was allocated by Comet's native code, so `allocated` already counts
Expand All @@ -151,10 +154,10 @@ Comet native memory usage: allocated 5412.3 MiB, reserved 3890.0 MiB (16 native
which allocator a buffer is charged to rather than where it was allocated, so treat the JVM's
part as an estimate.

Comet's untracked memory is what Comet holds outside the JVM heap that no pool reserves: `allocated`,
plus the JVM Arrow figure less the part imported from native, minus `reserved`. It is the part of
Comet's footprint that has to fit in `spark.executor.memoryOverhead`, alongside the JVM's own
non-heap memory. To size the overhead from it:
Comet's untracked memory is what Comet holds outside the JVM heap that Spark's off-heap pool does
not account for: `allocated`, plus the JVM Arrow figure less the part imported from native, minus
`reserved`. It is the part of Comet's footprint that has to fit in `spark.executor.memoryOverhead`,
alongside the JVM's own non-heap memory. To size the overhead from it:

1. Run a representative workload and find the line with the most untracked memory in each
executor's log. Take all the figures from the same line: they are sampled together, and figures
Expand Down
50 changes: 47 additions & 3 deletions native/core/src/execution/jni_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -106,7 +106,7 @@ use tokio::sync::mpsc;
use tokio::task::JoinHandle;

use crate::execution::memory_pools::{
create_memory_pool, parse_memory_pool_config, PlanMemoryPool,
create_memory_pool, overcommit, parse_memory_pool_config, PlanMemoryPool,
};
use crate::execution::operators::{ScanExec, ShuffleScanExec};
use crate::execution::shuffle::{
Expand Down Expand Up @@ -263,6 +263,17 @@ fn sum_reserved(pools: &[Arc<dyn MemoryPool>]) -> usize {
pools.iter().map(|pool| pool.reserved()).sum()
}

/// Bytes reserved across `pools`, less the part that Spark has not granted them; see
/// [`MemoryUsage::pools_reserved`].
fn sum_reserved_less_overcommit(pools: &[Arc<dyn MemoryPool>]) -> usize {
pools
.iter()
// The two figures are read at different moments, so a `grow` in between can leave the
// overcommit larger than the reservation read before it.
.map(|pool| pool.reserved().saturating_sub(overcommit(pool)))
.sum()
}

fn total_reserved_for_thread(thread_id: u64) -> usize {
sum_reserved(&snapshot_registry(Some(thread_id)).thread_pools)
}
Expand Down Expand Up @@ -310,7 +321,9 @@ struct MemoryUsage {
/// Bytes handed out by the Rust global allocator, process-wide.
native_allocated: usize,
/// Bytes reserved across every live Comet memory pool, counting each pool once however many
/// plans share it.
/// plans share it, less any the pools recorded beyond what Spark granted them; see
/// [`overcommit`]. Spark's off-heap pool does not account for those bytes, so the log counts
/// them with the native memory that no pool tracks.
pools_reserved: usize,
/// Live memory pools. With the task-shared pool types, which include both defaults, that is one
/// per task running native plans.
Expand All @@ -327,7 +340,7 @@ fn memory_usage() -> MemoryUsage {
let snapshot = snapshot_registry(None);
MemoryUsage {
native_allocated: crate::alloc_accounting::current_balance(),
pools_reserved: sum_reserved(&snapshot.all_pools),
pools_reserved: sum_reserved_less_overcommit(&snapshot.all_pools),
pools: snapshot.all_pools.len(),
plans: snapshot.plans,
}
Expand Down Expand Up @@ -2321,6 +2334,37 @@ mod tests {
drop(own_reservation);
}

/// The memory usage log leaves overcommit out of the reservations it reports, because Spark's
/// off-heap pool does not account for it, so the log counts it with the native memory that no
/// pool tracks. Tracing's process total still reports everything the pools recorded.
#[test]
fn memory_usage_leaves_out_what_spark_did_not_grant() {
use crate::execution::memory_pools::{
create_memory_pool_with_fake_spark, MemoryPoolConfig, MemoryPoolType,
};

let _guard = serial();
let before = memory_usage();
let traced_before = total_reserved_across_threads();
// A task's pool as `greedy_unified` creates it, where Spark grants at most 4096 bytes.
let config = MemoryPoolConfig::new(MemoryPoolType::GreedyUnified, 0);
let pool = create_memory_pool_with_fake_spark(&config, -6101, 4096);
let _registration = ThreadMemoryPoolRegistration::new(21, -6101, Arc::clone(&pool));
let reservation = MemoryConsumer::new("spill reader").register(&pool);

// A spilled batch read back from disk is recorded in full, although Spark grants only 4096
// of its 6144 bytes.
reservation.grow(6144);
assert_eq!(total_reserved_across_threads() - traced_before, 6144);
assert_eq!(memory_usage().pools_reserved - before.pools_reserved, 4096);

// Freeing memory repays the overcommit before anything goes back to Spark.
reservation.shrink(2048);
assert_eq!(memory_usage().pools_reserved - before.pools_reserved, 4096);
reservation.shrink(1024);
assert_eq!(memory_usage().pools_reserved - before.pools_reserved, 3072);
}

/// Stands in for a `CometFairMemoryPool` whose lock is held across a Spark acquire: it counts
/// its reservation reads, and notes whether the registry lock was held during any of them.
#[derive(Debug, Default)]
Expand Down
22 changes: 7 additions & 15 deletions native/core/src/execution/memory_pools/fair_pool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -18,11 +18,8 @@
use std::{
collections::HashMap,
fmt::{Debug, Display, Formatter, Result as FmtResult},
sync::Arc,
};

use jni::objects::{Global, JObject};

use super::spark_memory::SparkMemory;
use datafusion::common::resources_err;
use datafusion::execution::memory_pool::MemoryConsumer;
Expand Down Expand Up @@ -72,18 +69,7 @@ impl Debug for CometFairMemoryPool {
}

impl CometFairMemoryPool {
pub fn new(
task_memory_manager_handle: Arc<Global<JObject<'static>>>,
pool_size: usize,
task_attempt_id: i64,
) -> CometFairMemoryPool {
Self::with_spark(
SparkMemory::new(task_memory_manager_handle, task_attempt_id),
pool_size,
)
}

fn with_spark(spark: SparkMemory, pool_size: usize) -> CometFairMemoryPool {
pub(super) fn with_spark(spark: SparkMemory, pool_size: usize) -> CometFairMemoryPool {
Self {
spark,
pool_size,
Expand All @@ -93,6 +79,11 @@ impl CometFairMemoryPool {
}),
}
}

/// The part of [`MemoryPool::reserved`] that Spark has not granted; see [`SparkMemory`].
pub(super) fn overcommit(&self) -> usize {
self.spark.overcommit()
}
}

impl Display for CometFairMemoryPool {
Expand Down Expand Up @@ -213,6 +204,7 @@ impl MemoryPool for CometFairMemoryPool {
mod tests {
use super::super::spark_memory::fake::FakeSpark;
use super::*;
use std::sync::Arc;

#[test]
fn grow_past_the_fair_limit_is_recorded_and_refuses_the_next_try_grow() {
Expand Down
89 changes: 80 additions & 9 deletions native/core/src/execution/memory_pools/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ mod unified_pool;
use datafusion::execution::memory_pool::{MemoryPool, TrackConsumersPool, UnboundedMemoryPool};
use fair_pool::CometFairMemoryPool;
use jni::objects::{Global, JObject};
use spark_memory::SparkMemory;
use std::num::NonZeroUsize;
use std::sync::Arc;
use unified_pool::CometUnifiedMemoryPool;
Expand All @@ -42,6 +43,18 @@ pub(crate) fn create_memory_pool(
memory_pool_config: &MemoryPoolConfig,
comet_task_memory_manager: Arc<Global<JObject<'static>>>,
task_attempt_id: i64,
) -> Arc<dyn MemoryPool> {
create_pool(memory_pool_config, task_attempt_id, || {
SparkMemory::new(comet_task_memory_manager, task_attempt_id)
})
}

/// Creates the pool that [`create_memory_pool`] does, with `spark` connecting it to Spark's memory
/// manager, so that tests can connect it to a fake instead.
fn create_pool(
memory_pool_config: &MemoryPoolConfig,
task_attempt_id: i64,
spark: impl FnOnce() -> SparkMemory,
) -> Arc<dyn MemoryPool> {
const NUM_TRACKED_CONSUMERS: usize = 10;

Expand All @@ -57,18 +70,76 @@ pub(crate) fn create_memory_pool(

match pool_type {
MemoryPoolType::GreedyUnified => acquire_task_shared_pool(task_attempt_id, || {
tracked(CometUnifiedMemoryPool::new(
comet_task_memory_manager,
task_attempt_id,
))
tracked(CometUnifiedMemoryPool::with_spark(spark()))
}),
MemoryPoolType::FairUnified => acquire_task_shared_pool(task_attempt_id, || {
tracked(CometFairMemoryPool::new(
comet_task_memory_manager,
pool_size,
task_attempt_id,
))
tracked(CometFairMemoryPool::with_spark(spark(), pool_size))
}),
MemoryPoolType::Unbounded => Arc::new(UnboundedMemoryPool::default()),
}
}

/// The bytes that `pool` has recorded beyond what Spark granted it, which it carries as overcommit
/// until Spark grants them or the pool frees memory; see [`SparkMemory`]. This looks through the
/// wrappers that [`create_memory_pool`] puts around a Comet pool, and is zero for a pool that
/// takes nothing from Spark or that the function did not create.
pub(crate) fn overcommit(pool: &Arc<dyn MemoryPool>) -> usize {
let pool = task_shared::unwrap_task_shared(pool).unwrap_or(pool);
if let Some(tracked) = pool.downcast_ref::<TrackConsumersPool<CometUnifiedMemoryPool>>() {
tracked.inner().overcommit()
} else if let Some(tracked) = pool.downcast_ref::<TrackConsumersPool<CometFairMemoryPool>>() {
tracked.inner().overcommit()
} else {
0
}
}

/// [`create_memory_pool`], connected to a fake Spark that grants at most `limit` bytes.
#[cfg(test)]
pub(crate) fn create_memory_pool_with_fake_spark(
memory_pool_config: &MemoryPoolConfig,
task_attempt_id: i64,
limit: usize,
) -> Arc<dyn MemoryPool> {
let fake = spark_memory::fake::FakeSpark::with(limit);
create_pool(memory_pool_config, task_attempt_id, || fake.memory())
}

#[cfg(test)]
mod tests {
use super::*;
use datafusion::execution::memory_pool::MemoryConsumer;

#[test]
fn overcommit_is_read_through_the_wrappers_of_each_pool_type() {
// Task-shared pools are keyed by task attempt process-wide, so each gets its own id.
for (name, task_attempt_id, pool_type) in [
("greedy_unified", -3001, MemoryPoolType::GreedyUnified),
("fair_unified", -3002, MemoryPoolType::FairUnified),
] {
let config = MemoryPoolConfig::new(pool_type, 1000);
let pool = create_memory_pool_with_fake_spark(&config, task_attempt_id, 100);
let reservation = MemoryConsumer::new("spill reader").register(&pool);

// Spark grants 100 of the 150 bytes, and the pool records all of them.
reservation.grow(150);
assert_eq!(pool.reserved(), 150, "{name}");
assert_eq!(overcommit(&pool), 50, "{name}");

// Freeing memory repays the overcommit first.
reservation.shrink(30);
assert_eq!(overcommit(&pool), 20, "{name}");
drop(reservation);
assert_eq!(overcommit(&pool), 0, "{name}");
}
}

#[test]
fn a_pool_that_takes_nothing_from_spark_has_no_overcommit() {
let config = MemoryPoolConfig::new(MemoryPoolType::Unbounded, 0);
let pool = create_memory_pool_with_fake_spark(&config, -3003, 0);
let reservation = MemoryConsumer::new("sort").register(&pool);
reservation.grow(150);
assert_eq!(overcommit(&pool), 0);
}
}
6 changes: 6 additions & 0 deletions native/core/src/execution/memory_pools/task_shared.rs
Original file line number Diff line number Diff line change
Expand Up @@ -117,6 +117,12 @@ pub(crate) fn acquire_task_shared_pool(
memory_pool
}

/// The pool that `pool` wraps, if it is one that [`acquire_task_shared_pool`] returned.
pub(super) fn unwrap_task_shared(pool: &Arc<dyn MemoryPool>) -> Option<&Arc<dyn MemoryPool>> {
pool.downcast_ref::<TaskSharedMemoryPool>()
.map(|shared| &shared.inner)
}

#[cfg(test)]
mod tests {
use super::*;
Expand Down
24 changes: 8 additions & 16 deletions native/core/src/execution/memory_pools/unified_pool.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,18 +17,14 @@

use std::{
fmt::{Debug, Display, Formatter, Result as FmtResult},
sync::{
atomic::{AtomicUsize, Ordering::Relaxed},
Arc,
},
sync::atomic::{AtomicUsize, Ordering::Relaxed},
};

use super::spark_memory::SparkMemory;
use datafusion::{
common::{resources_datafusion_err, DataFusionError},
execution::memory_pool::{MemoryPool, MemoryReservation},
};
use jni::objects::{Global, JObject};
use log::warn;

/// A DataFusion `MemoryPool` implementation for Comet that delegates to
Expand All @@ -49,22 +45,17 @@ impl Debug for CometUnifiedMemoryPool {
}

impl CometUnifiedMemoryPool {
pub fn new(
task_memory_manager_handle: Arc<Global<JObject<'static>>>,
task_attempt_id: i64,
) -> CometUnifiedMemoryPool {
Self::with_spark(SparkMemory::new(
task_memory_manager_handle,
task_attempt_id,
))
}

fn with_spark(spark: SparkMemory) -> CometUnifiedMemoryPool {
pub(super) fn with_spark(spark: SparkMemory) -> CometUnifiedMemoryPool {
Self {
spark,
used: AtomicUsize::new(0),
}
}

/// The part of [`MemoryPool::reserved`] that Spark has not granted; see [`SparkMemory`].
pub(super) fn overcommit(&self) -> usize {
self.spark.overcommit()
}
}

impl Drop for CometUnifiedMemoryPool {
Expand Down Expand Up @@ -162,6 +153,7 @@ mod tests {
use super::super::spark_memory::fake::FakeSpark;
use super::*;
use datafusion::execution::memory_pool::MemoryConsumer;
use std::sync::Arc;

#[test]
fn grow_past_spark_is_recorded_and_refuses_try_grow_until_repaid() {
Expand Down
12 changes: 7 additions & 5 deletions spark/src/main/scala/org/apache/comet/CometExecIterator.scala
Original file line number Diff line number Diff line change
Expand Up @@ -584,15 +584,17 @@ object CometExecIterator extends Logging {
* A warning if the executor's native footprint exceeds `limitBytes`, the container's memory
* outside the JVM heap; see [[nativeMemoryLimit]].
*
* The footprint is the memory Comet holds outside the JVM heap that its pools do not track,
* The footprint is the memory Comet holds outside the JVM heap that Spark does not account for,
* `allocated + jvmArrow.allocatedByJvm - reserved`, plus `sparkOffHeapUsed`, everything in use
* in Spark's off-heap pool, which includes Comet's reservations as well as Spark's own off-heap
* execution and storage memory. The JVM's Arrow memory is added before the reservations are
* subtracted, because a native operator that holds on to a batch the JVM allocated, such as a
* sort buffering its input, reserves it: those bytes are in `reserved` and in the JVM figure
* but not in `allocated`, and are counted once. Comparing the sum rather than the untracked
* part against the overhead alone counts the part of `spark.memory.offHeap.size` that nothing
* has acquired at that moment, which untracked memory can occupy until Spark hands it out. The
* but not in `allocated`, and are counted once. Neither `reserved` nor `sparkOffHeapUsed`
* includes memory that a pool recorded beyond what Spark granted it, so that memory is counted
* as untracked; see [[Native.getMemoryUsage]]. Comparing the sum rather than the untracked part
* against the overhead alone counts the part of `spark.memory.offHeap.size` that nothing has
* acquired at that moment, which untracked memory can occupy until Spark hands it out. The
* limit also has to hold the JVM's own non-heap memory, so by the time the footprint exceeds it
* the executor has likely outgrown its container.
*/
Expand All @@ -604,7 +606,7 @@ object CometExecIterator extends Logging {
val untracked = math.max(usage(0) + jvmArrow.allocatedByJvm - usage(1), 0L)
val footprint = untracked + sparkOffHeapUsed
if (footprint > limitBytes) {
Some(s"Comet memory not tracked by any memory pool (${toMiB(untracked)}, native and JVM " +
Some(s"Comet memory that Spark does not account for (${toMiB(untracked)}, native and JVM " +
s"Arrow) plus Spark's off-heap memory in use (${toMiB(sparkOffHeapUsed)}, including " +
s"Comet's reservations) is ${toMiB(footprint)}, more than the ${toMiB(limitBytes)} the " +
"executor's container has outside the JVM heap (spark.memory.offHeap.size plus the " +
Expand Down
Loading
Loading