Skip to content

[fix](pipeline) Add streaming aggregation local exchange switch - #66222

Merged
Gabriel39 merged 1 commit into
apache:masterfrom
Gabriel39:agent/streaming-agg-local-exchange-master
Jul 29, 2026
Merged

[fix](pipeline) Add streaming aggregation local exchange switch#66222
Gabriel39 merged 1 commit into
apache:masterfrom
Gabriel39:agent/streaming-agg-local-exchange-master

Conversation

@Gabriel39

Copy link
Copy Markdown
Contributor

What

Add enable_local_exchange_before_streaming_agg, defaulting to false, and use it to control whether streaming aggregation requests a local hash exchange.

Why

Using the generic aggregation switch changes streaming aggregation whenever it is enabled. A dedicated switch preserves the existing streaming aggregation distribution by default. The master-specific non-hash child distribution guard remains in place so grouping keys are reshuffled when required.

Validation

  • Not run as requested.

@hello-stephen

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?

@Gabriel39

Copy link
Copy Markdown
Contributor Author

run buildall

@Gabriel39

Copy link
Copy Markdown
Contributor Author

/review

@Gabriel39
Gabriel39 marked this pull request as ready for review July 29, 2026 07:10
@hello-stephen

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

------ Round 1 ----------------------------------
============================================
q1	17733	4094	4021	4021
q2	2100	353	220	220
q3	10174	1454	816	816
q4	4683	489	339	339
q5	7521	1009	569	569
q6	192	178	140	140
q7	781	844	624	624
q8	9389	1555	1560	1555
q9	5645	4401	4367	4367
q10	6758	1734	1499	1499
q11	519	351	334	334
q12	772	595	460	460
q13	18086	3416	2713	2713
q14	273	261	241	241
q15	q16	785	778	704	704
q17	1031	1006	974	974
q18	6877	5829	5610	5610
q19	1302	1224	1126	1126
q20	825	682	546	546
q21	5954	2658	2275	2275
q22	453	348	304	304
Total cold run time: 101853 ms
Total hot run time: 29437 ms

----- Round 2, with runtime_filter_mode=off -----
============================================
q1	4377	4272	4306	4272
q2	293	319	211	211
q3	4583	4960	4407	4407
q4	2038	2134	1344	1344
q5	4398	4354	4296	4296
q6	237	178	124	124
q7	1754	1957	1811	1811
q8	2537	2157	2134	2134
q9	8097	8073	7750	7750
q10	4647	4639	4174	4174
q11	553	410	377	377
q12	761	787	563	563
q13	3344	3663	2875	2875
q14	318	300	291	291
q15	q16	700	715	642	642
q17	1368	1340	1306	1306
q18	7897	7347	7487	7347
q19	1170	1187	1150	1150
q20	2214	2206	1948	1948
q21	5196	4549	4378	4378
q22	518	432	392	392
Total cold run time: 57000 ms
Total hot run time: 51792 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
TPC-DS: Total hot run time: 177318 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 34d1fcab1ad29591e3f561a4110c357748917722, data reload: false

