ML Engineer MasterClass (October) | 4 seats left

Spark · Window frames
AmazonAmazon Analytics
Amazon · Window Functions

Window frames

Choose which neighboring rows or ordering values contribute to a window calculation.

Step 1 of 6 · Learn

Use a trailing row frame

rowsBetween uses inclusive offsets relative to the current row. (-1, 0) includes the previous row and current row; (-2, 0) adds one more preceding row. At the beginning of a partition, Spark uses only rows that exist. avg divides by the number of non-null values present, not the requested frame width. order_id gives a unique sequence here.

Lesson reference: PySpark and Scala

Use a trailing row frame

rowsBetween uses inclusive offsets relative to the current row. (-1, 0) includes the previous row and current row; (-2, 0) adds one more preceding row. At the beginning of a partition, Spark uses only rows that exist. avg divides by the number of non-null values present, not the requested frame width. order_id gives a unique sequence here.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import avg, round

result = orders.withColumn(
    "frame_value", round(avg("total").over(
        Window.orderBy("order_id").rowsBetween(-1, 0)
    ), 2)
).select("order_id", "total", "frame_value").orderBy("order_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import spark.implicits._

val result = orders.withColumn(
    "frame_value", round(avg("total").over(
        Window.orderBy("order_id").rowsBetween(-1, 0)
    ), 2)
).select("order_id", "total", "frame_value").orderBy("order_id")
result.show(false)

Include following rows

A positive upper bound includes rows after the current row. rowsBetween(0, 1) sums the current and next row. Following rows are available in this batch DataFrame; this is not a streaming promise about future events. Near the end, a frame contains fewer rows because no additional rows exist.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import sum

result = orders.withColumn(
    "frame_value", sum("total").over(
        Window.orderBy("order_id").rowsBetween(0, 1)
    )
).select("order_id", "total", "frame_value").orderBy("order_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import spark.implicits._

val result = orders.withColumn(
    "frame_value", sum("total").over(
        Window.orderBy("order_id").rowsBetween(0, 1)
    )
).select("order_id", "total", "frame_value").orderBy("order_id")
result.show(false)

Use value-based ranges

rangeBetween uses ordering values, not row positions. With item_count ordering, (0, 0) includes every row with the current item count. A lower bound of -1 also includes rows with one fewer item. Equal ordering values share the same range result. A bounded numeric range uses a single numeric ordering expression, so do not add an order_id tie-breaker here.

PySpark example

from pyspark.sql import Window
from pyspark.sql.functions import sum

result = orders.withColumn(
    "frame_value", sum("total").over(
        Window.orderBy("item_count").rangeBetween(0, 0)
    )
).select("order_id", "item_count", "frame_value").orderBy("order_id")
result.show(truncate=False)

Scala example

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
import spark.implicits._

val result = orders.withColumn(
    "frame_value", sum("total").over(
        Window.orderBy("item_count").rangeBetween(0, 0)
    )
).select("order_id", "item_count", "frame_value").orderBy("order_id")
result.show(false)
example.pyPySpark
1. Previous row1002: 149.00
2. Current row1003: 35.00
3. Frame average(149 + 35) / 2 = 92.00
Source orders6 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
1003101Cancelled351
1004103Delivered219.994
1005104Delivered49.991
1006105Shipped1202
Use a trailing row frame
Result6 rows
order_idtotalframe_value
100189.589.5
1002149119.25
10033592
1004219.99127.5
100549.99134.99
100612085
For order 1003, the example averages 149 and 35 to get 92. The first order uses only its own total.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.