Repository navigation
Hash aggregation produces batches reporting huge memory size #22526
Description
Activity
related work: #22164
@alamb I'd love to hear your thoughts on some feasible ways forward with this issue
- changed the title
[-]Partial hash aggregation produces batches reporting huge memory size[/-][+]Hash aggregation produces batches reporting huge memory size[/+]on May 28, 2026 Another proof of this issue:
First run with 100G memory limit reports:
partial_agg: output_rows=22.83 M, output_bytes=821.3 GB, output_batches=2.79 final_agg: output_rows=18.34 M, output_bytes=612.3 GB, output_batches=2.24 K ProjectionExec: output_rows=18.34 M, output_bytes=612.5 GB, output_batches=2.24 KWhile the second run with 600M memory limit reports:
partial_agg: output_rows=25.99 M, output_bytes=42.9 GB, output_batches=3.18 K final_agg: output_rows=43.53 M, output_bytes=19.1 GB, output_batches=6.17 K ProjectionExec: output_rows=18.34 M, output_bytes=5.1 GB, output_batches=3.09 KFor the same 18.34M rows, ProjectionExec reports:
- 612.5 GB output_bytes when run with a memory limit of 100 G
- 5.1 GB output_bytes when run with a memory limit of 600 M
The reason is that when OOM is not triggered (thus the partial aggregation doesn't emit early and the final aggregation doesn't spill), groups will continue accumulating as part of the large batch produced at the end of aggregation (the one that will be sliced as part of ProducingOutput execution state). Since each of these small batches reports the memory for the entire underlying arrow buffers, memory accounting ends up hugely inflated.
Full output:
$ cargo run --profile=release-nonlto --bin dfbench -- clickbench \ --path benchmarks/data/hits_partitioned \ --queries-path benchmarks/queries/clickbench/queries \ --query 34 \ -i 1 \ --debug \ --memory-limit 100G 2>&1 Running benchmarks with the following options: RunOpt { query: Some(34), pushdown: false, common: CommonOpt { iterations: 1, partitions: None, batch_size: None, mem_pool_type: "fair", memory_limit: Some(107374182400), sort_spill_reservation_bytes: None, debug: true, simulate_latency: false }, path: "benchmarks/data/hits_partitioned", queries_path: "benchmarks/queries/clickbench/queries", output_path: None, sorted_by: None, sort_order: "ASC", config_options: [] } Q34: -- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591 -- set datafusion.execution.parquet.binary_as_string = true SELECT 1, "URL", COUNT(*) AS c FROM hits GROUP BY 1, "URL" ORDER BY c DESC LIMIT 10; Query 34 iteration 0 took 1338.9 ms and returned 10 rows +-------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ | plan_type | plan | +-------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ | Plan with Metrics | SortPreservingMergeExec: [c@2 DESC], fetch=10, metrics=[output_rows=10, elapsed_compute=3.12µs, output_bytes=18.2 MB, output_batches=1] | | | SortExec: TopK(fetch=10), expr=[c@2 DESC], preserve_partitioning=[true], filter=[c@2 IS NULL OR c@2 > 55255], metrics=[output_rows=114, elapsed_compute=6.38ms, output_bytes=129.9 MB, output_batches=14, row_replacements=234] | | | ProjectionExec: expr=[1 as Int64(1), URL@0 as URL, count(Int64(1))@1 as c], metrics=[output_rows=18.34 M, elapsed_compute=2.84ms, output_bytes=612.5 GB, output_batches=2.24 K, expr_0_eval_time=1.82ms, expr_1_eval_time=280.68µs, expr_2_eval_time=156.67µs] | | | AggregateExec: mode=FinalPartitioned, gby=[URL@0 as URL], aggr=[count(Int64(1))], metrics=[output_rows=18.34 M, elapsed_compute=2.70s, output_bytes=612.3 GB, output_batches=2.24 K, spill_count=0, spilled_bytes=0.0 B, spilled_rows=0, peak_mem_used=9.47 B, aggregate_arguments_time=1.01ms, aggregation_time=78.08ms, emitting_time=70.33µs, time_calculating_group_ids=2.61s] | | | RepartitionExec: partitioning=Hash([URL@0], 14), input_partitions=14, metrics=[output_rows=22.83 M, elapsed_compute=2.09ms, output_bytes=5.6 GB, output_batches=2.79 K, spill_count=0, spilled_bytes=0.0 B, spilled_rows=0, fetch_time=11.11s, repartition_time=311.20ms, send_time=2.52s] | | | AggregateExec: mode=Partial, gby=[URL@0 as URL], aggr=[count(Int64(1))], metrics=[output_rows=22.83 M, elapsed_compute=6.08s, output_bytes=821.3 GB, output_batches=2.79 K, spill_count=0, spilled_bytes=0.0 B, spilled_rows=0, skipped_aggregation_rows=0, peak_mem_used=5.44 B, aggregate_arguments_time=21.59ms, aggregation_time=193.54ms, emitting_time=84.50µs, time_calculating_group_ids=5.84s, reduction_factor=22.83% (22.83 M/100.00 M)] | | | DataSourceExec: file_groups={14 groups: [[Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_0.parquet:0..122446530, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_1.parquet:0..174965044, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_10.parquet:0..101513258, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_11.parquet:0..118419888, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_12.parquet:0..149514164, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_15.parquet:88577877..103098894, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_16.parquet:0..101067219, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_17.parquet:0..116867853, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_18.parquet:0..133119589, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_19.parquet:0..103692598, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_23.parquet:73829085..79631107, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_24.parquet:0..78257049, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_25.parquet:0..144169728, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_26.parquet:0..156510916, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_27.parquet:0..166286210, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_30.parquet:67171810..124187913, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_31.parquet:0..123065410, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_32.parquet:0..94506004, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_33.parquet:0..78243765, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_34.parquet:0..119426616, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_39.parquet:94059938..103522954, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_4.parquet:0..140929275, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_40.parquet:0..142508647, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_41.parquet:0..290614269, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_42.parquet:0..288524057, ...], ...]}, projection=[URL], file_type=parquet, metrics=[output_rows=100.00 M, elapsed_compute=4.45ms, output_bytes=19.9 GB, output_batches=12.31 K, files_ranges_pruned_statistics=113 total → 113 matched, row_groups_pruned_statistics=325 total → 325 matched, row_groups_pruned_bloom_filter=325 total → 325 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, bytes_scanned=2.62 B, file_open_errors=0, file_scan_errors=0, files_opened=113, files_processed=113, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, bloom_filter_eval_time=3.90µs, metadata_load_time=2.11ms, page_index_eval_time=226ns, row_pushdown_eval_time=226ns, statistics_eval_time=226ns, time_elapsed_opening=8.92ms, time_elapsed_processing=4.33s, time_elapsed_scanning_total=11.09s, time_elapsed_scanning_until_data=313.52ms, output_rows_skew=2.18%, scan_efficiency_ratio=1637.04% (2.62 B/160.3 M)] | | | | +-------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ Query 34 avg time: 1338.93 ms$ cargo run --profile=release-nonlto --bin dfbench -- clickbench \ --path benchmarks/data/hits_partitioned \ --queries-path benchmarks/queries/clickbench/queries \ --query 34 \ -i 1 \ --debug \ --memory-limit 600M 2>&1 Finished `release-nonlto` profile [optimized] target(s) in 0.18s Running `target/release-nonlto/dfbench clickbench --path benchmarks/data/hits_partitioned --queries-path benchmarks/queries/clickbench/queries --query 34 -i 1 --debug --memory-limit 600M` Running benchmarks with the following options: RunOpt { query: Some(34), pushdown: false, common: CommonOpt { iterations: 1, partitions: None, batch_size: None, mem_pool_type: "fair", memory_limit: Some(629145600), sort_spill_reservation_bytes: None, debug: true, simulate_latency: false }, path: "benchmarks/data/hits_partitioned", queries_path: "benchmarks/queries/clickbench/queries", output_path: None, sorted_by: None, sort_order: "ASC", config_options: [] } Q34: -- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591 -- set datafusion.execution.parquet.binary_as_string = true SELECT 1, "URL", COUNT(*) AS c FROM hits GROUP BY 1, "URL" ORDER BY c DESC LIMIT 10; Query 34 iteration 0 took 7140.9 ms and returned 10 rows +-------------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ | plan_type | plan | +-------------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ | Plan with Metrics | SortPreservingMergeExec: [c@2 DESC], fetch=10, metrics=[output_rows=10, elapsed_compute=5.21µs, output_bytes=9.4 MB, output_batches=1] | | | SortExec: TopK(fetch=10), expr=[c@2 DESC], preserve_partitioning=[true], filter=[c@2 IS NULL OR c@2 > 55255], metrics=[output_rows=90, elapsed_compute=7.06ms, output_bytes=68.1 MB, output_batches=14, row_replacements=241] | | | ProjectionExec: expr=[1 as Int64(1), URL@0 as URL, count(Int64(1))@1 as c], metrics=[output_rows=18.34 M, elapsed_compute=3.06ms, output_bytes=5.1 GB, output_batches=3.09 K, expr_0_eval_time=2.02ms, expr_1_eval_time=195.81µs, expr_2_eval_time=82.64µs] | | | AggregateExec: mode=FinalPartitioned, gby=[URL@0 as URL], aggr=[count(Int64(1))], metrics=[output_rows=43.53 M, elapsed_compute=13.57s, output_bytes=19.1 GB, output_batches=6.17 K, spill_count=1.11 K, spilled_bytes=17.3 GB, spilled_rows=103.2 M, peak_mem_used=260.8 M, aggregate_arguments_time=2.26ms, aggregation_time=60.97ms, emitting_time=4.52ms, time_calculating_group_ids=2.49s] | | | RepartitionExec: partitioning=Hash([URL@0], 14), input_partitions=14, metrics=[output_rows=25.99 M, elapsed_compute=6.15ms, output_bytes=5.9 GB, output_batches=3.18 K, spill_count=10, spilled_bytes=45.2 MB, spilled_rows=237.6 K, fetch_time=9.46s, repartition_time=381.07ms, send_time=1.79s] | | | AggregateExec: mode=Partial, gby=[URL@0 as URL], aggr=[count(Int64(1))], metrics=[output_rows=25.99 M, elapsed_compute=4.32s, output_bytes=42.9 GB, output_batches=3.18 K, spill_count=0, spilled_bytes=0.0 B, spilled_rows=0, skipped_aggregation_rows=0, peak_mem_used=230.1 M, aggregate_arguments_time=23.73ms, aggregation_time=153.48ms, emitting_time=131.38ms, time_calculating_group_ids=3.95s, reduction_factor=25.99% (25.99 M/100.00 M)] | | | DataSourceExec: file_groups={14 groups: [[Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_0.parquet:0..122446530, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_1.parquet:0..174965044, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_10.parquet:0..101513258, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_11.parquet:0..118419888, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_12.parquet:0..149514164, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_15.parquet:88577877..103098894, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_16.parquet:0..101067219, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_17.parquet:0..116867853, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_18.parquet:0..133119589, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_19.parquet:0..103692598, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_23.parquet:73829085..79631107, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_24.parquet:0..78257049, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_25.parquet:0..144169728, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_26.parquet:0..156510916, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_27.parquet:0..166286210, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_30.parquet:67171810..124187913, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_31.parquet:0..123065410, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_32.parquet:0..94506004, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_33.parquet:0..78243765, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_34.parquet:0..119426616, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_39.parquet:94059938..103522954, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_4.parquet:0..140929275, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_40.parquet:0..142508647, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_41.parquet:0..290614269, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_42.parquet:0..288524057, ...], ...]}, projection=[URL], file_type=parquet, metrics=[output_rows=100.00 M, elapsed_compute=3.95ms, output_bytes=19.9 GB, output_batches=12.31 K, files_ranges_pruned_statistics=113 total → 113 matched, row_groups_pruned_statistics=325 total → 325 matched, row_groups_pruned_bloom_filter=325 total → 325 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, bytes_scanned=2.62 B, file_open_errors=0, file_scan_errors=0, files_opened=113, files_processed=113, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, bloom_filter_eval_time=5.47µs, metadata_load_time=2.15ms, page_index_eval_time=226ns, row_pushdown_eval_time=226ns, statistics_eval_time=226ns, time_elapsed_opening=10.48ms, time_elapsed_processing=4.12s, time_elapsed_scanning_total=17.84s, time_elapsed_scanning_until_data=348.59ms, output_rows_skew=3.54%, scan_efficiency_ratio=1637.04% (2.62 B/160.3 M)] | | | | +-------------------+-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ Query 34 avg time: 7140.90 msThe spilling issue occurs when there's an operator which accumulates the input partitions, such as
SortExec.
This can be achieved by removinglimit 10from q34:diff --git a/benchmarks/queries/clickbench/queries/q34.sql b/benchmarks/queries/clickbench/queries/q34.sql index fdb7edbb6..864294d7e 100644 --- a/benchmarks/queries/clickbench/queries/q34.sql +++ b/benchmarks/queries/clickbench/queries/q34.sql @@ -1,4 +1,4 @@ -- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591 -- set datafusion.execution.parquet.binary_as_string = true -SELECT 1, "URL", COUNT(*) AS c FROM hits GROUP BY 1, "URL" ORDER BY c DESC LIMIT 10; +SELECT 1, "URL", COUNT(*) AS c FROM hits GROUP BY 1, "URL" ORDER BY c DESC;As shown below, there are 735 spill counts for the SortExec operator, even though the memory limit for the execution was set to 50G.
spill_count=735, spilled_bytes=3.7 GB, spilled_rows=18.34 M
and the peak RSS doesn't go over 10 GB for this test:
$ cargo run --profile=release-nonlto --bin mem_profile -- clickbench \ --path benchmarks/data/hits_partitioned \ --queries-path benchmarks/queries/clickbench/queries \ --query 34 \ -i 1 \ --debug 2>&1 Finished `release-nonlto` profile [optimized] target(s) in 0.50s Running `target/release-nonlto/mem_profile clickbench --path benchmarks/data/hits_partitioned --queries-path benchmarks/queries/clickbench/queries --query 34 -i 1 --debug` Pre-building benchmark binary... Finished `release` profile [optimized] target(s) in 0.29s Benchmark binary built successfully. Query Time (ms) Peak RSS Peak Commit Major Page Faults ---------------------------------------------------------------- 34 1500.15 9.9 GB 10.2 GB 0$ cargo run --profile=release-nonlto --bin dfbench -- clickbench \ --path benchmarks/data/hits_partitioned \ --queries-path benchmarks/queries/clickbench/queries \ --query 34 \ -i 1 \ --debug \ --memory-limit 50G 2>&1 Finished `release-nonlto` profile [optimized] target(s) in 0.18s Running `target/release-nonlto/dfbench clickbench --path benchmarks/data/hits_partitioned --queries-path benchmarks/queries/clickbench/queries --query 34 -i 1 --debug --memory-limit 50G` Running benchmarks with the following options: RunOpt { query: Some(34), pushdown: false, common: CommonOpt { iterations: 1, partitions: None, batch_size: None, mem_pool_type: "fair", memory_limit: Some(53687091200), sort_spill_reservation_bytes: None, debug: true, simulate_latency: false }, path: "benchmarks/data/hits_partitioned", queries_path: "benchmarks/queries/clickbench/queries", output_path: None, sorted_by: None, sort_order: "ASC", config_options: [] } Q34: -- Must set for ClickBench hits_partitioned dataset. See https://github.com/apache/datafusion/issues/16591 -- set datafusion.execution.parquet.binary_as_string = true SELECT 1, "URL", COUNT(*) AS c FROM hits GROUP BY 1, "URL" ORDER BY c DESC; Query 34 iteration 0 took 2639.3 ms and returned 18342019 rows +-------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ | plan_type | plan | +-------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ | Plan with Metrics | SortPreservingMergeExec: [c@2 DESC], metrics=[output_rows=18.34 M, elapsed_compute=247.86ms, output_bytes=69.3 GB, output_batches=2.24 K] | | | SortExec: expr=[c@2 DESC], preserve_partitioning=[true], metrics=[output_rows=18.34 M, elapsed_compute=632.17ms, output_bytes=25.4 GB, output_batches=2.24 K, spill_count=735, spilled_bytes=3.7 GB, spilled_rows=18.34 M] | | | ProjectionExec: expr=[1 as Int64(1), URL@0 as URL, count(Int64(1))@1 as c], metrics=[output_rows=18.34 M, elapsed_compute=8.60ms, output_bytes=611.9 GB, output_batches=2.24 K, expr_0_eval_time=7.03ms, expr_1_eval_time=194.37µs, expr_2_eval_time=90.21µs] | | | AggregateExec: mode=FinalPartitioned, gby=[URL@0 as URL], aggr=[count(Int64(1))], metrics=[output_rows=18.34 M, elapsed_compute=2.67s, output_bytes=611.8 GB, output_batches=2.24 K, spill_count=0, spilled_bytes=0.0 B, spilled_rows=0, peak_mem_used=9.47 B, aggregate_arguments_time=1.05ms, aggregation_time=82.06ms, emitting_time=51.50µs, time_calculating_group_ids=2.58s] | | | RepartitionExec: partitioning=Hash([URL@0], 14), input_partitions=14, metrics=[output_rows=22.75 M, elapsed_compute=2.02ms, output_bytes=5.5 GB, output_batches=2.79 K, spill_count=0, spilled_bytes=0.0 B, spilled_rows=0, fetch_time=10.42s, repartition_time=364.31ms, send_time=2.04s] | | | AggregateExec: mode=Partial, gby=[URL@0 as URL], aggr=[count(Int64(1))], metrics=[output_rows=22.75 M, elapsed_compute=5.76s, output_bytes=827.3 GB, output_batches=2.78 K, spill_count=0, spilled_bytes=0.0 B, spilled_rows=0, skipped_aggregation_rows=0, peak_mem_used=5.39 B, aggregate_arguments_time=21.97ms, aggregation_time=189.26ms, emitting_time=78.41µs, time_calculating_group_ids=5.54s, reduction_factor=22.75% (22.75 M/100.00 M)] | | | DataSourceExec: file_groups={14 groups: [[Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_0.parquet:0..122446530, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_1.parquet:0..174965044, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_10.parquet:0..101513258, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_11.parquet:0..118419888, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_12.parquet:0..149514164, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_15.parquet:88577877..103098894, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_16.parquet:0..101067219, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_17.parquet:0..116867853, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_18.parquet:0..133119589, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_19.parquet:0..103692598, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_23.parquet:73829085..79631107, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_24.parquet:0..78257049, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_25.parquet:0..144169728, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_26.parquet:0..156510916, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_27.parquet:0..166286210, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_30.parquet:67171810..124187913, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_31.parquet:0..123065410, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_32.parquet:0..94506004, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_33.parquet:0..78243765, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_34.parquet:0..119426616, ...], [Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_39.parquet:94059938..103522954, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_4.parquet:0..140929275, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_40.parquet:0..142508647, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_41.parquet:0..290614269, Users/amiculas/work/datafusion/benchmarks/data/hits_partitioned/hits_42.parquet:0..288524057, ...], ...]}, projection=[URL], file_type=parquet, metrics=[output_rows=100.00 M, elapsed_compute=3.56ms, output_bytes=19.9 GB, output_batches=12.31 K, files_ranges_pruned_statistics=113 total → 113 matched, row_groups_pruned_statistics=325 total → 325 matched, row_groups_pruned_bloom_filter=325 total → 325 matched, page_index_pages_pruned=0 total → 0 matched, page_index_rows_pruned=0 total → 0 matched, limit_pruned_row_groups=0 total → 0 matched, batches_split=0, bytes_scanned=2.62 B, file_open_errors=0, file_scan_errors=0, files_opened=113, files_processed=113, num_predicate_creation_errors=0, predicate_evaluation_errors=0, pushdown_rows_matched=0, pushdown_rows_pruned=0, predicate_cache_inner_records=0, predicate_cache_records=0, bloom_filter_eval_time=3.89µs, metadata_load_time=1.92ms, page_index_eval_time=226ns, row_pushdown_eval_time=226ns, statistics_eval_time=226ns, time_elapsed_opening=9.13ms, time_elapsed_processing=4.09s, time_elapsed_scanning_total=10.39s, time_elapsed_scanning_until_data=260.32ms, output_rows_skew=2.95%, scan_efficiency_ratio=1637.04% (2.62 B/160.3 M)] | | | | +-------------------+------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+ Query 34 avg time: 2639.29 msThis also impacts performance:
50G memory limit: Query 34 avg time: 2911.36 ms
No memory limit: Query 34 avg time: 1505.30 ms@alamb I'd love to hear your thoughts on some feasible ways forward with this issue
From my perspective, the "right" solution is
The challenge is getting enough engineering effort to make it not (too much) slower than what we have today and break it up into managable (reviewable) pieces
Reacted by Ariel Miculas-TrifOk, please let me know if there are ways I could help get this PR merged
I filed an epic with a bunch of related memory items to try and coordinate some response
- marked Downstream consumers of AggregateExec significantly overcount memory usage #22757 as a duplicate of this issue
on Jun 4, 2026 I have another reproducer and explanation of this issue here, in case it helps - #22757
The solution we used internally for this was the alternate idea of Option 2:
Option 2. Fix memory accounting
Another idea:
extend this function to work across multiple batches and deduplicate the arrow buffersI think this will be easier to do and would also be a more general solution -- any operator that emits record batches with shared underlying buffers would be handled. I can open a draft PR doing this in hash join - we can discuss more there.
I think this will be easier to do and would also be a more general solution -- any operator that emits record batches with shared underlying buffers would be handled. I can open a draft PR doing this in hash join - we can discuss more there.
One problem with this approach is that the first sliced record batch would still report the memory of the entire original large record batch (the one being sliced).
It could lead to the following situation (observed in practice):
- the hash aggregation hits the OOM case, creates the large RecordBatch (which will be sliced) and releases the memory reservation
- the downstream operator consumes the first batch (which is a small one as a result of the
slice, but reporting the memory of all the underlying arrow buffers, i.e. large value returned by RecordBatch::get_array_memory_size) - the reservation for this batch fails (even though theoretically the hash aggregation released the memory reservation, the new reservation for this RecordBatch could require a bit more memory, meaning it would fail)
Conclusion: the downstream operator cannot even reserve memory for a single batch, resulting in spilling for every batch.
It would be nice if we could use something like:
impl GetSlicedSize for RecordBatch {
(but without breaking the functionality of the other operators)Reacted by Samyak Sarnayakthe reservation for this batch fails (even though theoretically the hash aggregation released the memory reservation, the new reservation for this RecordBatch could require a bit more memory, meaning it would fail)
Not sure I follow. Why would this sliced record batch require more memory than what agg released?
But the problem does make sense. Another example is
RepartitionExec. Even if only one record batch is actually buffered in the channels, we would be reserving size for the whole buffer (this happens today too).--
It would be nice if we could use something like:
datafusion/datafusion/physical-plan/src/spill/spill_manager.rs
Line 214 in e7f7fa9
impl GetSlicedSize for RecordBatch {(but without breaking the functionality of the other operators)
This also makes sense, but it would have the problem of undercounting in case of unreferenced data. One scenario is
StringViewArrays that wereslice()d/.take()n but not GC'd yet.Not sure I follow. Why would this sliced record batch require more memory than what agg released?
- the memory is reserved before the large record batch is created, so there's no guarantee that the released memory is sufficient for reserving the large record batch (well, a slice of this large record batch, but for memory accounting purposes it doesn't matter)
- it also depends on how memory is accounted:
RecordBatch::get_array_memory_size/ get_record_batch_memory_size or whatever the downstream operator decides to use - since there's no way to transfer a memory reservation from one operator to another, other memory pool operations could happen in-between, so there's no guarantee that if you free N bytes from a reservation in an operator you could reserve the same N bytes in another operator
- added a commit that references this issue
on Jun 5, 2026 FYI there’s now a PR for the fix I was imagining: #22862
Reacted by Ariel Miculas-Trif-
the memory is reserved before the large record batch is created, so there's no guarantee that the released memory is sufficient for reserving the large record batch (well, a slice of this large record batch, but for memory accounting purposes it doesn't matter)
-
since there's no way to transfer a memory reservation from one operator to another, other memory pool operations could happen in-between, so there's no guarantee that if you free N bytes from a reservation in an operator you could reserve the same N bytes in another operator
Very valid points. But all of these also apply to the current memory tracking, which is
get_record_batch_memory_size(used in HashJoin, Repartition, etc.). Currently, the downstream operator will try to reserve a lot more memory than what agg released. What I'm suggesting is strictly an improvement over the current behavior. Do you see any of these problems being made worse by a solution like #22862?-
As a short update here, in my mind the next steps are:
Refactor
@2010YOUY01 completes the refactor to separate more of the Hash aggregate operator
This will make it feasible to actually implement blocked state management in a way that will not reduce performance
Their POC looks promising
So we're moving away from #15591? Does it have performance issues?
So we're moving away from #15591? Does it have performance issues?
In my mind what @2010YOUY01 is doing is getting the aggregation code in shape so that something like #15591 could actually be feasible
Ok, that was my impression also, thanks for clarifying
Comet has hit this issue as well at the shuffle layer now that we don't deep copy there anymore (basically performing GC). My temporary solution is to dedupe the backing buffers as batches start accumulating at the shuffle writer. apache/datafusion-comet@b4cb826
Reacted by Ariel Miculas-Trif@mbutrovich basically #22862?
- added a commit that references this issue
on Jul 22, 2026 I filed a ticket / epic to track the idea of blocked state management here
Describe the bug
For aggregations (both partial and final) with
GroupOrdering::None, a huge batch is produced after consuming all the input RecordBatches, which is further sliced in order to produce batches ofbatch_sizelength.In row_hash.rs,
ExecutionState::ProducingOutput(batch)slices the large batch:Unfortunately
get_array_memory_sizefor each of these small RecordBatches returns the physical memory of the initial huge batch, causing unnecessary spills in the downstream operator.Operators such as RepartitionExec use
batch.get_array_memory_size()for deciding whether or not to spill.Related to #19481
To Reproduce
See the linked PR: #22527
Expected behavior
Option 1. Avoid producing the initial huge RecordBatch
Avoid producing a large RecordBatch and slicing it afterwards, instead emit only batch_size RecordBatches at a time, this would also fix: #18907
The issue with this approach is described here: #19906, i.e. the bottleneck of Vec::drain of shifting all the existing elements in the vector (though I'm not sure this is more significant then the performance impact of spilling)
Merging #15591 would probably make this solution more feasible, alternatively VecDeque could be considered: #19906 (comment)
Option 2. Fix memory accounting
Alternatively, fix the memory accounting for these small RecordBatches, so they only report the memory occupied by their slice instead of the entire underlying array.
Another idea:
extend this function to work across multiple batches and deduplicate the arrow buffers:
datafusion/datafusion/common/src/utils/memory.rs
Line 133 in 7c05b20
Additional context
No response