How do you run a sweep-line peak-concurrency job over a billion events under a memory ceiling?
answer
- the sort is the only expensive half
- two representations of one step function
- how big is the coordinate domain?
- counters per coordinate need no sort
- shards must be seeded with the carried-in count
basics
~20 sPick the representation the data justifies. Sort-then-scan still works out of core as an external sort plus one streaming pass. If coordinates are bounded small integers, count deltas into a dense array over the domain instead — no sort at all.
solid answer
~50 sThe sweep is two decoupled steps, and only the first is expensive. In sorted-event form you build `2n` records, order them and scan once; beyond memory that becomes an external sort plus a single streaming pass, so the cost lens stops being comparisons and becomes passes over storage. It shards by coordinate range only if each shard is seeded with the count carried in from earlier ranges. In dense-domain form — coordinates that are small bounded integers, such as per-second buckets over a day or stop numbers on a route — you accumulate deltas into an array indexed by coordinate in one streaming pass, then scan that array: `O(n + U)` time, `O(U)` memory, no sort. Choose on `U` versus `n log n`, on whether `U` fits the ceiling, and on whether the team can operate a sharded sort pipeline at all.
go deeper
Know that the sweep's cost is the sort, not the scan, and that the scan itself needs only a running sum. That single fact is what makes the technique survive data far larger than memory.
Explain the two representations and their bounds: sort-then-scan at O(n log n) time and O(n) memory, versus counting into a bounded coordinate domain at O(n + U) time and O(U) memory with no sort.
Work the constraint: what happens out of core, why the cost lens becomes passes over storage, and what a shard must be seeded with before its local maximum is meaningful.
Defend a choice under a memory ceiling and a batch window, naming the measurements you would take first, and be willing to pick the asymptotically worse option when predictable memory and operability serve the organisation better.
## Two representations of the same idea Peak concurrency is the maximum of a step function built from `+w`/`-w` deltas. Nothing forces you to compute that maximum by sorting events; sorting is one representation, and at a billion events it stops being the obvious one. **Representation A — sorted event stream.** Materialise `2n` `(coordinate, delta)` records, order them, scan once with a running sum. Time `O(n log n)`, working memory `O(n)`. **Representation B — dense coordinate domain.** If coordinates live in a bounded integer universe of size `U`, allocate a counter per coordinate, add each delta into its bucket in one streaming pass over the raw events, then walk the buckets accumulating a running sum. Time `O(n + U)`, memory `O(U)`, and **no sort at all**. This is the difference-array representation of the same step function; the mechanism belongs to that pattern, but the *choice* between the two is the decision here. ## The decision variables 1. **`U` versus `n log n`.** Per-second buckets over a day are 86,400 counters; over a year, about 31.5 million. Both are small next to a billion events. Nanosecond timestamps, or floating-point coordinates, give an effectively unbounded `U` and rule B out unless you quantise — and quantising is a product decision about the resolution the metric is reported at, not a free optimisation. 2. **Does `U` fit the ceiling?** `O(U)` memory is a hard, predictable number you can compute in advance and put in a capacity plan. `O(n)` is not: it grows with traffic, which is exactly the direction that hurts. 3. **Streaming or batch?** B reads each raw event once and never revisits it, so it works on a live stream and on data too large to store. A must see all events before the scan can begin, so it is inherently batch. 4. **Do you already own the sort?** If a pipeline already delivers events ordered by coordinate — many log pipelines do, or nearly do — representation A costs one linear pass and the whole argument evaporates. Free ordering is the strongest argument for A there is. ## When you are stuck with A and it does not fit Out-of-core, the sweep decomposes cleanly: sort the event file with an external sort, then make one sequential streaming pass with the running sum. The scan needs `O(1)` memory, so *the sort is the entire memory and I/O problem*. The relevant cost lens shifts from comparison counts to passes over storage, and the honest interview answer is that at this size you are budgeting I/O, not CPU. Sharding by coordinate range works and is the usual production shape, with one subtlety that is the actual interview point: **a shard's local running sum starts at zero, but the true count entering that range is whatever was open at its left boundary.** Each shard must be seeded with the carried-in count — the sum of all deltas in every earlier range — before its local maximum means anything. Since the sum of deltas per shard is a single number, one cheap pass over shard-level totals produces every seed, and only then may you combine the shard maxima. Taking the largest of unseeded shard maxima is the classic wrong answer, and it under-reports in exactly the busiest ranges. ## The organisational half of the decision Asymptotics rarely decide this alone. A sharded external sort with carried-in seeding is several moving parts, each with a failure mode, and someone is on call for it at 3 a.m. A bucket-counting job over a fixed 86,400-slot domain is one loop that a new engineer reads correctly on the first try, has memory usage that does not change when traffic doubles, and degrades in an obvious way if the coordinate range assumption breaks. If both meet the batch window, the maintainable one wins, and being able to say *why* — operability and predictable memory beat a better constant factor — is the principal-level part of this answer. The reverse case is real too: if the metric must be reported at full timestamp resolution, or the coordinate domain is unbounded because coordinates are opaque identifiers rather than clock ticks, then B is not available at any price and you build the sorted pipeline properly, including the seeding. ## What to state explicitly before choosing Say out loud which numbers you would go and measure: `n`, the coordinate domain size and its resolution, the memory ceiling per worker, the batch window, and whether ordering already exists upstream. A choice defended by those five is a decision; a choice defended by "`O(n + U)` beats `O(n log n)`" is a guess that happens to have notation attached.
- You shard the sorted sweep by coordinate range. What must each shard receive besides its own events?The carried-in count: the total of all deltas in every earlier coordinate range. A shard's local sum starts at zero, but the true occupancy entering its range is whatever was open at its left boundary. Sum the per-shard delta totals in one cheap pass to produce every seed, then combine the corrected maxima.
- When is the dense counter representation simply unavailable?When the coordinate domain is unbounded or absurdly large relative to the event count — full-resolution timestamps, floating-point positions, opaque identifiers. Quantising to coarser buckets restores it, but that changes the resolution the metric is reported at, so it is a product decision rather than a private optimisation.
- Why might you choose the asymptotically worse option here?Because predictable memory and operability often matter more than a better exponent. A fixed-size bucket job has memory that does not move when traffic doubles and reads correctly on first sight; a sharded external sort with seeded shards has more parts, more failure modes and a longer on-call story. If both fit the batch window, take the boring one.
saying these in an interview costs you the question
- Assumes the whole event list must fit in memory
- Combines shard maxima without seeding the carried-in count
- Counts comparisons when the real cost is storage passes
- Picks the dense domain without checking its size or resolution
- Judges only asymptotics, ignoring operability and memory ceilings