Skip to content

[feat](spill) Support multi-level partition spilling - #61212

Merged
yiguolei merged 4 commits into
apache:masterfrom
mrhhsg:spill_repartition
Mar 18, 2026
Merged

[feat](spill) Support multi-level partition spilling#61212
yiguolei merged 4 commits into
apache:masterfrom
mrhhsg:spill_repartition

Conversation

@mrhhsg

@mrhhsg mrhhsg commented Mar 11, 2026

Copy link
Copy Markdown
Member

What problem does this PR solve?

This PR adds support for spill repartition in partitioned operators.

When memory is still insufficient after a partition has been spilled, the spilled partition can be repartitioned into smaller sub-partitions, which further reduces peak memory usage. The repartitioner is made level-aware so repartition can be applied recursively when needed.

This PR also integrates the new spill repartition flow into partitioned hash join and partitioned aggregation, adds force-spill logic in the join sink operator, refactors spill file interfaces, improves spill metadata maintenance and memory-pressure handling, and fixes several spill-related issues such as revocable memory accounting and profile updates.

Related PR: #xxx

Problem Summary:

Release note

None

Check List (For Author)

  • Test

    • Regression test
    • Unit Test
    • Manual test (add detailed scripts or steps below)
    • No need to test or manual test. Explain why:
      • This is a refactor/code format and no logic has been changed.
      • Previous test can cover this change.
      • No code files have been changed.
      • Other reason
  • Behavior changed:

    • No.
    • Yes.
  • Does this need documentation?

    • No.
    • Yes.

Check List (For Reviewer who merge this PR)

  • Confirm the release note
  • Confirm test cases
  • Confirm document
  • Add branch pick label

@Thearas

Thearas commented Mar 11, 2026

Copy link
Copy Markdown
Contributor

Thank you for your contribution to Apache Doris.
Don't know what should be done next? See How to process your PR.

