Why does a stage that groups stream values into fixed-size batches usually also need an elapsed-time trigger?
answer
- each trigger bounds one side only
- size bounds bytes, time bounds waiting
- close on whichever fires first
- the oldest value sets the latency
- flush the partial group on completion
basics
~20 sA count-only group holds values until the count is reached, so when the source slows the partial group waits indefinitely and its oldest value goes stale. An elapsed-time trigger closes the group on age as well, bounding that wait.
solid answer
~40 sA grouping stage trades latency for batch efficiency, and each trigger bounds only one side of that trade. A size trigger bounds how big a group gets but says nothing about how long the first value of a group waits: at one value per minute, a group of a hundred takes over an hour to close. A time trigger bounds the wait but not the size: a fast burst can close an enormous group. Real pipelines therefore arm both and close on whichever fires first, so the group is at most `maxCount` values and at most `maxAge` old. The remaining edge is the tail - a partial group still open when the source ends must be emitted, not silently dropped.
code
pseudocode · 22 linesfunction batch(source, maxCount, maxAge):
buffer = []
timer = NONE
function flush():
if timer is not NONE:
cancel(timer)
timer = NONE
if buffer is not empty:
emit(copy_of(buffer))
buffer = []
on source value v:
if buffer is empty:
timer = schedule_after(maxAge, flush)
append v to buffer
if size(buffer) == maxCount:
flush()
on source end:
flush()
end_downstream()go deeper
Remember what a grouping stage costs: a value stops moving until its group closes, so something has to decide when that happens.
Explain that a size trigger bounds the group and a time trigger bounds the wait, and that arming both closes on whichever fires first.
Show the edges you have been bitten by: the partial group at completion, the buffered values lost on failure, the empty window, and the memory a sliding span holds.
Treat the two numbers as a published promise - maximum group size and maximum value age - and check them against measured arrival rates rather than defaults nobody revisits.
## What a grouping stage actually does A grouping stage turns a sequence of values into a sequence of collections of those values. Downstream sees fewer, larger items. That is worth doing when the consumer has a **per-item cost** it can amortise - a round trip, a transaction, a file handle - and it is never free, because a value that goes into a group stops moving until the group closes. So the only interesting question about such a stage is: **what closes the group?** ## The two triggers, and what each one bounds - **A size trigger** closes the group when it holds a fixed number of values. It bounds *group size* - and therefore the memory the stage holds and the work the consumer does per item. It bounds *nothing at all* about time: the wait for the first value of a group is `(count - 1)` further arrivals away, whatever that takes. - **An elapsed-time trigger** closes the group a fixed duration after it opened. It bounds *the age of the oldest value in the group*, which is the latency the consumer feels. It bounds nothing about size: a burst arriving inside one window produces one very large group. Each trigger leaves exactly the hazard the other covers. That is why they are usually armed together: | trigger armed | group size | wait for the oldest value | fails when | |---|---|---|---| | size only | bounded | unbounded | the source slows or stops mid-group | | elapsed time only | unbounded | bounded | the source bursts inside one window | | both, first to fire | bounded | bounded | neither; you have chosen both numbers deliberately | With both armed, the promise to the consumer is precise and testable: *no group exceeds N values, and no value waits more than T before its group closes.* ## The tail nobody tests Three edges around a grouping stage account for most of its bugs: 1. **Source completion with a partial group.** The values are already accepted; ending the sequence without emitting them silently loses the tail. A correct stage flushes on completion. 2. **Source failure with a partial group.** A failure signal normally terminates the sequence, and values buffered in the open group go with it. A grouping stage is a latency device, not durable storage - if losing the tail is unacceptable, the durability has to live somewhere else. 3. **An empty window.** If the elapsed-time trigger fires and nothing has arrived, emitting an empty group makes the consumer handle an item that carries no work. Decide deliberately whether empty groups are emitted or skipped, and write it down. ## When consecutive groups overlap Grouping is usually **disjoint**: every value lands in exactly one group, so a total computed per group and then added up equals the total over the sequence. A **sliding window** deliberately breaks that. Each group covers a span, and consecutive groups start closer together than that span, so a value near a boundary appears in more than one group. That is exactly what you want for a moving average or a rolling count, where every point should be computed over the last T of history rather than over an arbitrary slice of it. It costs two things: - **Retention.** The stage must hold values for the whole span, not only since the last emission, so its memory is set by the span and the arrival rate. - **Double counting.** A downstream consumer that sums across groups counts overlapped values once per group they appear in. Overlap is a property the consumer must know about; it cannot be inferred from the group contents. ## Choosing the numbers The age bound comes from the consumer's latency budget: it is the worst-case wait a value inherits from this stage alone, so it has to fit inside the end-to-end budget alongside everything else in the pipeline. The size bound comes from the consumer's per-item cost and the memory you are willing to hold: large enough that the amortised per-value cost is small, small enough that one group does not dominate the stage's footprint. Measure the arrival rate before picking either. If the typical rate fills a group well before the age bound, the size trigger is doing the work and the age bound is a safety net for quiet periods; if the age bound fires first most of the time, the group size you chose is aspirational and the consumer is not getting the amortisation you designed for.
- What changes for a downstream total when consecutive groups overlap instead of being disjoint?A value near a boundary appears in more than one group, so adding the group totals counts it once per group. Overlap also forces the stage to retain values for the whole window span rather than only since the last emission. The consumer must be told the groups overlap; it cannot tell from their contents.
- The source ends while a group is half full - what should the stage do?Emit the partial group, then end downstream. Those values were already accepted, and dropping them loses the tail of the sequence silently. Failure is different: a failure signal terminates the sequence and the buffered values usually go with it, which is why a grouping stage should not be relied on as durable storage.
- How do you pick the age bound?From the consumer's latency budget. The age bound is the worst-case delay this stage alone adds to the oldest value in a group, so it has to fit inside the end-to-end budget alongside every other stage. Then check it against the measured arrival rate to see which of the two triggers actually fires in practice.
saying these in an interview costs you the question
- Thinks a size trigger alone bounds how long a value waits
- Assumes an elapsed-time trigger also caps the group size
- Drops the partial group when the source completes
- Treats a grouping stage as durable storage for buffered values
- Sums across overlapping groups as if every value appeared once