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


##########
native/core/src/execution/memory_pools/fair_pool.rs:
##########
@@ -72,7 +97,99 @@ impl CometFairMemoryPool {
         Self {
             spark,
             pool_size,
-            state: Mutex::new(CometFairPoolState { used: 0, num: 0 }),
+            state: Mutex::new(CometFairPoolState {
+                used: 0,
+                num: 0,
+                anchor_held: false,
+            }),
+        }
+    }
+
+    /// Whether an anchor request came back covered. A declined anchor is a 
zero grant, so
+    /// there is nothing to hand back either way.
+    fn anchor_granted(acquired: i64) -> bool {
+        let granted = usize::try_from(acquired).unwrap_or(0);
+        if granted > ANCHOR_BYTES {
+            warn!("Requested {ANCHOR_BYTES} bytes from the JVM but it reports 
{granted} granted");
+        }
+        granted >= ANCHOR_BYTES
+    }
+
+    /// Takes the anchor on the first grow and retries it while Spark declines 
it, as a
+    /// request of its own that never rides on a real grow. Spark declines it 
only while the
+    /// task sits at its share, so the extra JNI call is paid on that path 
alone and never
+    /// once the anchor is held. A `try_grow` rolls back its reservation if 
this fails.
+    fn take_missing_anchor(&self) -> CometResult<()> {

Review Comment:
   Reproduced: 26 warnings reading `: 1` in `CometExecSuite`, in the bucketed 
table and TakeOrderedAndProject tests, and every other warning one byte high.
   
   The anchor now goes through two new methods on `CometTaskMemoryManager`, 
`acquireAnchor` and `releaseAnchor`. They acquire from Spark's 
`TaskMemoryManager` through the same `NativeMemoryConsumer` but count into a 
separate field, so the task's balance and the consumer's usage include the byte 
while `getUsed`, which `close` checks, does not. The unified pool never calls 
them. I went with that over reporting the pool total at `releasePlan` because a 
task-shared pool's total is the whole task's, so a plan closing while a sibling 
still holds memory would warn with the sibling's bytes, and it would change the 
warning's source for the unified pool as well. It also avoids a constant on the 
JVM side that would hide a real one byte leak.
   
   `CometExecIteratorLifecycleSuite` has a two plan task under `fair_unified` 
that checks close stays quiet with the anchor held, and that a reservation 
leaked on purpose is reported at its exact size. `CometTaskMemoryManagerSuite` 
pins the two counters, and the fair pool test double keeps them apart too. With 
the fix the `: 1` lines are gone from `CometExecSuite`. The anchor paragraph in 
`memory_management.md` says where the byte counts and why.
   
   One thing worth knowing from the same run: 83 other close warnings in 
`CometExecSuite` are the same on main, same tests and same values, from a plan 
closing while a sibling in the task still holds live reservations. No pool ever 
dropped with bytes reserved. Not from this PR, so left alone.
   



-- 
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