dwsmith1983 commented on code in PR #5613:
URL: https://github.com/apache/datafusion-comet/pull/5613#discussion_r4133984221
##########
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:
> 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?
Coordinated with completed teardown, in ac0a691. The registry entry used to
go away before the fair pool's drop handed the anchor back, so a replacement
could be created, park its first acquire in Spark and wake to the removed task
entry, as your test showed. Now the entry stays, with an expired `Weak`, until
the old pool has finished dropping, anchor release included, and
`acquire_task_shared_pool` waits on a condvar while such an entry exists. A
guard field declared after the pool removes the entry and notifies, so the JNI
release still runs with no Rust lock held and the registry lock is only taken
for the map update.
I did not share the anchor across generations. The byte is charged to the
`CometTaskMemoryManager` that acquired it, and a replacement created between
"no successor, release" and the release landing would hit the same window, so
it needs this ordering anyway.
Two regressions, both failing on the previous head:
`a_replacement_waits_for_the_previous_pool_to_finish_dropping` holds the old
pool inside its drop and checks no replacement is handed out until it finishes,
and `the_next_plan_of_a_task_waits_for_the_dying_pool_to_hand_back_its_anchor`
is your scenario on the Spark stub (100 bytes, other tasks holding 74 and 25,
the dying pool holding only its anchor).
The cost is that `createPlan` can wait for one release call made by another
thread.
The branch also picked up main's #6261, which runs Spark acquires in
`block_in_place`. The anchor acquire now goes through the same call (cd19f9a),
so a parked anchor request no longer holds its Tokio worker.
--
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]