dwsmith1983 commented on code in PR #5613:
URL: https://github.com/apache/datafusion-comet/pull/5613#discussion_r4105532543


##########
native/core/src/execution/memory_pools/fair_pool.rs:
##########
@@ -148,37 +295,90 @@ impl MemoryPool for CometFairMemoryPool {
         additional: usize,
     ) -> Result<(), DataFusionError> {
         if additional > 0 {
-            let mut state = self.state.lock();
-            let num = state.num;
-            let limit = self
-                .pool_size
-                .checked_div(num)
-                .expect("overflow in checked_div");
-            // We use state.used instead of reservation.size() because 
DataFusion 53+
-            // calls pool.try_grow() before incrementing the reservation's 
atomic size,
-            // so reservation.size() would not include prior grows.
-            let used = state.used;
-            if limit < used + additional {
-                return resources_err!(
-                    "Failed to acquire {additional} bytes where {used} bytes 
already reserved ({} bytes overcommitted) and the fair limit is {limit} bytes, 
{num} registered",
-                    self.spark.overcommit()
-                );
-            }
-
-            // A partial grant is handed back and refused, which triggers 
spilling in the caller.
-            if let Err(refusal) = self.spark.try_acquire(additional)? {
-                return resources_err!(
-                    "Failed to acquire {} bytes plus {} bytes overcommitted, 
only got {} bytes. Reserved: {} bytes",
-                    additional,
-                    refusal.overcommit,
-                    refusal.granted,
-                    state.used
-                );
-            }
-            state.used = state
-                .used
-                .checked_add(additional)
-                .expect("overflow in checked_add");
+            // Checking the fair limit and reserving the bytes is one atomic 
step, so concurrent
+            // grows can never jointly exceed pool_size / num. The blocking 
JVM calls then run
+            // without any lock held, and the reservation rolls back if the 
JVM does not back it.
+            {
+                let mut state = self.state.lock();
+                let num = state.num;
+                let limit = self
+                    .pool_size
+                    .checked_div(num)
+                    .expect("overflow in checked_div");
+                // The pool tracks one total across every consumer and checks 
the fair limit
+                // against that total, not against this reservation's own size.
+                let used = state.used;
+                match used.checked_add(additional) {
+                    Some(total) if total <= limit => state.used = total,
+                    _ => {
+                        return resources_err!(
+                            "Failed to acquire {additional} bytes where {used} 
bytes already reserved ({} bytes overcommitted) and the fair limit is {limit} 
bytes, {num} registered",
+                            self.spark.overcommit()
+                        );
+                    }
+                }
+            }
+
+            // The anchor comes after the local limit check, so a grow the 
pool rejects itself
+            // never makes a JVM call, and before the real request, so the 
byte is held before
+            // the balance can reach zero. The JVM call can panic inside its 
JNI frame; the
+            // optimistic reservation must not outlive either call, or the 
leaked bytes poison
+            // the task-shared pool for every other consumer.
+            match std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
+                self.take_missing_anchor()
+            })) {
+                Ok(Ok(())) => {}
+                Ok(Err(e)) => {
+                    self.settle_acquire(additional, 0);
+                    return Err(e.into());
+                }
+                Err(panic) => {
+                    self.settle_acquire(additional, 0);
+                    std::panic::resume_unwind(panic);
+                }
+            }
+            // Spark is asked for the request plus any outstanding overcommit, 
and a full grant
+            // repays the overcommit. A short grant stays with Spark until 
this pool hands it
+            // back below, so the bytes can stay charged meanwhile.
+            let refusal = match 
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
+                self.spark.try_acquire_leaving_a_short_grant(additional)

Review Comment:
   > Could the bridge avoid reacquiring the task monitor after obtaining the 
grant, or cover acquisition and diagnostics with one monitor scope while 
keeping releases independent?
   
   It now avoids it. `CometTaskMemoryManager.acquireMemory` no longer calls 
`showMemoryUsage` on a short grant. The warning still logs the task, the 
request, the grant, this manager's total and `getMemoryConsumptionForThisTask`. 
That last call takes only the memory manager's monitor, which a waiting acquire 
gives up in `lock.wait()`, so it cannot close the cycle.
   
   One monitor scope would also work, since the monitor is reentrant. It would 
mean Comet locking on `internal` and relying on `TaskMemoryManager` using 
`this` as its lock. That holds from 3.4 through 4.1 but is not an API. It would 
also keep a thread that has bytes to hand back inside the task monitor for the 
whole dump. Dropping the dump avoids both. Releases were already independent. 
`releaseExecutionMemory` takes the memory manager's monitor, and on 4.x 
`offHeapMemoryLock`, but never the task's.
   
   `CometTaskMemoryManagerSuite` now runs your interleaving against a real 
`UnifiedMemoryManager` with a 100 byte off-heap pool. Other tasks hold 82 bytes 
and 1 byte, and this task holds the anchor. A 30 byte request is granted 16, 
the 1 byte task leaves, and a 10 byte request waits in `ExecutionMemoryPool`. A 
`TaskMemoryManager` subclass holds the first thread right after its grant until 
the second has parked. With the old bridge the second acquire never finishes, 
and the test fails with the first thread BLOCKED and the second WAITING on both 
Spark 3.5 and 4.1. With the change the first thread hands its 16 bytes back and 
the second gets its 10.
   
   Main makes the same call. There the fair pool's mutex keeps two native 
acquires of a task apart. `greedy_unified` has no such lock, so the same 
interleaving is open to it whenever two threads of one task acquire at once. 
The change is in the bridge, so it covers both pools.
   
   The memory management guide now states the rule next to the one for 
`getUsed` and `spill`.
   



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to