Skip to main content
Version: Current

Aggregations in Time Windows

Mode: Streaming

What are aggregations in time windows?​

Aggregations in time windows let you calculate values such as count, sum, average, minimum or maximum over events that occurred during a period of time.

For example, you may want to calculate:

  • the number of transactions made by a customer during the last hour,
  • the total value of payments during each 10-minute period,
  • the number of actions performed during a user session.

The period of time over which events are collected for an aggregation is called a window.

For example, if you want to calculate how much a customer has spent during the last hour, the window contains that customer's transactions from the relevant one-hour period. One or more aggregation functions, such as Sum or Count, can then be applied to the events in that window.

At a high level, all aggregation components behave similarly:

  • they collect events into a time window and, when groupBy is configured, keep separate aggregates for each groupBy value,
  • they apply one or more aggregation calculations to the events belonging to each window and group,
  • they either create an aggregation event containing the calculated values, or enrich incoming events with the current calculated values.
note

In this documentation, record is a neutral term used across all processing modes. In streaming contexts, we typically use the term event — a record with a timestamp that lets Nussknacker apply time-based processing logic.

Which window to choose?​

Nussknacker supports three types of time windows. They differ mainly in how the time boundaries of a window are determined.

Window typeThink of it asExample question
SlidingA continuously moving period, such as the last 60 minutes; triggered by the current eventHow much has this customer spent during the last hour?
TumblingConsecutive, non-overlapping periods, such as 10:00–10:10, 10:10–10:20, etc.How many transactions occurred in each 10-minute period?
SessionA period of activity that ends after a configured period of inactivityHow much did this customer spend during one continuous period of activity?

Event time​

Events are assigned to time windows based on their event time. This means that the timestamp associated with an event determines where it belongs, rather than simply the moment when Nussknacker processes it. Read Notion of time to learn more.

Common aggregation configuration​

Regardless of the window type, you need to define which events should be considered together and what should be calculated from them.

A single time-window aggregation component can calculate several aggregate values over the same window and groups of events. For example, for each customer and the same one-hour window, you can calculate at once:

  • the sum of transaction amounts,
  • the average transaction amount,
  • the number of transactions.

The window and groupBy configuration are common to all these calculations. Each individual aggregate calculation is defined by a name, an aggregator input expression and an aggregator function.

ParameterQuestion it answersDescription
groupByWhich events should be aggregated together?One or more expressions that together identify the group to which an event belongs. A separate set of aggregate values is maintained for each combination of groupBy values.
For example, using #input.customerId calculates separate results for each customer.
Whenever a result is emitted, the values of the expressions forming groupBy are available as a list in the #key variable - see Example 2 below.
UI tip: After entering each groupBy expression, press Enter to add it. To edit an existing expression, click it.
output variable nameWhere should the results be stored?Name of the variable containing all aggregation results. For example, if it is set to myAggr, the results are available under #amyAggr.
Parameters for each aggregation calculation
nameWhat should this result be called?Name of one aggregation result within that variable. For example, if name is txTotal, the result is available as #myAggr.txTotal.
aggregator inputWhat value from each event should be used in this calculation?An expression evaluated for every event included in the window. Its result is passed to the selected aggregator. Each aggregate calculation can use a different aggregator input expression.
aggregatorWhat calculation should be performed?The function applied to the aggregator input values, for example Sum, Average, Count, or Max. Each aggregate calculation has its own aggregator.

Aggregation components also provide the #key variable, which contains the result of the groupBy expression.

Examples​

Sample data​

{"subscriberId":"1","transactionId":11,"operation":"RECHARGE","amount":500.00}
{"subscriberId":"2","transactionId":12,"operation":"RECHARGE","amount":200.00}
{"subscriberId":"1","transactionId":13,"operation":"TRANSFER","amount":5000.00}
{"subscriberId":"1","transactionId":14,"operation":"TRANSFER","amount":1000.00}

For the examples below, assume that the window for each group contains exactly the events shown in the sample data. We focus only on the aggregate values calculated from those events. The tables do not represent complete output events or specify when these values become available. That depends on the window type and its configuration.

Example 1: multiple aggregation calculations in one node​

Let's calculate three values for each customer: the total transaction amount, the number of transactions, and the list of transactions.