query5	4310	637	491	491
query6	469	226	213	213
query7	4865	593	349	349
query8	346	187	170	170
query9	8780	4094	4095	4094
query10	491	348	312	312
query11	5946	2328	2116	2116
query12	161	100	109	100
query13	1290	609	441	441
query14	6174	5224	4882	4882
query14_1	4241	4236	4202	4202
query15	219	205	182	182
query16	1082	519	481	481
query17	1143	735	590	590
query18	2740	477	347	347
query19	216	195	156	156
query20	111	111	107	107
query21	240	155	138	138
query22	13676	13537	13377	13377
query23	17618	16497	16063	16063
query23_1	16217	16268	16202	16202
query24	7572	1763	1276	1276
query24_1	1319	1300	1283	1283
query25	570	472	386	386
query26	1346	371	215	215
query27	2537	630	386	386
query28	4453	2047	2066	2047
query29	1095	615	508	508
query30	353	264	230	230
query31	1120	1090	995	995
query32	105	65	61	61
query33	527	331	258	258
query34	1174	1142	654	654
query35	785	787	658	658
query36	1029	1029	894	894
query37	156	108	100	100
query38	1886	1712	1678	1678
query39	873	861	857	857
query39_1	831	830	826	826
query40	261	173	142	142
query41	63	66	63	63
query42	97	91	89	89
query43	317	318	277	277
query44	1443	770	770	770
query45	195	192	174	174
query46	1045	1270	738	738
query47	2146	2122	1970	1970
query48	407	405	313	313
query49	580	427	308	308
query50	1012	420	343	343
query51	10894	10540	10363	10363
query52	83	85	73	73
query53	261	279	199	199
query54	287	234	220	220
query55	74	73	66	66
query56	317	302	297	297
query57	1346	1301	1230	1230
query58	286	249	244	244
query59	1614	1667	1463	1463
query60	312	271	256	256
query61	155	146	152	146
query62	539	496	425	425
query63	246	207	205	205
query64	2794	1067	872	872
query65	4714	4665	4618	4618
query66	1757	491	426	426
query67	29228	29181	29024	29024
query68	3010	1619	1001	1001
query69	411	302	270	270
query70	887	821	797	797
query71	377	337	329	329
query72	3049	2649	2409	2409
query73	858	812	417	417
query74	5118	4910	4723	4723
query75	2546	2509	2143	2143
query76	2329	1158	750	750
query77	335	375	281	281
query78	12048	11937	11315	11315
query79	1408	1175	743	743
query80	1284	541	461	461
query81	530	338	287	287
query82	645	158	120	120
query83	373	318	298	298
query84	272	161	136	136
query85	983	607	539	539
query86	416	241	226	226
query87	1811	1827	1746	1746
query88	3741	2823	2813	2813
query89	444	370	332	332
query90	1925	201	198	198
query91	205	197	165	165
query92	62	62	58	58
query93	1624	1509	1029	1029
query94	729	369	310	310
query95	793	497	455	455
query96	1104	805	339	339
query97	2649	2613	2495	2495
query98	211	209	215	209
query99	1088	1123	962	962
Total cold run time: 264271 ms
Total hot run time: 177318 ms

@hello-stephen

Copy link
Copy Markdown
Contributor
ClickBench: Total hot run time: 25.47 s
machine: 'aliyun_ecs.c7a.8xlarge_32C64G'
scripts: https://github.com/apache/doris/tree/master/tools/clickbench-tools
ClickBench test result on commit 34d1fcab1ad29591e3f561a4110c357748917722, data reload: false

query1	0.01	0.00	0.01
query2	0.18	0.09	0.08
query3	0.39	0.24	0.24
query4	1.62	0.24	0.24
query5	0.33	0.32	0.31
query6	1.16	0.67	0.67
query7	0.04	0.01	0.01
query8	0.10	0.08	0.07
query9	0.52	0.38	0.37
query10	0.57	0.58	0.56
query11	0.32	0.18	0.17
query12	0.32	0.18	0.18
query13	0.52	0.51	0.51
query14	0.93	0.91	0.92
query15	0.67	0.58	0.59
query16	0.38	0.39	0.38
query17	1.01	0.99	0.99
query18	0.30	0.27	0.28
query19	1.91	1.77	1.81
query20	0.02	0.02	0.02
query21	15.40	0.37	0.31
query22	4.86	0.13	0.13
query23	15.78	0.48	0.30
query24	2.44	0.60	0.43
query25	0.16	0.10	0.10
query26	0.75	0.27	0.21
query27	0.10	0.10	0.10
query28	3.40	0.95	0.53
query29	12.48	4.26	3.35
query30	0.40	0.28	0.26
query31	2.77	0.63	0.33
query32	3.22	0.59	0.47
query33	3.07	2.96	2.94
query34	15.69	4.06	3.36
query35	3.25	3.25	3.25
query36	0.65	0.53	0.50
query37	0.12	0.10	0.09
query38	0.08	0.07	0.06
query39	0.08	0.06	0.06
query40	0.20	0.18	0.17
query41	0.12	0.08	0.07
query42	0.09	0.06	0.06
query43	0.07	0.06	0.07
Total cold run time: 96.48 s
Total hot run time: 25.47 s

@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.

Request changes.

