intermediate 3 min answer

A live dashboard's "views in the last 30 minutes" reads about 30 times the true number. The job emits a 30-minute window advancing every minute; individual window values look correct when spot-checked; the sink is a time-series table and the panel sums the rows in the selected range. Lag and CPU are normal. What is happening and what do you change?

windowinghopping windowsdouble countingsinksreconciliation
Show the full answer Hide the answer

The first three things I would look at, and why in that order

  1. The ratio of window size to advance. An error factor that is a round number is almost never arithmetic drift; 30 is the number of one-minute hops inside a 30-minute window.
  2. How many rows the sink holds per key per minute. One row per window end means the writer is fine and the reader is wrong.
  3. The panel's aggregation, which here is a sum over rows rather than a read of the latest row.

The diagnosis

A hopping window with size 30 minutes and advance 1 minute places every input record into 30 distinct windows, because that is what overlap means. Each window's value is correct: it counts the views in its own 30 minutes. But the windows are not disjoint, so summing rows counts each event once per window it belongs to — exactly size / advance times. The dashboard is not reporting a corrupted number; it is adding up 30 overlapping answers to the same question.

The second, quieter effect is volume. At 20k events/s, a 30x overlap turns the output into 600k window updates per second and multiplies the job's window state by the same factor. A 30x bill for a value the reader wanted once.

The misleading signal

Spot-checking a single window always passes, which is what makes this survive review. The per-window numbers are right, the pipeline is healthy, and the fault lives in a contract nobody wrote down: whether the rows in the sink are disjoint facts or overlapping views.

The fix, in order

  1. Change the read, not the pipeline. Select the row for the most recent window end rather than summing. This is a one-line fix and it restores correctness today.
  2. Then change the pipeline, because the overlap is still paying 30x for one value. Emit one-minute tumbling windows and let the reader sum 30 disjoint rows, or compute the rolling figure in the serving layer. Tumbling output is disjoint, so a naive sum is correct by construction and the state per key drops to one window.
  3. Make the grain explicit in the sink: write window_start and window_end on every row, so a future reader cannot mistake overlapping rows for disjoint ones.

The alert that would have caught it earlier

An assertion on emitted rows per input record. For a tumbling window it is 1, for a hopping window it is size / advance, and any drift means the topology changed. A daily reconciliation against a batch count over the same period would also have flagged a 30x gap in the first day, which is the more general lesson: a streaming aggregate with no periodic comparison to an independent count will be wrong for as long as nobody happens to look.

When not to remove the overlap

When the consumer genuinely needs a smoothed value evaluated frequently — an alerting rule that must fire within a minute on a 30-minute condition — the hopping window is the correct tool and the 30x output is the cost of the alert's responsiveness. Choose it deliberately, and never let a naive SUM downstream of it.