comphead commented on code in PR #6544:
URL: https://github.com/apache/datafusion-comet/pull/6544#discussion_r4179690756


##########
native/core/src/execution/memory_pools/fair_pool.rs:
##########
@@ -84,6 +85,41 @@ impl CometFairMemoryPool {
     pub(super) fn overcommit(&self) -> usize {
         self.spark.overcommit()
     }
+
+    /// Records `additional` bytes for `reservation` whatever the fair and 
pool limits, and carries
+    /// what Spark doesn't grant as overcommit. See [`SparkMemory`].
+    fn record(

Review Comment:
   #5613 is open and reworks these same functions. It charges under the state 
lock, calls Spark without it, settles afterwards, and `consumer_used` there 
takes the consumer id. The heads conflict in `fair_pool.rs` and the memory 
guide (checked with `git merge-tree`). `record` takes the locked state and 
calls Spark while it is held, and all three refusals go through `refuse`. In 
#5613 the Spark refusal is handled after the lock is dropped. I expect 
whichever PR lands second to need more than a textual rebase.
   
   That, and the fact that #6583 will have to undo edits in both pools and 
their tests, makes me wonder about a `MemoryPool` wrapper in `spill_replay.rs`, 
like `PlanMemoryPool` and `TaskSharedMemoryPool`. On a `ResourcesExhausted` 
from a replaying final aggregate it would call the inner `grow`, which already 
skips the limits and carries the shortfall as overcommit, so neither pool would 
change. The cost is a per-consumer map of its own for `fair_unified` and one 
more layer for `overcommit()` to look through. I have not compiled this, so 
please read it as an option.



##########
native/core/src/execution/memory_pools/fair_pool.rs:
##########
@@ -160,31 +192,34 @@ impl MemoryPool for CometFairMemoryPool {
                 .expect("overflow in checked_div");
             let consumer_used = *state.consumer_used(reservation);
             if limit < consumer_used.saturating_add(additional) {
-                return resources_err!(
+                let err = resources_datafusion_err!(
                     "Failed to acquire {additional} bytes where this consumer 
already holds {consumer_used} bytes and the fair limit is {limit} bytes, {num} 
registered ({} bytes overcommitted)",
                     self.spark.overcommit()
                 );
+                return self.refuse(&mut state, reservation, additional, err);
             }
             // The shares alone do not bound the pool's total, because a 
consumer keeps what it
             // reserved before another consumer registered.
             let used = state.used;
             if self.pool_size < used.saturating_add(additional) {
-                return resources_err!(
+                let err = resources_datafusion_err!(
                     "Failed to acquire {additional} bytes where {used} bytes 
already reserved ({} bytes overcommitted) and the pool limit is {} bytes",
                     self.spark.overcommit(),
                     self.pool_size
                 );
+                return self.refuse(&mut state, reservation, additional, err);

Review Comment:
   As far as I can tell, this is the one `refuse` call site that none of the 
new tests reach. The fair pool test trips the fair limit first, and the shared 
tests keep `pool_size` at 1000 so only Spark refuses. Putting `return Err(err)` 
back here would not fail any of them. A sketch that should reach it (not run): 
a pool of 100 over a Spark that grants everything, `other` registers and grows 
to 70, then `FinalHashAggregateStream[0]` registers, grows `merge` by 25 and 
grows `merge.new_empty()` by 20. The share is 50, so 25 + 20 passes the fair 
check, and 70 + 25 + 20 is over the pool size.



##########
native/core/src/execution/memory_pools/fair_pool.rs:
##########
@@ -84,6 +85,41 @@ impl CometFairMemoryPool {
     pub(super) fn overcommit(&self) -> usize {
         self.spark.overcommit()
     }
+
+    /// Records `additional` bytes for `reservation` whatever the fair and 
pool limits, and carries
+    /// what Spark doesn't grant as overcommit. See [`SparkMemory`].
+    fn record(
+        &self,
+        state: &mut CometFairPoolState,
+        reservation: &MemoryReservation,
+        additional: usize,
+    ) {
+        self.spark.acquire(additional);
+        state.used = state.used.saturating_add(additional);
+        let consumer_used = state.consumer_used(reservation);
+        *consumer_used = consumer_used.saturating_add(additional);
+    }
+
+    /// Refuses a `try_grow` with `err`, unless it comes from a final 
aggregate reading its spill
+    /// files back, which can't spill. That request is recorded instead; see 
[`spill_replay`].
+    fn refuse(

Review Comment:
   `refuse` returns `Ok(())` after recording the request, so `return 
self.refuse(...)` at the call sites reads as the opposite of what it can do. A 
name like `refuse_unless_replay` would say it.



##########
native/core/src/execution/memory_pools/spill_replay.rs:
##########
@@ -0,0 +1,158 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements.  See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership.  The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License.  You may obtain a copy of the License at
+//
+//   http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied.  See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+//! Lets a final aggregate read its spill files back past its share of memory.
+//!
+//! Once one of DataFusion 55's final aggregates has spilled, it merges its 
sorted spill files and
+//! replays them through an `OrderedFinalAggregateStream` that has no way to 
spill, so a refused
+//! memory request there fails the task (#6254). `FinalHashAggregateStream` 
does this, and so does
+//! `OrderedFinalAggregateStream` itself, which DataFusion uses when the input 
is sorted on some of
+//! the grouping keys. The merge reserves read buffers for as many spill files 
as fit, and those
+//! buffers belong to the same consumer, so the replay often finds the 
consumer's share already
+//! taken. The replay only asks for memory once it has aggregated a batch, so 
like a `grow`, the
+//! request is for memory that already exists. The pools record it the way 
they record a `grow`,
+//! carrying what Spark doesn't grant as overcommit. The replay emits every 
finished group after
+//! each batch, so it holds about one batch of groups, and releasing memory 
repays the overcommit
+//! first.

Review Comment:
   The module doc and the new guide section repeat each other closely (the 
replay can't spill, the merge buffers take the share, the request is recorded 
like a `grow`), and the doc on `is_spill_replay` repeats the guide's second 
paragraph about when a sibling holds memory. For overcommit, the guide keeps a 
short paragraph and leaves the detail to `spark_memory.rs`. It might be worth 
doing the same here. The DataFusion 55.1 facts that an upgrader has to re-check 
could stay in this module, and the guide could say what the pools do and link 
to it. The module doc could also name #6583 next to the removal condition so 
that the removal is easy to find.



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