adriangb commented on PR #23565:
URL: https://github.com/apache/datafusion/pull/23565#issuecomment-5779637711

   ## Summary: behavior under default settings and under memory pressure
   
   ### Default settings (no memory limit)
   
   This PR has no effect. The changed code runs only from 
`InProgressSpillFile::append_batch` / `append_batch_async`, and operators call 
these only when a memory reservation fails. The default memory pool is 
unbounded, so no operator spills.
   
   The three default-setting runs (tpch, tpcds, clickbench_partitioned) all 
report `Peak spill: 0 B`. Their per-query differences are noise: TPC-H shows 
"13 faster" and TPC-DS / ClickBench show "5 slower" for a code path that did 
not execute.
   
   ### Under memory pressure
   
   On each spilled batch, for each `Utf8View` / `BinaryView` column (also when 
nested in `List`, `Map`, `Dictionary`, ...) with more than 10 KB of data 
buffers:
   
   - **On `main`:** `gc()` copies the bytes of each non-inline view into new 
buffers. When views share bytes, `gc()` writes one copy per row. Parquet 
dictionary pages decode to views that all point into one dictionary buffer, so 
this is the normal case for low-cardinality string columns read from Parquet. 
This is the regression from https://github.com/apache/datafusion/pull/21633 
that https://github.com/apache/datafusion/issues/23564 shows (1.3 MB Parquet 
file → 83.2 MB spilled).
   - **With this PR:** `GenericByteViewBuilder::with_deduplicate_strings()` 
hashes each non-inline value (> 12 bytes) and writes each distinct value one 
time.
   
   Consequences:
   
   - **Disk and read-back memory:** smaller when a batch has repeated 
non-inline values. The dedup scope is one batch (default `batch_size` = 8192 
rows), not the full spill file.
   - **CPU:** one hash and one probe per non-inline value, on every spilled 
batch. When the values in a batch are unique, this cost gives no benefit.
   - **Transient memory:** the hash table and the builder blocks. The memory 
pool does not track this memory, but it is small (hundreds of KB per 8192-row 
batch). The builder blocks have `capacity > len`, which is why the test now 
measures `len`. This does not change the IPC bytes that are written to disk.
   
   ### Measured results (sort_tpch SF1, `c4a-highmem-16`)
   
   | Config | Result |
   | --- | --- |
   | 256 MiB, 12 partitions | Q3, Q7–Q11 fail on **both** sides (pre-existing). 
Other queries: no change. |
   | 512 MiB, 12 partitions | Q7, Q10 fail on both sides. Q3 +4% (586.9 → 611.3 
ms, low stddev). Others within ±2%. |
   | 512 MiB, 4 partitions | Q3 **+15%** (613.9 → 706.2 ms), Q11 +6%, Q8 +5%, 
Q9 +4%, Q10 +3%. Total +4%, CPU user +4%. |
   
   The slower queries all include `l_comment` (4.5M distinct values, longer 
than 12 bytes). This is the worst case for this PR: each value is hashed and no 
value is deduplicated. This agrees with the concern in the review thread about 
unique strings.
   
   Q1, Q2, Q4–Q6 have no non-inline strings, and they show no change, as 
expected. Q7 is the only query with a repeated non-inline column 
(`l_shipinstruct`: 4 distinct values, 2 of them > 12 bytes). It shows +1%, but 
the bot does not report spill bytes, so we do not know the disk effect.
   
   Memory pool peaks do not change (±2%). This is also expected, because 
sort_tpch has no repeated long strings in its hot columns.
   
   ## Benchmarks that are missing
   
   1. **The target case, with spill bytes measured.** No run measures the claim 
of this PR. The bot reports `Peak spill: 0 B` also for the sort_tpch runs that 
spilled, so we have no disk numbers. I recommend a main-vs-PR sweep that uses 
the #23564 repro and records `spilled_bytes`, `spill_count`, wall time, pool 
peak and RSS from `EXPLAIN ANALYZE`:
      - distinct values per batch: 1, 100, 1000, 8192 (unique)
      - string length: 16, 64, 256 bytes
      - source: dictionary-encoded Parquet (views share bytes) and 
plain-encoded or computed strings (views do not share bytes)
      - memory limit: a light spill (2–3 spill files) and a heavy spill
   
      This gives the break-even cardinality, and it shows the size of the win.
   2. **Real data.** Sort ClickBench `hits` under a memory limit. The 
`data_sort_pushdown` step in `bench.sh` already does `COPY (SELECT * FROM hits 
ORDER BY "EventTime")` with `memory_limit` set. `hits` has a realistic mix of 
low- and high-cardinality long strings. Measure wall time and `spilled_bytes`.
   3. **Criterion microbenchmark** of `gc()` against `gc_dedup_view` on one 
batch, for the same cardinality and length matrix. For example, extend 
`datafusion/physical-plan/benches/spill_io.rs`. This isolates the per-batch CPU 
cost with low noise.
   4. **Other spill producers.** Aggregation spills write group keys, and the 
keys are distinct in each spilled batch. For these columns, dedup cannot make 
the output smaller, so the hash cost is always overhead. `external_aggr` uses 
only integer keys, so it cannot show this. A string-key external aggregation 
(for example `SELECT l_comment, count(*) FROM lineitem GROUP BY l_comment` with 
a memory limit) is necessary. Repartition (`spill_pool`) and sort-merge join 
spills also use this path.
   5. **A/A control** for the 512 MiB / 4 partitions sort_tpch run, to confirm 
that Q3 +15% is real.
   
   ## Possible mitigation (for discussion)
   
   A cheap check can select between `gc()` and dedup without hashing the 
strings. For example, compare the sum of the non-inline view lengths with the 
total size of the data buffers. If the sum is larger, the views share bytes (as 
with Parquet dictionary pages), and `gc()` will make the array larger. Only 
then is dedup necessary. Another option is to preserve the sharing that exists: 
remap views by `(buffer_index, offset)` instead of by value. This does not hash 
the string bytes, and it keeps the output no larger than the input. Both 
options keep the #23564 fix and remove most of the cost on high-cardinality 
columns such as `l_comment`.
   


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