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]