CertSafari
    Databricks Certified Associate Developer for Apache Spark· Lessons

    Domain 5 · Lesson 27/32

    Windowed Aggregations on Streaming DataFrames with window()

    Perform basic operations on Streaming DataFrames and Streaming Datasets, such as selection, projection, window and aggregation.

    9 min read
    3.12% of exam
    2 sources
    Published 3 Oct 2026
    Docs as of 30 Sep 2026

    What you will be able to do

    • Use window() inside groupBy to aggregate a stream over event-time intervals, and read the resulting window struct
    • Choose between tumbling, sliding and session windows and write each one
    • Apply the window constraints the exam tests: half-open intervals, slide no longer than the window, timestamp-typed time column, no month durations
    • Bound windowed-aggregation state with a watermark and pick an output mode that fits

    1.window(): grouping a stream by time

    A running count for each key never closes. Business questions are usually tied to time instead: orders per hour, clicks per five minutes. window() makes that possible. It is a grouping expression that you put inside groupBy, and it assigns every row to one or more time buckets based on a timestamp column.

    Parameters of window()
    ParameterMeaning
    timeColumnThe column or expression to use as the timestamp; must be TimestampType or TimestampNTZType
    windowDurationWidth of each window, e.g. 10 minutes; a fixed length, not a calendar unit
    slideDuration (optional)A new window starts every slideDuration; must be less than or equal to windowDuration; omit it for tumbling windows
    startTime (optional)Offset from 1970-01-01 00:00:00 UTC, e.g. 15 minutes for hourly windows running 12:15-13:15

    Window starts are inclusive and window ends are exclusive. By default the grouping column in the result is a struct named window with nested start and end fields, both timestamps. Downstream you therefore refer to window.start and window.end, not to a single bucket label.

    Checkpoint 1 of 6· Check yourself

    After groupBy(window("dt", "5 seconds")).agg(sum("v")), what does the window column in the result look like by default?

    Sources1

    2.Tumbling, sliding and session windows

    The window shape depends on the arguments you pass. With only a duration you get tumbling windows: fixed-size, non-overlapping intervals where each input row falls into exactly one window. That fits discrete totals such as hourly sales:

    Tumbling window: total sales per non-overlapping hourpython
    from pyspark.sql.functions import window, sum
    
    hourly_sales = (orders
      .withWatermark("timestamp", "1 hour")
      .groupBy(window("timestamp", "1 hour"))
      .agg(sum("amount").alias("total_sales"))
    )

    Adding a slide duration gives sliding windows. They are still fixed-size, but they overlap, so one row can belong to several windows. A 6-hour window that slides every hour produces intervals such as 5–11 AM and 6 AM–12 PM, which gives you a rolling aggregate. The slide must be less than or equal to the window duration.

    Checkpoint 2 of 6· Fill the gap

    Which keyword argument turns this into a rolling 6-hour total that advances every hour?

    from pyspark.sql.functions import window, sum
    
    rolling_sales = (orders
      .withWatermark("timestamp", "1 hour")
      .groupBy(window("timestamp", "6 hours",  ? ="1 hour"))
      .agg(sum("amount").alias("total_sales"))
    )

    Session windows are a separate function, session_window, and they have no fixed size. A window opens when a row arrives, each later row within the gap duration extends it, and it closes once a full gap passes with no new rows. They suit bursts of activity separated by idle periods:

    Session window: page views per user, closing after 30 idle minutespython
    from pyspark.sql.functions import session_window, sum
    
    sessionized_page_views = (activity
      .withWatermark("timestamp", "1 hour")
      .groupBy("user_id", session_window("timestamp", gapDuration="30 minutes"))
      .agg(sum("page_views").alias("total_page_views"))
    )

    Checkpoint 3 of 6· Match them up

    Match each window type to how it assigns rows

    Tap a term, then the definition that fits it.

    Sources2

    3.The time column and duration rules

    A few rules decide whether a windowed aggregation is valid at all. First, the time column passed to window() or session_window() must be of type TimestampType or TimestampNTZType. Normally that is the event's own timestamp. To bucket by arrival time instead, Databricks says to use current_timestamp(), which defines windows on processing time.

    Second, durations are fixed lengths of time, not calendar units: one day always means 86,400,000 milliseconds. Window durations can range from microseconds up to days, and month durations or longer are not supported.

    Third, when you combine groupBy() with window(), refer to columns by name, as "<colName>" or col("<colName>"). Databricks says this preserves the event-time marker that a watermark depends on.

    Checkpoint 4 of 6· Check yourself

    A team wants to aggregate a stream into monthly revenue buckets with window("ts", "1 month"). What happens?

    Checkpoint 5 of 6· Exam question

    A monitoring pipeline reads a streaming DataFrame `logins` with an event-time column `event_time` and needs 10-minute tumbling windows that report a fresh count for every 10-minute period with no overlap between consecutive windows. Which line correctly defines this windowed aggregation? ```python result = logins.groupBy( ___ ).count() ```

    Sources2

    4.Bounding windowed state with a watermark

    Every open window is state. A stream that runs indefinitely keeps creating new windows, so something has to decide when an old window's aggregate can be dropped. A watermark makes that decision. You declare it with withWatermark on the same timestamp column before grouping:

    5-minute tumbling counts per id, with a 10-minute watermark on event_timepython
    from pyspark.sql.functions import window
    
    (df
      .withWatermark("event_time", "10 minutes")
      .groupBy(
        window("event_time", "5 minutes"),
        "id")
      .count()
    )

    That choice ties back to output mode. Complete mode keeps all window state indefinitely, because it rewrites every result each trigger. Append mode with a suitable watermark bounds state growth and prevents memory problems on large data. The trade-off is that a window's result is written only once the watermark has passed it.

    Checkpoint 6 of 6· Check yourself

    A windowed aggregation over a high-volume, never-ending stream is running out of memory. Which change does Databricks recommend for bounding state growth?

    Sources2

    Exam traps

    Each one states something that sounds right. Open it to see what is actually true.

    1. 1.An event stamped exactly on a window boundary (e.g. 12:05) is counted in both adjacent 5-minute windows.Why is that wrong?

      Window starts are inclusive and ends are exclusive, so 12:05 belongs only to [12:05,12:10).

      Covered in window(): grouping a stream by time

    2. 2.A sliding window can use a slide longer than the window, e.g. window(ts, "5 minutes", "10 minutes"), to sample the stream.Why is that wrong?

      slideDuration must be less than or equal to windowDuration. If it is omitted, the windows are tumbling.

      Covered in Tumbling, sliding and session windows

    3. 3.A "1 day" window follows calendar days, so it adjusts for daylight-saving changes, and "1 month" windows are fine as well.Why is that wrong?

      Durations are fixed lengths: 1 day is always 86,400,000 milliseconds. Month-scale durations are not supported.

      Covered in The time column and duration rules

    Sources

    Every claim above is drawn from one of these pages, quoted as it was written on the date shown.

    1. 1.
      “If the slideDuration is not provided, the windows will be tumbling windows.”
      ↩︎ window(): grouping a stream by time
      “Window starts are inclusive but the window ends are exclusive”
      ↩︎ window(): grouping a stream by time
      “Window starts are inclusive but the window ends are exclusive”
      ↩︎ Exam trap 1
      “For example, 1 day always means 86,400,000 milliseconds, not a calendar day.”
      ↩︎ Exam trap 3
      “12:05 will be in the window [12:05,12:10) but not in [12:00,12:05)”
      ↩︎ Prediction
      “The output column will be a struct called 'window' by default with the nested columns 'start' and 'end'”
      ↩︎ Checkpoint
    2. 2.
      “Tumbling windows are fixed-size with non-overlapping intervals. Each input row belongs to exactly one window.”
      ↩︎ Tumbling, sliding and session windows
      “A window opens when a row arrives and closes after a gap duration that contains no new rows.”
      ↩︎ Tumbling, sliding and session windows
      “The timeColumn argument for window() and session_window() must be of TimestampType or TimestampNTZType.”
      ↩︎ The time column and duration rules
      “Use current_timestamp() to define windows based on processing time rather than event time.”
      ↩︎ The time column and duration rules
      “reference columns by name, "<colName>" or col("<colName>"), to ensure the event time marker is preserved”
      ↩︎ The time column and duration rules
      “State information is maintained for each count until the end of the window is 10 minutes older than the latest observed event_time.”
      ↩︎ Bounding windowed state with a watermark
      “Use complete output mode with windowed aggregations to keep all window state indefinitely.”
      ↩︎ Bounding windowed state with a watermark
      “slideDuration must be less than or equal to the windowDuration.”
      ↩︎ Exam trap 2
      “Sliding windows are fixed-size with intervals that can overlap. A single row can belong to multiple windows.”
      ↩︎ Checkpoint
      “You can set window durations from microseconds up to days. Month durations and longer are not supported.”
      ↩︎ Checkpoint
      “Use append output mode with an appropriate watermark to bound state growth and prevent memory issues for large data sets.”
      ↩︎ Checkpoint

    Spotted a mistake, or was something unclear? Tell us.