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]