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


##########
native/core/src/execution/memory_pools/fair_pool.rs:
##########
@@ -90,9 +144,110 @@ impl CometFairMemoryPool {
             state: Mutex::new(CometFairPoolState {
                 used: 0,
                 consumers: HashMap::new(),
+                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<()> {
+        if self.state.lock().anchor_held {
+            return Ok(());
+        }
+        // The lock is not held across the call.
+        if 
!Self::anchor_granted(self.spark.manager().acquire_anchor(ANCHOR_BYTES)?) {
+            return Ok(());
+        }
+        {
+            let mut state = self.state.lock();
+            if !state.anchor_held {
+                state.anchor_held = true;
+                return Ok(());
+            }
+        }
+        // A grow or a release on another thread took the anchor meanwhile. 
This byte was never
+        // booked, so a failed return only leaves Spark holding it until the 
task ends, as on
+        // drop.
+        if let Err(e) = self.spark.manager().release_anchor(ANCHOR_BYTES) {
+            warn!("Failed to return a duplicate memory pool anchor byte: 
{e:?}");
+        }
+        Ok(())
+    }
+
+    /// Hands `bytes` that Spark granted this pool back to it. While the 
anchor is missing, one
+    /// of them is kept as the anchor instead: this pool holds at least those 
bytes from Spark,
+    /// so the release cannot take the task's balance to zero under an acquire 
parked there.
+    /// Claiming the anchor under the lock means one release keeps it, and an 
anchor retry that
+    /// lands afterwards hands its byte back as a duplicate. The JVM releases 
to Spark before it
+    /// moves its counters, so a failed call has moved nothing: the claim is 
rolled back and the
+    /// pool is unanchored again, with Spark still holding the bytes, and 
later grows retry the
+    /// anchor as usual. A retry that returned its byte while the claim stood 
holds nothing.
+    fn release_to_spark(&self, bytes: usize) -> CometResult<()> {
+        let keep_anchor = !std::mem::replace(&mut 
self.state.lock().anchor_held, true);
+        if !keep_anchor {
+            return self.spark.manager().release(bytes);
+        }
+        let released = self.spark.manager().release_keeping_anchor(bytes);
+        if released.is_err() {
+            self.state.lock().anchor_held = false;
+        }
+        released
+    }
+
+    /// Settles a release the JVM has accepted: the bytes come off the pool's 
total and the
+    /// consumer's share only now, so a grow is never admitted on bytes Spark 
still holds.
+    fn settle_release(&self, consumer: usize, bytes: usize) {
+        self.state.lock().settle(consumer, bytes, 0);
+    }
+
+    /// Settles a finished JVM acquire. `charged` bytes went on the pool's 
total and the
+    /// consumer's share before the call and `held` is what Spark granted and 
still holds, which
+    /// stays charged until it is handed back. The difference is rolled back.
+    fn settle_acquire(&self, consumer: usize, charged: usize, held: usize) {
+        self.state.lock().settle(consumer, charged, held);
+    }
+}
+
+impl Drop for CometFairMemoryPool {
+    /// The last plan of the task letting go of the pool runs this, on 
whatever thread that
+    /// happens on; `with_env` attaches the thread and the JVM object handle 
outlives the pool.
+    /// Nothing may panic out of a drop, so failures are only logged, and 
Spark frees the task's
+    /// whole balance when the task ends anyway.
+    fn drop(&mut self) {
+        let state = self.state.get_mut();
+        if state.used != 0 {
+            warn!(
+                "Task {} dropped CometFairMemoryPool with {} bytes still 
reserved ({} bytes overcommitted)",
+                self.spark.task_attempt_id(),
+                state.used,
+                self.spark.overcommit()
+            );
+        }
+        if !state.anchor_held {
+            return;
+        }
+        let released = 
std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
+            self.spark.manager().release_anchor(ANCHOR_BYTES)

Review Comment:
   [P2] Coordinate the final anchor release with replacement pools for the same 
task. `TaskSharedMemoryPool` permits replacement once its `Weak` expires, 
before its inner fair pool finishes dropping. In a 100-byte Spark pool, let 
other tasks hold 74 and 25 bytes and the dying pool hold only its anchor. Pause 
this release, create a replacement through `acquire_task_shared_pool`, and 
start its first grow. The replacement's anchor request parks. Resuming the old 
release removes the task's `memoryForTask` entry, so the replacement wakes with 
`NoSuchElementException` instead of acquiring memory or returning a spillable 
refusal. This can fail a task when final pool destruction overlaps creation of 
its next native plan. Base has no teardown anchor release. Could the anchor be 
shared across pool generations, or could replacement acquisition be coordinated 
with completed teardown without holding the global registry lock across JNI?
   
   Evidence: The disposable native test 
`review_replacement_pool_survives_the_old_pools_anchor_release` exercised the 
actual head's registry and fair-pool implementation and failed with `key not 
found: 0`. A separate probe compiled with the unchanged head's Java bridge 
against real Spark 4.1.3 produced `releaseOldAfterStart=true 
failure=java.util.NoSuchElementException: key not found: 0`, at 
`ExecutionMemoryPool.scala:115`. Releasing the old anchor before starting the 
replacement produced `failure=null`. Logs: 
`/tmp/comet-5613-review-xdjznr_c/native-probes.log` and 
`/tmp/comet-5613-review-xdjznr_c/jvm-probes.log`.



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