advanced

Windowed Aggregations

10 min readLast updated: 2026-07-09

Overview

Learn how to aggregate streaming data over time intervals using Tumbling (non-overlapping) and Sliding (overlapping) window queries.

What You Will Learn

In this lesson, you will learn:
  • Tumbling Windows: Fixed, contiguous time ranges.
  • Sliding Windows: Overlapping time intervals.
  • Window Aggregations: Calculating metrics per time window.

Detailed Concept Explanation

Structured Streaming supports aggregations over time-based windows, grouping events based on when they occurred.

Window Types

  1. Tumbling Windows: Fixed-size, non-overlapping time intervals (e.g. every 10 minutes).
  2. Sliding Windows: Fixed-size, overlapping time intervals (e.g. 10-minute windows sliding every 5 minutes). A single event can fall into multiple sliding windows.

Code Examples

Python (PySpark) Implementation

python
from pyspark.sql import SparkSession
from pyspark.sql.functions import window

spark = SparkSession.builder.appName("Windowing").getOrCreate()

# Stream sensor data
sensor_stream = spark.readStream.format("rate").load()

# Calculate sliding window counts
windowed_counts = sensor_stream.groupBy(
    window(sensor_stream.timestamp, "10 minutes", "5 minutes")
).count()

query = windowed_counts.writeStream \
    .format("console") \
    .outputMode("complete") \
    .start()

query.awaitTermination(5)
query.stop()

Expected Output

text
+------------------------------------------+-----+
|                                    window|count|
+------------------------------------------+-----+
|[2026-07-09 19:20:00, 2026-07-09 19:30:00]|   10|
+------------------------------------------+-----+

Execution Plan Diagram (Python & Scala)

Execution Plan Diagram
SparkSession.builder
readStream
groupBy(window(timestamp 10m 5m)).count()
writeStream.outputMode(complete)
start()

Scala Implementation

scala
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.window

val spark = SparkSession.builder().appName("WindowingScala").getOrCreate()

val sensorStream = spark.readStream.format("rate").load()

val windowedCounts = sensorStream.groupBy(
  window($"timestamp", "10 minutes", "5 minutes")
).count()

val query = windowedCounts.writeStream
  .format("console")
  .outputMode("complete")
  .start()

query.awaitTermination(5000)
query.stop()

Expected Output

text
+------------------------------------------+-----+
|                                    window|count|
+------------------------------------------+-----+
|[2026-07-09 19:20:00, 2026-07-09 19:30:00]|   10|
+------------------------------------------+-----+

Common Mistakes

  • Failing to use Complete/Update Mode: Running windowed aggregations with the default Append mode without setting a watermark. Grouped aggregations must use Complete or Update mode to output incremental changes.

Interview Perspective

What is the difference between tumbling and sliding windows in Spark?

Tumbling windows are contiguous, fixed-size, and non-overlapping. For example, a 10-minute window will group data from 12:00-12:10, 12:10-12:20, etc. Sliding windows are fixed-size but overlap based on a slide interval. For example, a 10-minute window sliding every 5 minutes will group data from 12:00-12:10, 12:05-12:15, etc., meaning a single record can fall into multiple windows.


Related Topics