wirybeaver opened a new issue, #19727:
URL: https://github.com/apache/pinot/issues/19727

   # MSE hash aggregation disk spill: roadmap and design
   
   ## Purpose and scope
   
   This issue is the design and delivery roadmap for the MSE keyed 
hash-aggregation spill work originally implemented in #19469. The monolithic 
review is being decomposed by functionality; this issue keeps the complete 
design in one place so each implementation PR can describe only its own change.
   
   Source snapshot: #19469 at `2e40e9402d3069aa86e22cf7e6858b904795b57f` 
(+2,045/-41 across 18 files). This is a description of the reviewed proposal, 
not a claim that spill is merged or enabled in a release.
   
   - In scope: opt-in local-disk spill for keyed MSE aggregation, 
intermediate-state round trips, partition-at-a-time merge, resource limits, 
cleanup, and observability.
   - Out of scope for this first version: SSE spill/parallel combine, global 
aggregation, leaf-final-result and group-trim modes, remote spill, asynchronous 
I/O, recursive repartitioning, and a retained-memory-byte trigger.
   - #19667 is the independent `ANY_VALUE(BYTES)` correctness prerequisite and 
is still open. Its latest type-consistency fix must be included; the older copy 
in #19469 is not authoritative.
   - #19666 remains the follow-up for a memory-based trigger; do not close it 
with this group-count-based implementation.
   - This work only partially addresses #12080. The SSE and parallel-combine 
requests remain open.
   
   ## Delivery and review rules
   
   1. Split by invariant and behavior, not by file or an arbitrary line count. 
Each PR carries the tests for the behavior it introduces.
   2. Keep production spilling unreachable until disk budgets, cleanup, and the 
complete operator lifecycle are present. New library seams may be directly 
tested before activation.
   3. Every child must compile and pass its focused tests on its declared 
parent without needing a later child.
   4. Use the personal fork for stacked review: each child PR targets its 
direct predecessor branch, not fork `master`, so the diff is incremental. 
Promote to `apache/pinot` only after its prerequisites are available on 
upstream `master`.
   5. Preserve #19469's reviewed source SHA and review discussions before 
replacing its code with a narrow slice. Do not silently drop existing 
correctness or regression coverage.
   
   ## Acceptance criteria
   
   - Below resource/group limits, spilled and non-spilled execution have the 
same group identity and aggregate results; intermediate states are not 
finalized twice.
   - `numGroupsLimit` continues to bound each input/merge hash table and the 
final result admission policy. Exceeding it preserves the documented 
truncation/error behavior.
   - Spill location and query/server-JVM disk budgets are server-owned; client 
metadata cannot bypass them.
   - Success, failure, cancellation, early termination, and close clean up 
owned spill files; startup handles files left by a crashed server instance.
   - Unsupported modes report why spill was skipped, and malformed options are 
rejected independently of eligibility.
   - New statistics cannot make the query fail when a peer cannot decode them.
   - The group-count trigger is a v1 memory proxy, not a strict heap-byte 
bound. Disk quotas are not memory budgets. Keep internal disk-quota metadata 
distinct from the future memory-trigger option in #19666.
   
   ## Implementation breakdown
   
   Sizes are approximate planning budgets for additions plus deletions, 
including focused tests. Storage budgets are roughly 330–400 lines; the 
operator slice is 500–650. Prefer complete behavior and its tests over 
artificial 300-line cuts. Do not pad tiny correctness fixes. The one larger 
operator slice keeps result production and cleanup atomic; use shared test 
fixtures instead of deleting coverage or creating a tests-only tail.
   
   | PR | Phase | What it adds | Est. Lines | Depends On | Target / branch |
   |----|-------|-------------|------------|------------|-----------------|
   | Existing | BYTES correctness | Keep aggregation and serde BYTES values 
consistently `ByteArray` | Existing PR | None | `apache/pinot#19667` |
   | 1 | Executor state contract | Export/remerge intermediate rows; capacity 
