Non-streaming scatter ordering resorts already ordered shard runs

perfloop/vitess · INEFFICIENT ALGORITHM

https://perfloop.ai/t/oss/case_7jxxapcjj5

Verdict

VERIFIED · settled 2026-08-04

What happened: The paired measurements met the required improvement.

Hypothesis

On the accepted Execute path, executeShards calls ExecuteMultiShard and then, on every non-streaming request where len(route.OrderBy)>0 and len(rss)>1, calls route.sort before returning. Route.sort makes a shallow result copy and Comparison.SortResult ultimately calls slices.SortFunc over all result rows. The streaming counterpart takes the same ordered multi-shard condition through route.mergeSort; MergeSort's source keeps one row per sorted stream in its merge tree. ScatterConn.ExecuteMultiShard currently AppendResult's each shard result into one flat result, so the non-streaming path has lost those runs and performs a global comparison sort after every qualifying scatter response. The proposed delta is removal of that all-row post-scatter sort in favor of merging preserved shard runs, not removal of the required ordering semantics. Its cost shape is R combined returned rows versus S shards: a full comparison sort grows roughly with R log R, whereas a k-way merge grows with R log S. It is only a candidate for result-heavy ordered scatters (for example hundreds to thousands of rows across several shards), and no profile establishes that this CPU dominates RPC time. I ran `go test ./go/vt/vtgate/engine -run '^(TestRouteSort|TestMergeSortDescending)$' -count=1`; it passed, confirming existing sort and merger behavior but measuring neither CPU nor production cardinality. A case should benchmark genuine planner-produced scatter ORDER BY plans across at least 2, 8, and 32 shards and increasing combined row counts, collect CPU profiles attributing time to Route.sort/Comparison.Sort, and compare output to the current full sort. The correctness differential must cover ascending and descending keys, duplicate keys, collation/weight-string keys, empty and partially failed shards, warnings, and truncation, while verifying the planner's shard-local ordering guarantee before enabling the merge path.

Change to test: Add a planner-marked ordered non-streaming scatter path that preserves each shard's result run and materializes it with a heap/k-way merge analogous to MergeSort. Keep the existing flat-result full-sort fallback unless the route is proven to send a compatible local ORDER BY to every shard.

Where it lives

perfloop/vitess · go/vt/vtgate/grpcvtgateservice/server.go

Evidence

ordered non-streaming scatter, 2 shards, 1024 rows · 10 sample pairs

metric baseline candidate paired median change confidence range required result
ns/op 278665 142959 −49% (−136590) −142512 to −124407 < 0 PASSED

ordered non-streaming scatter, 2 shards, 8192 rows · 10 sample pairs

metric baseline candidate paired median change confidence range required result
ns/op 2315936 829756 −64.4% (−1491880) −1738104 to −1384960 < 0 PASSED

ordered non-streaming scatter, 8 shards, 1024 rows · 10 sample pairs

metric baseline candidate paired median change confidence range required result
ns/op 425939 253278 −39.6% (−168622) −179042 to −160483 < 0 PASSED

ordered non-streaming scatter, 8 shards, 8192 rows · 10 sample pairs

metric baseline candidate paired median change confidence range required result
ns/op 3147273 1247165 −60.5% (−1903363) −2033476 to −1782870 < 0 PASSED

ordered non-streaming scatter, 32 shards, 1024 rows · 10 sample pairs

metric baseline candidate paired median change confidence range required result
ns/op 578104 437193 −24.2% (−140140) −151490 to −125243 < 0 PASSED

ordered non-streaming scatter, 32 shards, 8192 rows · 10 sample pairs

metric baseline candidate paired median change confidence range required result
ns/op 3166109 1796222 −44.3% (−1401791) −1513499 to −1345731 < 0 PASSED

Checks: 8 of 8 passed. Verification: no defect found.

Timeline