skip to content

Six months of stored records are pushed through a job that groups by event time in two hours - how do the results differ from the live run's?

level: seniorimportance: should knowfreq 42%

answer

  1. the claim moves with the data
  2. six months of it in seconds
  3. late records are no longer late
  4. parallel pieces read out of time order
  5. count what fell outside its group

basics

~20 s

The job's running claim about how far time has advanced is derived from the data, so it races through six months in seconds. Groups close back to back, and the replay can end up either more complete than the live run or less, depending on the order history is read in.

solid answer

~60 s

Grouping by event time - the timestamp carried on the record, which the pipeline assigned from a field, from arrival or from source metadata - means the job maintains a running claim that nothing older than some time is still expected. That claim advances with the data, not with the clock, so a replay advances it at reading speed: months pass in seconds and groups close in rapid succession. Two opposite outcomes follow. If history is read roughly in time order, the replay is **more** complete than the live run, because records that were too late to be counted live are simply present. If history is read by many pieces in parallel, each covering a different period, the claim can run ahead of records another piece has not reached yet, and those are treated as too old - the replay is then **less** complete in a way the live run never was. The remedy is to control the read order or to widen the tolerance for the duration of the replay.

go deeper

for a junior

Grasp that a job grouping by the timestamp on the record advances through time as fast as it reads, so a replay compresses months into minutes and closes groups far faster than the live job ever did.

for a middle

Explain both directions of divergence: records too late to count live are present in the archive, while pieces read in parallel out of time order can push the claim past records not yet read.

for a senior

Control the read order or widen the tolerance deliberately for the replay, count rejected records as evidence, and predict the emission and memory differences before comparing anything to the live output.

for a principal

Decide whether the platform stores history in a layout that can be replayed in time order at all, and what it is worth paying for that, given that every future correction and reprocessing depends on it.

## Why the claim races A job grouping by event time cannot know when a group is finished, so it maintains a claim: *no further input older than time T is expected*. Everything downstream of that claim - when a group closes, whether a record counts, when a result is emitted - moves when it moves. The crucial property for a replay is that the claim is derived from the data the job has read, not from the machine's clock. Read six months of history in two hours and the claim traverses six months in two hours. Read one piece of it in ten seconds and the claim traverses that period in ten seconds. Nothing is broken by this on its own. Groups close far faster than they did live, results appear in a burst, and the job finishes. The differences from the live run come from two directions at once. ## Direction one: the replay is more complete While live, a record that arrived after its group had closed was late: diverted, dropped, or applied as a correction depending on what the job was written to do. In a replay that record is in the stored history like any other, positioned by its own timestamp, present before the claim passes it. Groups that were short live come out whole. This is the archive being better than the original, and it has a consequence worth stating plainly: **a correct replay is expected to disagree with the live output**, and treating any disagreement as a defect will send you hunting a bug that does not exist. ## Direction two: the replay is less complete Stored history is normally laid out for efficient reading rather than for time order, and a replay reads many pieces of it at once. Each piece covers some period. A piece covering last spring may be read in seconds while the piece covering the same days from another source has not been opened yet. If the claim is allowed to advance to the fastest piece, it passes records the slower piece has not produced, and those records are then out of bounds - dropped or diverted, in volumes the live run never saw. This is the single most common surprise in a historical replay, and it is entirely an artefact of the read. Live, records arrived roughly in order because time itself sequenced them. In a replay nothing sequences them but the layout of the files. ## Remedies, in order of preference 1. **Read history in an order aligned with time.** One piece per period, periods processed in order. It costs parallelism within a period boundary and it removes the problem at the source. 2. **Slow the claim to the slowest piece.** Where the runtime derives one job-wide claim as the minimum across all readers, this happens naturally; where it does not, a replay can advance far ahead of a lagging reader. 3. **Widen the tolerance for the replay only.** Allowing a much larger amount of disorder costs memory and delays closing, but for a bounded replay both are affordable. Return it to the live setting afterwards. 4. **Count what was rejected.** A counter that every worker adds to and the coordinating process sums, reporting how many records fell outside their group, turns a silent loss into a number you can compare against the live run's. ## What else compression changes - **Far more is held at once.** Many periods are open simultaneously rather than one or two, so the memory a worker holds during a replay is not comparable to the live job's, and a replay may write its working set to local disk where the live job never did. - **Emission counts differ.** Where a runtime emits an updated running result on every arrival, the live run produced many rows per group over hours and the replay produces very few, because the group completed almost at once. Where a runtime emits once at close, the count matches and only the timing differs. Compare final values per key, not row counts. - **The three execution models behave differently.** A continuous job assembled from repeated small finite runs will pull vastly more into each slice during a replay than it ever did live, which changes its memory profile and its emission cadence. A record-at-a-time runtime with key-bound state runs the same code but with every period's state live at once. The two-phase disk-to-disk model has no running claim at all: the boundary is whatever the code computes from the records in the input, which is why a replay there reproduces the original grouping exactly and the other two do not. The honest summary to give an interviewer: replaying history does not reproduce the past, it recomputes it under different conditions, and the differences are predictable enough to be enumerated before the run rather than discovered after it.

  • During a replay you see a large number of records rejected as too old. What is the first thing to check?
    The order in which history is being read. If pieces covering different periods are read in parallel, the claim advances to whichever piece is fastest and passes records the slower ones have not produced yet. Check whether the claim is derived as the minimum across readers, and if reading in time order is not practical, widen the disorder tolerance for the duration of the replay and count the rejections.
  • Why does a replay of a job grouping by event time often need much more memory than the live job did?
    Live, only the currently open periods are held, because time advances slowly. In a replay months of periods are opened within minutes, so the job can hold state for a great many groups at once. Expect a working set that spills to local disk where the live job never did, and either widen the run or bound the replay to shorter periods processed in sequence.

saying these in an interview costs you the question

  • A replay should reproduce the live output exactly if the code is unchanged
  • The completeness claim advances at the same pace regardless of reading speed
  • Records rejected as too old during a replay indicate corrupt stored data
  • Assumes memory use during a replay resembles the live job's
  • Compares row counts between the two runs as a correctness check