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


##########
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 problem I'm worried about: with a blocked `count`, partial aggregation 
now fails queries that main finishes in the same pool. Every failure below 
reports `0.0 B already allocated`: the blocks were already allocated by 
`aggregate_batch` before the check failed, so returning the error here frees 
nothing, it only discards the work.
   #24852 settled that this path must not fail (@Weijun-H: "fail with OOM after 
the batch is already allocated ... preserve progress"), and its test 
`test_partial_hash_stream_emits_whole_batch_when_held_batch_does_not_fit` 
("must not fail with a resources exhausted error") fails on this PR.
   
   ```sql
   -- datafusion-cli -m 32M (also fails at 64M and 128M; main passes)
   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)
   );
   ```
   
   | `count(*)` query, 12 partitions | `PartialHashAggregateStream` errors: PR 
| main |
   |---|---|---|
   | TPC-H SF1 `ROLLUP(l_orderkey, l_suppkey)`, 32–128M | 9/9 | 0/9 |
   | TPC-H SF1 `GROUPING SETS ((l_orderkey), (l_partkey))`, 32–64M | 6/6 | 0/6 |
   | TPC-H SF1 same ROLLUP joined with `orders`, 128–256M | 5/6 | 0/6 |
   | TPC-H SF10 same ROLLUP, 512M and 1G | 4/4 | 0/4 |
   
   How would you like to handle it?



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