Accumulate through the current row
sum over an ordered window adds values from the start through the current row. rowsBetween(unboundedPreceding, currentRow) makes that frame explicit. The example processes order_id ascending. Reversing the window order changes which rows have already contributed, even if the final table remains sorted ascending. A global window is appropriate here only because the teaching dataset is small.
Lesson reference: PySpark and Scala
Accumulate through the current row
sum over an ordered window adds values from the start through the current row. rowsBetween(unboundedPreceding, currentRow) makes that frame explicit. The example processes order_id ascending. Reversing the window order changes which rows have already contributed, even if the final table remains sorted ascending. A global window is appropriate here only because the teaching dataset is small.
PySpark example
from pyspark.sql import Window
from pyspark.sql.functions import col, sum
result = orders.withColumn(
"running_total", sum("total").over(
Window.orderBy(col("order_id").asc())
.rowsBetween(Window.unboundedPreceding, Window.currentRow)
)
).select("order_id", "status", "customer_id", "total", "running_total").orderBy("order_id")
result.show(truncate=False)
Scala example
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
val result = orders.withColumn(
"running_total", sum("total").over(
Window.orderBy(col("order_id").asc())
.rowsBetween(Window.unboundedPreceding, Window.currentRow)
)
).select("order_id", "status", "customer_id", "total", "running_total").orderBy("order_id")
result.show(false)
Restart for each group
partitionBy(status) keeps a separate accumulator for each status. Interleaved rows of other statuses do not contribute. Within each group, order_id determines the sequence and the explicit row frame includes the current row. Changing the partition key changes where totals restart.
PySpark example
from pyspark.sql import Window
from pyspark.sql.functions import col, sum
result = orders.withColumn(
"running_total", sum("total").over(
Window.partitionBy("status").orderBy(col("order_id").asc())
.rowsBetween(Window.unboundedPreceding, Window.currentRow)
)
).select("order_id", "status", "customer_id", "total", "running_total").orderBy("order_id")
result.show(truncate=False)
Scala example
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
val result = orders.withColumn(
"running_total", sum("total").over(
Window.partitionBy("status").orderBy(col("order_id").asc())
.rowsBetween(Window.unboundedPreceding, Window.currentRow)
)
).select("order_id", "status", "customer_id", "total", "running_total").orderBy("order_id")
result.show(false)
Count qualifying rows so far
count(*) over the same row frame counts every row encountered. A conditional sum can count only matching rows without removing the other rows from the result. Use 1 for matches and 0 otherwise. Filtering the DataFrame first would remove rows that should remain visible.
PySpark example
from pyspark.sql import Window
from pyspark.sql.functions import count
result = orders.withColumn(
"running_count", count("*").over(
Window.orderBy("order_id").rowsBetween(Window.unboundedPreceding, Window.currentRow)
)
).select("order_id", "status", "running_count").orderBy("order_id")
result.show(truncate=False)
Scala example
import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._
val result = orders.withColumn(
"running_count", count("*").over(
Window.orderBy("order_id").rowsBetween(Window.unboundedPreceding, Window.currentRow)
)
).select("order_id", "status", "running_count").orderBy("order_id")
result.show(false)