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

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
    ],
)

Scala

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
partitionBy Required. One or more columns to group events by.
orderBy Required. The timestamp column that orders events. Must be a TimestampType column.
frontierSpec The frontier definition. Set delay to control how far the frontier lags behind processing time, and alias to name the output frontier column.
measures A list of aggregate expressions. Add .over() for a bounded window, or omit it for a lifetime window.

Supported aggregate functions

Rolling windows support the following aggregate functions:

  • sum
  • count
  • avg
  • min
  • max
  • stddev, including sample (stddev_samp) and population (stddev_pop) variants
  • variance, including sample (var_samp) and population (var_pop) variants
  • approx_count_distinct
  • approx_percentile
  • Scala Aggregator UDAFs registered with functions.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

Note

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 update output mode. The append and complete output 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 Aggregator UDAFs 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.