A data engineer is working with a PySpark DataFrame representing sensor readings. The DataFrame `sensor_df` has columns `device_id`, `timestamp`, and `temperature`. They need to calculate a 7-day rolling average of `temperature` for each `device_id`. Which PySpark window function concept is most appropriate for this task?
- ALead window function
- BRolling average window function
- CRow_number window function
- DLag window function
Show answer & explanationAnswer & explanation
Correct answer: B. Rolling average window function
A rolling average (or moving average) is a specific type of window function that computes the average over a defined range of rows (a 'window') that slides with each row. This is precisely what's needed for a '7-day rolling average'. While not a direct PySpark function name, it describes the *concept* of the window function required, typically implemented using `avg()` with `ROWS BETWEEN` or `RANGE BETWEEN` in a window specification.
Why the other options are wrong
- A. Lead retrieves a value from a subsequent row, not an average over a range.
- C. Row_number assigns a sequential integer to rows within each partition, not an average.
- D. Lag retrieves a value from a previous row, not an average over a range.
PySpark Rolling Average Window
A rolling average (or moving average) in PySpark is calculated using a window function that partitions data, orders it, and defines a frame (e.g., 7 preceding rows) over which an aggregate function (like `avg`) is applied.
- Requires `Window.partitionBy()`, `Window.orderBy()`.
- Uses `Window.rowsBetween()` or `Window.rangeBetween()` for the frame.
- Calculates aggregate (e.g., `avg`, `sum`) over the defined frame.
- Ideal for trend analysis over a moving period.
Memory trick: Partition, Order, Frame, Aggregate to get your window results.