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


##########
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:
   Thanks, I went with the wrapper in cc60fe783. `SpillReplayPool` sits between 
`TaskSharedMemoryPool` and `TrackConsumersPool`, and on a `ResourcesExhausted` 
from a replaying final aggregate it calls the inner `grow`. Both pools are back 
to main's versions, and `git merge-tree` against #5613 now only conflicts on 
the pool stack diagram in the guide, where both PRs add a line. The costs are 
the ones you listed: its own map, which only holds final aggregates, and one 
more line in `overcommit()`. It only matches `ResourcesExhausted`, so a failed 
JNI call during the replay still goes back to the aggregate, and there is a 
test for that.



##########
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:
   Added your scenario as 
`a_final_aggregate_reading_its_spill_files_back_may_pass_the_pool_limit`. With 
the wrapper all three refusals take the same path now, but I checked that the 
test does reach the pool limit: with the wrapper taken out it fails with 
`Failed to acquire 20 bytes where 95 bytes already reserved (0 bytes 
overcommitted) and the pool limit is 100 bytes`.



##########
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:
   Agreed. The guide section is now one paragraph on what `SpillReplayPool` 
does, and it points to `spill_replay.rs` for the rest. The module doc keeps the 
DataFusion 55.1 facts an upgrade has to re-check as a list, and names #6583 
next to the removal condition. I also moved the guide section to after the 
task-shared pools section, because #5613 adds its anchor text where it used to 
be.



##########
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` is gone with the wrapper. Its `try_grow` matches the refusal and 
only calls `grow` for the replay, so there is nothing left to rename.



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