zhuqi-lucas opened a new issue, #24742:
URL: https://github.com/apache/datafusion/issues/24742

   ### Is your feature request related to a problem or challenge?
   
   `first_value` / `last_value` pick a per-type state strategy 
(`first_last.rs:114/171/188`):
   
   | value type | state | per winning row |
   |---|---|---|
   | primitive | `PrimitiveValueState<T>` — `Vec<T::Native>` + null bitmap | 
one store, no allocation |
   | utf8 / binary | `BytesValueState` — `Vec<Option<Vec<u8>>>` | one memcpy 
into that group's `Vec` |
   | struct / list / map | `GenericValueState` — `Vec<Option<ScalarValue>>` | 
build a 1-row array, `compact()`, two recursive `size()` walks |
   
   The nested branch (`first_last/state.rs:329`) runs this for every row that 
beats the current winner:
   
   ```rust
   self.total_size -= v.size();                                // walk the old 
value
   let mut scalar = ScalarValue::try_from_array(array, idx)?;  // build a 1-row 
StructArray, wrap in Arc
   scalar.compact();                                           // copy the 
referenced bytes
   self.total_size += scalar.size();                           // walk it again
   self.vals[group_idx] = Some(scalar);                        // store, drop 
the old one
   ```
   
   against the primitive path (`first_last/state.rs:80`):
   
   ```rust
   self.vals[group_idx] = array.value(idx);
   self.nulls.set_bit(group_idx, !array.is_null(idx));
   ```
   
   So the cost scales with how often a row wins, and that scaling only bites at 
the high end. Two points on the curve, both from `run benchmark first_last`:
   
   **~7% win rate** — `update_bench` feeds a random `ORDER BY` key over 65536 
rows / 1024 groups, so wins are the running minima of a random sequence:
   
   ```
   first_value update_bench struct(i64,utf8,f64) nulls=0%    34.9 ms
   first_value update_bench nulls=0%, filter=false           30.4 ms   
(primitive)
   ```
   
   15%. Entirely reasonable, which is why nothing has flagged this before.
   
   **100% win rate** — `coalesce_peers ... (winner changes)` feeds a strictly 
decreasing key, so every row wins:
   
   ```
   first_value coalesce_peers(i64,utf8,f64) coalesced struct (winner changes)   
411.3 ms
   first_value coalesce_peers(i64,utf8,f64) separate x3      (winner changes)   
 77.8 ms
   ```
   
   5.3x. Holding the accumulator and data volume fixed and varying only win 
frequency isolates it:
   
   ```
                           winner stable    winner changes
   coalesced struct           34.6 ms          411.3 ms      12x
   separate x3                90.4 ms           77.8 ms      flat
   ```
   
   A monotonically increasing `ORDER BY` key over time-ordered data is the 
ordinary shape that lands here.
   
   ### Describe the solution you'd like
   
   `BytesValueState` is the precedent: byte values are variable-length too, but 
rather than falling back to `ScalarValue` they got a purpose-built state. A 
struct could be decomposed the same way — one child state per field, each 
picking its own strategy:
   
   ```
   struct(i64, utf8, f64)
        ├─ field 0 → PrimitiveValueState<i64>
        ├─ field 1 → BytesValueState
        └─ field 2 → PrimitiveValueState<f64>
   ```
   
   A winning row then costs three ordinary field updates, which is what 
`separate x3` already measures at 77.8 ms. Note what that implies for the 
coalescing rewrite in #23682: it keeps the "N compares become 1" saving *and* 
loses the retain-path penalty, so it would win in both regimes rather than 
trading one for the other.
   
   Two cheaper things that stand on their own:
   
   1. **Drop one of the `size()` walks.** Each update walks the value twice — 
once to subtract the old size, once to add the new — and `ScalarValue::size()` 
recurses for nested types. Having `compact()` return the size, or recomputing 
lazily, removes one traversal.
   2. **Relax when `compact()` runs.** The comment there explains the tradeoff: 
without it a single stored winner pins its whole source batch. But that could 
be driven by how many distinct batches are currently pinned rather than copying 
on every row.
   
   ### Describe alternatives you've considered
   
   Storing `(Arc<ArrayRef>, row_idx)` and materializing once in `take()` would 
remove the per-row work entirely, but it reintroduces exactly the batch-pinning 
problem `compact()` exists to avoid, so it would need the same threshold logic 
as (2) above.
   
   Leaving it as-is is reasonable for the current default configuration — 
nothing regresses today. It matters because the cost is paid by any 
`first_value(<nested> ORDER BY ...)` regardless of #23682, and because it is 
the blocker for ever enabling `optimizer.enable_coalesce_first_last` by default.
   
   ### Additional context
   
   - `GenericValueState` was added in #23628
   - Numbers from the benchmark run on #23682; `coalesce_peers` cases come from 
#24559
   


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