From 236efa250d13ad9b69f53306dbf342ecd687b6b0 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Sun, 27 Sep 2026 09:36:49 -0600 Subject: [PATCH] fix: count pool overcommit as untracked memory in the native memory usage log The executor's memory usage log reports allocated and reserved, and both the memory tuning guide and the container warning read the difference as native memory that has to fit outside spark.memory.offHeap.size. Since #6128 a pool's reserved() also includes overcommit, the bytes it records when Spark grants less than a grow asked for. Spark's off-heap pool does not account for those bytes, so while a pool was overcommitted the untracked figure and the warning's footprint both came out low by the amount of the overcommit. The log's reserved figure now leaves out each pool's overcommit, read through the wrappers that create_memory_pool puts around the Comet pools. Tracing's comet_memory_reserved_total is unchanged. Closes #6260. --- .../contributor-guide/memory_management.md | 4 +- .../source/user-guide/latest/tuning/memory.md | 7 +- native/core/src/execution/jni_api.rs | 50 ++++++++++- .../src/execution/memory_pools/fair_pool.rs | 22 ++--- native/core/src/execution/memory_pools/mod.rs | 89 +++++++++++++++++-- .../src/execution/memory_pools/task_shared.rs | 6 ++ .../execution/memory_pools/unified_pool.rs | 24 ++--- .../org/apache/comet/CometExecIterator.scala | 10 ++- .../main/scala/org/apache/comet/Native.scala | 4 +- 9 files changed, 165 insertions(+), 51 deletions(-) diff --git a/docs/source/contributor-guide/memory_management.md b/docs/source/contributor-guide/memory_management.md index 9768698b0d8..d9cb338235e 100644 --- a/docs/source/contributor-guide/memory_management.md +++ b/docs/source/contributor-guide/memory_management.md @@ -422,7 +422,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 diff --git a/docs/source/user-guide/latest/tuning/memory.md b/docs/source/user-guide/latest/tuning/memory.md index 0fb4288abac..d8bd23fa25c 100644 --- a/docs/source/user-guide/latest/tuning/memory.md +++ b/docs/source/user-guide/latest/tuning/memory.md @@ -154,8 +154,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`. The difference between the two, `allocated - reserved`, is Comet's untracked native memory. It is the part of Comet's footprint that has to fit in `spark.executor.memoryOverhead`, alongside the diff --git a/native/core/src/execution/jni_api.rs b/native/core/src/execution/jni_api.rs index e199d5282fd..a39de53610d 100644 --- a/native/core/src/execution/jni_api.rs +++ b/native/core/src/execution/jni_api.rs @@ -104,7 +104,7 @@ use std::{ use tokio::runtime::{Handle, Runtime}; use tokio::sync::mpsc; -use crate::execution::memory_pools::{create_memory_pool, parse_memory_pool_config}; +use crate::execution::memory_pools::{create_memory_pool, overcommit, parse_memory_pool_config}; use crate::execution::operators::{ScanExec, ShuffleScanExec}; use crate::execution::shuffle::{ decode_remote_shuffle_batch, read_ipc_compressed, CompressionCodec, ShuffleWriterExec, @@ -260,6 +260,17 @@ fn sum_reserved(pools: &[Arc]) -> 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]) -> 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) } @@ -307,7 +318,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. @@ -324,7 +337,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, } @@ -2232,6 +2245,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)] diff --git a/native/core/src/execution/memory_pools/fair_pool.rs b/native/core/src/execution/memory_pools/fair_pool.rs index 2fbd8224bc3..56b09da7365 100644 --- a/native/core/src/execution/memory_pools/fair_pool.rs +++ b/native/core/src/execution/memory_pools/fair_pool.rs @@ -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; @@ -72,18 +69,7 @@ impl Debug for CometFairMemoryPool { } impl CometFairMemoryPool { - pub fn new( - task_memory_manager_handle: Arc>>, - 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, @@ -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 { @@ -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() { diff --git a/native/core/src/execution/memory_pools/mod.rs b/native/core/src/execution/memory_pools/mod.rs index 69311698b2e..0b62d2f7e17 100644 --- a/native/core/src/execution/memory_pools/mod.rs +++ b/native/core/src/execution/memory_pools/mod.rs @@ -25,6 +25,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; @@ -40,6 +41,18 @@ pub(crate) fn create_memory_pool( memory_pool_config: &MemoryPoolConfig, comet_task_memory_manager: Arc>>, task_attempt_id: i64, +) -> Arc { + 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 { const NUM_TRACKED_CONSUMERS: usize = 10; @@ -55,18 +68,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) -> usize { + let pool = task_shared::unwrap_task_shared(pool).unwrap_or(pool); + if let Some(tracked) = pool.downcast_ref::>() { + tracked.inner().overcommit() + } else if let Some(tracked) = pool.downcast_ref::>() { + 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 { + 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); + } +} diff --git a/native/core/src/execution/memory_pools/task_shared.rs b/native/core/src/execution/memory_pools/task_shared.rs index b5b4da61f9f..4b13083b491 100644 --- a/native/core/src/execution/memory_pools/task_shared.rs +++ b/native/core/src/execution/memory_pools/task_shared.rs @@ -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) -> Option<&Arc> { + pool.downcast_ref::() + .map(|shared| &shared.inner) +} + #[cfg(test)] mod tests { use super::*; diff --git a/native/core/src/execution/memory_pools/unified_pool.rs b/native/core/src/execution/memory_pools/unified_pool.rs index f023d51af23..f32b2d61858 100644 --- a/native/core/src/execution/memory_pools/unified_pool.rs +++ b/native/core/src/execution/memory_pools/unified_pool.rs @@ -17,10 +17,7 @@ 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; @@ -28,7 +25,6 @@ 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 @@ -49,22 +45,17 @@ impl Debug for CometUnifiedMemoryPool { } impl CometUnifiedMemoryPool { - pub fn new( - task_memory_manager_handle: Arc>>, - 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 { @@ -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() { diff --git a/spark/src/main/scala/org/apache/comet/CometExecIterator.scala b/spark/src/main/scala/org/apache/comet/CometExecIterator.scala index 5b70fbaaf24..7471fdf8ea7 100644 --- a/spark/src/main/scala/org/apache/comet/CometExecIterator.scala +++ b/spark/src/main/scala/org/apache/comet/CometExecIterator.scala @@ -545,10 +545,12 @@ 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 native memory Comet's pools do not track, `allocated - reserved`, plus + * The footprint is the native memory Spark does not account for, `allocated - 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. Comparing the sum - * rather than the untracked part against the overhead alone counts the part of + * reservations as well as Spark's own off-heap execution and storage memory. 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 @@ -562,7 +564,7 @@ object CometExecIterator extends Logging { val footprint = untracked + sparkOffHeapUsed if (footprint > limitBytes) { Some( - s"Comet native memory not tracked by any memory pool (${toMiB(untracked)}) plus " + + s"Comet native memory that Spark does not account for (${toMiB(untracked)}) plus " + s"Spark's off-heap memory in use (${toMiB(sparkOffHeapUsed)}, including Comet's " + 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 " + diff --git a/spark/src/main/scala/org/apache/comet/Native.scala b/spark/src/main/scala/org/apache/comet/Native.scala index 664cab7959c..89215e8a52f 100644 --- a/spark/src/main/scala/org/apache/comet/Native.scala +++ b/spark/src/main/scala/org/apache/comet/Native.scala @@ -272,7 +272,9 @@ class Native extends NativeBase { * @return * `[nativeAllocated, poolsReserved, pools, plans]`. `nativeAllocated` is the bytes the native * allocator has handed out. `poolsReserved` is the bytes reserved across every Comet memory - * pool, counting a pool shared by several plans once. `pools` is the number of live pools, + * pool, counting a pool shared by several plans once, less any that a pool recorded beyond + * what Spark granted it. Spark's off-heap pool does not account for those bytes, so they are + * counted with the native memory that no pool tracks. `pools` is the number of live pools, * which with the default task-shared pool types is one per task running native plans, and * `plans` is the number of native plans created and not yet released. */