A window running-total over full history reruns nightly; how do you make cost grow with new data only?
answer
- same work every night for one new day
- cost grows quadratically with elapsed time
- associativity lets a seed replace history
- the watermark buys correctness against late data
- reconcile periodically or do not trust it
basics
~20 sCheckpoint the state. Persist the last cumulative value per key at a watermark, then run the window only over rows newer than it and offset each partition by its stored seed. Cost then tracks the daily delta instead of all history.
solid answer
~50 sA full-history window reshuffles and re-sorts everything every night, so cost grows with the square of time rather than with the data you added. The fix is a seed-and-delta design: store, per key, the cumulative value as of a watermark; each run reads only rows after that watermark, computes the window within the delta, and adds the stored seed to each partition. This works because running sums and counts are associative — a local prefix plus an offset equals the global prefix. Three things decide whether it is worth building. **Late-arriving data**: choose a bounded restatement window (recompute the last N days) and keep the watermark behind it. **Non-decomposable functions**: medians, exact distinct counts and global ranks cannot be seeded this way and either need recomputation or an approximation. **Operational cost**: if a full recompute finishes cheaply inside its schedule, keep it — the incremental version adds state, backfill logic and a reconciliation job.
code
sql · 14 lines-- Delta run: window over new rows only, offset by the stored seed per key
WITH delta AS (
SELECT l.account_id, l.event_ts, l.amount
FROM ledger l
JOIN checkpoint c ON c.account_id = l.account_id
WHERE l.event_ts > c.as_of
)
SELECT d.account_id,
d.event_ts,
c.running_balance
+ sum(d.amount) OVER (PARTITION BY d.account_id ORDER BY d.event_ts)
AS running_balance
FROM delta d
JOIN checkpoint c ON c.account_id = d.account_id;go deeper
Understand that recomputing a running total over all history every night repeats yesterday's work, and that starting from a saved total plus the new rows gives the same answer for far less work.
Explain why the shortcut is legal: a running sum is associative, so a stored cumulative value plus the local running value over new rows equals the true cumulative value. Note that this fails for medians and exact distinct counts.
Own the mechanics end to end — keyed checkpoint, data-derived watermark, idempotent delta writes, per-key repair for corrections — and be able to say which columns you would refuse to make incremental.
Frame it as a cost-versus-operability trade. Quantify the growth curve, name reconciliation as the price of trusting incremental state, and be explicit about the condition under which you would keep the simple full recompute instead.
## The cost curve you are fixing A nightly job that computes `sum(amount) OVER (PARTITION BY account_id ORDER BY event_ts)` across the entire ledger does the same work on day 400 as on day 399, plus a day. Every run shuffles the whole table, sorts every partition end to end, and rewrites every row — while the business only added one day of events. Total cost over a year is quadratic in elapsed time, which is why these pipelines quietly become the largest line on a warehouse bill and eventually stop fitting in their schedule. The target is a run whose cost is proportional to the delta: read new rows, compute over new rows, write new rows. ## Seed and delta The enabling property is that a running aggregate is *decomposable*. For an associative accumulator, the cumulative value at a row equals the cumulative value at some earlier checkpoint plus the local running value computed from the checkpoint onward. So: 1. Maintain a checkpoint table keyed by the partition key, holding the cumulative value and the watermark timestamp it is valid as of. 2. Each run selects rows with `event_ts > watermark`. 3. Compute the window over that delta only, partitioned and ordered the same way. 4. Add the seed for each key to every delta row's local running value. 5. Write the new rows and update the checkpoint to the new maximum per key. The shuffle and sort now touch delta-sized data, and the checkpoint join is a small dimension-style join keyed on the partition key. Keys with no new activity are not touched at all, which in most real datasets is the large majority. ## Which functions survive the treatment - **Fully decomposable**: SUM, COUNT, MIN, MAX, and anything built from them (AVG as sum and count). Seeding is exact. - **Row numbering per key**: also seedable — carry the last row number per key as the offset — as long as the ordering never changes retroactively. - **Not decomposable**: exact `COUNT(DISTINCT)`, medians and other exact quantiles, and any global ranking whose value depends on rows that may still arrive. Some of these have mergeable approximate forms, but choosing one is a precision decision to make explicitly with the consumer, not a silent substitution. - **Look-ahead frames**: a value that depends on *future* rows cannot be finalised at the watermark at all. Those rows must stay provisional until the frame closes. ## Late data is the real design problem Incremental state is only correct if the past is immutable. It usually is not: events arrive late, corrections are issued, a source replays a day. Two mechanisms handle it. A **bounded restatement window** keeps the watermark deliberately behind wall-clock time — say three days — so anything arriving inside that lag is picked up by the normal delta run. Rows older than the lag are treated as final, and the pipeline publishes that guarantee as an SLA rather than pretending it is unbounded. A **repair path** handles the rest: when a correction lands outside the window, invalidate the affected keys' checkpoints and recompute those keys from the beginning. Because the recompute is per key rather than per table, it stays cheap. This is why the checkpoint should be keyed and independently repairable rather than a single global marker. ## Correctness discipline Incremental state drifts. Anything that makes the delta run twice, skip a batch, or read a partially loaded partition corrupts the seed, and — unlike a full recompute — the error persists forever. Three safeguards are standard: make the run idempotent (delete-and-rewrite the delta partition rather than blind-append), derive the watermark from the data rather than from the clock, and schedule a periodic full recompute that reconciles against the incremental output and alarms on divergence. If the reconciliation job cannot be afforded, the incremental design cannot be trusted either. ## When not to do this The honest default is to keep the full recompute. It is one query, obviously correct, self-healing after any upstream fix, and needs no state. Move to incremental when at least one of these is true: the run no longer fits its schedule; its cost is a material share of the platform bill; or the data volume is growing fast enough that the first two are a quarter away. Otherwise the incremental version trades a compute bill you can measure for an operational surface — state, backfills, reconciliation, on-call understanding — that you cannot. A middle option often wins: keep the full recompute but bound its input. Many running totals are only consumed for a recent horizon, in which case anchoring at a periodic snapshot (an opening balance per key per month) caps the scanned history at one period without any incremental state at all. It is less elegant and much easier to operate. ## What to present in an interview Name the cost curve, name decomposability as the property that makes seeding legal, name late data as the constraint that decides the watermark, and name the reconciliation job as the price of admission. Then say plainly under what condition you would not build it — that last part is what separates a considered design from an enthusiastic one.
- Which window computations cannot be seeded from a checkpoint, and what do you do about them?Exact distinct counts, medians and other exact quantiles, and any ranking whose position depends on rows that may still arrive. Options are to recompute those columns over a bounded horizon, to substitute a mergeable approximate sketch with an agreed error bound, or to publish them at a coarser refresh cadence than the decomposable ones. The choice is a conversation with the consumer, not a silent substitution.
- How do you keep an incremental window correct when a correction arrives for a month-old row?Keep the checkpoint keyed by partition so repair is per key: invalidate that key's seed and recompute it from the beginning, leaving every other key untouched. Routine lateness is absorbed differently — hold the watermark a few days behind wall clock so ordinary late arrivals fall inside the normal delta run, and publish that lag as the pipeline's guarantee.
- When would you keep the full recompute despite the cost curve?When it still finishes comfortably inside its schedule and its cost is not a material share of the bill. A full recompute is self-healing after any upstream fix and carries no state; the incremental version adds checkpoints, backfill paths and a reconciliation job that someone must own. Convert when the run stops fitting or the growth rate says it will within a quarter.
A bank does not re-add every transaction since the account opened to print this month's statement; it starts from last month's closing balance and adds what came since.
saying these in an interview costs you the question
- Adds incremental state before checking the run still fits its schedule
- Assumes all window functions can be seeded, including medians and distinct counts
- Treats history as immutable and has no plan for late or corrected rows
- Derives the watermark from wall-clock time rather than from the data
- Ships incremental state with no periodic reconciliation against a full recompute