cetra3 commented on PR #23565: URL: https://github.com/apache/datafusion/pull/23565#issuecomment-5867838821
## Results for 9c2f8f19f (keep shared buffers instead of deduplicating values) 9c2f8f19f replaces the value dedup with a check on the buffers: merge repeated entries of the same buffer, then run `gc()` only if it writes less than the merged buffers hold. No values are hashed or sampled. Details are in the commit message. ### The issue's repro 1M rows, 1000 distinct labels, `-m 64M`, default partitions. `ORDER BY id DESC`, because `main` now skips the sort for `ORDER BY id` on this file. Three runs each, identical. | Build | Spilled bytes | | --- | --- | | `main` | 83.2 MB | | previous version (7effeab5c) | 83.2 MB | | 9c2f8f19f | **31.0 MB** | The previous version did not reduce this at all: sorted by `id`, the labels come out in cycling order, so the sample saw 256 distinct values and fell back to `gc()` in every batch. ### Multi-level merge: the repro at 10M rows, 1 partition Median of 5 interleaved runs. Spilled rows above 10M are rows re-spilled by extra merge passes. | Limit | Spilled: main → 9c2f8f19f | Spilled rows | Time | | --- | --- | --- | --- | | 24M | 1959 → 1152 MB | 23.5M → 20.0M | 832 → 751 ms | | 32M | 1477 → 334 MB | 17.8M → **10.0M** (one pass) | 715 → **490 ms** | | 48M | 832 → 353 MB | 10.0M → 10.0M | 571 → 501 ms | | 64M | 832 → 371 MB | 10.0M → 10.0M | 574 → 494 ms | | 128M | 832 → 444 MB | 10.0M → 10.0M | 614 → 578 ms | Smaller spilled batches lower the per-file reservation in the multi-level merge, so more files fit in each pass. At 32M the merge needs one pass instead of about 1.8. ### spill_views and mixed data, local `datafusion-cli`, 1M rows, 4 partitions, 15 rounds with the build order shuffled each round. Ratios are the median per-round time over `main`. "N% repeated" sets N% of the rows to one value, with the rest distinct. | Case | Spilled: main / previous / 9c2f8f19f | previous / main | 9c2f8f19f / main | | --- | --- | --- | --- | | q01 sort, 1 distinct value | 84.2 / 23.2 / 23.8 MB | 0.97x | **0.85x** | | q02 sort, 1000 distinct values | 84.2 / 30.6 / 31.1 MB | 1.05x | **0.88x** | | q03 sort, BinaryView | 84.2 / 30.6 / 31.1 MB | 1.01x | **0.85x** | | q04 sort, all distinct | 84.2 / 84.2 / 84.2 MB | 1.02x | 1.00x | | q05 GROUP BY, all distinct | 84.2 / 84.2 / 84.2 MB | 0.94x | 1.01x | | 5% repeated | 84.2 / 81.1 / 84.2 MB | **1.53x** | 1.04x | | 20% repeated | 84.2 / 72.0 / 84.2 MB | **1.56x** | 0.99x | | 50% repeated | 84.2 / 53.7 / 84.2 MB | **1.48x** | 1.04x | | 80% repeated | 84.2 / 35.4 / 84.2 MB | **1.31x** | 1.00x | | 1000 labels in cycling order | 84.2 / **84.2** / 31.1 MB | 1.03x | **0.79x** | - Shared-dictionary data spills 2.7–3.6x less and runs 12–21% faster. - Data without shared buffers takes the same `gc()` path as `main`: the same bytes, and times within noise (slower than `main` in 6–9 of 15 rounds). - The trade-off: repeated values stored as separate copies (the 50% and 80% cases) are no longer deduplicated. The previous version spilled 36–58% less there, at 1.31–1.48x the time. ### Benchmark bot Two A/B runs of 9c2f8f19f against (merge-base), , 4 partitions. The A/A (main vs main) runs are from the previous round, on the same merge-base. Ratios are branch / base, mean times. **spill_views** (each query sets its own limit: 40M for q01–q03, 96M for q04–q05) | Query | A/B run 1 | A/B run 2 | A/A (previous round) | Previous version (7effeab5c) | | --- | --- | --- | --- | --- | | q01 sort, 1 distinct value | 109.7 → 73.3 ms (**0.67x**) | 105.4 → 71.5 ms (**0.68x**) | 0.99x | 0.78–0.79x | | q02 sort, 1000 distinct values | 104.7 → 73.7 ms (**0.70x**) | 100.5 → 71.5 ms (**0.71x**) | 0.99x | 0.92x | | q03 sort, BinaryView | 103.3 → 74.0 ms (**0.72x**) | 100.2 → 71.8 ms (**0.72x**) | 0.98x | 0.94x | | q04 sort, all distinct | 1.00x | 0.99x | 0.99x | 1.00–1.02x | | q05 GROUP BY, all distinct | 1.01x | 0.97x | 0.98x | 0.99–1.01x | q05 passed at the default 96M limit in both runs. **sort_tpch** (SF1, 512M, 4 partitions): no query changes beyond A/A noise. Q3 fails on both sides, as before. | Query | A/B run 1 | A/B run 2 | A/A (previous round) | | --- | --- | --- | --- | | Q8 () | 509 → 488 ms (0.96x) | 487 → 480 ms (0.98x) | 0.99x | | Q9 () | 536 → 540 ms (1.01x) | 529 → 527 ms (1.00x) | 1.00x | | Q11 () | 379 → 371 ms (0.98x) | 377 → 376 ms (1.00x) | 0.99x | | Total | 4324 → 4254 ms (0.98x) | 4290 → 4267 ms (0.99x) | 0.99x | -- 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]
