Repository navigation
GC & Deduplicate String View on Spill - #23565
Conversation
|
run benchmarks |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (25cc81b) to f755cb4 (merge-base) diff using: clickbench_partitioned File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (25cc81b) to f755cb4 (merge-base) diff using: tpcds File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (25cc81b) to f755cb4 (merge-base) diff using: tpch File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usagetpch — base (merge-base)
tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: CPU Details (lscpu)Details
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
25cc81b to
d6079c7
Compare
|
Those benchmarks don't actually stress this path at all. In fact I don't think there are any standard benchmarks that will sweat this particular change. There is a |
| fn gc_dedup_view<T: ByteViewType>( | ||
| array: &GenericByteViewArray<T>, | ||
| ) -> GenericByteViewArray<T> { | ||
| let mut builder = GenericByteViewBuilder::<T>::with_capacity(array.len()) |
There was a problem hiding this comment.
Every qualifying Utf8View and BinaryView now allocates a hash table sized for all rows and hashes every non-inline value. This is valuable for repeated values, but unique strings or large binary values get no disk reduction over gc() while paying extra CPU and peak memory during spilling.
There was a problem hiding this comment.
Yes this is the trade off. However, I believe the size of this is actually dictated by the configured batch size so it's never going to be more than that.
So with 16384 as the default batch size this could use ~ 16384 * 64 = 1MiB <+ hash internals> RAM to write this out + CPU overhead of hash lookups which I think is ~ log(n). So yeah high cardinality string views will probably cause increased memory/cpu but not by much. Keeping in mind that we are spilling here to reduce RAM usage, and so when we consider the re-hydrated spilled file, the memory usage is much more reduced for low cardinality views.
In high cardinality scenarios this will add overhead, but it's hard to measure without doing a second pass to work out what the cardinality is. Maybe there is a better way to gauge the cardinality of the input view and use that to decide? Not sure there is a good method for this.
But, in production we have seen a good reduction in file usage/memory pressure with gc & dedup (i.e, we have applied this as a patch to our version of DF), so I feel that for real use cases the trade off is worth it.
| @@ -1345,7 +1359,7 @@ mod tests { | |||
| let size_without_gc: usize = array_without_gc | |||
| .data_buffers() | |||
| .iter() | |||
| .map(|buffer| buffer.capacity()) | |||
| .map(|buffer| buffer.len()) | |||
There was a problem hiding this comment.
Keep coverage for both length and capacity.
we can compare the repeated-label reproduction with a forced-spill high-cardinality case such as sort-TPCH Q3. |
d6079c7 to
20bc3e1
Compare
|
run benchmark sort_tpch env: |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (20bc3e1) to a9b61ab (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "256M"Results will be posted here when complete File an issue against this benchmark runner |
|
run benchmark sort_tpch env: |
Codecov Report❌ Patch coverage is
Additional details and impacted files@@ Coverage Diff @@
## main #23565 +/- ##
========================================
Coverage 82.73% 82.73%
========================================
Files 1147 1147
Lines 449389 449540 +151
Branches 449389 449540 +151
========================================
+ Hits 371788 371924 +136
+ Misses 54943 54940 -3
- Partials 22658 22676 +18 ☔ View full report in Codecov by Harness. 🚀 New features to boost your workflow:
|
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (20bc3e1) to a9b61ab (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (20bc3e1) to a9b61ab (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "256M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagesort_tpch — base (merge-base)
sort_tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (20bc3e1) to a9b61ab (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagesort_tpch — base (merge-base)
sort_tpch — branch
File an issue against this benchmark runner |
|
run benchmark sort_tpch env: |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (20bc3e1) to a9b61ab (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_EXECUTION_TARGET_PARTITIONS: "4"
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (20bc3e1) to a9b61ab (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_EXECUTION_TARGET_PARTITIONS: "4"
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagesort_tpch — base (merge-base)
sort_tpch — branch
File an issue against this benchmark runner |
20bc3e1 to
cad94ba
Compare
cad94ba to
ca22a41
Compare
|
Local results with the Time: 3 batches × 10 interleaved rounds (random arm order; arms main, main, PR, PR), 100 timings per arm per batch, median. Spill bytes from
Bot results with the suite will follow after #25625 merges and this PR is rebased. |
ca22a41 to
cbdecca
Compare
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_EXECUTION_TARGET_PARTITIONS: "4"
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagesort_tpch — base (merge-base)
sort_tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark sort_tpch
env:
DATAFUSION_EXECUTION_TARGET_PARTITIONS: "4"
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagesort_tpch — base (merge-base)
sort_tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark spill_views
env:
DATAFUSION_EXECUTION_TARGET_PARTITIONS: "4"
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagespill_views — base (merge-base)
spill_views — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark spill_views
env:
DATAFUSION_EXECUTION_TARGET_PARTITIONS: "4"
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagespill_views — base (merge-base)
spill_views — branch
File an issue against this benchmark runner |
Results for 9c2f8f1 (keep shared buffers instead of deduplicating values)9c2f8f1 replaces the value dedup with a check on the buffers: merge repeated entries of the same buffer, then run The issue's repro1M rows, 1000 distinct labels,
The previous version did not reduce this at all: sorted by Multi-level merge: the repro at 10M rows, 1 partitionMedian of 5 interleaved runs. Spilled rows above 10M are rows re-spilled by extra merge passes.
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
Benchmark botTwo A/B runs of 9c2f8f1 against spill_views (each query sets its own limit: 40M for q01–q03, 96M for q04–q05)
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.
|
|
run benchmarks env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M" |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark clickbench_partitioned
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark tpcds
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark running (GKE) | trigger CPU Details (lscpu)Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark tpch
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"Results will be posted here when complete File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark tpch
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagetpch — base (merge-base)
tpch — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark tpcds
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usagetpcds — base (merge-base)
tpcds — branch
File an issue against this benchmark runner |
|
🤖 Benchmark completed (GKE) | trigger Instance: Comparing spill-dedup-view-arrays (9c2f8f1) to ed43e69 (merge-base) diff Run configurationrun benchmark clickbench_partitioned
env:
DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"CPU Details (lscpu)Details
Memory Pool PeaksPeak Base:
Pool accounting vs. process RSS Max pool peak is the largest reservation any single query in the run reached; peak RSS covers the whole invocation, including data loading and allocator retention, and the two high-water marks need not coincide in time. The gap is therefore an upper bound on what the pool did not account for, not a measurement of it.
Resource Usageclickbench_partitioned — base (merge-base)
clickbench_partitioned — branch
File an issue against this benchmark runner |
adriangb
left a comment
There was a problem hiding this comment.
Benchmarks and implementation look good to me, I'll leave this open for @kumarUjjawal to weigh in
# Which issue does this PR close? Fixes #11280 And also helps with: apache/datafusion#23565 # Rationale for this change See the original issue, but, basically the interleave kernel can make sub-optimal string views especially when buffers are shared. # What changes are included in this PR? We use a btreemap on the seen buffers (not individual values) and use that to prevent us from cloning the same buffer multiple times. # Are these changes tested? Yes tested, and locally benchmarks have shown a moderate speed increase in some of the interleave benchmarks, on my dev machine which has an **AMD Ryzen 9 9950X3D** processor: | benchmark | main | this PR | change | |-----------------------------------------------------------------------|----------|----------|--------| | interleave str_view(0.0) 100 [0..100, 100..230, 450..1000] | 284.5 ns | 268.6 ns | -5.6% | | interleave str_view(0.0) 400 [0..100, 100..230, 450..1000] | 706.9 ns | 563.9 ns | -20.2% | | interleave str_view(0.0) 1024 [0..100, 100..230, 450..1000] | 1.717 µs | 1.250 µs | -27.2% | | interleave str_view(0.0) 1024 [0..100, 100..230, 450..1000, 0..1000] | 1.783 µs | 1.274 µs | -28.6% | # Are there any user-facing changes? No user facing changes, just an optimization.
# Which issue does this PR close? Fixes apache#11280 And also helps with: apache/datafusion#23565 # Rationale for this change See the original issue, but, basically the interleave kernel can make sub-optimal string views especially when buffers are shared. # What changes are included in this PR? We use a btreemap on the seen buffers (not individual values) and use that to prevent us from cloning the same buffer multiple times. # Are these changes tested? Yes tested, and locally benchmarks have shown a moderate speed increase in some of the interleave benchmarks, on my dev machine which has an **AMD Ryzen 9 9950X3D** processor: | benchmark | main | this PR | change | |-----------------------------------------------------------------------|----------|----------|--------| | interleave str_view(0.0) 100 [0..100, 100..230, 450..1000] | 284.5 ns | 268.6 ns | -5.6% | | interleave str_view(0.0) 400 [0..100, 100..230, 450..1000] | 706.9 ns | 563.9 ns | -20.2% | | interleave str_view(0.0) 1024 [0..100, 100..230, 450..1000] | 1.717 µs | 1.250 µs | -27.2% | | interleave str_view(0.0) 1024 [0..100, 100..230, 450..1000, 0..1000] | 1.783 µs | 1.274 µs | -28.6% | # Are there any user-facing changes? No user facing changes, just an optimization. (cherry picked from commit 007894e)
…iew columns (apache#25625) ## Which issue does this PR close? - Refers to apache#23564. - Precursor for apache#23565. This PR must merge first, so the benchmark bot can compare `main` and that PR with `run benchmark spill_views`. ## Rationale for this change apache#23565 changes how spilled `StringView` and `BinaryView` columns are compacted. No current benchmark shows the effect: - The `sort_tpch` string columns are either high-cardinality (`l_comment`) or have a dictionary buffer that is too small to be compacted (`l_shipinstruct`). - Most suites spill only when the environment sets a low memory limit. ## What changes are included in this PR? A new SQL benchmark suite `spill_views` (in `benchmarks/sql_benchmarks/spill_views/`) and a `bench.sh` entry for it. There is no change to the spill code. Each query reads 1M rows from a Parquet file that the suite writes with `COPY`. The data must come from Parquet, because then the views of a dictionary-encoded column point into one shared buffer. Each query sets its own memory limit and `target_partitions = 4`, so it spills in any environment. | Query | Subgroup | What it does | Memory limit | |---|---|---|---| | `q01_sort_string_1_distinct` | `repeated` | `ORDER BY` a shuffled key, `StringView` payload with 1 distinct value (64 bytes) | 40M | | `q02_sort_string_1000_distinct` | `repeated` | Same, with 1000 distinct values | 40M | | `q03_sort_binary_1000_distinct` | `repeated` | Same as q02, with a `BinaryView` payload | 40M | | `q04_sort_string_all_distinct` | `distinct` | Same, with all-distinct values (compaction cannot remove data) | 96M | | `q05_group_by_string_all_distinct` | `distinct` | `GROUP BY` an all-distinct string key (compaction cannot remove data) | 96M | The asserts check the row count, the number of distinct values, that the column is read as a view type, and that the memory limit and `target_partitions` settings are applied. ## What is the testing strategy for this PR? I ran the suite with `benchmark_runner` and with `./bench.sh run spill_views` (the `cargo bench --bench sql` path that the bot uses). All asserts pass and each query takes less than 150 ms per iteration. I also ran the same queries with `EXPLAIN ANALYZE`, on `main` and with this commit on top of apache#23565. All queries spill on both sides. Spilled bytes are from the sort or final aggregation operator. Times are the median of 40 iterations of `benchmark_runner` on a laptop (Apple M4 Pro, local SSD), with runs of both builds alternated: | Query | Spilled (`main`) | Spilled (PR) | Time (`main`) | Time (PR) | |---|---|---|---|---| | q01 sort, 1 distinct | 84.2 MB | 23.2 MB | 63.5 ms | 57.0 ms | | q02 sort, 1000 distinct | 84.2 MB | 30.6 MB | 57.0 ms | 65.3 ms | | q03 sort, 1000 distinct binary | 84.2 MB | 30.6 MB | 54.0 ms | 65.9 ms | | q04 sort, all distinct | 84.2 MB | 84.2 MB | 71.0 ms | 95.7 ms | | q05 group by, all distinct | 84.2 MB | 84.2 MB | 75.2 ms | 77.7 ms | The machine was under load, so the times are only approximate. On a local SSD, spill I/O is cheap. The benchmark bot will give better numbers. ## Are there any user-facing changes? No. This PR only adds a benchmark. 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
…nment (apache#25634) ## Which issue does this PR close? - Follow-up to apache#25625. ## Rationale for this change Each `spill_views` query hardcodes its memory limit, and a benchmark-bot trigger comment can only set environment variables. So the only way to change a limit today is another DataFusion PR. At the current defaults every query spills on `main` and passes (q05, the GROUP BY, passed 30 of 30 executions at every limit from 96M to 192M). But apache#23565 fails q05 at 96M on every bot run (4 of 4), which aborts the whole suite run, so that PR gets no numbers. ## What changes are included in this PR? - `SPILL_VIEWS_LIMIT_REPEATED` (default `40M`) and `SPILL_VIEWS_LIMIT_DISTINCT` (default `96M`) set the limit of each subgroup, also as `--limit-repeated` / `--limit-distinct` on `benchmark_runner`. The defaults do not change. - The docs suggest `128M` as an override: q05 stops spilling at about `192M`, so larger values no longer measure the spill path. ## Are these changes tested? Ran the suite through `benchmark_runner` at the defaults and with `SPILL_VIEWS_LIMIT_DISTINCT=128M`. The suite's asserts read back the limit each query ran with, so a wrong override fails the run. ## Are there any user-facing changes? No, benchmarks only. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 5 <noreply@anthropic.com>
…ess (apache#25644) ## Which issue does this PR close? - None. This is follow-up work to apache#23985, which added per-query `pool_peak_bytes` to the `dfbench` results JSON. ## Rationale for this change The suites that `bench.sh` runs through the Criterion SQL harness (`cargo bench --bench sql`: `spill_views`, `wide_schema`, `predicate_eval`, and others) do not report a per-query memory pool peak. The harness prints one number to stdout and writes no results JSON. For example, this benchmark bot run of `spill_views` with `DATAFUSION_RUNTIME_MEMORY_LIMIT: "512M"` has no pool peak section: apache#23565 (comment) A second problem: a suite that sets its own limit with `SET datafusion.runtime.memory_limit` records nothing, even if the output were written. `spill_views` does this in its `init` script. The `SET` makes `SessionContext` build a new `RuntimeEnv` with a new pool, and that drops the `PeakRecordingPool` that `CommonOpt::runtime_env_builder` installed. ## What changes are included in this PR? Two commits: 1. **Keep the recorder across a SQL `SET`.** After a benchmark's `load` and `init` steps, `prepare_benchmark` puts a `PeakRecordingPool` in front of the session's pool again if the pool has a finite limit and no recorder. The new `RuntimeEnv` is installed the same way `SET datafusion.runtime.*` does it (`SessionStateBuilder::from(state).with_runtime_env(..)`), so tables and config are kept. The simple runner (`benchmark_runner -o`) also goes through `prepare_benchmark`, so it gets this fix too. 2. **Write a results JSON from the Criterion harness.** When `BENCH_RESULTS_FILE` is set, the harness writes the same JSON format as `dfbench -o`. Each case is the Criterion id (for example, `spill_views/q01_sort_string_1_distinct_repeated`, the same as in `critcmp`). `pool_peak_bytes` covers every execution Criterion makes of the query, including warm-up. The `load`, `init` and `assert` steps are not included. `iterations` is empty, because Criterion keeps the timings. `bench.sh` passes `results/<name>/criterion/<file>.json` at each SQL-harness call site. The subdirectory keeps the file out of `bench.sh compare`, which reads `results/<name>/*.json` as timing results. This means no second timing table. The benchmark bot reads the new file in adriangb/datafusion-benchmarking#38. The bot takes `bench.sh` from `main`, so it shows these tables only after this PR is merged. Until then, and for any side that is older than this PR, the bot shows a note that the pool peaks are not available. Not changed here: `SET datafusion.runtime.memory_limit` always builds a `GreedyMemoryPool` and drops any wrapper, so `--mem-pool-type` has no effect on suites that set their own limit. The fix in this PR stays in the benchmarks crate. A core change that keeps a pool wrapper across `SET` is possible, but it is out of scope for this PR. ## What is the testing strategy for this PR? New unit tests in `benchmarks/src/sql_benchmark_runner.rs`: - `sql_memory_limit_keeps_the_peak_recorder`: a SQL `SET` replaces the harness pool, and the recorder is put back and records the next query. - `no_memory_limit_gets_no_recorder`: without any limit, nothing changes. - `criterion_harness_writes_pool_peak_per_case`: an end-to-end Criterion run of a suite whose `init` sets the limit writes a JSON with a non-zero peak and empty `iterations`. If the rewrap is disabled, the first and third tests fail. `./bench.sh run spill_views`, the command that the bot runs (`SQL_CARGO_COMMAND="cargo bench --bench sql -- --save-baseline ..."`), writes `results/<name>/criterion/spill_views.json`. The peaks are the same with `DATAFUSION_RUNTIME_MEMORY_LIMIT=512M` and with no env limit, because the suite's own limits apply (40M for q01 to q03, 96M for q04 and q05): | Query | `pool_peak_bytes` | | --- | --- | | `spill_views/q01_sort_string_1_distinct_repeated` | 42614784 (40.6 MiB) | | `spill_views/q02_sort_string_1000_distinct_repeated` | 41709504 (39.8 MiB) | | `spill_views/q03_sort_binary_1000_distinct_repeated` | 41709504 (39.8 MiB) | | `spill_views/q04_sort_string_all_distinct_distinct` | 107544576 (102.6 MiB) | | `spill_views/q05_group_by_string_all_distinct_distinct` | 100633968 (96.0 MiB) | The q01 and q04 peaks are a little above their limits. The pool allows this: infallible `grow` calls are not limited, and the recorder counts what the pool grants. `./bench.sh compare` with two such result sets reads no file. `sort_tpch` (dfbench, queries 1 and 2, `DATAFUSION_RUNTIME_MEMORY_LIMIT=512M`) is unchanged. It writes `results/<name>/sort_tpch1.json` with 5 iterations and `pool_peak_bytes` for each query, and it writes no `criterion/` directory. ## Are there any user-facing changes? Only for benchmark tooling. There is a new optional `BENCH_RESULTS_FILE` env var for `cargo bench --bench sql` (documented in `benchmarks/sql_benchmarks/README.md`). `bench.sh run` now writes `results/<name>/criterion/*.json` for SQL-harness suites. A SQL-set memory limit now reports a pool peak. There are no public API changes outside the benchmarks crate. 🤖 Generated with [Claude Code](https://claude.com/claude-code) --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com>
## Which issue does this PR close? - Closes apache#23564 ## Rationale for this change Spill files inflate string views via GC. Arrow's `gc()` copies the bytes of every view separately, so a highly repeated string view column (for example a dictionary-encoded Parquet column) spills one copy per row, causing more memory and disk pressure. ## What changes are included in this PR? The spill path now compacts `StringView` / `BinaryView` arrays by copying each distinct value once: - Values are deduplicated with a hash table over the views, and null views are zeroed. - The first 256 non-inline values are sampled. If they contain almost no repeats, compaction falls back to plain `gc()`, so all-distinct data does not pay for hashing. ## Are these changes tested? Yes: - `test_gc_copies_repeated_values_once` checks that repeated values are written once, both when views share one buffer (as with Parquet dictionary columns) and when they reference separate copies. - `test_gc_distinct_values` checks the fallback to `gc()` for distinct values. - Existing spill tests cover the rest of the path. Benchmarks (`spill_views`, `sort_tpch` with a memory limit) show 0.78–0.94x on repeated-value spills, and no change for all-distinct data or `sort_tpch`. See the results in the comments. ## Are there any user-facing changes? No, this is an internal method change. --------- Co-authored-by: Claude Opus 5.5 <noreply@anthropic.com> Co-authored-by: Adrian Garcia Badaracco <1755071+adriangb@users.noreply.github.com>
Which issue does this PR close?
Rationale for this change
Spill files inflate string views via GC. Arrow's
gc()copies the bytes of every view separately, so a highly repeated string view column (for example a dictionary-encoded Parquet column) spills one copy per row, causing more memory and disk pressure.What changes are included in this PR?
The spill path now compacts
StringView/BinaryViewarrays by copying each distinct value once:gc(), so all-distinct data does not pay for hashing.Are these changes tested?
Yes:
test_gc_copies_repeated_values_oncechecks that repeated values are written once, both when views share one buffer (as with Parquet dictionary columns) and when they reference separate copies.test_gc_distinct_valueschecks the fallback togc()for distinct values.Benchmarks (
spill_views,sort_tpchwith a memory limit) show 0.78–0.94x on repeated-value spills, and no change for all-distinct data orsort_tpch. See the results in the comments.Are there any user-facing changes?
No, this is an internal method change.