Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Triggers and accumulation modes

Batch & Streaming · Streaming Design

Triggers and accumulation modes

Hardbatch-streaming-40
triggerswindowingapache-beamflink

Question

In Beam or Flink, what are triggers and accumulation modes?

Solution

In frameworks like Apache Beam and Apache Flink, a trigger determines the precise moment when a window evaluates its accumulated records and emits an output pane. Accumulation modes govern whether successive panes retain previously computed state to emit cumulative totals or discard past state to emit only incremental deltas. Because triggers can emit several panes for a single window, downstream systems must be architected to handle multiple updates per time boundary without double-counting records.

How triggers control window emission

A window defines the logical grouping of records in event time, but it does not dictate when outputs appear. Triggers provide this control across three operational stages:

  • Early speculative firings: Emitting partial results at periodic processing-time intervals (such as every sixty seconds) allows operational dashboards to display emerging trends before a long window concludes.
  • On-time firings: Firing when the watermark crosses the window end boundary ensures that all on-time data contributes to an authoritative baseline calculation.
  • Late firings: Emitting whenever straggling records arrive after the watermark but within allowed lateness allows the pipeline to incorporate delayed arrivals.

For example, on a one-hour tumbling window, a pipeline can emit speculative updates every minute, an on-time result when the watermark passes, and occasional late corrections whenever stragglers appear.

Accumulating versus discarding panes

When a trigger fires multiple times for one window, the accumulation mode determines how consecutive outputs represent state:

  • Accumulating mode: The window retains all previous data. Each emitted pane represents the full running aggregate from the window start (for example: 10 orders, then 25 orders, then 40 orders).
  • Discarding mode: The window resets its state after each firing. Each emitted pane represents only the increment that arrived since the prior trigger (for example: +10 orders, then +15 orders, then +15 orders).

Downstream integration requirements

Multiple pane emissions place specific requirements on downstream storage sinks.

If you emit accumulating panes, your destination sink must be update-capable, using primary key upserts into a database or merge operations in a lakehouse table to overwrite previous window totals. If you emit discarding panes into an append-only log or message bus, downstream consumers must apply delta summation logic to construct correct overall totals.

PreviousNext