spark/core
intermediate
Connecting…

Window functions: keep the row, add the answer

Rank, lag, and running totals keep every row. partitionBy is a shuffle — and forgetting it funnels everything to one task.

Lesson 17 of 29 · Spark path

Explain it at my level

  1. groupBy
  2. Window
  3. partitionBy
  4. orderBy
  5. Frame
Watch the canvas:input rowscollapsed groupswindow columnpartition laneexchange / frameLive simulation
Input rowsUS · t1amount 100US · t2amount 140EU · t1amount 90US · t3amount 160EU · t2amount 120groupBy(region).sum(amount)USsum 400EUsum 2105 rows in → 2 rows outts and amount are gonesum(amount).over(w)US · t1no column yetUS · t2no column yetEU · t1no column yetUS · t3no column yetEU · t2no column yetpartitionBy(region) · Exchange hashpartitioningPartition: USt2 140 · t1 100 · t3 160 (no order)🔒Partition: EUt2 120 · t1 90🔒Aggregation is lossy by design. The detail columns have nowhere to go.
groupBy collapses rows; a window keeps them. partitionBy is the shuffle, orderBy gives position, the frame gives the number.