jayzhan211 commented on code in PR #25877:
URL: https://github.com/apache/datafusion/pull/25877#discussion_r4207746687


##########
datafusion/physical-plan/src/aggregates/hash_stream.rs:
##########
@@ -316,20 +316,44 @@ impl PartialHashAggregateStream {
                         break;
                     }
                     HandleInputResult::OOM => {
-                        let materialized_group_states = 
hash_table.take_state_batch()?.ok_or_else(|| {
-                            internal_datafusion_err!(
+                        let materialized_group_states =
+                            hash_table.take_state_batches()?;
+                        if materialized_group_states.is_empty() {
+                            return Err(internal_datafusion_err!(
                                 "Partial hash aggregate ran out of memory with 
no aggregated groups"
-                            )
-                        })?;
+                            ));
+                        }
 
                         self.early_emit_count.add(1);
                         timer.done();
-                        self.emit_on_memory_pressure(
-                            materialized_group_states,
-                            &mut emitter,
-                            hash_table.memory_size(),
-                        )
-                        .await?;
+
+                        // Blocked storage returns one batch per block, moved 
out of
+                        // the table without copying. Emit them in turn, 
keeping the
+                        // batches not emitted yet in the reservation.
+                        let pending_memory: usize = materialized_group_states
+                            .iter()
+                            // Don't include the first batch since we will 
emit it right away and release its memory
+                            // if it needs slicing then a we will try to hold 
on that reservation while slicing
+                            .skip(1)
+                            .map(|b| b.memory_size)
+                            .sum();
+
+                        // Make sure we can hold on the hash tables and all 
the batches that need to be emitted (except the first one)
+                        // if we can't hold it than we can't do anything about 
it.
+                        self.reservation.try_resize(

Review Comment:
   the fix below never calls `grow()`. But erroring fails queries main finishes 
in the same pool: the PR fails this 6/6 at both 64M and 128M 
(`PartialHashAggregateStream[1] … with 0.0 B already allocated`), main 0/6.
   
   ```sql
   -- datafusion-cli -m 128M
   SET datafusion.execution.target_partitions = 12;
   SELECT count(*), sum(c) FROM (
     SELECT a, b, count(*) AS c
     FROM (SELECT v % 1500000 AS a, v % 7 AS b FROM generate_series(1, 6000000) 
AS t(v))
     GROUP BY ROLLUP(a, b)
   );
   ```
   
   Same on TPC-H `ROLLUP(l_orderkey, l_suppkey)` + `count(*)`: SF1 at 32–128M 
fails 15/15 (main 0/15), SF10 at 512M fails 3/3 (main 0/3). Keeping the pre-OOM 
reservation and shrinking it as blocks drain, as 
`PartialReduceHashAggregateStream` does, passes all of them with peak RSS at or 
below main (SF10 512M: 1658 vs 1720 MB):
   
   ```diff
   -                        let pending_memory: usize = 
materialized_group_states
   -                            .iter()
   -                            // Don't include the first batch since we will 
emit it right away and release its memory
   -                            // if it needs slicing then a we will try to 
hold on that reservation while slicing
   -                            .skip(1)
   -                            .map(|b| b.memory_size)
   -                            .sum();
   -
   -                        // Make sure we can hold on the hash tables and all 
the batches that need to be emitted (except the first one)
   -                        // if we can't hold it than we can't do anything 
about it.
   -                        self.reservation.try_resize(
   -                            hash_table.memory_size()
   -                              // The batches that need to be held while 
emitting each one
   -                              + pending_memory,
   -                        )?;
   +                        let mut held = hash_table.memory_size()
   +                            + materialized_group_states
   +                                .iter()
   +                                // The first batch is emitted right away
   +                                .skip(1)
   +                                .map(|b| b.memory_size)
   +                                .sum::<usize>();
   +
   +                        // The blocks already exist, so failing would free 
nothing.
   +                        // If the pool can't cover them, keep the pre-OOM 
reservation
   +                        // and only shrink it as they drain, like
   +                        // `PartialReduceHashAggregateStream`
   +                        match self.reservation.try_resize(held) {
   +                            Ok(()) | 
Err(DataFusionError::ResourcesExhausted(_)) => {}
   +                            Err(e) => return Err(e),
   +                        }
   
                            for (i, batch) in
                                
materialized_group_states.into_iter().enumerate()
                            {
                                if i != 0 {
   -                                
self.reservation.try_shrink(batch.memory_size)?;
   +                                held -= batch.memory_size;
   +                                let reserved = self.reservation.size();
   +                                if held < reserved {
   +                                    self.reservation.shrink(reserved - 
held);
   +                                }
                                }
   ```
   
   The two `test_partial_hash_stream_*` tests assume one flat state batch. With 
a `Utf8` key astream_under_memory_limit`, and the pointer check on `column(1)` 
as `Int64Type`, all 301`aggregates::` tests pass. 



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