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]