The dedicated option is wired through Thrift and BE safely, but two required execution paths are incomplete: default FE local-shuffle planning still uses the generic aggregation flag, and observer-to-master forwarding drops the new session value. The variable therefore does not reliably control streaming aggregation.

Critical checkpoint conclusions:

  • Goal and proof: not achieved on the default FE-planned or forwarded-query paths; no new test proves either switch value.
  • Scope and parallel paths: the four-file change is focused but omits the FE AggregationNode planner path and forwarding metadata.
  • Concurrency and lifecycle: no new concurrency, lock, allocation, or lifecycle hazard was found. StreamingAggOperatorX::init() rejects merge/finalize before distribution planning, so removing _needs_finalize from the early return is safe on reachable BE paths.
  • Conditions and data correctness: BE preserves hash-join passthrough, empty-key behavior, serial-child reshuffle behavior, grouping expressions, and the non-hash-child reshuffle guard. Final merge aggregation preserves rows; the defects are broken session semantics plus avoidable or suppressed local-exchange work.
  • Configuration and compatibility: this is a per-query session variable. Thrift ID 228, the false default, BE __isset, and rolling unknown-field behavior are compatible, but both FE and BE planner paths must consume the same setting and follower forwarding must preserve it.
  • Tests and results: no test file changed; existing coverage references only the old option, and the PR states tests were not run. Add direct BE branch tests, FE plan-shape/parity tests, and forwarding coverage.
  • Observability, persistence, transactions, and writes: no new observability is required; reflective persistence and cloning are sound; no transaction, journal, or data-write behavior changes.
  • User focus: no additional focus was provided; the whole PR was reviewed.

tResult.setEnableDistinctStreamingAggregation(enableDistinctStreamingAggregation);
tResult.setEnableStreamingAggHashJoinForcePassthrough(enableStreamingAggHashJoinForcePassthrough);
tResult.setEnableLocalExchangeBeforeAgg(enableLocalExchangeBeforeAgg);
tResult.setEnableLocalExchangeBeforeStreamingAgg(enableLocalExchangeBeforeStreamingAgg);

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.

[P1] Make the FE planner honor this dedicated switch

With enable_local_shuffle_planner=true (the default), NereidsPlanner inserts local exchanges in FE and RuntimeState::plan_local_shuffle() prevents BE from replanning them. AggregationNode.enforceAndDeriveLocalExchange() still gates the useStreamingPreagg branch on enableLocalExchangeBeforeAgg, so new=false/old=true still inserts a hash exchange and new=true/old=false suppresses it—the new variable has no effect on the default path. Please use enableLocalExchangeBeforeStreamingAgg in that branch, retain the old flag for the other aggregation operators, preserve the non-hash-child reshuffle invariant, and add opposite-value FE/BE planner tests plus direct BE branch coverage.

@VarAttrDef.VarAttr(name = ENABLE_LOCAL_EXCHANGE_BEFORE_AGG, fuzzy = true)
public boolean enableLocalExchangeBeforeAgg = true;

@VarAttrDef.VarAttr(name = ENABLE_LOCAL_EXCHANGE_BEFORE_STREAMING_AGG, fuzzy = true)

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.

[P1] Forward this query-affecting variable to the master

getForwardVariables() only includes fields marked needForward or one of the query-result-affect flags, and observer queries sent by FEOpExecutor are planned under a fresh master context. Because this annotation sets none of them, SET enable_local_exchange_before_streaming_agg=true on a follower is omitted and the master keeps false, so the knob silently changes behavior depending on the connected FE. Please add needForward=true and cover follower-to-master restoration.

@github-actions github-actions Bot added the approved Indicates a PR has been approved by one committer. label Jul 29, 2026
@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.00% (5/5) 🎉

Increment coverage report
Complete coverage report

Category Coverage
Function Coverage 75.70% (31957/42218)
Line Coverage 60.36% (356205/590174)
Region Coverage 56.82% (298150/524771)
Branch Coverage 58.22% (134251/230573)

@Gabriel39
Gabriel39 merged commit 2c4b048 into apache:master Jul 29, 2026
33 of 34 checks passed
@hello-stephen

Copy link
Copy Markdown
Contributor

FE Regression Coverage Report

Increment line coverage 2.17% (3/138) 🎉
Increment coverage report
Complete coverage report

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. reviewed

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants