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]