Please clearly describe your PR:

  1. What problem was fixed (it's best to include specific error reporting information). How it was fixed.
  2. Which behaviors were modified. What was the previous behavior, what is it now, why was it modified, and what possible impacts might there be.
  3. What features were added. Why was this function added?
  4. Which code was refactored and why was this part of the code refactored?
  5. Which functions were optimized and what is the difference before and after the optimization?

@mrhhsg

mrhhsg commented Mar 11, 2026

Copy link
Copy Markdown
Member Author

run buildall

@doris-robot

Copy link
Copy Markdown

Cloud UT Coverage Report

Increment line coverage 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 79.31% (1798/2267)
Line Coverage 64.62% (32273/49944)
Region Coverage 65.55% (16163/24659)
Branch Coverage 55.95% (8610/15388)

@doris-robot

Copy link
Copy Markdown

BE UT Coverage Report

Increment line coverage 84.95% (1801/2120) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 52.59% (19638/37342)
Line Coverage 36.20% (183323/506451)
Region Coverage 32.39% (142046/438579)
Branch Coverage 33.56% (62023/184838)

@doris-robot

Copy link
Copy Markdown
TPC-H: Total hot run time: 28022 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit cfcd7ad768bfd59f259de091c3cd51802cdc6662, data reload: false

------ Round 1 ----------------------------------
============================================
q1	17643	4491	4329	4329
q2	q3	10746	788	519	519
q4	4725	383	256	256
q5	8178	1220	1025	1025
q6	245	178	145	145
q7	849	875	670	670
q8	10882	1502	1360	1360
q9	6799	4816	4781	4781
q10	6679	1927	1672	1672
q11	469	280	248	248
q12	767	578	472	472
q13	18084	2940	2183	2183
q14	233	230	217	217
q15	938	817	830	817
q16	740	721	689	689
q17	736	876	430	430
q18	6155	5576	5315	5315
q19	1187	1030	627	627
q20	495	499	626	499
q21	4654	1944	1507	1507
q22	385	299	261	261
Total cold run time: 101589 ms
Total hot run time: 28022 ms

----- Round 2, with runtime_filter_mode=off -----
============================================
q1	4807	4643	4598	4598
q2	q3	3942	4372	3821	3821
q4	877	1215	815	815
q5	4165	4414	4363	4363
q6	187	171	144	144
q7	1796	1660	1516	1516
q8	2538	2780	2597	2597
q9	7788	7542	7350	7350
q10	3803	4098	3573	3573
q11	587	496	420	420
q12	490	582	454	454
q13	2829	3392	2307	2307
q14	283	301	283	283
q15	891	861	827	827
q16	756	768	723	723
q17	1166	1599	1400	1400
q18	7248	7013	6844	6844
q19	964	938	978	938
q20	2092	2208	2069	2069
q21	3956	3458	3364	3364
q22	450	445	373	373
Total cold run time: 51615 ms
Total hot run time: 48779 ms

@doris-robot

Copy link
Copy Markdown
TPC-DS: Total hot run time: 153851 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpcds-tools
TPC-DS sf100 test result on commit cfcd7ad768bfd59f259de091c3cd51802cdc6662, data reload: false

query5	4322	657	523	523
query6	333	237	193	193
query7	4212	492	264	264
query8	332	245	231	231
query9	8703	2761	2740	2740
query10	515	448	369	369
query11	7389	5883	5645	5645
query12	185	128	123	123
query13	1269	469	347	347
query14	5663	3820	3593	3593
query14_1	2801	2821	2800	2800
query15	203	200	176	176
query16	970	476	441	441
query17	904	737	605	605
query18	2441	459	363	363
query19	228	217	181	181
query20	136	129	133	129
query21	232	148	124	124
query22	5049	5005	5004	5004
query23	16711	16190	15911	15911
query23_1	16021	15967	15741	15741
query24	7672	1814	1312	1312
query24_1	1271	1324	1324	1324
query25	575	512	417	417
query26	1287	291	155	155
query27	2871	492	289	289
query28	4512	1854	1863	1854
query29	828	579	518	518
query30	317	253	220	220
query31	1371	1296	1204	1204
query32	87	74	72	72
query33	513	337	284	284
query34	944	952	552	552
query35	645	703	599	599
query36	1069	1095	932	932
query37	143	97	85	85
query38	2941	2971	2863	2863
query39	898	886	861	861
query39_1	832	858	825	825
query40	240	151	138	138
query41	60	60	59	59
query42	304	290	296	290
query43	243	255	232	232
query44	
query45	199	191	189	189
query46	933	1018	610	610
query47	2104	2127	2016	2016
query48	316	327	228	228
query49	637	458	375	375
query50	708	283	217	217
query51	4167	4192	4008	4008
query52	290	299	283	283
query53	293	337	281	281
query54	296	265	266	265
query55	97	90	81	81
query56	313	310	303	303
query57	1338	1349	1282	1282
query58	287	276	271	271
query59	1363	1436	1299	1299
query60	338	332	328	328
query61	147	147	146	146
query62	639	593	543	543
query63	313	277	276	276
query64	4944	1272	1014	1014
query65	
query66	1462	494	355	355
query67	16413	16419	16572	16419
query68	
query69	405	352	284	284
query70	1010	977	984	977
query71	339	302	311	302
query72	2800	2707	2476	2476
query73	550	606	323	323
query74	10040	9926	9833	9833
query75	2926	2771	2483	2483
query76	2272	1127	687	687
query77	371	400	308	308
query78	11251	11245	10690	10690
query79	2301	803	614	614
query80	1767	644	566	566
query81	567	301	261	261
query82	1017	155	117	117
query83	340	281	258	258
query84	307	134	110	110
query85	994	549	517	517
query86	429	313	320	313
query87	3154	3177	3003	3003
query88	3611	2685	2631	2631
query89	429	368	344	344
query90	2028	184	164	164
query91	164	158	141	141
query92	86	79	70	70
query93	1036	863	508	508
query94	638	333	247	247
query95	595	401	312	312
query96	638	538	232	232
query97	2510	2523	2426	2426
query98	237	221	223	221
query99	998	992	923	923
Total cold run time: 236257 ms
Total hot run time: 153851 ms

@hello-stephen

Copy link
Copy Markdown
Contributor

BE Regression && UT Coverage Report

Increment line coverage 85.78% (1809/2109) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 73.31% (26802/36559)
Line Coverage 56.57% (285603/504901)
Region Coverage 53.75% (237944/442722)
Branch Coverage 55.55% (102989/185414)

@hello-stephen

Copy link
Copy Markdown
Contributor

BE Regression && UT Coverage Report

Increment line coverage 85.78% (1809/2109) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 73.33% (26808/36559)
Line Coverage 56.59% (285720/504901)
Region Coverage 53.76% (237996/442722)
Branch Coverage 55.56% (103023/185414)

yiguolei
yiguolei previously approved these changes Mar 12, 2026
@mrhhsg

mrhhsg commented Mar 12, 2026

Copy link
Copy Markdown
Member Author

/review

@mrhhsg

mrhhsg commented Mar 13, 2026

Copy link
Copy Markdown
Member Author

run buildall

@mrhhsg

mrhhsg commented Mar 13, 2026

Copy link
Copy Markdown
Member Author

/review

@doris-robot

Copy link
Copy Markdown
TPC-H: Total hot run time: 27280 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit bedb76c2363416672171c03d30a2e8aed75b4ff7, data reload: false

------ Round 1 ----------------------------------
orders	Doris	NULL	NULL	0	0	0	NULL	0	NULL	NULL	2023-12-26 18:27:23	2023-12-26 18:42:55	NULL	utf-8	NULL	NULL	
============================================
q1	17641	4564	4410	4410
q2	q3	10742	840	528	528
q4	4730	379	259	259
q5	8193	1202	1017	1017
q6	234	177	158	158
q7	821	850	699	699
q8	10709	1509	1354	1354
q9	6810	4817	4737	4737
q10	6357	1918	1662	1662
q11	483	266	251	251
q12	739	568	476	476
q13	18073	2966	2190	2190
q14	235	227	221	221
q15	q16	751	742	677	677
q17	746	853	460	460
q18	5932	5425	5224	5224
q19	1123	979	638	638
q20	537	493	390	390
q21	4659	2004	1645	1645
q22	379	305	284	284
Total cold run time: 99894 ms
Total hot run time: 27280 ms

----- Round 2, with runtime_filter_mode=off -----
orders	Doris	NULL	NULL	150000000	42	6422171781	NULL	22778155	NULL	NULL	2023-12-26 18:27:23	2023-12-26 18:42:55	NULL	utf-8	NULL	NULL	
============================================
q1	4696	4567	4537	4537
q2	q3	3906	4471	3910	3910
q4	946	1195	809	809
q5	4070	4375	4365	4365
q6	180	176	142	142
q7	1754	1698	1578	1578
q8	2521	2806	2626	2626
q9	7465	7357	7631	7357
q10	3751	3971	3585	3585
q11	559	527	446	446
q12	527	610	465	465
q13	2796	3131	2306	2306
q14	273	304	285	285
q15	q16	748	768	750	750
q17	1146	1368	1436	1368
q18	7260	6822	6719	6719
q19	950	948	928	928
q20	2063	2169	2141	2141
q21	4033	3502	3446	3446
q22	479	438	387	387
Total cold run time: 50123 ms
Total hot run time: 48150 ms

@doris-robot

Copy link
Copy Markdown
TPC-DS: Total hot run time: 168100 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpcds-tools
TPC-DS sf100 test result on commit bedb76c2363416672171c03d30a2e8aed75b4ff7, data reload: false

query5	4338	651	522	522
query6	329	225	201	201
query7	4220	485	278	278
query8	349	256	241	241
query9	8722	2792	2761	2761
query10	534	391	350	350
query11	7026	5094	4874	4874
query12	187	127	123	123
query13	1282	467	360	360
query14	5756	3733	3530	3530
query14_1	2846	2892	2854	2854
query15	206	199	180	180
query16	1003	498	453	453
query17	1123	719	618	618
query18	2441	462	364	364
query19	220	211	184	184
query20	134	132	127	127
query21	215	134	125	125
query22	13293	13446	13250	13250
query23	15903	15541	15971	15541
query23_1	15818	15671	15854	15671
query24	8028	1695	1327	1327
query24_1	1362	1315	1325	1315
query25	713	522	466	466
query26	1266	283	173	173
query27	3224	514	313	313
query28	4953	1974	1927	1927
query29	872	574	474	474
query30	296	216	192	192
query31	984	943	876	876
query32	82	70	74	70
query33	525	326	291	291
query34	891	878	528	528
query35	635	680	603	603
query36	1055	1137	992	992
query37	136	93	88	88
query38	3021	2985	2873	2873
query39	864	829	810	810
query39_1	799	786	804	786
query40	229	150	133	133
query41	64	59	60	59
query42	256	248	256	248
query43	253	246	233	233
query44	
query45	195	191	186	186
query46	886	1002	625	625
query47	2564	2104	2026	2026
query48	321	318	224	224
query49	637	460	389	389
query50	696	275	211	211
query51	4142	4054	4011	4011
query52	261	272	251	251
query53	288	335	287	287
query54	296	268	274	268
query55	93	88	93	88
query56	314	332	319	319
query57	1961	1879	1520	1520
query58	287	277	273	273
query59	2798	2970	2745	2745
query60	348	337	319	319
query61	155	148	153	148
query62	633	586	534	534
query63	313	274	279	274
query64	5167	1266	1022	1022
query65	
query66	1451	465	360	360
query67	24247	24289	24279	24279
query68	
query69	415	313	291	291
query70	956	1003	948	948
query71	342	310	307	307
query72	2852	2670	2437	2437
query73	547	555	315	315
query74	9595	9617	9397	9397
query75	2865	2788	2474	2474
query76	2292	1022	686	686
query77	379	379	329	329
query78	11034	11184	10436	10436
query79	1130	795	582	582
query80	1299	647	531	531
query81	561	261	231	231
query82	1011	149	119	119
query83	332	268	251	251
query84	251	121	95	95
query85	905	493	436	436
query86	417	341	297	297
query87	3151	3167	3058	3058
query88	3580	2700	2654	2654
query89	428	369	350	350
query90	2008	186	170	170
query91	170	161	136	136
query92	76	73	74	73
query93	981	856	521	521
query94	647	314	260	260
query95	601	402	328	328
query96	630	514	233	233
query97	2454	2474	2402	2402
query98	235	224	216	216
query99	998	1013	916	916
Total cold run time: 252159 ms
Total hot run time: 168100 ms

@doris-robot

Copy link
Copy Markdown

BE UT Coverage Report

Increment line coverage 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 52.70% (19722/37424)
Line Coverage 36.27% (184376/508295)
Region Coverage 32.45% (142805/440050)
Branch Coverage 33.67% (62492/185615)

@hello-stephen

Copy link
Copy Markdown
Contributor

BE Regression && UT Coverage Report

Increment line coverage 100% (0/0) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 73.29% (26841/36624)
Line Coverage 56.63% (286842/506562)
Region Coverage 53.94% (239523/444093)
Branch Coverage 55.69% (103633/186097)

@doris-robot

Copy link
Copy Markdown

Cloud UT Coverage Report

Increment line coverage 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 79.24% (1798/2269)
Line Coverage 64.50% (32267/50026)
Region Coverage 65.42% (16160/24702)
Branch Coverage 55.81% (8603/15416)

@zclllyybb

Copy link
Copy Markdown
Contributor

/review

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Code Review: feat Support multi-level partition spilling

Reviewed ~80 files (9124 additions, 3594 deletions). The overall architecture is sound — the new SpillFile/SpillFileWriter/SpillFileReader abstraction is well-designed, the SpillRepartitioner cleanly handles both partitioner (join) and column-index (aggregation) modes, and the recursive repartition approach is a solid solution for large partitions that don't fit in memory.

Critical Checkpoints (Part 1.3)

Goal & correctness: Multi-level spill repartitioning for hash join and aggregation. The core goal is achieved. Test coverage is strong for join probe repartition (level transitions, max depth, row conservation) and good for aggregation source repartition. SpillFile infrastructure and SpillRepartitioner have comprehensive unit tests.

Concurrency: Single-writer pattern for SpillFile is correct — no concurrent mutation of _part_count/_total_written_bytes. Pipeline task execution is single-threaded per task. Dependency block/set_ready patterns are properly paired.

Lifecycle: SpillFileWriter uses weak_ptr back-reference to SpillFile (no circular reference). SpillFile gc() handles disk cleanup on destruction. No SIOF concerns.

Configuration: New spill_file_part_size_bytes config is mutable (good). New session variables properly forwarded via thrift. One default value issue found (see inline comment).

Incompatible changes: New thrift field IDs (203-208) are additive, backward-compatible.

Parallel code paths: Sort operators correctly adapted to new SpillFile API. Multi-cast data streamer changes look correct.

Observability: Profile counters are present in SpillFileWriter and SpillFileReader. Some old counters (SpillTotalTime, etc.) were removed but replaced by more granular per-file counters — acceptable.

Test coverage: 10 test files with ~150+ test cases. Notable gap: no end-to-end multi-level escalation test for aggregation (level 0→1→2). Hash join multi-level is tested.

Performance: One O(n²) vector erasure pattern found (minor, small n). No other hot-path concerns.

Issues Found: 3 bugs, 1 minor

See inline comments below.

description = {"一旦触发 spill 后,aggregation sink 的 revocable memory 超过该阈值就主动落盘(字节)。默认 64MB。",
"After spill is triggered, aggregation sink will proactively spill when revocable memory "
+ "exceeds this threshold (in bytes). Default is 64GB."})
public long spillAggregationSinkMemLimitBytes = 64L * 1024L * 1024L * 1024L;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bug (Critical): Default value is 64L * 1024L * 1024L * 1024L = 64GB, not 64MB.

