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]
