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


##########
datafusion/functions-aggregate/src/array_agg.rs:
##########


Review Comment:
   Each coalesced payload is copied twice: `copy_array_data` + 
`compact_payload` detach it, then `concat` copies tail + payload again. On the 
PR's own bench (interleaved runs, base 696eaf58d vs this branch) that is a 
regression on exactly the inputs this PR targets:
   
   | case | main | PR |
   |---|---|---|
   | i64 ordered, 1 row/update | 645 µs | 903 µs (+40%) |
   | i64 random, 1 row/update | 710 µs | 970 µs (+35%) |
   | i64 ordered, 8 rows/update | 90 µs | 122 µs (+35%) |
   
   For primitive, boolean and offset byte types, `concat` already produces 
fresh buffers, so the detach copy can be skipped when coalescing. A prototype 
of that gives 391 µs / 465 µs on the 1-row cases and 64 µs / 136 µs on 8 rows, 
faster than main, and all `array_agg` tests pass. View and dictionary types 
still need the detach first: their raw input buffers over-report 
`get_buffer_memory_size`, so they would never pass the byte cap.
   
   ```rs
   let row_count = values.len();
   let concat_copies = values.data_type().is_primitive()
       || matches!(
           values.data_type(),
           DataType::Boolean
               | DataType::Utf8
               | DataType::LargeUtf8
               | DataType::Binary
               | DataType::LargeBinary
       );
   let can_coalesce = concat_copies
       && self.batches.last().is_some_and(|last| {
           last.len() + row_count <= ORDERED_ARRAY_AGG_COALESCE_ROWS
               && last.get_buffer_memory_size() + 
values.get_buffer_memory_size()
                   <= ORDERED_ARRAY_AGG_COALESCE_BYTES
       });
   // `concat` copies fixed-layout payloads, so detaching first is only
   // needed when the payload is stored as-is or may share buffers.
   let values = if can_coalesce {
       values
   } else {
       compact_payload(make_array(copy_array_data(&values.to_data())))?
   };
   ```
   Then branch on `can_coalesce` instead of re-checking the tail. Can you add 
the before/after numbers to the PR description?



##########
datafusion/functions-aggregate/src/array_agg.rs:
##########
@@ -1531,11 +1539,37 @@ impl OrderSensitiveArrayAggAccumulator {
         }
 
         let start = self.entries.len();
-        let batch_idx = self.batches.len();
-        self.batches.push(values);
-        self.entries.extend(
-            (0..row_count).map(|row_idx| OrderedArrayAggEntry { batch_idx, 
row_idx }),
-        );
+        let (batch_idx, row_offset) = match self.batches.last() {
+            Some(last_batch)
+                if last_batch.len() + row_count <= 
ORDERED_ARRAY_AGG_COALESCE_ROWS
+                    && last_batch.get_buffer_memory_size()
+                        + values.get_buffer_memory_size()
+                        <= ORDERED_ARRAY_AGG_COALESCE_BYTES =>
+            {
+                let merged =

Review Comment:
   For Utf8View/BinaryView, `concat` → `GenericByteViewBuilder::append_array` 
clones every input's data buffers instead of copying the bytes. After 64 
one-row updates the coalesced tail is 1 array with **64 data buffers**, so each 
row still pays for its own `Buffer` + `Arc<Bytes>` allocation. That is the 
fixed per-row overhead this PR is meant to remove, and `size()` can't see it. 
Running `compact_payload` on the merged array gets it down to 1 buffer. The 
extra copy is bounded by the 4 KiB cap.
   
   ```diff
   -                let merged =
   -                    arrow::compute::concat(&[last_batch.as_ref(), 
values.as_ref()])?;
   +                // View concat shares the inputs' data buffers; GC folds 
them
   +                // into one so the tail does not keep a buffer per update.
   +                let merged = compact_payload(arrow::compute::concat(&[
   +                    last_batch.as_ref(),
   +                    values.as_ref(),
   +                ])?)?;
   ```
   
   The Utf8View test could also assert this:
   
   ```rs
   assert_eq!(acc.batches[0].as_string_view().data_buffers().len(), 1);
   ```
   
   Fine to handle in a follow-up if you'd rather keep this PR focused.
   
   I checked this with a local probe (64 single-row Utf8View updates): 64 data 
buffers on the PR branch, 1 with compact_payload after the concat. I haven't 
run the suggested extra assertion as part of the PR's own Utf8View test.



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