Repository navigation
Benchmark harness for the streaming selection ORDER BY combine - #1
Draft
rohityadav1993 wants to merge 5 commits into
Draft
rohityadav1993 wants to merge 5 commits into
rohityadav1993 wants to merge 5 commits into
Conversation
rohityadav1993
force-pushed
the
oss/pr1-streaming-selection-combine
branch
from
September 26, 2026 18:41
f3f2db3 to
d605e80
Compare
rohityadav1993
force-pushed
the
rohity/pr1-arm4-memory-crossover-bench
branch
from
September 27, 2026 03:28
97151cd to
968d36f
Compare
rohityadav1993
force-pushed
the
oss/pr1-streaming-selection-combine
branch
from
September 28, 2026 17:47
d605e80 to
ff800d6
Compare
1 task
rohityadav1993
force-pushed
the
oss/pr1-streaming-selection-combine
branch
from
September 29, 2026 07:07
ff800d6 to
8271cf9
Compare
rohityadav1993
force-pushed
the
rohity/pr1-arm4-memory-crossover-bench
branch
from
September 29, 2026 08:48
968d36f to
0143a87
Compare
Adds the JMH harness used to measure the streaming selection ORDER BY
combine against the existing MinMaxValueBased combine. Confined to
pinot-perf; no production code is touched.
SortedMergeFixture builds and caches the segment fixture, with
DISJOINT/PARTIAL/FULL overlap and a
keyCardinality knob giving runs of K rows
per distinct tsCol value
SortedMergeDriver builds the query context and runs one
configuration on either arm
SortedMergeMemoryProbe the correctness gate: runs both arms on
every configuration and requires agreement
SortedMergeCpuMeter per-thread CPU accounting
BenchmarkSortedMergeMemoryCrossover
the JMH entry point
sorted-merge-bisect-xmx.sh bisecting -Xmx search for the minimum heap
a configuration survives on
The query shape is selectable: ORDER BY tsCol, valCol keeps a total
order, and ORDER BY tsCol alone drops the tiebreaker so the leaf's scan
block tightens to the limit. Under the latter the order is not total, so
the gate compares a commutative digest over tsCol values rather than
whole rows: which tied rows come back is free, but how many rows carry
each value is not.
SortedMergeGateTest adds a test source root to pinot-perf and covers the
gate's comparison primitives with negative controls, since running the
gate by hand only ever exercises the pass path.
This is the harness as run for the published results, committed before
any follow-up changes so the numbers stay reproducible against it.
The -Xmx bisection forks a fresh JVM per attempt, so two attempts can race to build the same fixture key. Nothing serialized them: the in-JVM SEGMENT_CACHE only orders callers within one process. Guard each key with a FileLock held across validation and build, and route destroyAll() and purgeAll() through the same per-key lock instead of deleting the base directory wholesale. Without that the teardown path bypassed the lock entirely and could delete a directory another process was still building. FileLock is scoped to the JVM rather than the thread, so a teardown colliding with a build inside one JVM throws OverlappingFileLockException. That extends IllegalStateException, not IOException, so it escaped the existing handler -- out of a shutdown hook, abandoning the rest of the cleanup loop. Both sides now handle it: the teardown skips the key being built and continues, and a build that collides with a teardown fails with an explanation rather than a bare exception. SortedMergeFixtureTest covers the teardown half, which is the half that regressed. The build half needs two JVMs and is not reachable from a unit test. Every test points the fixture at a temporary directory and asserts the override took effect before anything destructive runs, because purgeAll() deletes every key under the base directory and would otherwise destroy the fixtures that published results were measured against. The base directory override exists only for that guard and is unset on every non-test path, so the on-disk cache built by previous runs is unchanged.
gate() decided whether two arms agreed and formatted the result line in one body, so the decision could only be tested by running a full merge on both arms. The decision is the part worth pinning: under TS_ONLY the ORDER BY is not a total order, ties can straddle the LIMIT boundary, and two correct operators may legitimately return different rows -- so the comparison is deliberately weaker there than under TS_VAL, and nothing tested that the relaxation is exactly as wide as intended. cellVerdict() now returns the failure reason or null, and gate() keeps only the formatting. The printf output is unchanged. SortedMergeGateTest gains nine cases over the verdict itself, including the pair that fixes the shape of the relaxation: the same two arms that TS_ONLY must accept because they picked different tied rows, TS_VAL must reject. Also corrects a comment on the TS_ONLY branch claiming the full-row digest was never computed. gate() computes it for both arms and derives armsPickedDifferentTiedRows from it. The real reason that digest cannot decide the verdict is that two correct arms may differ on which tied rows fall inside the LIMIT.
TS_VAL and TS_ONLY are not two spellings of one query. TS_ONLY makes sortedColumnsPrefixSize equal the ORDER BY length, which tightens the leaf's maxDocsPerCall to min(limit + offset, 10000); TS_VAL leaves it pinned at the block size. That difference is the axis two suites were run to measure, and it rested on an untested assumption about what buildQueryContext emits. SortedMergeDriverTest asserts the ORDER BY of each shape, that the select list is identical and three wide under both, that the no-arg overload still means TS_VAL, and that the direction is ascending. No production behavior changes.
The -Xmx search measures the smallest heap a configuration completes on. A JVM fails at live set plus GC headroom plus fragmentation, and headroom scales with allocation rate, which is the one thing the two arms differ in by about 5x. So the figure is a provisioning requirement, not a retained set, and a ratio read off the grid is not a ratio of retention. The script now states that at the top, and the output is named accordingly. Six defects found in review of the v3 run: - The floor was never probed. A cell that already fit in LOW_MB converged to 95 MB and printed it, a heap size that was never tried. It now probes LOW_MB first and reports AT_OR_BELOW_LOW. - Neither the collector nor -Xms was pinned, so the figures were specific to the JDK defaults that produced them and could not be read alongside the 4 GB JMH runs. Now -Xms = -Xmx with G1 explicit. - Survival near the boundary is a sigmoid, and requiring N consecutive survivals converges on an unstated quantile that depends on the draws it got. Each cell is now measured REPLICATES times, every replicate is its own row, and the result is the half-open bracket (low_mb, high_mb] where both bounds were actually run. - REPEATS_ON was 1, justified by identical segments-processed counts. That shows the workload is deterministic; what varies at the OOM boundary is GC and allocation timing, which is not. Now 2. - No timeout, so a JVM thrashing without tripping the GC overhead limit hung the queue. Now wrapped, with a distinct outcome. The kill backstop and the kernel OOM killer share exit 137, so they are separated by elapsed time: recording a wedged JVM as an OOM would have manufactured evidence for the claim under test. - No per-attempt record. Roughly 900 launches collapsed into 48 rows, so a non-monotone response was unrecoverable afterwards. Each attempt is now logged with its outcome, duration, and segments processed. The tolerance is relative, since a flat 32 MB was half the floor and 0.8% of the ceiling. TOLERANCE_MB is refused rather than ignored, because a knob someone deliberately set must not be silently dropped. sorted-merge-bisect-xmx-selftest.sh stubs java and injects the failure modes a real run cannot produce on demand. The grid axes became overridable so it can do that in seconds rather than hours; the banner prints the grid actually used, so a run says what it measured. The CSV schema changed and v3 output no longer parses. The peak live set the PR's claim is really about still needs a separate instrument, which this is deliberately not.
rohityadav1993
force-pushed
the
rohity/pr1-arm4-memory-crossover-bench
branch
from
October 5, 2026 18:43
0143a87 to
e38986b
Compare
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
This PR is the benchmark harness and the measurement report for #19120. apache#19120 adds the operator. This PR adds no operator code. It stacks on
oss/pr1-streaming-selection-combineand changespinot-perfonly.TL;DR
ON= the streaming combine from apache#19120.OFF= the existing MinMax combine (10 threads). 100 segments × 50,000 rows, ASC.LIMIT≥ 100,000,ONneeds 2x to 14x less heap. AtLIMIT1,000,000,OFFneeds about 930 MB;ONneeds less than 64 MB (DISJOINT) to 237 MB (FULL).ONstays flat asLIMITgrows.LIMIT≥ 100,000,ONis faster: 0.07 to 0.75 ofOFF(single-shot timing). One 4 GB average-time run showsONup to 1.15x slower atLIMIT100,000. At DISJOINT and PARTIAL,ONis also faster at most smallerLIMITvalues.LIMIT< 100,000:ONis 1.8x to 7x slower (up to 10x with a one-columnORDER BY).ONreads all tied segments on one thread;OFFuses 10. At FULL /LIMIT10,000,ONneeds about 1.27x more heap.ONuses less CPU in 17 of 18 cases.6df6b87de5) took one case from 77.8x slower to 0.82x.Details, methods and raw sources follow.
Which PR1 code was measured. The rerun with the fix used operator code identical to
6df6b87de5.git diff 4f79f49f78 6df6b87de5is empty on the combine, the leaf,SelectionPlanNode,CombinePlanNode,InstancePlanMakerImplV2andQueryContext. Four suites ran on code before the fix. The FULL K=1 rerun ran on642af405b8(deferral only, no frontier fix). The rerun with the fix ran on4f79f49f78.The earlier figures still apply. Tie deferral is off under
TS_VAL. Frontier pruning can only keep closed a segment that would have opened before. At FULL underTS_VALit cannot fire, and the FULL rerun matched the earlier figures within noise. At DISJOINT and PARTIAL the earlier figures can only be pessimistic forON.OFFis not changed by the fix.The four later commits cannot change a result:
19235491468271cf9b0e3b92607347ba2b193b32MutableSegment. All fixture segments are immutable. A rerun is not necessary.What this adds
All paths are under
pinot-perf. The Java files are insrc/main/java/org/apache/pinot/perf/sortedmerge/. The scripts are insrc/main/scripts/.SortedMergeFixtureSortedMergeDriverQueryContextand the plan nodes. RunsCombinePlanNodeand drains the blocks. Asserts the operator class (G1). Holds the digest helpers and theOrderByshapes.BenchmarkSortedMergeMemoryCrossovermergeTopK). Reports wall time, CPU through@AuxCounters, and allocation through-prof gc.SortedMergeCpuMeter-Darm4.cpuTime=true).SortedMergeMemoryProbegateis the correctness gate (G2).rundoes one merge per JVM for the heap search.sorted-merge-bisect-xmx.sh-XmxoverProbe run.sorted-merge-bisect-xmx-selftest.shjava. 14 cases, all pass.SortedMergeGateTest(15),SortedMergeDriverTest(5),SortedMergeFixtureTest(3). Total 23.The rest of this description is the measurement report.
How to run
Run all commands from the repository root. Use the shaded
benchmarks.jar. Themain()in the benchmark class ignores its arguments.The flag
-Ddevelocity.cache.local.enabled=falsestops the build cache from reporting success when no tests ran. Run the heap search from the repository root, because it uses relative paths.1. Objective
apache#19120 adds
StreamingSelectionOrderByCombineOperator. Query optionsortedSelectionMergeMode=ONselects it. This report calls it the streaming combine (ON). The existingMinMaxValueBasedSelectionOrderByCombineOperatoris the MinMax combine (OFF). The target isORDER BYplusLIMITon a leading column that is sorted in each segment.The benchmark answers two questions. Memory: does
ONneed less heap thanOFFasLIMITgrows? Crossover: at whichLIMITand overlap doesONstop costing more wall time and CPU? Memory and latency need separate instruments, because allocation and elapsed time do not show live heap.2. Definitions
tsColrange contains the global midpoint.DISJOINT=1,PARTIAL≈10,FULL=100. Depth is not measured at theLIMITboundary, so it is exact only at FULL.OFF,ONOFF: MinMax combine, work on 10 pool threads.ON: streaming combine, merge on the calling thread.TS_VALORDER BY tsCol, valCol.valColis unique, so the order is total.TS_ONLYORDER BY tsCol. Ties can cross theLIMITboundary.tsCol.6df6b87de5). See section 8.ON÷OFFwall time (below 1:ONfaster).ONminusOFFCPU in ms. A/A control:OFFcases rerun later, to measure noise.segmentsProcessedlog line).-Xmxat which a run completes.(low, high]: failed atlow, survived athigh. An upper bound that includes GC headroom.3. Setup
Pipeline:
SortedMergeFixturefeedsSortedMergeDriver. The driver feeds JMH,SortedMergeCpuMeterandProbe run.Probe runfeeds the heap search.Probe gatecompares both arms in a separate step.LIMIT10 to 1,000,000 in JMH, 1,000 to 1,000,000 in the heap search.tsColLONG (sorted, dictionary).valColINT (unique).payloadColSTRING (raw, 54 to 59 characters).maxExecutionThreads=10.sortedSelectionMergeBlockSize=10,000. The block ismin(limit + offset, 10,000)rows underTS_ONLY.rohity-pinot-oss1), 96 cores, 377 GB RAM, Temurin 25.0.3, G1. Hot page cache.4. Results
Memory:
OFF÷ONminimum surviving heapTS_VAL, K=1. Each figure is a lower bound: the heap whereOFFfailed divided by the heap whereONsurvived. "≥" meansONsurvived at the 64 MB floor, so the true ratio is not resolved.LIMIT1,000LIMIT10,000LIMIT100,000LIMIT1,000,000ONneeds ≥ 1.27x moreOFFgrows withLIMIT.ONdoes not. AtLIMIT1,000,000,OFFneeds about 914 to 946 MB at every overlap. Section 7 has the brackets.LIMIT10,000 is the one adverse case. UnderTS_VAL, all 100 segments with the same minimum must open.LIMIT≥ 100,000,OFFallocates about 4.5x more thanONat FULL/100k and DISJOINT/1M. At FULL/1M it allocates about 2.7x more (1.62 GB against 606 MB).Wall clock:
ON÷OFFBelow 1 means
ONis faster. Bold meansONis slower beyond noise.TS_VAL, K=1:TS_ONLY, FULL, with the fix. A blank cell has no measurement. The 0.64, 0.82 and 0.90 entries have error bars that include 1. Read them as parity.ONis faster at everyLIMITinSingleShotTime(0.50 at 100k, 0.32 at 1M). The 4 GBavgtrun (raw/jmh-two-column-grid/r4-heap4g.json) showsON1.10x slower at DISJOINT/100k and 1.15x slower at PARTIAL/100k. At PARTIAL,ONis slower only at 10,000 in the 16 GB runs. At FULL,ONis faster only from 100,000 (up to 9x underTS_VAL, 13x underTS_ONLY).ONis slower. It reads every tied segment on one thread.OFFuses 10.ONis at parity or faster, except K=100/LIMIT1,000 (1.76). K=10,000/LIMIT10,000 went from 77.8x to 0.82x.LIMIT≥ 100 stays slow with the fix. The output needs key 0 from every segment, so all 100 open.CPU and GC
ONuses less CPU in 17 of 18TS_VALcases (12AverageTime, 6SingleShotTime). The exception is FULL/LIMIT10,000, at +25.6 ms. AtLIMIT≥ 100,000 the CPU ratio is 0.09 to 0.20.OFFand 0.57% forON(section 7).5. Harness
Instruments
@Threads(1). Allocation:-prof gc(B/op), in a separate run.SortedMergeCpuMetersumsThreadMXBeanCPU over the calling thread and its pool threads. GC and JIT threads are excluded, soOFFis understated.LIMIT. Use the CPU difference there.CombinePlanNodewith aResultsBlockStreamer. The mode is explicitlyONorOFF, neverAUTO, and it also selects the leaf. The streamer is a no-op. The executor is a fixed pool of 10 daemon threads.JMH settings
AverageTime, msLIMIT10 to 10,000:-bm avgt.LIMIT100,000 and 1,000,000:-bm ss(SingleShotTime).avgt. 10 forss. 1 for the allocation run. The rerun with the fix used 1 fork and 3 iterations, except K=10,000/LIMIT100,000 (ss, 10 forks).avgt: 3 × 2 s (-wi 3 -w 2s).ss: 5 (-wi 5).avgt: 5 × 2 s (-i 5 -r 2s).ss: 1 (-i 1).The allocation run uses
-f 1 -wi 1 -i 3 -r 5s -w 5s -prof gc. The 4 GB run uses-Xms4g -Xmx4g -XX:+UseG1GC -prof gc. The default grid has 72 configs (the class Javadoc says 36).OFFis bimodal, because its segment count races across 10 threads.Guards
ON, MinMax forOFF). Runs in JMH setup, everyrunand the gate.Probe gatestep. Not in JMH.The gate covers 3 overlaps × 6
LIMITvalues,ONandOFF. Row counts must be equal. It also requires:TS_VAL: equal full-row digest. AtLIMIT≤ 10,000, also equal row multiset andtsColsorted.TS_ONLY: equal digest oftsColvalues only. AtLIMIT≤ 10,000, also equaltsColmultiset andtsColsorted.Under
TS_ONLYthe arms can return different tied rows. The gate proves nothing aboutvalColor payload there. AboveLIMIT10,000 it proves nothing about order.SortedMergeGateTesthas negative controls.Fixture cache and locking
/tmp/pinot-arm4-sorted-merge/<key>on disk. The harness reuses a fixture only with a.arm4-completemarker, after it reloads and recounts every segment.<key>.lockfile lock covers validate-or-build, because each heap-search attempt is a new JVM. A build takes 15 to 25 s. The run drivers also usedflock /tmp/arm4-bench.lock, which is not in the committed code.Heap search
flowchart TD A["Probe the floor (64 MB)"] --> B{"Survives?"} B -->|yes| R1["AT_OR_BELOW_LOW"] B -->|no| C["Probe the ceiling (4,096 MB)"] C --> D{"Survives?"} D -->|no| R2["ABOVE_HIGH"] D -->|yes| E["Pick the midpoint of (low, high]"] E --> F["Run N consecutive probes<br/>(N=3 for OFF, N=2 for ON)"] F --> G{"All N survive?"} G -->|yes| H["high = midpoint"] G -->|no| I["low = midpoint"] H --> J{"high - low <= max(16 MB, 5% of high)?"} I --> J J -->|no| E J -->|yes| K["Output bracket (low, high]"]OFFuses N=3, because its retention depends on thread scheduling. This raises theOFFminimum, so it is conservative against the PR.-Xms=-Xmx, G1,-XX:+ExitOnOutOfMemoryError. Fixtures are prebuilt at-Xmx8g.Exit 3 (or 137/143 before the timeout) is an OOM. Exit 124, or 137 after the 900 s timeout, is
ERROR_TIMEOUT. No[arm4] builtorreusedline isERROR_FIXTURE_NOT_READY. Any other exit isERROR_PROBE_FAILED.An
ERROR_*outcome never steers the search, so a harness bug is never recorded as "needs more heap". The harness abandons that case. Every attempt goes to an audit CSV. Each case has 2 replicates. The script hardcodesTS_VAL.Suites
TS_VALwall, CPU, 4 GB heap, A/ATS_ONLYagainstTS_VALOFFbaselines, before-fixON642af405b8)ONonly4f79f49f78ONfigures for cases the fix changes6. Fixtures
All fixtures come from
SortedMergeFixture. ThetsColof each segment starts atbase. In each segmenttsColascends, so the leaf scans forward.baseof segment ii × 50,000i × 5,000tsCol = (base + j) / K. Each key covers K rows. At FULL, every segment starts with the same K rows of key 0. This is the tie in section 8.Measured depth is 1 / 10 / 100 at every K, except PARTIAL at K=10,000 (11).
7. Details
Heap search results (
TS_VAL, MB)Each cell is one heap-search result.
≤64means survival at the floor. AtLIMIT≤ 10,000 every cell is ≤64 except FULL/10,000. If the two replicates disagree, the bracket spans both.LIMITOFFK=1ONK=1OFFK=10,000ONK=10,000LIMITand K. At FULL,ONneeds at most about 205 to 237 MB at everyLIMITfrom 10,000 to 1,000,000. At FULL all 100 segments open, each with a live 10,000-row block.OFFholds a small capped queue for each.OFFatLIMIT100,000, and differ by one adjacent bracket.CPU and GC detail
TS_VAL, K=1,LIMIT≤ 10,000: -0.09 to -7.8 ms at DISJOINT and PARTIAL. At FULL: -2.0, -1.9, -4.1, then +25.6 ms at 10,000.LIMIT≥ 100,000:OFF247 MB to 1.62 GB for each operation.ON34 to 606 MB. AtLIMIT1,000:OFF1.8 MB,ON0.98 MB.ONtook 135.7, 123.6 and 96.3 ms in three runs. Do not quote it alone.Shape, K and precision
TS_VALdoes not get worse forONas K rises. FULL/LIMIT10,000: 127.9, 120.7 and 105.5 ms at K=1, 100, 10,000.TS_ONLYat DISJOINT and PARTIAL is within noise ofTS_VAL.scoreError/score atLIMIT≤ 10,000: 11.8%. CV atLIMIT≥ 100,000: 3.3% to 49.2%. A/A CPU: 2.65% median, 8.74% worst. The A/A control covers 12OFFcases atLIMIT≤ 10,000, 90 minutes later.8. Finding fed back to apache#19120
How it surfaced
The heap search found the one case where
ONneeds more heap thanOFF: FULL,LIMIT10,000,ON(221,237] againstOFF(158,174]. At FULL every segment has minimumtsCol0. The tie defeats the lazy opening of segments. All 100 open, and each pins a 10,000-row block.The key-cardinality suite showed the cost under
TS_ONLY: at FULL, K=10,000,LIMIT10,000,OFFopened 10 segments andONopened 100 (78x slower).The issue and the fix
ONvisits segments in order of their minimum. It keeps a segment closed only while that minimum sorts strictly after the merge frontier. A segment whose minimum ties the frontier always opened.OFFalready skips such a segment underTS_ONLY.A second flaw hid behind the first. The frontier is the smaller of two heads (the leader's and the heap top's), but the check used the larger. This kept K=100 at 100 segments.
The fix has two parts. It defers tied segments under a single-column ORDER BY only. With two or more columns, a tie on the first column can hide an earlier row on the second. It also takes the frontier as the smaller of the two heads, for all shapes.
A deferred segment is postponed, never skipped. It opens when no other segment can supply a row. Its remaining rows tie, so they are interchangeable under a single-column sort key.
Effect
TS_ONLY, FULL,ON. "Before" is from the key-cardinality suite (K=1/LIMIT10: query-shape suite). The wall ratio with the fix is in section 4. Allocation is in MB per operation.LIMITON/OFFallocation before / with fixLIMIT10,000: segments 100 to 1, wall 121.6 to 1.29 ms (94x),ONallocation 285.06 to 7.14 MB (OFF25.46).LIMIT10,000 is unchanged. The output is the 100 key-0 rows of each segment, so all 100 must open.OFFalso opens all 100.OFFbaselines predate a rebase onto newer master. They matched within noise in the FULL K=1 rerun.Under
TS_VAL,ONstill opens all tied segments. The adverse case in section 4 remains.9. Limits and open items
Limits
TS_ONLY. The expected 1-segment gain at K=10,000 is unmeasured.AUTOratio (0.8).maxExecutionThreads. Segment count and rows. Payload width. Concurrent queries.Open items
ORDER_BYknob. Run a before-and-after pair at FULL underTS_ONLY. Rerun withLOW_MB=16.TS_VAL. It was not rerun. It can only helpON.depth(L). The fixture measures depth at the midpoint, exact only at FULL. It does not compute depth at theLIMITboundary.OFFallocation at FULL but not inON. The cause is unexplained.