and trigger seams | ~300 | None | Retain `apache/pinot#19469`; 
`pinot-agg-spill` |
   | 2 | Stats compatibility | Unknown-key decode failures degrade instead of 
failing queries | ~50 | None logically | Fork `review/agg-spill/02-stat-compat` 
|
   | 3 | Round-trip contracts | Honest ANY_VALUE intermediate schema; 
content-based single-array-key grouping | ~180 | #19667; PR 1 for round-trip 
tests | Fork `review/agg-spill/03-value-key-contracts` |
   | 4 | Framed file storage | Append/read partition records; file ownership 
and cleanup | ~360 | Baseline APIs | Fork `review/agg-spill/04-framed-storage` |
   | 5 | Partitioned buffering | Deep hashing, bounded records/buffers, writer 
reuse | ~330 | PR 4 | Fork `review/agg-spill/05-partition-buffering` |
   | 6 | Disk safeguards/config | Query/JVM budgets, owned root, orphan sweep, 
trusted options | ~370 | PR 5 | Fork `review/agg-spill/06-disk-safeguards` |
   | 7 | Operator lifecycle | Complete consume/spill/restore/output/cleanup 
path and spill stats | ~600 | PRs 1, 2, 3, 6 | Fork 
`review/agg-spill/07-operator-lifecycle` |
   | 8 | Server activation | Enable the trusted gate; multi-server 
result/stats/cleanup proof | ~180 | PR 7 | Fork 
`review/agg-spill/08-activation-e2e` |
   
   ### Execution checklist
   
   - [ ] Preserve the full #19469 source as an archive branch before any 
rewrite.
   - [ ] **PR 1 / #19469:** export/remerge intermediate rows and keep the 
existing group-limit backstop. Add executor-only SUM/AVG, row-heap/serialized, 
layout-remapping, and finalization tests. No disk manager, QueryRunner 
activation, or new stat keys; it does not depend on #19667.
   - [ ] **PR 2:** harden `StatMap` and `MultiStageStatsTreeDecoder` with 
unknown-ordinal/decode-failure regressions. Do not reintroduce the 
SAFE/homogeneous-cluster spill gate.
   - [ ] **PR 3:** fix ANY_VALUE intermediate type reporting and 
single-array-key identity, with actual serde round trips. Use the latest #19667 
changes instead of the stale BYTES fix embedded in the original monolith.
   - [ ] **PR 4:** explicit-partition, length-prefixed file storage and 
idempotent cleanup. Test multiple records, truncated framing, consumer 
failures, and deletion retry. No production caller yet.
   - [ ] **PR 5:** iterator-driven hash partitioning, bounded batching, and 
writer reuse. Test 1,024-row records, more than eight partitions, the aggregate 
buffer limit, and writer closure. Keep this storage API unexposed to queries.
   - [ ] **PR 6:** disk reservation/release, instance-root lifecycle and 
startup sweep, validated options, and server-owned metadata. Test sibling 
operators, different queries, reuse after close, spoofed metadata, and 
unrelated-file preservation. Force the internal enable metadata to `false` 
until final activation.
   - [ ] **PR 7:** complete consume/spill/restore/output/cleanup and stats 
emission. Test real spill across repeated keys, supported modes/types/filters, 
group limits, and all terminal paths. Normal QueryRunner still forces the gate 
off; operator tests can use direct test contexts.
   - [ ] **PR 8:** activate only via the real default-off server gate, and 
prove SUM/AVG and DISTINCTCOUNT across two query servers with result equality, 
positive stats, and directory cleanup. Keep both request and stage metadata 
unable to bypass server controls.
   
   The operator slice is intentionally the largest: consumption, restoration, 
and cleanup are one correctness contract. Reduce repeated test setup with 
focused fixtures, not by dropping cases or postponing them to a tests-only PR. 
All sizes are estimates to check against the actual parent-to-head diff during 
extraction, not already-validated child diffs.
   
   ### Logical dependencies
   
   ```mermaid
   flowchart LR
       A[Executor state] --> D[Operator lifecycle]
       B[Value and key contracts] --> D
       C[Guarded partition storage] --> D
       D --> E[Server activation]
       F[Stats compatibility] --> D
   ```
   
   ### Fork review and upstream promotion
   
   1. Keep #19469 open and narrow it to PR 1 only after the narrow code and 
