andygrove commented on code in PR #2518:
URL: 
https://github.com/apache/datafusion-ballista/pull/2518#discussion_r4160693723


##########
ballista/core/src/utils.rs:
##########
@@ -220,25 +227,36 @@ pub async fn write_stream_to_disk(
 
         let mut writer =
             StreamWriter::try_new_with_options(file, schema.as_ref(), 
options)?;
+        let mut coalescer = BatchCoalescer::new(schema, batch_size)
+            .with_biggest_coalesce_batch_size(Some(batch_size / 2));
+        let mut num_batches = 0;
 
         while let Some(batch) = rx.blocking_recv() {
             let timer = write_metric.timer();
-            writer.write(&batch)?;
+            coalescer.push_batch(batch)?;

Review Comment:
   [P2] This call can exceed the documented `batch_size`-row per-partition 
bound. In arrow-select 59.3, if a small partial batch is buffered and the next 
batch is larger than `batch_size / 2`, the large-batch bypass is not taken 
while `buffered_rows <= limit`; `push_batch` instead splits the entire input 
and queues every completed batch internally before returning. The loop below 
cannot drain those completions until then, so a small-then-large sequence can 
duplicate/materialize essentially the whole large input at once. The current 
large-batch test only exercises the empty-coalescer bypass. Please add that 
transition case and flush/drain the partial batch before pushing an oversized 
input (or use an API that drains incrementally).



##########
ballista/core/src/utils.rs:
##########
@@ -220,25 +227,36 @@ pub async fn write_stream_to_disk(
 
         let mut writer =
             StreamWriter::try_new_with_options(file, schema.as_ref(), 
options)?;
+        let mut coalescer = BatchCoalescer::new(schema, batch_size)

Review Comment:
   [P1] Please bound or account for aggregate coalescer memory across output 
partitions. `ShuffleWriterExec::execute_shuffle_write` invokes this once for 
every child output partition and starts all K drains concurrently. 
Range-repartition children expose `execution.target_partitions` here, so one 
task can have K independent coalescers; this is not limited to the task's 
active input/vcore count. Each populated coalescer may retain/copy roughly 
`batch_size` rows in addition to the existing per-partition channels, and none 
of this is registered with the task memory pool. With the defaults, K=1000 and 
1 KiB rows is about 7.8 GiB of new buffering before channel/input overhead. The 
channel capacity does not cap this state because every output has its own 
channel and coalescer. Please enforce a shared per-task byte budget / memory 
reservation and flush partial coalescers under pressure (or otherwise make the 
total independent of K), and cover a high-K, wide-row case.



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