arXiv:2609.06982v1 [cs.DC] 7 Sep 2026
TreeRedux: Separating Concerns in Spark’s Distributed Tree Aggregation David A. G. Harrison*
Ivan Cao
Computer and Information Science University of Mississippi Oxford, MS, USA [email protected]
Computer and Information Science University of Mississippi Oxford, MS, USA [email protected]
Abstract—By default, Apache Spark’s tree aggregation primitives place the tree root on the driver, requiring the driver to participate in the same aggregation computation over intermediate aggregation state as executor nodes. For large aggregates, this can expose the single coordinator to substantial computation and memory requirements. Recent Spark versions optionally move the root to an executor, but the completed aggregate must still be returned to and materialized on the driver. We demonstrate this limitation using exact quantile computation and heavy-hitter identification, where the intermediate aggregation state can be substantially larger than the desired final result. We propose TreeRedux, a minimal extension that adds a terminal finalize operation U → V executed on an executor, allowing the compact result V rather than the potentially large aggregation state U to be materialized on the driver. For exact quantile computation, applying TreeRedux to GK Select removes the driver’s εn memory term, reducing driver memory requirements to the same asymptotic order as Spark’s GK Sketch. In our experiments, the default GK Select implementation encountered a driver OOM at 2.5 billion elements. Spark’s executor-side final aggregation option extended this limit to approximately 16–18 billion elements but still required the final aggregation state to be materialized on the driver. Redux Select completed through 28 billion elements without a driver OOM. TreeRedux allowed Space-Saving sketches with up to 32× the capacity of the largest configuration that materializes a full sketch on the driver.
I. I NTRODUCTION Modern distributed data-processing systems are commonly built around a shared-nothing architecture [1], in which each worker has its own CPU, memory, and storage, and communication between workers occurs through explicit data exchange rather than shared memory. This architecture enables horizontal scaling when computation can be decomposed and moved toward its data. A defining characteristic of shared-nothing systems is the separation between workers and coordinators. Workers perform data-intensive operations using local computation and * Corresponding author.
This is the extended version of a paper submitted for peer review. This work has been submitted to the IEEE for possible publication. Copyright may be transferred without notice, after which this version may no longer be accessible. Please cite the published conference version of this work once available, unless referring to material that appears only in this extended version.
memory, while coordinators manage execution and collect results. This asymmetry enables systems to scale by adding workers without requiring the coordinator to process data at the same scale. Typically worker failure affects only tasks and state associated with that worker, and the coordination layer reassigns the work to another worker. Shared-nothing architectures have become particularly prevalent in cloud-scale data processing because their independent-resource model maps naturally onto cloud infrastructure, where computation is provisioned as networked virtual machines, and wherein each virtual machine has its own local processing, memory and storage. Apache Spark is a prominent distributed data-processing system that exemplifies the shared-nothing architecture. In Spark terminology, each data processing application has a driver which coordinates execution and schedules work on executors. Executors are processes running on worker nodes. Each executor has its own CPU cores and memory allocation and operates on one or more data partitions. Spark’s fault-tolerance model treats executors as replaceable data-processing resources whose state can often be reconstructed through lineage: if an executor fails, its work is reassigned to another executor and computation continues. The driver, however, remains the central coordinator for the lifetime of the application. If the driver crashes, executors receive no further work and cannot report results, so execution halts as soon as the executors finish their currently assigned tasks; recovering requires an external mechanism to restart the driver. We also observed the driver hang under memory pressure without crashing outright — a second failure mode with the same halting effect, but one that lacks the clear processexit signal of a crash and so may go undetected for longer. We confirmed this failure mode directly, on independent runs across both of our compute backends; see the appendix for the full classification. Both failure modes underscore why the driver should not be placed under the same memory and computational pressure as executors. Table I itemizes what is actually lost when the driver crashes or hangs without deliberate checkpointing. Unfortunately, some Spark operations do exactly this. In this paper, we focus on tree aggregation, a distributed computation pattern in which partition-local intermediate
TABLE I: State lost when the Spark driver crashes or hangs without intentional checkpointing. Executor failures are recoverable because lost partitions are reconstructed from lineage on another worker; the driver has no equivalent fallback, since it is the sole holder of the rows marked “Yes” below. Component RDD lineage DAG DAGScheduler state Task completion state Shuffle map output tracking Accumulator state Cached RDD partitions Shuffle data files Input data Checkpointed RDD data
Location / owner Driver JVM Driver JVM Driver JVM Driver JVM Driver JVM Executors Executors / shuffle service External storage Reliable storage
Lost on driver failure? Yes Yes Yes Yes Yes Maybe Maybe No No
states are recursively merged through an aggregation tree until a final aggregate is produced. This pattern maps naturally onto shared-nothing systems. Workers compute locally, aggregate, and return the result to the coordination layer represented by the driver. Spark provides tree aggregation primitives such as treeReduce and treeAggregate. Before SPARK36419 [2] was included in v3.3.0, the root of the aggregation tree was placed on the driver rather than an executor. Even after this change, Spark by default still places the root of the aggregation tree on the driver. Only when finalAggregateOnExecutor is set to true does Spark move the root. Even when true, Spark continues to return a result of the same type as the intermediate aggregated state. For algorithms that construct mergeable sketches, such as GK Sketch [3], Misra-Gries [4], or Count-Min Sketch [5], the final aggregation state may be a hash table, a 2-D array, a vector– possibly with supporting data structures. For such algorithms, the requested answer can be much smaller than the aggregation state. Spark nevertheless returns the full state to the driver, where the final answer is computed. TreeRedux demonstrates the advantages of avoiding materializing the final aggregation state on the driver. It removes a source of driver hangs and crashes, and we show that it can improve asymptotic complexity as well as practical performance for two illustrative use cases. II. R ELATED W ORK MapReduce and Hadoop established large-scale batch processing on shared-nothing clusters [6]–[8], but neither provides a distributed tree aggregation primitive. Dryad and DryadLINQ generalized this model to directed acyclic graphs of operators [9], [10]; Spark adopted the same dataflow model while exposing explicit tree aggregation [11], [12]. Spark’s treeReduce and treeAggregate suffice when intermediate state and result share one type, as in scalar reductions. Mergeable summaries generalize this pattern: independent local summaries are combined while preserving their guarantees [13]. TreeRedux addresses the case in which the completed summary remains much larger than the answer ultimately needed by the driver. Table II lists examples where the desired result is often much smaller than the intermediate state.
Consequence Lineage must be reconstructed; lost RDD dependency graph. Cannot resume jobs, stages, or dependencies. Completed tasks are not known to the new driver. Shuffle locations may be lost; stages may rerun. Application metrics and counters may be lost. Data may remain, but reuse is not guaranteed. Data may survive, but metadata may be unavailable. Can be reread to recompute results. Provides a recovery point and limits recomputation.
In this paper we pick two representative examples where the final state is significantly smaller than the intermediate state. The first returns a scalar, the second an aggregate stripped of overhead: • exact quantile computation using GK Select, and • heavy-hitter identification using a Space-Saving sketch. III. S PARK API E XTENSION With Spark’s tree aggregation, each executor operates on one or more partitions of the data. For each partition, an executor constructs some aggregate state of type U. The aggregate state is repeatedly merged through the tree and the final U is returned to the driver. This is appropriate for scalar reductions, but it is unnecessarily restrictive for mergeable summaries whose completed state can be reduced to a smaller answer. TreeRedux extends the tree aggregation pattern by adding an explicit terminal finalize : U → V operation, allowing the completed summary to be converted to the desired result before it is returned to the driver. Spark exposes treeReduce as a familiar reduction primitive: def treeReduce(f: (T, T) => T, depth: Int = 2): T It recursively applies f to objects of type T until only one object of type T remains. It does this using an aggregation tree distributed across the Spark executors with the root on the driver. This interface is appropriate when the elements being reduced, the intermediate aggregation state, and the final result all have the same type. Internally, Spark implements treeReduce using the more general treeAggregate primitive, which separates the input element type T from the aggregation state type U:
TABLE II: Examples of mergeable summaries where the intermediate state U may be substantially larger than the final result V . The scalar column indicates the useful final-result form in the TreeRedux setting. Problem
Sketch / Summary
Intermediate State U
Final Result V
Scalar?
Quantiles Quantiles Quantiles Quantiles Heavy hitters Heavy hitters Frequency estimation Cardinality estimation Cardinality estimation Frequency moments Similarity estimation Set membership Set membership Range counting
GK sketch [3] GK Select [14] KLL sketch [15] t-Digest [16] Misra–Gries [4], [13] Space-Saving [17] Count-Min sketch [5] HyperLogLog [18] HyperLogLog++ [19] AMS sketch [20] MinHash [21] Bloom filter [22] Counting Bloom filter [23] ε-approximation [13]
Quantile summary Summary, rank-interval candidates Compactor hierarchy Centroids Candidate counter table Candidate counter table Counter matrix Register array Register array Randomized projections Signature vector Bit vector Counter vector Sample / coreset
Quantile value(s) Quantile value(s) Quantile value(s) Quantile value(s) Frequent-item list Frequent-item list Estimated count† Distinct-count estimate Distinct-count estimate Moment estimate Jaccard estimate Membership query Membership query Range-count estimate
✓ ✓ ✓ ✓ – – ✓ ✓ ✓ ✓ ✓ Batch† Batch† Batch†
† A single query produces a scalar result but does not justify constructing and merging the summary. TreeRedux is useful when a
batch of queries is evaluated against the merged summary on an executor, producing a collection of answers; a single query can instead be answered exactly by a direct scan.
be substantially smaller than U. For example, in exact quantile computation, U may be a large mergeable sketch while V is a single numeric quantile. In heavy-hitter identification, U may contain auxiliary hash tables and counters, while V contains only the compact list of reported heavy hitters. A natural question is whether combOp alone could discard unneeded state at the root, thus making finalize unnecesThe type U represents both the intermediate aggregation sary. For a mergeable summary, the contract of combOp is state and the value returned to the driver. This coupling is the more than its type U × U → U : merging states representing source of the bottleneck addressed in this paper: even when the two disjoint input subsets must produce another valid state desired answer is small, Spark must return the full aggregate representing their union. This closure property permits Spark state U to the driver. to rearrange the aggregation tree while preserving the sumThe finalAggregateOnExecutor parameter moves mary’s meaning. Finalization, by contrast, has type U → V the root of the aggregation tree to an executor. This avoids and may intentionally discard information needed by a later performing the final merge on the driver, but the returned value merge. Spark provides combOp no direct indication that an is still of type U. Thus, if U is a large sketch, hash table, or invocation is terminal, so an algorithm that attempts to finalize other intermediate representation, that large state must still be inside combOp must reconstruct that condition itself. materialized on the driver. Nevertheless, an algorithm-specific workaround can encode We propose treeAggRedux, which separates the aggre- terminality inside the aggregation state. For example, suppose gation state from the final returned value by adding a terminal that each partial state carries the number of input elements finalization function: it represents, and that the total cardinality n is known. The algorithm can replace U with a tagged state such as def treeAggRedux[U: ClassTag, V: ClassTag]( Partial(U, count) or Final(V). When a combOp zeroValue: U, invocation combines partial states whose counts sum to n, depth: Int = 2 it can apply the algorithm’s finalization logic and return )( Final(V). With finalAggregateOnExecutor=true, seqOp: (U, T) => U, this terminal state can be formed on the executor holding the combOp: (U, U) => U, root before it is sent to the driver. finalize: U => V This workaround preserves the nominal signature U × U → ): V U only by widening it to a tagged sum type over both The functions seqOp and combOp are identical to those cases — in Scala 2, a sealed trait such as Partial(U, used by treeAggregate. They construct and merge inter- count) / Final(V) (Scala 3 offers native union types for mediate states of type U. The new function finalize is the same purpose). Because only one variant is populated per executed on an executor after the final aggregate state has instance, a Final(V) value need not itself be as large as been produced but before the result is returned to the driver. U ; the workaround’s real cost is structural rather than spatial. This allows the driver to receive a value of type V, which may combOp must branch on the tag on every merge, not only
def treeAggregate[U]( zeroValue: U, seqOp: (U, T) => U, combOp: (U, U) => U, depth: Int, finalAggregateOnExecutor: Boolean )(implicit arg0: ClassTag[U]): U
the terminal one, and it must define a result for combining a Final(V) with a Partial(U) or another Final(V) — exactly the case the closure property above rules out as generally meaningful, since a finalized result usually cannot be re-merged with further raw state. Cardinality metadata must also be threaded through every seqOp/combOp call rather than checked once, and any residual driver-side fold (e.g., combination with the zero value) must special-case the Final tag. TreeRedux instead makes this boundary explicit: finalize is invoked exactly once after aggregation completes, returns V directly, and avoids algorithm-specific terminal detection and tagging machinery. For convenience, we also define a treeRedux wrapper analogous to treeReduce: def treeRedux[V: ClassTag]( depth: Int = 2 )( f: (T, T) => T, finalize: T => V ): V This wrapper handles the special case where the aggregation state is the same type as the input element type. The more general primitive is treeAggRedux. The proposed API could also be implemented as overloads of Spark’s existing treeReduce and treeAggregate methods. We use the name Redux in this paper only to distinguish the proposed behavior from the current Spark API. Our prototype implementation reuses nearly all of Spark’s existing tree aggregation machinery. The appendix describes the implementation and presents a noninferiority study showing that treeAggRedux introduces no measurable overhead relative to treeAggregate with finalAggregateOnExecutor=true. IV. E XAMPLE U SE C ASES To evaluate TreeRedux we consider two representative applications whose intermediate aggregation state may be substantially larger than the final result returned to the user. The first workload is exact quantile computation using GK Select [14]. The second workload is heavy-hitter identification using a Space-Saving sketch [17]. TABLE III: Definitions used in algorithm descriptions and analysis. Symbol n ni P q k ∆k ε d ϕ b r
Description Total number of elements across all partitions. Total elements in the ith partition. Number of partitions. Quantile queried (e.g., 0.5 for median). Target rank k = nq Target rank minus the approximate rank. GK sketch relative error parameter. Requested depth of the aggregation tree. Maximum actual fan-in over aggregation nodes. Target fan-in. d = ⌈logb P ⌉. The number of samples collected for splitter selection.†
† Relevant only to Spark Full Sort.
A. Exact Quantile Computation with GK Select GK Select [14] computes exact quantiles by combining a mergeable Greenwald–Khanna (GK) sketch [3] with a secondstage exact selection. First, each executor constructs a local GK sketch for each of the partitions assigned to that executor. The local sketches are merged through Spark’s tree aggregation to produce an approximate quantile with rank error |∆k| ≤ ⌈εn⌉. ε sets the allowed rank error. The algorithm then identifies the exact quantile by computing the |∆k| candidate values local to each partition within the rank interval induced by the approximate quantile, and then tree aggregating across partitions to find the global set of |∆k| candidates. The exact quantile is the appropriate extremum of this global set. This set requires O(εn) space on the driver, even though the final output is a single scalar quantile. This second aggregation phase provides an illustrative workload for evaluating TreeRedux because its intermediate aggregation state grows with the dataset size while the desired result remains constant in size. B. Heavy-Hitters with the Space-Saving Sketch Space-Saving [17] is a streaming heavy-hitter sketch that maintains a bounded table of candidate items and a counter for each candidate used to estimate its frequency. When an observed item is already present in the table, its counter is incremented; when a new item arrives and capacity remains, it is inserted; when the table is full, the item with the smallest count is replaced. The algorithm assumes that any previously unseen item could have occurred no more frequently than the current least frequent candidate. The new item is therefore inserted with the replaced item’s count plus one, yielding an error-bounded estimate. In a distributed setting, each partition constructs a local Space-Saving sketch, and these sketches are merged through tree aggregation. The intermediate state consists of the sketch itself, including the candidate table and supporting data structures required for efficient updates and merges, while the final result is simply the reported heavy hitters and their estimated frequencies. Shedding the supporting data structures and tuning k to return the reliable subset of the heaviest hitters, as we show, dramatically reduces the size of the state. V. E XPERIMENT M ETHODOLOGY All experiments were conducted on Amazon EMR 7.9.0 running Spark 3.5.5. Experiments used a cluster consisting of one primary node and 30 core nodes. The driver ran on the primary node, and each core node ran one executor. All experiments, in both Example I and Example II, use r6g.xlarge nodes, which provide 4 vCPUs and 32 GiB of memory per node. Increasing instance memory can avert a particular OOM, but it merely delays the problem to larger data sets or clusters; it neither solves the underlying problem nor changes the asymptotic behavior. We leave evaluation of remedies for executor bottlenecks and memory limitations to future work.
Fig. 1: GK Select vs. full sort as dataset size n grows (P = 120, ε = 0.01, 30×r6g.xlarge, executor memory 18 GiB, driver memory 512 MiB). (a) Runtime vs. n (log-scale x-axis): GK Select is faster than full sort throughout the range shown, but OOMs at n = 2.5B (dashed vertical line); no GK Select data exists at or beyond that point. (b) Driver peak heap: climbs from 433 MiB at n = 1B to 505 MiB at n = 2B, approaching the 512 MiB driver heap limit (dotted line) before the driver OOMs at n = 2.5B. Full Sort’s driver memory stays flat (346–355 MiB) across the entire range, since it never collects an εn-sized candidate slice on the driver. (c) Across the test range, executor peak heap grows slower for GK Select, and of more importance, GK Select does not cause an executor OOM, confirming that the problem is isolated to the driver.
All reported runtimes represent end-to-end Spark job execution times. Driver heap usage is sampled every 100 ms; the reported peak is the largest observed sample. This is used JVM heap, not committed heap, non-heap memory, process RSS, or container memory. The driver is configured with a 512 MiB ceiling. VI. E XAMPLE I: E XACT Q UANTILE C OMPUTATION In order to demonstrate the need for TreeRedux, we start with GK Select applied without TreeRedux and explore its scaling limitations. We show that with each measure taken to increase scale, we ran into a new barrier until we were convinced that the proper measure was to compute the final quantile on the executor and only report the scalar quantile to the driver. A. The scaling limitations of GK Select with respect to n Full sort as implemented by sortByKey or sortBy is the traditional means to obtain an exact quantile, although establishing a full ordering is excessive if the objective is only to obtain one or more quantiles. In Figure 1, we include a comparison against Spark’s full sort to explain why someone might use GK Select to compute a quantile. It is clearly faster than a full sort and uses less executor memory. The limiting factor is the driver memory. As n increases, the driver’s peak memory climbs until the driver experiences an Out-of-Memory (OOM). B. Analysis of GK Select We use the terminology from Table III. We begin by demonstrating how GK Select encounters scaling limitations when treeAggregate returns intermediate aggregation state to the driver. The problem becomes
increasingly pronounced as the dataset size n grows. The intermediate state has two substantive components: 1) the GK sketch, which produces an approximate quantile whose rank error is bounded by |∆k| ≤ ⌈εn⌉, and 2) the set of |∆k| candidate values within the rank interval induced by the approximate quantile, which is used during the second-stage exact selection phase. GK Select is algorithmically identical to the implementation presented by Cao et al. [14]. However, the analysis in that work assumes that the driver aggregates GK sketches from all P partitions. This assumption was valid for Spark versions prior to Spark 2.3, but no longer accurately describes current Spark implementations. Beginning with Spark 2.3, Spark’s GK Sketch implementation uses treeAggregate at a depth of 2. At that depth, treeAggregate performs the first aggregation on executors and then the second aggregation on the driver. Note that tree aggregation for GK Sketch is separate from tree aggregation for the |∆k| candidates. The GK Sketch tree’s depth is fixed at two and is not exposed through Spark’s public API, so no experiment in this paper varies it. We are, however, free to set the depth of the |∆k| aggregation tree, and do so in the following subsection. Spark’s treeAggregate determines the fan-in at each aggregation level from the number of partitions P and the requested depth d: scale = max(⌈P 1/d ⌉, 2). It then repeatedly reduces the number of intermediate partitions by this scale using integer (truncating) division until another tree level would no longer provide benefit. Specifically, Spark applies the recurrence Pi Pi+1 = scale
TABLE IV: Nominal fan-in vs. actual fan-in at the driver, by depth (P = 120).
while Pi > scale +
Pi . scale
The number of intermediate partitions remaining after this recurrence determines the number of partial summaries that must be merged by the driver’s final fold(). Section VI-C analyzes this recurrence in detail and demonstrates that the actual surviving partition count can differ from the idealized P 1/d due to Spark’s integer arithmetic. For depth 2, the idealized fan-in is √ scale = ⌈ P ⌉, so √ Spark reduces the original P partition summaries to Θ( P ) pre-merged summaries before they reach the driver. This differs from the driver-only aggregation model analyzed in Cao et al [14]. GK Select uses substantial driver-resident state in two successive phases. During the first phase, the driver receives and merges the GK sketch state that remains after executor-side tree aggregation. After extracting the approximate quantile, GK Select releases this sketch state. During the second phase, the driver aggregates the |∆k| candidate values used for exact selection. Because the sketch and candidate set reach their respective peaks at different times, the peak logical driver memory is the larger of the two phase-specific requirements: √ O max
!! P n log ε √ , ϕ · εn ε P
(1)
where ϕ denotes the fan-in when aggregating |∆k| candidates. A sequential combOp intrinsically requires only O(εn) candidate space, independent of ϕ; the factor ϕ captures worstcase Spark materialization when multiple incoming results remain buffered. For the candidate aggregation, Spark’s construction gives a nominal fan-in of O(P 1/d ). To keep fan-in bounded as P grows, we choose a constant target b > 1 and set d = ⌈logb P ⌉. The resulting fan-in is ϕ = O(b) = O(1). For P = 120, choosing b = 4 gives d = 4. Section VI-C deliberately varies d to evaluate the effects of other tree shapes. Both arguments of the maximum can limit scalability, as demonstrated in the following sections. However, the second term is particularly problematic because it grows linearly with the dataset size n and proportionally with the fan-in ϕ. Higher fan-in ϕ requires a node to perform proportionally more merge operations. Spark does not provide backpressure based on the rate at which these folds execute. Consequently, incoming results may be deserialized before the aggregation function consumes them and may remain buffered while the fold progresses. Therefore, the worst-case memory requirement is the sum of all arriving |∆k| sets.
Depth
Fan-in (scale)
Actual fan-in at driver
Real driver combines
1 2 4 6
120 11 4 3
120 10 1 4
119 9 0 3
C. Reducing fan-in helps partially Section VI-B derives how Spark’s treeAggregate sets fan-in from the partition count P and depth d via scale = max(⌈P 1/d ⌉, 2) and the recurrence Pi+1 = ⌊Pi /scale⌋ while Pi > scale + ⌈Pi /scale⌉. Increasing d decreases this fan-in at every level of the tree, including the driver. We hold P = 120 and ε = 0.01 fixed and sweep d ∈ {1, 2, 4, 6}, increasing n until driver OOM occurs for each setting. Figure 2a shows driver peak heap versus n for each depth. Every curve rises with n and eventually hits the driver’s 512 MB ceiling, but the OOM threshold is not monotonic in depth: depths 1, 2, and 6 fail at n ≈ 1.5B, 3B, and 5B respectively, while depth 4 survives to n ≈ 17B—more than 3× further than depth 6, despite depth 6 having a smaller nominal fan-in. Table IV, “Nominal fan-in vs. actual fan-in at the driver, by depth,” shows why: what matters is not the nominal fan-in per tree level, but how many partial results actually survive to reach the driver. Whatever partition count survives the recurrence from Section VI-B is exactly what the driver’s final .fold() call must merge — the “actual fan-in at the driver” column above. Since .fold() combines these one at a time starting from an empty zero value, the number of real driver-side combine operations is (actual fan-in at driver −1). For d = 4 at P = 120: scale = ⌈1201/4 ⌉ = ⌈3.31⌉ = 4, and the recurrence runs 120 → 30 → 7 → 1, terminating at exactly one remaining partition. The driver’s final fold has nothing to combine: the single already-merged |∆k|-sized slice passes straight through, so depth 4 incurs zero driver combines. This is an artifact of the arithmetic at P = 120, not a general property of moderate depths: depth 6’s recurrence (120 → 40 → 13 → 4) terminates at four remaining partitions (three real combines), which is enough to force an OOM at roughly a third of depth 4’s n, despite depth 6 having a smaller average fan-in (3 vs. 4). Modifying the depth to minimize the number of partial results arriving at the driver can increase the maximum number n that can be accommodated prior to a driver experiencing an OOM, but even when only 1 result arrives at the driver, when n becomes large enough the driver still experiences an OOM.
FAE eliminates the ϕ factor from the second argument to max, causing the driver memory requirement in Equation 1 to become !! √ P n log ε √ , εn (2) O max ε P
(a) Without Final Aggregate on Executor (FAE). Depth 4 survives to n ≈ 17B while depths 1, 2, and 6 fail at n ≈ 1.5B, 3B, and 5B respectively. A separate, later trial at n = 25B (same cluster and executor configuration) produced a driver hang rather than a clean OOM crash, confirmed for all three seeds: the driver log ran cleanly through the full treeReduce, survived one executor OOM and retry, reached the final single-task combine stage (consistent with depth 4’s fan-in-1 behavior), received a task result of approximately 134.5 MB on the executor, and then went silent with no exception — the same dag-scheduler-event-loop silent-OOM pattern described in the appendix, where the full classification appears.
(b) With FAE, OOM boundaries cluster tightly at n ≈ 16–18B for every depth, in sharp contrast to (a).
Fig. 2: Driver peak heap vs. n across depths 1, 2, 4, and 6 (P = 120, ε = 0.01, 30×r6g.xlarge, 18 GB executor heap, 512 MB driver heap), without FAE (a) and with FAE (b). Dashed verticals mark each depth’s confirmed driver-OOM boundary; the dotted vertical in (a) marks a confirmed driver hang, a different failure mode discussed in the Introduction. FAE delays failure substantially but does not eliminate it.
D. When executor-side final aggregation fails Spark’s treeAggregate offers a finalAggregateOnExecutor (FAE) flag [2] that moves the root of the aggregation tree from the driver to an executor. The executor at the root of the aggregation tree then sends a single aggregated result to the driver. When FAE is applied to the aggregation of the |∆k| candidates within the rank interval induced by the approximate quantile,
We reran the experiments in Section VI-C but with finalAggregateOnExecutor turned on. The results appear in Figure 2(b). FAE shifts all driver OOMs into the same range as the depth-4 result in Section VI-C, because in both cases a single aggregate of size |∆k| is returned to the driver. When |∆k| exceeds the driver memory, the driver OOMs. If we want to push beyond this limitation, the binding constraint appears to be the interface itself : for GK Select, any aggregation primitive whose return type equals its accumulator type U must materialize O(εn) state at the driver. Moving the root of the candidate aggregation tree to the executors cannot circumvent this. E. Finalize We reran the depth-2 and depth-4 configurations from Section VI-D with Redux Select, increasing n until a resource limit was reached. Unlike every preceding configuration, the first failures were executor OOMs rather than driver failures. This distinction is consequential. TreeRedux prevents the |∆k| candidate array from being materialized by the driver, removing a growing memory demand from Spark’s single, non-replaceable coordinator. The terminal executor must still hold that array while running quickSelect, so TreeRedux does not eliminate the O(εn) capacity limit; it relocates that demand to the data plane. An executor failure is eligible for Spark’s task-retry and rescheduling mechanisms, and a retry may succeed when the original failure resulted from transient memory pressure or contention with other tasks on the same instance. In contrast, a driver OOM halts the application and loses its coordinator state. If the final candidate array alone exceeds the memory available to any executor, retries cannot resolve the underlying capacity limit. Characterizing executor OOM behavior—including the effects of task contention, retry, and rescheduling—is outside the scope of this paper and left for future work. Across the tested range, driver peak heap remained nearly level, and full sort and Redux Select exhibited almost identical peak memory usage. The measured baseline before either computation was stable across runs, as shown by the 95% bootstrap confidence interval in Figure 3, but it lies well below both peaks. Thus, the measurements do not identify the source of the remaining, apparently shared driver allocation. We leave that question outside the scope of this paper. Figure 3 confirms the expected change in the driver-side bound: TreeRedux removes the linear εn candidate term, leaving only the GK-sketch and Spark bookkeeping state. No driver OOM occurred at either depth through n = 28B, well beyond the range reached by the methods in Sections VI-C
Fig. 3: Redux Select at depths 2 and 4, compared against Spark full sort, on the same cluster used throughout this section (30×r6g.xlarge, 18 GB executor, cores=2, 512 MB driver, P = 120, ε = 0.01; horizontal-axis sizes are labeled in billions (B)). (a) Runtime: both depths track closely and are roughly an order of magnitude faster than full sort (9.0× and 9.5× at n = 9B, for depth 2 and depth 4 respectively); full sort’s own sweep did not complete beyond n = 9B, exceeding its 30-minute per-trial cap at n = 12B. (b) Driver peak heap is essentially flat across the entire swept range for both depths, in contrast to the FAE curves in Figure 2(b), which continue climbing toward the 512 MB ceiling regardless of depth. (c) Executor peak heap becomes the binding constraint instead: depth 2 survives to n = 28B before failing at n = 30B; depth 4 survives to n = 25B before failing at n = 28B. Both failures were confirmed executor-side (ExecutorLostFailure / YARN container kill, or an executor heartbeat timeout, each followed by a caught SparkException and a graceful SparkContext shutdown) rather than the driver’s silent self-kill signature seen throughout Sections VI-C and VI-D.
and VI-D. The final executor nevertheless remains subject to the |∆k| = O(εn) candidate-memory requirement, which accounts for the eventual executor OOMs shown in the figure. The runtime comparison is also important because full sort is Spark’s conventional mechanism for exact quantile computation. At n = 9B, Redux Select completed in 156.0 s at depth 2 and 147.4 s at depth 4, compared with 1398.0 s for full sort: respective speedups of 9.0× and 9.5×. Full sort did not complete the n = 12B trial within its 30-minute cap, so we do not claim a measured speedup beyond n = 9B. Nevertheless, the completed points establish a substantial practical advantage for Redux Select over the traditional exact baseline. F. Redux Analysis The preceding sections derive the costs of GK Select and the Spark GK sketch it uses. Redux Select changes only the candidate-aggregation phase: finalize runs QuickSelect on the root executor and returns a scalar. This removes the εn candidate term from both driver time and memory, leaving √ O εP log(ε √nP ) for each, determined by the sketch phase. The root executor still processes and stores the |∆k| ≤ ⌈εn⌉ candidates. Without an assumption relating ε and P , its time is O Pn log 1ε + Pn log log(ε Pn ) + εn and its memory is O max( Pn , εn)
Under the ε ≲ 1/P regime assumed by Cao et al. [14], εn ≲ n/P , so these reduce to the executor bounds shown for Redux Select in Table V. VII. E XAMPLE II: H EAVY H ITTERS The heavy-hitter problem seeks to identify the most frequently occurring items in a stream or dataset while using substantially less memory than would be required to maintain exact counts for all distinct items. Throughout this section, we use the term heavy hitters or top-k frequent items to refer to the k items with the largest frequencies. This differs from the top-k selection problem used in the appendix, where top-k refers to the k largest values in a numeric dataset. The heavy-hitter problem provides a useful second case study because the mergeable sketch state can be substantially larger than the final result. During aggregation, the sketch must retain a large candidate frontier in order to distinguish items with similar frequencies and avoid premature elimination of potential heavy hitters. However, the final output consists only of the top-k frequent items. This separation between the intermediate representation U and the desired result V makes heavy hitters a natural application of treeAggRedux, which performs the final extraction of the top-k frequent items on an executor and returns only the compact result to the driver. We use a distribution with a wide frontier because this is the regime in which heavy-hitter sketches are most stressed. When many items have similar frequencies near the reporting threshold, small count errors can change which items appear in the reported top-k. Such frontiers arise naturally in saturated or capped processes, where many entities reach similar upperend frequencies: for example, products with similar sales under
TABLE V: Asymptotic executor and driver complexity for quantile methods. Executor timea n n Spark Full Sort O P log P n n n Classical GK Sketch O P log 1ε + P log log(ε P ) c n 1 n 1 n Spark GK Sketch O P log B + ε P B log(ε P ) n n n GK Selectd O P log 1ε + P log log(ε P ) n n n Redux Select O P log 1ε + P log log(ε P )
Driver time
Executor memory
Driver memory
GK depthb
|∆k| depth
O(rP log(rP )) n/a (streaming)
n O( P ) 1 n log(ε ε P ) 1 n O(B + ε log(ε √ )) P n O P
O(rP ) n/a (streaming)
n/a n/a
n/a n/a
Algorithm
O
√
O
P ε log(εn)
O √
P √n εn ε log(ε P ) + √ P √n ) O log(ε ε P
O
n P
√
O( εP log(ε √n ))
√
P
O max( εP log(ε √n ), εn) P √ P √n ) O ε log(ε P
2
n/a
2
⌈logb P ⌉
2
⌈logb P ⌉
a Executor-time entries assume P = O(E), where E is available executor task parallelism. With fixed E and increasing P , task waves add scheduling time. b Depth of the GK sketch aggregation tree. As of Spark 3.5.5, approxQuantile (GK Sketch) implements a fixed depth of 2, which is not exposed and
thus not settable. We made it settable by slightly modifying the Spark’s GK Sketch implementation, but we move such exploration out-of-scope for this paper. c B in the Spark GK Sketch time complexities refers to the size of headSampled which buffers samples before flushing them to the sketch. The size of this buffer is fixed, which modifies the complexities from the classical GK Sketch. Improving Spark’s GK Sketch itself is outside the scope of this paper and is the subject of an upcoming paper. d For GK Select, we modified the GK Sketch to set the size of headSampled a factor larger than the post-compress sketch size, which restores the classical GK Sketch time complexity.
TABLE VI: Communication and synchronization complexity for quantile methods. Algorithm
Network volume
Full Shuffles
Rounds
E/A
Spark Full Sort GK Sketch GK Select Redux Select
O(n) n P ε log(ε P ) n P ε log(ε P ) + εnP P n ε log(ε P ) + εnP
1 0 0 0
1 1 3 3
Exact Approx. Exact Exact
O
O O
inventory limits, popular posts constrained by recommendation exposure, network flows limited by rate caps, or sensor and event streams where many devices emit near a maximum reporting rate. In these cases, the difficulty is not identifying a single dominant heavy hitter, but resolving a crowded set of contenders near the cutoff. We evaluate treeAggRedux on the heavy-hitter problem using a Space-Saving sketch. We generate n=830,472,175 integers spread evenly across P =100 partitions using the default aggregation tree depth 2. We use a logistic function because it provides an easily controlled plateau of high-frequency items together with a tunable transition region. The parameter Cmax determines the saturation level, r0 determines the location of the transition, and α controls the sharpness of the drop-off. This allows us to construct a large candidate frontier while maintaining a compact parameterization and a smooth count-versus-rank curve. Item frequencies are assigned according to the logistic function c(r) =
Cmax , 1 + eα(r−r0 )
with Cmax = 30, r0 = 25,165,824, and α = 5 × 10−7 . The resulting distribution concentrates many items near the topk frequent-item boundary, a regime where sketch quality is most sensitive to sketch capacity. Figure 4 shows the resulting count-versus-rank curve. To demonstrate the difficulty of this regime we considered two different heavy hitter sketches: 1) Misra-Gries (MG) [4]
Fig. 4: Item frequency profile used in the heavyhitter experiments (Cmax =30, α=5×10−7 , r0 =25,165,824, n=830,472,175 elements across 108,818,368 distinct labels, of which 31,054,701 fall in the transition region with count > 1). The logistic function Cmax /(1 + eα(r−r0 ) ) (blue) determines the frequency of each item at rank r; the red step function is the discrete band approximation used in the experiments. The dashed line marks r0 , where count = Cmax /2 = 15. Items to the left of r0 form a dense plateau of high-frequency counts; the steep transition near r0 is the frontier where many items compete for the top-k positions, making sketch accuracy most sensitive to k in that region.
using mergeable summaries from Agarwal et al [13], and the Space-Saving (SS) sketch [17]. Although MG and SS are similar, we found MG to be far more erratic in the regime with a wide frontier. We include partial results for MG in Table VII to demonstrate erratic behavior. We varied the sketch size and ran 20 trials at each sketch size while shuffling the items within each partition in between trials. We evaluate sketch quality with three metrics. Precision is the fraction of returned items whose true rank is at most k: a value of 1 means every returned item genuinely belongs in the top-k. Recall is the fraction of the true top-k items that appear in the result: a value of 1 means nothing from the true top-k was missed. Rank MAE (mean absolute error) is the average absolute difference between each returned item’s position in
TABLE VII: Misra-Gries sketch vs. sketch size k. Sketch performance metrics show mean ± 1.96×SE (95% CI) over 20 trials. mean ± 95% CI
k
100K 250K 400K 500K 750K 1M
Precision
Recall
Rank MAE‡
0.589 ± 0.002 0.520 ± 0.000 0.598 ± 0.026 0.538 ± 0.000 collapse† collapse†
0.001 ± 0.000 0.015 ± 0.000 0.014 ± 0.005 0.031 ± 0.000 0.000 ± 0.000 0.000 ± 0.000
4,765K ± 30K 6,163K ± 7K 4,717K ± 454K 5,796K ± 4K — —
† All 20 trials returned an empty sketch. ‡
Rank MAE is computed over the returned items only, comparing each item’s position in the result against its true rank in the dataset. Values can substantially exceed k when returned items are drawn from outside the true top-k: an item at estimated rank 50,000 whose true rank is 5,000,000 contributes an error of 4,950,000.
the result and its true rank in the dataset, computed only over the items that were returned. Because rank MAE is restricted to returned items, it can substantially exceed k when the sketch surfaces items far outside the true top-k. Table VII shows a collapse to empty results at k=750K and k=1M, but not at smaller values of k. This is counterintuitive: one might expect larger sketches to be more robust, yet it is the intermediate sizes that fail completely. The explanation lies in a resonance between sketch capacity and the density of the logistic frontier. The Agarwal merge compresses a merged sketch of size 2C back to capacity C by subtracting the (C+1)-th largest count from every entry and removing items whose count reaches zero. The amount subtracted per merge step is therefore governed by the minimum count among the items that must be evicted. At k=750K–1M (capacity 3.75M–5M), the sketch fills predominantly with items from the logistic transition region, where millions of items have counts clustered tightly around Cmax /2 = 15. Because the minimum count of the excess items is close to the maximum count of the retained items, each merge step subtracts nearly as much as the items hold. After approximately 25 merge steps in the reduction tree (P =100, depth 2), this compounding effect empties the sketch entirely. At small k (100K–500K, capacity 500K–2.5M), the sketch is too small to accommodate the frontier: only items that accumulated the highest per-partition evidence survive the partition-level seqOp, and these carry counts well above the per-merge decrement threshold. The result is inaccurate (rank MAE in the millions, recall near zero) because those survivors are not the true top-k, but the sketch is not empty. We did not experience quite such erratic behavior from the SS sketch, so we used it for the remaining results in this section. With the SS sketch we compare three aggregation strategies: standard treeAggregate (Agg), treeAggregate with finalAggregateOnExecutor=true (FAE), and treeAggRedux (Redux). Each sketch has capacity 5k. Table VIII reports driver peak memory for all three strategies
and, for Redux, executor peak memory, precision, recall, and rank MAE. Table VIII reveals a sharp progression across aggregation strategies. Standard treeAggregate (Agg) encounters a driver OOM at k=75K and cannot be used for larger sketches. treeAggregate with finalAggregateOnExecutor=true (FAE) defers the final merge to an executor, reducing the driver peak at k=50K (298 MB vs. 488 MB) and extending the feasible range to k=250K, but it too fails by OOM at k=400K because the completed sketch must still be serialized and shipped to the driver afterward. treeAggRedux (Redux), by contrast, operates across the full range up to k=8M while holding the driver peak within 219–412 MB, and is the only strategy to survive to a sufficiently large k to reach a precision of 1.0, first achieved at k=2M. Redux ultimately fails at k=9M, but on the executor rather than the driver: the collapse step’s combOp momentarily holds two capacity-5k sketches before compressing them, and at k=9M this exceeds the executor heap even though the driver itself never approaches its limit. In other words, Redux does not eliminate the memory pressure created by a large mergeable sketch — it relocates it from the single, non-fungible driver to the many, horizontally-scalable executors. This is the outcome TreeRedux is designed for, but it also means executor memory becomes the new limiting resource; characterizing and relaxing that limit is a direction we leave for future work. The large gap in reachable k between FAE and Redux is explained by what each strategy sends to the driver. To accommodate the broad frontier of contenders described above, the sketch is allocated 5k counters rather than k; the output size is only the inner k items. Beyond the raw counter storage, our Space-Saving implementation uses the classical StreamSummary structure of Metwally et al. [17]: a HashMap[Int, SSCounter] maps each label to a counter object, and counters are additionally threaded into a doubly linked list of buckets ordered by count, so that both incrementing a counter and evicting the global-minimum counter are O(1) pointer operations rather than requiring a scan or a re-sort. Each counter object carries its label, its count, and three references — to its bucket and to the previous and next counter within that bucket’s list — well above the 12 bytes required for the raw (label, count) pair, and the bucket objects themselves add further overhead once amortized across however many counters currently share a given count. The result is that the live sketch for a capacity of 5k can occupy substantially more than 5k × 12 bytes in the driver heap. With FAE, the final merged sketch of capacity 5k must be fully materialized on the driver before any extraction can occur. With Redux, finalize runs on the executor and extracts only the top-k result into two primitive Array[Int] — one for labels, one for counts — costing 8k bytes before the result crosses the network: 4 bytes per label plus 4 bytes per count, narrowing the sketch’s internal 64-bit count (needed so a single counter cannot overflow for arbitrarily large n) to a 32-bit count in the exported result, which is safe here because every count is
TABLE VIII: Space-Saving sketch driver/executor memory and Redux accuracy vs. sketch size k (n=830,472,175, P =100, depth 2, logistic Cmax =30, r0 =25,165,824; cluster as in Section V). “OOM” marks the k that first experienced an out-ofmemory failure. For AGG and FAE, the driver experienced OOMs. Redux survived until an executor OOM. k
Driver Peak (MB)
Executor Peak (MB)
Heavy-Hitter Accuracy
Agg
FAE
Redux
Redux
Precision
Recall
Rank MAE
50K 488 75K OOM 100K — 250K — 400K — 500K — 750K — 1M — 2M — 3M — 5M — 8M — 9M —
298 301 298 506 OOM — — — — — — — —
219 220 221 237 238 242 259 273 326 345 348 412 —
2,418 3,884 3,520 2,837 6,208 5,998 5,372 6,492 9,519 10,753 12,602 15,727 OOM
0.7444 0.0022 0.7500 0.0033 0.7579 0.0045 0.7871 0.0116 0.8025 0.0189 0.8034 0.0236 0.8165 0.0360 0.8306 0.0488 1.0000 0.1176 1.0000 0.1764 1.0000 0.2939 1.0000 0.4703 — —
5,104K 4,977K 4,774K 4,152K 3,849K 3,869K 3,651K 3,419K 0 0 0 0 —
bounded by n ≪ 231 . The driver receives two flat primitive arrays with none of the boxing, hash-map, or linked-structure overhead described above. At k=250K — FAE’s last surviving point — the Redux result costs only ≈1.9 MB, while the observed driver peaks are 506 MB for FAE and 237 MB for Redux. The observed completion and OOM outcomes establish the feasible ranges: FAE’s last viable sketch size is k=250K, whereas Redux reaches k=8M, a 32× extension. In a separate set of smaller controlled experiments on our locally instrumented cluster — distinct from the EMR sweep reported in Table VIII, where every observed failure was a clean driver crash — we also observed the silent driver-hang failure mode described in the Introduction during heavy-hitters aggregation, using a Misra–Gries sketch under memory pressure at k=1,000,000 with the driver heap capped at 512 MB. See the appendix for the captured evidence and full classification. VIII. C ONCLUSION Spark’s tree-aggregation interface couples the intermediate aggregation state with the value returned to the driver. This coupling is harmless for scalar reductions but becomes a scalability and reliability problem when a large mergeable state is needed only to derive a small final answer. TreeRedux breaks that coupling through executor-side finalization, preserving Spark’s aggregation structure while returning only the final result to the driver. Across exact quantiles and heavy hitters, the experiments show that this small API change removes an important driver-side bottleneck and extends the usable range of otherwise driver-limited computations. The remaining executor-memory limits and the behavior of Spark’s GK Sketch implementation are important directions for future work. R EFERENCES [1] M. Stonebraker, “The case for shared nothing,” Database Engineering, vol. 9, no. 1, pp. 4–9, 1986.
[2] A. Patnam, “SPARK-36419: Optionally move final aggregation in rdd.treeaggregate to executor,” Apache Spark GitHub Pull Request #33644, 2021, associated with Apache Spark JIRA SPARK-36419, accessed 2026-06-19. [Online]. Available: https://github.com/apache/ spark/pull/33644 [3] M. Greenwald and S. Khanna, “Space-efficient online computation of quantile summaries,” in Proceedings of the 2001 ACM SIGMOD International Conference on Management of Data, ser. SIGMOD ’01. New York, NY, USA: Association for Computing Machinery, 2001, p. 58–66. [Online]. Available: https://doi-org.umiss.idm.oclc.org/10.1145/ 375663.375670 [4] J. Misra and D. Gries, “Finding repeated elements,” Science of Computer Programming, vol. 2, no. 2, pp. 143–152, 1982. [Online]. Available: https://www.sciencedirect.com/science/article/pii/0167642382900120 [5] G. Cormode and S. Muthukrishnan, “An improved data stream summary: The count-min sketch and its applications,” Journal of Algorithms, vol. 55, no. 1, pp. 58–75, 2005. [6] J. Dean and S. Ghemawat, “Mapreduce: Simplified data processing on large clusters,” in Proceedings of the 6th Symposium on Operating Systems Design and Implementation, 2004. [7] Apache Software Foundation, “Apache hadoop,” https://hadoop.apache. org, accessed 2026-08-13. [8] K. Shvachko, H. Kuang, S. Radia, and R. Chansler, “The hadoop distributed file system,” 2010 IEEE 26th Symposium on Mass Storage Systems and Technologies, 2010. [9] M. Isard, M. Budiu, Y. Yu, A. Birrell, and D. Fetterly, “Dryad: Distributed data-parallel programs from sequential building blocks,” in Proceedings of the 2nd ACM SIGOPS/EuroSys European Conference on Computer Systems (EuroSys ’07). New York, NY, USA: ACM, 2007, pp. 59–72. [10] Y. Yu, M. Isard, D. Fetterly, M. Budiu, Ú. Erlingsson, P. K. Gunda, and J. Currey, “DryadLINQ: A system for General-Purpose distributed Data-Parallel computing using a High-Level language,” in 8th USENIX Symposium on Operating Systems Design and Implementation (OSDI 08). San Diego, CA: USENIX Association, 2008. [Online]. Available: https://www.usenix.org/conference/osdi-08/dryadlinqsystem-general-purpose-distributed-data-parallel-computing-using-high [11] M. Zaharia, M. Chowdhury, M. J. Franklin, S. Shenker, and I. Stoica, “Spark: Cluster computing with working sets,” in 2nd USENIX Workshop on Hot Topics in Cloud Computing (HotCloud 10). Boston, MA: USENIX Association, Jun. 2010. [Online]. Available: https://www.usenix.org/conference/hotcloud10/spark-cluster-computing-working-sets [12] M. Zaharia, M. Chowdhury, T. Das, A. Dave, J. Ma, M. McCauly, M. J. Franklin, S. Shenker, and I. Stoica, “Resilient distributed datasets: A Fault-Tolerant abstraction for In-Memory cluster computing,” in 9th USENIX Symposium on Networked Systems Design and Implementation (NSDI 12). San Jose, CA: USENIX Association, Apr. 2012, pp. 15–28. [Online]. Available: https://www.usenix.org/conference/nsdi12/ technical-sessions/presentation/zaharia
[13] P. K. Agarwal, G. Cormode, Z. Huang, J. M. Phillips, Z. Wei, and K. Yi, “Mergeable summaries,” ACM Trans. Database Syst., vol. 38, no. 4, Dec. 2013. [Online]. Available: https://doi.org/10.1145/2500128 [14] I. Cao, J. Saloni, and D. Harrison, “A quick and exact method for distributed quantile computation,” in 2025 IEEE International Conference on Big Data (BigData), 2025. [15] Z. Karnin, K. Lang, and E. Liberty, “Optimal quantile approximation in streams,” in Proceedings of the 57th IEEE Symposium on Foundations of Computer Science (FOCS). IEEE, 2016, pp. 71–78. [16] T. Dunning and O. Ertl, “Computing extremely accurate quantiles using t-digests,” 2019. [17] A. Metwally, D. Agrawal, and A. El Abbadi, “Efficient computation of frequent and top-k elements in data streams,” in Database Theory ICDT 2005, T. Eiter and L. Libkin, Eds. Berlin, Heidelberg: Springer Berlin Heidelberg, 2005, pp. 398–412. [18] P. Flajolet, É. Fusy, O. Gandouet, and F. Meunier, “Hyperloglog: The analysis of a near-optimal cardinality estimation algorithm,” in Proceedings of the 2007 Conference on Analysis of Algorithms (AofA), ser. Discrete Mathematics and Theoretical Computer Science Proceedings, 2007, pp. 127–146. [19] S. Heule, M. Nunkesser, and A. Hall, “Hyperloglog in practice: Algorithmic engineering of a state of the art cardinality estimation algorithm,” in Proceedings of the 16th International Conference on Extending Database Technology (EDBT), 2013, pp. 683–692. [20] N. Alon, Y. Matias, and M. Szegedy, “The space complexity of approximating the frequency moments,” in Proceedings of the 28th Annual ACM Symposium on Theory of Computing (STOC). ACM, 1996, pp. 20–29. [21] A. Z. Broder, “On the resemblance and containment of documents,” in Proceedings of Compression and Complexity of Sequences 1997. IEEE, 1997, pp. 21–29. [22] B. H. Bloom, “Space/time trade-offs in hash coding with allowable errors,” Communications of the ACM, vol. 13, no. 7, pp. 422–426, 1970. [23] L. Fan, P. Cao, J. Almeida, and A. Z. Broder, “Summary cache: A scalable wide-area web cache sharing protocol,” IEEE/ACM Transactions on Networking, vol. 8, no. 3, pp. 281–293, 2000.
A PPENDIX I MPLEMENTATION D ETAILS AND P ERFORMANCE VALIDATION Because treeAggRedux reuses nearly all of Spark’s existing tree aggregation machinery, any performance difference should be attributable to the additional semantic capability rather than to changes in the aggregation algorithm itself. We use a non-inferiority test to evaluate whether treeAggRedux introduces measurable overhead relative to treeAggregate with finalAggregateOnExecutor=true. Specifically, non-inferiority is established if the upper bound of the 95% bootstrap confidence interval on the median runtime ratio remains below a 10% margin. The observed median runtime ratios were slightly below 1.0 for this workload, although the purpose of this experiment is to establish non-inferiority rather than superiority. We implemented treeAggRedux by copying the structure of Spark’s source and replacing the terminal partiallyAggregated.fold( copiedZeroValue)(cleanCombOp) call, which ships U to the driver, with a SinglePartitioner foldByKey followed by finalize executed on the executor. The differences are small and intentional, but we felt it was necessary to demonstrate that our implementation was no slower using different aggregation sizes and different tree depths.
Since we are only comparing tree aggregation performance, we use a finalize that immediately returns the passed value. The finalize is effectively the identity function and thus should not appreciably affect aggregation performance. To perform the non-inferiority test, we assume a simple problem that generates equal-length arrays for aggregation and at each aggregation performs a linear operation and discards down to the size of a single array. This merge-and-trim process repeats until only k elements remain. A natural problem that has this aggregation structure is the top-k selection problem (as distinguished from the top-k heavy-hitter problem studied in Section VII). It differs from the more general selection problem in that we assume k is small compared to n and thus may use a binary min-heap to keep the top-k as we sweep within each partition, pop all elements in the heap to order it, and then each combOp performs a linear in-order merge and returns the top k, discarding the remainder. Finding the smallest k within each partition using a min-heap has higher time complexity O(ni log k) than performing a quick select O(ni ), where ni is the number of elements in the i-th partition, but avoids the necessity of materializing the entire partition in memory. Workload. We use the same number of nodes and experimental configuration as the other experiments in this paper, except for the instance type: 30 EMR core nodes of type m6g.xlarge. We populate the cluster with n = 150 M integers distributed evenly across P = 128 partitions. Both methods return the full size-k array to the driver, so driver memory and network transfer are identical. Trials were run across four cells: two result sizes (k ∈ {100 K, 1 M}) × two tree depths (d ∈ {1, 7}). Protocol. Each trial submits both methods in randomly interleaved order to control for time-varying cluster state. Run 1 of each method is discarded as a cold JVM warmup; runs 2–10 are used for inference (9 warm pairs per cell). Between each trial, we randomly shuffle the elements in each partition to reduce the effect of order-specific behaviors. Table IX summarizes the paired bootstrap results. All four cells satisfy the non-inferiority criterion: the upper 97.5th percentile of the bootstrap distribution for the median paired runtime ratio remains at or below 1.034, well below the 10% margin. The observed warm median ratios range from 0.99 to 1.01, indicating that treeAggRedux introduces no measurable aggregation overhead relative to treeAggregate with finalAggregateOnExecutor=true for these workloads. To confirm the result is not an artifact of input order, the experiment was repeated with a per-trial Fisher-Yates shuffle applied independently to each partition before each submission. All four cells again satisfy the non-inferiority criterion. C LASSIFICATION OF S ILENT D RIVER FAILURES The Introduction describes a second driver failure mode observed during our experiments: the driver becomes unresponsive under memory pressure without producing the clean process-exit signal of a crash. This appendix describes the
TABLE IX: Bootstrap 95% CI on the median paired ratio (Redux / FAE) for each experimental cell (shuffled replication, redux_vs_fae-20260621-183029-nodes30). Non-inferiority margin δ = 1.10. npairs = 9 warm runs per cell. k
d
npairs
Redux (ms)
FAE (ms)
Ratio
CI 2.5%
CI 97.5%
Result
100K 100K 1M 1M
1 7 1 7
9 9 9 9
308.5 412.0 1257.5 756.5
298.5 402.5 1261.5 762.5
0.9889 0.9949 1.0093 0.9877
0.9698 0.9698 0.9866 0.9784
1.0335 1.0065 1.0295 1.0252
PASS PASS PASS PASS
underlying mechanism, how we detect and classify it, and the instances we directly confirmed. a) Mechanism.: Spark runs several background driverside threads independently of the main application thread – notably the dag-scheduler-event-loop thread, which processes all task completion and stage-transition events, and the task-result-getter-N threads, which deserialize incoming task results. An OutOfMemoryError is a Throwable, not an Exception; when one of these threads throws it, Spark’s normal error handling does not intercept it, and the thread dies silently while the driver process itself remains alive. Because these threads are singletons responsible for delivering completion events, the driver is left holding open a blocking call (e.g., a fold or collect) that can never be satisfied. No further log output is produced, and the process never exits on its own. This is mechanistically distinct from the clean, self-terminating driver crash (-XX:OnOutOfMemoryError-triggered) used elsewhere in this paper to positively identify ordinary driver OOMs. b) Detection.: On our heavily instrumented local cluster, we built an automated watchdog that polls each driver’s log for these specific thread-death signatures and, failing a match, enforces a hard wall-clock timeout (600 s) as a fallback. Both trigger paths kill the driver process and record a result distinct from an ordinary completion, so the failure is captured automatically as part of the normal experiment pipeline rather than requiring manual intervention. EMR has no equivalent automated detector; the EMR instances described below were identified by manually inspecting driver logs retrieved from S3 after the fact. c) Representative example.: The clearest documented instance occurred on 2026-06-16, during a heavy-hitters sweep (Misra–Gries, k = 1,000,000, sketch capacity 5k, driver heap 512 MB). The driver log shows: Exception in thread "dag-scheduler-event-loop" java.lang.OutOfMemoryError: Java heap space at MisraGries.trim (MisraGries.scala:64) at MisraGries.merge (MisraGries.scala:52) followed by no further output; the process required an external
TABLE X: Directly confirmed instances of the silent driverthread-death failure mode, by workload and compute backend. Workload Quantile (GK Select fan-in) Heavy hitters (Space-Saving / Misra–Gries)
Local cluster
EMR
10 2
3 0
kill. This incident directly motivated building the watchdog described above. d) Confirmed instances.: Table X summarizes every instance we directly confirmed via captured exception text or an unambiguous silent-termination signature (a substantial, otherwise error-free log that stops mid-stream with no exception, crash banner, or shutdown hook). These counts are a lower bound: they reflect the specific experiment logs we searched, not an exhaustive census of every trial we ever ran, and our local cluster’s automated watchdog additionally classified a number of wide-k heavy-hitters sweep trials as driver failures via this same mechanism without our independently reverifying each one’s captured exception text. The quantile instances span depths 1, 2, and 6 across independent seeds (local cluster) and all three seeds of the corrected depth-4, n = 25B configuration (EMR). We found no heavy-hitters instances among the 102 EMR driver logs we inspected across five wide-k sweeps, so on the evidence available this failure mode is confirmed on EMR for the quantile workload but not, so far, for heavy hitters.