tests exist; do not label the current monolithic diff as already split.
   2. Use `review/agg-spill/base` as a staging baseline in `wirybeaver/pinot`. 
If #19667 or narrowed #19469 are still pending, include their exact head 
commits in this baseline and explicitly record them; do not duplicate those 
patches in child diffs.
   3. Open the remaining PRs in `wirybeaver/pinot`. The first targets the 
staging baseline; each later child targets the preceding fork branch. Comparing 
every child to fork master would recreate the cumulative-diff problem.
   4. Promote one ready functional slice at a time to `apache/pinot:master` 
after its true prerequisites have landed upstream. GitHub cannot use a 
fork-only branch as the base of an upstream PR.
   5. Refresh the staging root as upstream changes land and drop duplicate 
prerequisite patches. Do not force-reset fork master or merge the staging 
baseline wholesale into upstream.
   6. Preserve the old review threads and map their topics to the relevant 
child. Moving code is not, by itself, a reason to mark a correctness concern 
resolved.
   
   **Accounting scope:** the query disk limit is shared across a query's local 
operators within each server JVM, plus a JVM-wide limit across queries. It is 
not a cluster-wide reservation.
   
   
   ## Design
   
   ```text
                            input MseBlock.Data
                                    |
                                    v
                       +--------------------------+
                       | MultistageGroupByExecutor|
                       | in-memory hash table     |
                       +--------------------------+
                            |               |
           groups < min(maxGroups,limit)     | groups >= min(maxGroups,limit)
                            |               v
                            |    export [group keys + intermediate states]
                            |               |
                            |               v
                            |       hash(group keys) % P
                            |        /        |        \
                            |       v         v         v
                            |   part-0     part-1 ... part-P-1
                            |   append-only, length-prefixed DataBlocks
                            |               |
                            +------- continue with a fresh hash table
                                            |
                                            v
                                   end of input (EOS)
                                            |
                                            v
                       +------------------------------------+
                       | restore one partition at a time    |
                       | merge intermediate aggregation     |
                       | states into a bounded hash table   |
                       +------------------------------------+
                            | emit rows | delete partition
                            v
                       next partition ... -> EOS
   
    cleanup: success / upstream error / cancellation / early termination / close
   ```
   
   Key properties:
   
   - Spill files are scoped to one `AggregateOperator` under a configurable 
server spill directory (`pinot-aggregation-spill-*`); startup cleans orphaned 
spill files for that server instance.
   - Rows are partitioned using a deep hash of all group keys, including arrays 
and nulls.
   - Each serialized spill record contains at most 1,024 rows; aggregate 
buffering is bounded at 8,192 rows.
   - Aggregation intermediate states are serialized instead of finalized 
values, preserving DIRECT, INTERMEDIATE, and FINAL aggregation semantics.
   - Partitions are restored and emitted sequentially. The restored hash table 
obeys `numGroupsLimit` even if a partition exceeds the spill trigger; reaching 
the limit preserves the existing truncation/error behavior.
   - The input hash table also obeys `numGroupsLimit` as a hard backstop, while 
the global result limit remains in force across spill runs. Plans using group 
trimming are not eligible for spill and log why it was skipped.
   - Spill files are removed after consumption and on success, failure, 
cancellation, early termination, or operator close.
   - Operator stats expose `spillCount`, `spilledRows`, and `spilledBytes`. 
Server-owned per-query and per-process spill byte limits protect local disk 
usage.
   
   ## Enabling spill
   
   Spill requires both a server-side feature gate and per-query options.
   
   ### 1. Enable the server gate
   
   Set the following properties on every server participating in multi-stage 
