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)