We want to maintain separate results for each subscriber, so we configure:

  • groupBy: #input.subscriberId
  • output variable name: aggr

Within the same aggregation node, we define several aggregate calculations:

nameaggregator inputaggregatorResult available as
txTotal#input.amountSum#aggr.txTotal
txCountN/ACount#aggr.txCount
txList{"tid": #input.transactionId, "val": #input.amount}List#aggr.txList

For the events shown above, these calculations produce the following values for each groupBy value:

#key (json notation)#aggr.txTotal#aggr.txList#aggr.txCount
["1"]6500.0[{"tid":11,"val":500.0}, {"tid":13,"val":5000.0}, {"tid":14,"val":1000.0}]3
["2"]200.0[{"tid":12,"val":200.0}]1

Example 2: grouping by more than one event property​

Using the same input events, suppose we want to calculate the maximum transaction amount separately for each subscriber and operation type.

We configure:

  • groupBy: #input.subscriberId #input.operation
  • output variable name: aggr2

and define the following aggregate calculation:

nameaggregator inputaggregatorResult available as
maxAmount#input.amountMax#aggr2.maxAmount

For the events shown above, the calculation produces the following values for each groupBy value:

#key (json notation)#aggr2.maxAmount
["1","RECHARGE"]500.0
["1","TRANSFER"]5000.0
["2","RECHARGE"]200.0

This example shows that groupBy does not have to identify only a single entity such as a subscriber. It can also define more specific groups, for example by combining the subscriber with the operation type.

Available aggregator functions​

  • ApproximateSetCardinality — estimates the number of distinct values using HyperLogLog. null is treated as a distinct value.
  • Average
  • Count
  • CountWhen — counts events for which aggregator input evaluates to true
  • First
  • Last
  • List — collects aggregator input values into a list
  • Max
  • Median
  • Min
  • Set — collects distinct aggregator input values into a set; large sets may consume significant memory
  • StddevPop — population standard deviation
  • StddevSamp — sample standard deviation
  • Sum
  • VarPop — population variance
  • VarSamp — sample variance

Analogy to SQL GROUP BY​

If you are familiar with SQL, the way groupBy, aggregator input, and aggregator work together is roughly analogous to a GROUP BY query.

SELECT
customerId,
SUM(amount),
AVG(amount)
FROM transactions
GROUP BY customerId

Here:

  • groupBy corresponds to the SQL GROUP BY,
  • each aggregator input + aggregator pair corresponds to an aggregate expression such as SUM(amount) or AVG(amount).

Unlike a regular GROUP BY query, Nussknacker performs these calculations only over events belonging to the relevant time window.

Additional considerations​

Window precision​

Aggregations in time windows are optimized to reduce the amount of state and computation required during stream processing. As a consequence, the boundaries of some windows are not evaluated with millisecond precision.

By default, Aggregate Sliding and Aggregate Session use a precision of 60 seconds. If seconds are specified in windowLength or sessionTimeout, respectively, the precision is 1 second. This means that an event may effectively remain in the aggregation for up to about one minute less than the configured window length. For a 60-minute window, the effective window length may be anywhere between just over 59 minutes and 60 minutes. For this reason, the configured window length should not be shorter than the precision used by the component.

The default slice lengths are:

aggregationTypeslice length
sliding60 seconds
session60 seconds
join60 seconds
tumblingwindowLength

If seconds are specified in windowLength or sessionTimeout, the precision is set to 1 second.* With 1-second precision, the amount of state can become significantly larger, requiring substantially more memory.

Internally, Nussknacker reduces resource consumption by precomputing aggregates in slices rather than maintaining every event individually until it leaves the window. A slice represents a short interval of event time. When a slice falls outside the current window, its contribution can be removed from the aggregate as a whole.

Handling out-of-order and late events​

In event-time processing, events may arrive out of order, so the system must decide how long to wait before considering a time window complete. Waiting too long delays results, while closing a window too early may cause late events to be ignored; Flink uses watermarks to make this trade-off, and understanding them helps explain why some aggregation or join results may be delayed or not produced as expected. A good description of the problem can be found here.

caution

Late events are dropped - data carried by the late events will not be included in aggregation computations.