query execution, then restart or roll the cluster:
   
   ```properties
   pinot.server.query.executor.mse.aggregation.spill.enabled=true
   # Optional: dedicated, writable local filesystem path; default is <instance 
dataDir>/aggregation-spill.
   
pinot.server.query.executor.mse.aggregation.spill.dir=/mnt/pinot/aggregation-spill
   # Optional: per-query and per-server-JVM byte limits (defaults: 1 GiB and 8 
GiB).
   pinot.server.query.executor.mse.aggregation.spill.max.bytes=1073741824
   pinot.server.query.executor.mse.aggregation.spill.server.max.bytes=8589934592
   ```
   
   The spill gate defaults to `false`. No `SAFE` stats mode is required: the 
mailbox and stream stats decoders discard an undecodable stat buffer during a 
mixed-version rollout, rather than failing the query.
   
   The user-facing query cannot override this server gate. 
`mseAggregationSpillEnabled` is internal metadata and any query-supplied value 
is overwritten by the server.
   
   Choose a dedicated local disk path with enough capacity for concurrent 
queries; avoid tmpfs. An instance-specific subdirectory is used under the 
configured root. Startup removes orphaned spill directories in that instance 
subdirectory, which should not be shared by concurrent server processes.
   
   ### 2. Configure each query
   
   Submit the query through the multi-stage engine and set a positive group 
threshold. The partition count is optional:
   
   ```sql
   SET mseAggregationSpillMaxGroups = 100000;
   SET mseAggregationSpillPartitions = 16;
   
   SELECT dimension, SUM(metric)
   FROM myTable
   GROUP BY dimension;
   ```
   
   - `mseAggregationSpillMaxGroups`: positive group-count spill trigger (a v1 
proxy for memory); `numGroupsLimit` remains a hard per-table ceiling. If 
absent, spilling is disabled for the query.
   - `mseAggregationSpillPartitions`: hash partition count, from 1 through 64; 
defaults to 8.
   
   Choose enough partitions to keep restored partition memory reasonable. A 
partition may exceed `mseAggregationSpillMaxGroups` but remains subject to 
`numGroupsLimit`.
   
   ## End-to-end test
   
   `QueryRunnerTest.testMSEAggregationSpill` exercises the complete in-process 
MSE path with two query servers, distributed realtime segments, mailbox 
exchange, and broker-side reduction.
   
   It first runs the baseline query without spilling:
   
   ```sql
   SELECT col1, SUM(col3), AVG(col3)
   FROM a
   GROUP BY col1
   ORDER BY col1;
   ```
   
   It then forces multiple spill cycles with:
   
   ```sql
   SET mseAggregationSpillMaxGroups = 2;
   SET mseAggregationSpillPartitions = 8;
   
   SELECT col1, SUM(col3), AVG(col3)
   FROM a
   GROUP BY col1
   ORDER BY col1;
   ```
   
   The test verifies the same schema and ordered rows as the non-spilled query, 
positive `spillCount`, `spilledRows`, and `spilledBytes` stats, and removal of 
spill directories after success. A second two-server end-to-end query covers 
`DISTINCTCOUNT` intermediate states.
   
   ## Verification and rollout
   
   The full proposal already has operator tests and an in-process 
two-query-server MSE test for SUM/AVG and DISTINCTCOUNT. These are the baseline 
requirements to preserve during extraction; each split PR must be revalidated 
on its own actual parent. Previously reported green results for the monolith 
are not a substitute for per-child validation.
   
   For each changed module, run focused Maven tests plus `spotless:apply`, 
`checkstyle:check`, `license:format`, and `license:check`. The final activation 
slice must prove spill actually happened, compare the result with the 
non-spilling query, and check cleanup after a successful drain. This is 
in-process multi-server testing, not a deployed-cluster or performance 
validation claim.
   
   Roll out only after all safety prerequisites land; keep the feature 
default-off. Disable the server gate (using the normal server 
configuration/restart process) to stop new queries from spilling, or revert the 
activation slice while leaving independently useful correctness fixes in place.
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to