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.
| Parameter | Meaning |
|---|---|
| timeColumn | The column or expression to use as the timestamp; must be TimestampType or TimestampNTZType |
| windowDuration | Width 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?
window() produces a struct column called window by default, with start and end fields of TimestampType.
“The output column will be a struct called 'window' by default with the nested columns 'start' and 'end'”Source: docs.databricks.com
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:
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"))
)slideDuration sets how often a new window begins. startTime only offsets window boundaries, and gapDuration belongs to session_window. windowDuration is already given as "6 hours".
Source: docs.databricks.comSession 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:
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.
Tumbling and sliding windows both come from window(); the slide duration decides whether they overlap. Session windows come from session_window and are sized by inactivity gaps.
“Sliding windows are fixed-size with intervals that can overlap. A single row can belong to multiple windows.”Source: docs.databricks.com
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?
Window durations can range from microseconds up to days. Month-scale durations are explicitly unsupported.
“You can set window durations from microseconds up to days. Month durations and longer are not supported.”Source: docs.databricks.com
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() ```
Correct answer: A — `window(logins.event_time, "10 minutes")`
- A. This is correct because omitting the slide duration argument to `window()` makes the slide equal to the window duration, which produces non-overlapping tumbling 10-minute windows.
- B. This is incorrect because passing a 5-minute slide duration alongside the 10-minute window duration creates overlapping sliding windows, not the non-overlapping tumbling windows the scenario requires.
- C. This is incorrect because `session_window` produces variable-length windows based on gaps between events, not the fixed-size tumbling windows the pipeline needs.
- D. This is incorrect because writing the slide duration explicitly as equal to the window duration is redundant with omitting it, and a mismatched literal like this is easy to get wrong; the cleaner, standard tumbling-window call omits the third argument entirely.
- E. This is incorrect because a 5-minute window with a 10-minute slide leaves 5-minute gaps between windows where no data is captured, rather than producing continuous non-overlapping coverage.
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:
from pyspark.sql.functions import window
(df
.withWatermark("event_time", "10 minutes")
.groupBy(
window("event_time", "5 minutes"),
"id")
.count()
)Until the end of that window is 10 minutes older than the latest event_time observed. After that, its state can be dropped.
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?
Complete mode keeps all window state forever. Append mode plus a watermark lets old windows be finalised and their state dropped.
“Use append output mode with an appropriate watermark to bound state growth and prevent memory issues for large data sets.”Source: docs.databricks.com
Sources2
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
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.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.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.
“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.
“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