Rolling windows in Structured Streaming
A rolling window computes time-based aggregations over a continuously advancing frontier. Unlike tumbling or sliding windows, which assign each event to a fixed time bucket, a rolling window emits updated aggregates as the frontier moves forward and older events fall out of range. Use rolling windows for real-time feature serving and low-latency monitoring, where you need the current value of an aggregate over a recent lookback period.
Rolling windows support the following:
- Bounded windows that aggregate over a fixed lookback duration, such as the last 10 minutes.
- Lifetime windows that accumulate all data since the stream started.
- Multiple windows in one operator, so you can compute different lookback durations in a single call.
Requirements
- Databricks Runtime 19 and above.
- Micro-batch mode: classic compute (Standard or Dedicated access mode) or serverless compute.
- Real-time mode: classic compute (Standard or Dedicated access mode). In Lakeflow pipelines, real-time mode also runs on serverless compute through pipeline configuration. See Use real-time mode in Lakeflow pipelines.
For more information, see Execution modes.
How rolling windows work
A rolling window query groups events by one or more partition columns, orders them by a timestamp column, and computes aggregates relative to a frontier. Each output row includes the frontier value as a timestamp column.
Frontier
The frontier defines the current "now" for the operator. It advances based on processing time (wall clock), and can optionally lag behind by a configurable delay. A delay is useful to absorb slight ingestion delays so that in-flight events are counted before the frontier passes them.
Define the frontier with FrontierSpec:
# Frontier with a 5-second delay, emitted as a column named "window_time"
RollingWindow.FrontierSpec(delay="5 seconds", alias="window_time")
Bounded windows
A bounded window aggregates data within a fixed duration, looking backward from the frontier. Define the lookback with .over():
# Sum of "value" over the last hour
F.sum("value").over(RollingWindow.preceding(RollingWindow.Range("1 hour")))
Lifetime windows
A lifetime window accumulates data from the start of the stream with no expiration. To create a lifetime window, omit .over():
# Running count of non-null values since the stream started
F.count("value")
Execution modes
Rolling windows run in both micro-batch mode and real-time mode. The frontier behaves differently in each:
- Micro-batch mode: The frontier is fixed for each micro-batch. Each partition key emits at most one output row per batch.
- Real-time mode: The frontier advances continuously during a batch. A partition key can emit multiple output rows within a single batch as the frontier moves.
Late events and watermarks
A rolling window sets the window boundary from the processing-time frontier, but you can apply an event-time watermark to drop late-arriving events. Add .withWatermark() before rollingWindow(), using the same column you pass to orderBy:
df.withWatermark("event_time", "5 seconds").rollingWindow(
partitionBy="user_id",
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(delay="0 seconds"),
measures=[F.sum("value").over(RollingWindow.preceding(RollingWindow.Range("10 minutes")))],
)
An event whose event time is at or below the watermark is dropped, even when it falls within the window's time range. Without a watermark, every event that falls within the window is aggregated, including late arrivals. Watermarks apply in both micro-batch and real-time mode, and the watermark position is checkpointed, so it's restored when a query restarts.
Syntax
Call rollingWindow() on a streaming DataFrame:
- Python
- Scala
from pyspark.sql.streaming.rolling_window import RollingWindow
from pyspark.sql import functions as F
df.rollingWindow(
partitionBy="user_id",
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(
delay="0 seconds",
alias="frontier",
),
measures=[
F.sum("value").over(RollingWindow.preceding(RollingWindow.Range("10 minutes"))),
F.count("value").alias("non_null_value_count"), # No .over() = lifetime window
],
)
import org.apache.spark.sql.streaming.RollingWindow
import org.apache.spark.sql.streaming.RollingWindow.{FrontierSpec, Range, preceding}
import org.apache.spark.sql.functions._
df.rollingWindow(
partitionBy = Seq("user_id"),
orderBy = "event_time",
frontierSpec = FrontierSpec(delay = "0 seconds", alias = "frontier"),
measures = Seq(
sum("value").over(preceding(Range("10 minutes"))),
count("value")
)
)
The rollingWindow() operator takes the following parameters:
Parameter | Description |
|---|---|
| Required. One or more columns to group events by. |
| Required. The timestamp column that orders events. Must be a |
| The frontier definition. Set |
| A list of aggregate expressions. Add |
Supported aggregate functions
Rolling windows support the following aggregate functions:
sumcountavgminmaxstddev, including sample (stddev_samp) and population (stddev_pop) variantsvariance, including sample (var_samp) and population (var_pop) variantsapprox_count_distinctapprox_percentile- Scala
AggregatorUDAFs registered withfunctions.udaf. Supported on classic compute only. For an example, see Example 3: Scala UDAF.
Examples
Example 1: Per-user rolling revenue over the last hour
The following example computes a per-user rolling revenue sum over the last hour and writes updates to Kafka:
from pyspark.sql.streaming.rolling_window import RollingWindow
from pyspark.sql import functions as F
result = events_df.rollingWindow(
partitionBy="user_id",
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(delay="0 seconds"),
measures=[
F.sum("revenue")
.over(RollingWindow.preceding(RollingWindow.Range("1 hour")))
.alias("rolling_revenue"),
],
)
(
result.select(
F.col("user_id").cast("string").alias("key"),
F.to_json(F.struct("__frontier", "rolling_revenue")).alias("value"),
)
.writeStream.format("kafka")
.option("kafka.bootstrap.servers", "<host1:port1,host2:port2>")
.option("topic", "<topic-name>")
.outputMode("update")
.start()
)
Example 2: Mixed bounded and lifetime windows
The following example computes two bounded aggregates and one lifetime aggregate in a single operator:
result = events_df.rollingWindow(
partitionBy=["region", "product"],
orderBy="event_time",
frontierSpec=RollingWindow.FrontierSpec(delay="10 seconds", alias="window_ts"),
measures=[
F.sum("revenue").over(RollingWindow.preceding(RollingWindow.Range("1 hour"))),
F.avg("price").over(RollingWindow.preceding(RollingWindow.Range("10 minutes"))),
F.count("value"), # Lifetime: non-null values since the stream started
],
)
The output schema is (region, product, window_ts, sum(revenue), avg(price), count(value)).
Example 3: Scala UDAF
Scala Aggregator UDAFs are supported only on classic compute.
The following example defines a custom average as a Scala Aggregator, registers it with udaf, and uses it as a measure over a 10-minute bounded window:
import org.apache.spark.sql.{Encoder, Encoders}
import org.apache.spark.sql.expressions.Aggregator
import org.apache.spark.sql.functions._
import org.apache.spark.sql.streaming.RollingWindow.{FrontierSpec, Range, preceding}
case class AvgBuffer(sum: Long, count: Long)
object CustomAverage extends Aggregator[Long, AvgBuffer, java.lang.Double] {
override def zero: AvgBuffer = AvgBuffer(0L, 0L)
override def reduce(buffer: AvgBuffer, value: Long): AvgBuffer =
AvgBuffer(buffer.sum + value, buffer.count + 1)
override def merge(left: AvgBuffer, right: AvgBuffer): AvgBuffer =
AvgBuffer(left.sum + right.sum, left.count + right.count)
override def finish(buffer: AvgBuffer): java.lang.Double =
if (buffer.count == 0) null else buffer.sum.toDouble / buffer.count
override def bufferEncoder: Encoder[AvgBuffer] = Encoders.product
override def outputEncoder: Encoder[java.lang.Double] = Encoders.DOUBLE
}
val customAvg = udaf(CustomAverage)
val result = events.rollingWindow(
partitionBy = Seq("user_id"),
orderBy = "event_time",
frontierSpec = FrontierSpec(delay = "0 seconds", alias = "frontier"),
measures = Seq(
customAvg(col("value").cast("long"))
.over(preceding(Range("10 minutes")))
.as("custom_avg_10m")
)
)
Limitations
- Use
updateoutput mode. Theappendandcompleteoutput modes aren't supported. - Each measure can have at most one
.over()clause. You can't chain window specifications. - You can't combine rolling windows with regular Spark window functions. A single column expression can't use both
.over(WindowSpec)and.over(RollingWindow...). - The frontier advances by processing time only. Event-time frontiers aren't supported.
- Scala
AggregatorUDAFs are supported only on classic compute. - Databricks recommends using rolling window durations less than or equal to 7 days. Durations greater than 7 days might require substantial executor storage and cause slow performance or cluster failures.