All sibling variables use 64MB:

  • spillJoinBuildSinkMemLimitBytes = 64L * 1024L * 1024L (64MB)
  • spillSortSinkMemLimitBytes = 64L * 1024L * 1024L (64MB)
  • spillSortMergeMemLimitBytes = 64L * 1024L * 1024L (64MB)

The thrift default is also 64MB (67108864). The Chinese description says "默认 64MB" but the English says "Default is 64GB".

The BE accessor clamps to [1MB, 4GB], so the effective value becomes 4GB — 62x higher than the intended 64MB. This effectively disables proactive spilling for aggregation sinks, causing excessive memory usage or OOM in aggregation spill scenarios.

Fix:

public long spillAggregationSinkMemLimitBytes = 64L * 1024L * 1024L;

Also fix the English description to say "Default is 64MB."

Comment thread be/src/runtime/runtime_state.h Outdated
int spill_repartition_max_depth() const {
if (_query_options.__isset.spill_repartition_max_depth) {
// Clamp to a reasonable range: [1, 128]
return std::min(_query_options.spill_repartition_max_depth, 128);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bug (Major): Comment says "Clamp to a reasonable range: [1, 128]" but only the upper bound is enforced. If a user sets spill_repartition_max_depth = 0, the function returns 0.

With depth=0, the repartition check if (new_level >= _repartition_max_depth) (in both partitioned_hash_join_probe_operator.cpp:424 and partitioned_aggregation_source_operator.cpp:450) will ALWAYS fail because new_level (which is partition.level + 1 >= 1) is always >= 0. This causes spill repartitioning to immediately return InternalError("exceeded max depth"), making queries fail under memory pressure.

Fix:

return std::max(1, std::min(_query_options.spill_repartition_max_depth, 128));

// Estimate rows that will land in the hash table so we can reserve
// enough for JoinHashTable::first[] + JoinHashTable::next[].
size_t rows = std::max(static_cast<size_t>(state->batch_size()),
static_cast<size_t>(local_state._recovered_build_block->rows()));

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Bug (Edge case): _recovered_build_block can be nullptr when about_to_build is true, causing a null dereference here.

Scenario:

  1. At line 749, _recovered_build_block is reset when a new partition is popped
  2. recover_build_blocks_from_partition is called (line 763) and reads the build file
  3. If ALL blocks from the file happen to be empty (line 359 continue), the function reaches EOS without ever creating _recovered_build_block (lines 366-371 never execute)
  4. build_file is reset (line 379), function returns OK
  5. Back in _pull_from_spill_queue, the function returns OK to the pipeline task
  6. In the NEXT pipeline iteration, get_reserve_mem_size is called. At this point: is_valid()=true, build_finished=falseabout_to_build=true, but _recovered_build_block is null
  7. Line 909: _recovered_build_block->rows()null dereference

While an all-empty-blocks spill file is an edge case, the code should be defensively correct. The simplest fix:

if (about_to_build && local_state._recovered_build_block) {

Or add a null check:

size_t rows = local_state._recovered_build_block
    ? std::max(static_cast<size_t>(state->batch_size()),
               static_cast<size_t>(local_state._recovered_build_block->rows()))
    : static_cast<size_t>(state->batch_size());

while (!local_state._blocks.empty()) {
auto blk = std::move(local_state._blocks.front());
merged_rows += blk.rows();
local_state._blocks.erase(local_state._blocks.begin());

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Minor (Performance): _blocks.erase(_blocks.begin()) on a std::vector is O(n) per call (shifts all elements), making this loop O(n²) overall.

With typical block sizes and 8MB buffer limits, n is usually small (single digits), so the practical impact is minimal. Still, a cleaner pattern would be index-based iteration:

for (size_t i = 0; i < local_state._blocks.size(); ++i) {
    auto blk = std::move(local_state._blocks[i]);
    merged_rows += blk.rows();
    status = _agg_source_operator->merge_with_serialized_key_helper(...);
    RETURN_IF_ERROR(status);
}
local_state._blocks.clear();

// number of spill partitions configured for this operator
size_t _partition_count = 0;
// max repartition depth (configured from session variable in FE)
size_t _repartition_max_depth = SpillRepartitioner::MAX_DEPTH;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Should be type int, as in be/src/exec/operator/partitioned_hash_join_probe_operator.h

!state->is_cancelled()) {
const auto& key = iter.template get_key<typename HashTableType::key_type>();
auto partition_index = Base::_shared_state->get_partition_index(hash_table.hash(key));
auto partition_index = hash_table.hash(key) % parent._partition_count;

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

What if parent._partition_count == 0?

Comment thread be/src/exec/spill/spill_file_reader.cpp Outdated
void SpillFileReader::seek(size_t block_index) {
auto st = _seek_to_block(block_index);
DCHECK(st.ok()) << "SpillFileReader::seek failed, block_index=" << block_index
<< ", error=" << st.to_string();

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

status should not be ignored.

void update_profile(RuntimeProfile* child_profile);
/// Flush the current in-memory hash table by draining it as blocks and routing
/// each block through the repartitioner into the output sub-spill-files.
Status flush_hash_table_to_sub_spill_files(RuntimeState* state);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Declared but no definition, delete it.

/// unread spill files from `remaining_spill_files`, and push resulting sub-partitions into
/// `_partition_queue`. After this call the hash table is reset and
/// `remaining_spill_files` is cleared.
Status flush_and_repartition(RuntimeState* state);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Declared but no definition, delete it.

@mrhhsg

mrhhsg commented Mar 16, 2026

Copy link
Copy Markdown
Member Author

run buildall

@doris-robot

Copy link
Copy Markdown

Cloud UT Coverage Report

Increment line coverage 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 79.24% (1798/2269)
Line Coverage 64.50% (32266/50026)
Region Coverage 65.38% (16150/24702)
Branch Coverage 55.81% (8604/15416)

@hello-stephen

Copy link
Copy Markdown
Contributor

FE UT Coverage Report

Increment line coverage 100.00% (12/12) 🎉
Increment coverage report
Complete coverage report

@doris-robot

Copy link
Copy Markdown
TPC-H: Total hot run time: 27246 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit 96582422b67e2461d89ffff786f88c57ee3f3228, data reload: false

------ Round 1 ----------------------------------
orders	Doris	NULL	NULL	0	0	0	NULL	0	NULL	NULL	2023-12-26 18:27:23	2023-12-26 18:42:55	NULL	utf-8	NULL	NULL	
============================================
q1	17628	4494	4419	4419
q2	q3	10763	829	550	550
q4	4736	360	271	271
q5	8190	1228	1022	1022
q6	242	176	147	147
q7	814	865	676	676
q8	10743	1533	1351	1351
q9	6717	4747	4749	4747
q10	6362	1940	1651	1651
q11	471	252	244	244
q12	702	583	468	468
q13	18082	2965	2163	2163
q14	228	235	218	218
q15	q16	740	767	676	676
q17	727	870	436	436
q18	5989	5462	5324	5324
q19	1110	991	669	669
q20	538	493	379	379
q21	4697	2133	1576	1576
q22	402	344	259	259
Total cold run time: 99881 ms
Total hot run time: 27246 ms

----- Round 2, with runtime_filter_mode=off -----
orders	Doris	NULL	NULL	150000000	42	6422171781	NULL	22778155	NULL	NULL	2023-12-26 18:27:23	2023-12-26 18:42:55	NULL	utf-8	NULL	NULL	
============================================
q1	5723	4739	4615	4615
q2	q3	3915	4395	3859	3859
q4	879	1216	780	780
q5	4096	4429	4373	4373
q6	185	174	145	145
q7	1850	1709	1583	1583
q8	2539	2757	2632	2632
q9	7545	7379	7470	7379
q10	3882	4046	3674	3674
q11	533	483	442	442
q12	527	699	459	459
q13	2855	3175	2292	2292
q14	280	297	276	276
q15	q16	722	763	731	731
q17	1221	1360	1327	1327
q18	7347	6795	6735	6735
q19	915	918	899	899
q20	2079	2243	2002	2002
q21	4072	3484	3349	3349
q22	458	445	380	380
Total cold run time: 51623 ms
Total hot run time: 47932 ms

@github-actions

Copy link
Copy Markdown
Contributor

PR approved by at least one committer and no changes requested.

@github-actions

Copy link
Copy Markdown
Contributor

PR approved by anyone and no changes requested.

@hello-stephen

Copy link
Copy Markdown
Contributor

BE Regression && UT Coverage Report

Increment line coverage 100% (0/0) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 71.54% (26179/36595)
Line Coverage 54.35% (275318/506606)
Region Coverage 51.61% (228751/443235)
Branch Coverage 53.04% (98631/185946)

@hello-stephen

Copy link
Copy Markdown
Contributor

BE Regression && UT Coverage Report

Increment line coverage 100% (0/0) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 71.54% (26179/36595)
Line Coverage 54.34% (275313/506606)
Region Coverage 51.60% (228692/443235)
Branch Coverage 53.04% (98622/185946)

@yiguolei yiguolei closed this Mar 17, 2026
@yiguolei yiguolei reopened this Mar 17, 2026
yiguolei and others added 4 commits March 17, 2026 16:59
[c7b56dbe6a5][2026-03-05][Hu Shenggang] fix ut
[b3a486d14b9][2026-03-05][Hu Shenggang] tiny mod
[6411b66][2026-03-05][yiguolei    ] refactor code and add unit test
[23854d8][2026-03-05][yiguolei    ] add force spill logic in join sink operator
[a631dac][2026-03-04][yiguolei    ] update dir and file meta realtime
[a3fd36e][2026-03-04][yiguolei    ] fix agg profile bug
[c0380ad][2026-03-04][yiguolei    ] change spill file to shared ptr
[05670c5][2026-03-04][yiguolei    ] fix probe revokeable memory size bug
[5595d90][2026-03-03][yiguolei    ] fix compile bug
[36e3004][2026-03-03][yiguolei    ] simplify code
[9094b6b][2026-03-03][yiguolei    ] simplify agg code
[fdb355f][2026-03-03][yiguolei    ] simplify probe code enhancement probe operator
[82e94bb][2026-03-03][yiguolei    ] refactor spill file interface
[7548fe8][2026-03-02][Hu Shenggang] some tiny fix
[2329442][2026-03-01][Hu Shenggang] fix agg revocable mem size
[2c23a07][2026-02-28][Hu Shenggang] disbale distinct streaming agg when spill enabled
[5c27350][2026-02-28][Hu Shenggang] Using spill_buffer_size_bytes as read limit when recovering data
[35f2c55][2026-02-28][Hu Shenggang] [pipeline] Proactively pause query for spill under memory pressure in PipelineTask
[4df2277][2026-02-28][Hu Shenggang] Make spill stream RAII
[6a99170][2026-02-28][Hu Shenggang] Clear revoked data in agg & avoid null pointer in join
[995aeec][2026-02-27][Hu Shenggang] Make repartitioner level-aware
[87cbf87][2026-02-27][yiguolei    ] fix compile
[8a250cd][2026-02-27][yiguolei    ] fix compile
[a92b150][2026-02-27][yiguolei    ] remove some codes
[37df1fd][2026-02-27][yiguolei    ] f
[dd6d1f0][2026-02-27][yiguolei    ] f
[eb8a827][2026-02-27][yiguolei    ] f
[f73c0e9][2026-02-27][yiguolei    ] f
[2c989cb][2026-02-27][yiguolei    ] repartitioner
@mrhhsg

mrhhsg commented Mar 17, 2026

Copy link
Copy Markdown
Member Author

run buildall

@github-actions github-actions Bot removed the approved Indicates a PR has been approved by one committer. label Mar 17, 2026
@hello-stephen

Copy link
Copy Markdown
Contributor

Cloud UT Coverage Report

Increment line coverage 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 79.15% (1796/2269)
Line Coverage 64.48% (32285/50072)
Region Coverage 65.40% (16172/24727)
Branch Coverage 55.81% (8614/15434)

@hello-stephen

Copy link
Copy Markdown
Contributor

FE UT Coverage Report

Increment line coverage 100.00% (12/12) 🎉
Increment coverage report
Complete coverage report

@doris-robot

Copy link
Copy Markdown
TPC-H: Total hot run time: 27155 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpch-tools
Tpch sf100 test result on commit 3991576dc69cb2b4a1e358459c7876a42b148283, data reload: false

------ Round 1 ----------------------------------
orders	Doris	NULL	NULL	0	0	0	NULL	0	NULL	NULL	2023-12-26 18:27:23	2023-12-26 18:42:55	NULL	utf-8	NULL	NULL	
============================================
q1	17627	4562	4332	4332
q2	q3	10648	827	534	534
q4	4675	365	258	258
q5	7558	1219	1008	1008
q6	186	177	156	156
q7	797	865	678	678
q8	9307	1520	1402	1402
q9	4951	4727	4735	4727
q10	6261	1926	1689	1689
q11	486	264	261	261
q12	755	582	481	481
q13	18055	2958	2196	2196
q14	233	233	222	222
q15	q16	741	724	674	674
q17	739	861	454	454
q18	5991	5412	5136	5136
q19	1125	992	645	645
q20	549	498	387	387
q21	4428	1900	1626	1626
q22	451	387	289	289
Total cold run time: 95563 ms
Total hot run time: 27155 ms

----- Round 2, with runtime_filter_mode=off -----
orders	Doris	NULL	NULL	150000000	42	6422171781	NULL	22778155	NULL	NULL	2023-12-26 18:27:23	2023-12-26 18:42:55	NULL	utf-8	NULL	NULL	
============================================
q1	4869	4625	4581	4581
q2	q3	3927	4354	3839	3839
q4	958	1228	844	844
q5	4128	4443	4353	4353
q6	190	187	155	155
q7	1817	1690	1514	1514
q8	2545	2762	2571	2571
q9	7816	7646	7462	7462
q10	3759	4008	3739	3739
q11	546	446	420	420
q12	503	596	451	451
q13	2713	3153	2326	2326
q14	300	311	294	294
q15	q16	751	812	735	735
q17	1181	1329	1321	1321
q18	7178	6859	6580	6580
q19	977	944	948	944
q20	2148	2276	2024	2024
q21	4081	3687	3549	3549
q22	448	433	384	384
Total cold run time: 50835 ms
Total hot run time: 48086 ms

@doris-robot

Copy link
Copy Markdown
TPC-DS: Total hot run time: 170013 ms
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/tpcds-tools
TPC-DS sf100 test result on commit 3991576dc69cb2b4a1e358459c7876a42b148283, data reload: false

query5	4327	653	514	514
query6	352	237	216	216
query7	4234	473	273	273
query8	358	250	246	246
query9	8762	2800	2794	2794
query10	511	415	344	344
query11	7027	5102	4949	4949
query12	192	136	130	130
query13	1279	496	355	355
query14	5768	3831	3582	3582
query14_1	2946	2915	2936	2915
query15	208	202	177	177
query16	1034	426	469	426
query17	1122	744	631	631
query18	2476	463	370	370
query19	224	215	193	193
query20	144	136	132	132
query21	216	140	116	116
query22	13237	14102	14809	14102
query23	16603	15949	15885	15885
query23_1	15610	15372	15396	15372
query24	7194	1656	1248	1248
query24_1	1253	1256	1264	1256
query25	564	485	423	423
query26	1254	279	146	146
query27	2768	491	301	301
query28	4458	1877	1860	1860
query29	855	575	484	484
query30	307	232	194	194
query31	1041	959	878	878
query32	80	74	71	71
query33	524	342	283	283
query34	891	889	550	550
query35	655	696	618	618
query36	1076	1135	977	977
query37	137	97	82	82
query38	2970	2967	2830	2830
query39	859	833	806	806
query39_1	798	790	785	785
query40	230	155	137	137
query41	61	60	60	60
query42	265	262	257	257
query43	259	275	226	226
query44	
query45	198	192	183	183
query46	875	986	626	626
query47	2132	2110	2060	2060
query48	318	346	248	248
query49	636	461	393	393
query50	704	280	220	220
query51	4128	4026	4086	4026
query52	269	267	264	264
query53	293	337	291	291
query54	309	278	269	269
query55	93	92	85	85
query56	321	364	331	331
query57	1960	1774	1807	1774
query58	298	284	279	279
query59	2826	2974	2760	2760
query60	345	344	337	337
query61	158	157	152	152
query62	626	596	542	542
query63	315	288	284	284
query64	5090	1290	1022	1022
query65	
query66	1456	478	361	361
query67	24482	24475	24362	24362
query68	
query69	411	319	291	291
query70	966	977	992	977
query71	357	332	311	311
query72	2849	2723	2574	2574
query73	553	561	334	334
query74	9716	9666	9498	9498
query75	2961	2838	2493	2493
query76	2295	1068	698	698
query77	393	412	334	334
query78	11028	11189	10486	10486
query79	1128	793	597	597
query80	744	680	585	585
query81	499	269	239	239
query82	1369	161	129	129
query83	398	277	261	261
query84	260	124	104	104
query85	892	492	467	467
query86	381	308	293	293
query87	3143	3157	3013	3013
query88	3652	2753	2707	2707
query89	438	380	350	350
query90	1984	193	178	178
query91	170	165	138	138
query92	81	80	73	73
query93	925	874	509	509
query94	463	334	297	297
query95	604	413	328	328
query96	659	541	229	229
query97	2445	2469	2437	2437
query98	246	230	223	223
query99	1021	978	909	909
Total cold run time: 250003 ms
Total hot run time: 170013 ms

@jacktengg jacktengg left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM

@github-actions github-actions Bot added the approved Indicates a PR has been approved by one committer. label Mar 18, 2026
@github-actions

Copy link
Copy Markdown
Contributor

PR approved by at least one committer and no changes requested.

@yiguolei

Copy link
Copy Markdown
Contributor

skip check_coverage

@yiguolei
yiguolei merged commit a0bc783 into apache:master Mar 18, 2026
28 of 31 checks passed
mrhhsg added a commit to mrhhsg/doris that referenced this pull request Mar 24, 2026
This PR adds support for spill repartition in partitioned operators.

When memory is still insufficient after a partition has been spilled,
the spilled partition can be repartitioned into smaller sub-partitions,
which further reduces peak memory usage. The repartitioner is made
level-aware so repartition can be applied recursively when needed.

This PR also integrates the new spill repartition flow into partitioned
hash join and partitioned aggregation, adds force-spill logic in the
join sink operator, refactors spill file interfaces, improves spill
metadata maintenance and memory-pressure handling, and fixes several
spill-related issues such as revocable memory accounting and profile
updates.

Related PR: #xxx

Problem Summary:

None

- Test <!-- At least one of them must be included. -->
    - [ ] Regression test
    - [ ] Unit Test
    - [ ] Manual test (add detailed scripts or steps below)
    - [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
        - [ ] Previous test can cover this change.
        - [ ] No code files have been changed.
        - [ ] Other reason <!-- Add your reason?  -->

- Behavior changed:
    - [ ] No.
    - [ ] Yes. <!-- Explain the behavior change -->

- Does this need documentation?
    - [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
apache/doris-website#1214 -->

- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->

---------

Co-authored-by: yiguolei <guolei@selectdb.com>
yiguolei added a commit that referenced this pull request Mar 25, 2026
…61677)

Pick #61212

This PR adds support for spill repartition in partitioned operators.

When memory is still insufficient after a partition has been spilled,
the spilled partition can be repartitioned into smaller sub-partitions,
which further reduces peak memory usage. The repartitioner is made
level-aware so repartition can be applied recursively when needed.

This PR also integrates the new spill repartition flow into partitioned
hash join and partitioned aggregation, adds force-spill logic in the
join sink operator, refactors spill file interfaces, improves spill
metadata maintenance and memory-pressure handling, and fixes several
spill-related issues such as revocable memory accounting and profile
updates.

Related PR: #xxx

Problem Summary:

None

- Test <!-- At least one of them must be included. -->
    - [ ] Regression test
    - [ ] Unit Test
    - [ ] Manual test (add detailed scripts or steps below)
    - [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
- [ ] Previous test can cover this change. - [ ] No code files have been
changed. - [ ] Other reason <!-- Add your reason? -->

- Behavior changed:
    - [ ] No.
    - [ ] Yes. <!-- Explain the behavior change -->

- Does this need documentation?
    - [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
apache/doris-website#1214 -->

- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->

---------

### What problem does this PR solve?

Issue Number: close #xxx

Related PR: #xxx

Problem Summary:

### Release note

None

### Check List (For Author)

- Test <!-- At least one of them must be included. -->
    - [ ] Regression test
    - [ ] Unit Test
    - [ ] Manual test (add detailed scripts or steps below)
    - [ ] No need to test or manual test. Explain why:
- [ ] This is a refactor/code format and no logic has been changed.
        - [ ] Previous test can cover this change.
        - [ ] No code files have been changed.
        - [ ] Other reason <!-- Add your reason?  -->

- Behavior changed:
    - [ ] No.
    - [ ] Yes. <!-- Explain the behavior change -->

- Does this need documentation?
    - [ ] No.
- [ ] Yes. <!-- Add document PR link here. eg:
apache/doris-website#1214 -->

### Check List (For Reviewer who merge this PR)

- [ ] Confirm the release note
- [ ] Confirm test cases
- [ ] Confirm document
- [ ] Add branch pick label <!-- Add branch pick label that this PR
should merge into -->

Co-authored-by: yiguolei <guolei@selectdb.com>
@morningman morningman mentioned this pull request Jul 3, 2026
76 tasks
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

approved Indicates a PR has been approved by one committer. dev/4.1.0-merged reviewed

Projects

None yet

Development

Successfully merging this pull request may close these issues.

7 participants