Interview/Spark

Spark interview questions: Window functions: keep the row, add the answer

Interactive Spark interview questions on Window functions: keep the row, add the answer. Practice with the matching lesson. Rank, lag, and running totals keep every row. partitionBy is a shuffle — and forgetting it funnels everything to one task.

Lesson · Simulation

You need a per-customer running total but the report still needs every transaction row. Window or groupBy?

Answer it out loud, then reveal. Play steps through like the simulators.

Production scenario

Window functions: keep the row, add the answer

A 'rolling 30-day average' column turned a 6-minute job into a 90-minute job

Symptoms

  • One new metric column shipped last sprint
  • A single task never finishes while the rest of the cluster is idle
  • Runtime grows a little every day as history accumulates

All questions on this page

Indexed as FAQ. Open any item if you prefer a list to Play.

beginner

You need a per-customer running total but the report still needs every transaction row. Window or groupBy?

Window functions: keep the row, add the answer · tap to open the answer

Short: Window — it computes over related rows but returns a value per row, so nothing collapses.

Detailed: groupBy + agg reduces each key to one row; a window function keeps input cardinality and adds a column. In the plan you get a WindowExec with a Sort under it instead of a HashAggregate, and no join-back step to reattach the aggregate.

Common mistake: Doing groupBy then joining the aggregate back to the detail — two shuffles for what a window does in one.

Follow-up: How would you use a window to grab the latest row per key?

Lesson · Simulation

intermediate

rowsBetween(-2, 0) and rangeBetween(-2, 0) on the same ordered window return different numbers. Why?

Window functions: keep the row, add the answer · tap to open the answer

Short: rowsBetween counts physical rows; rangeBetween compares values of the ordering column.

Detailed: rowsBetween is positional — the two preceding rows regardless of their values. rangeBetween is a value frame: every row whose ordering value falls inside the offset is included, so gaps shrink the frame and ties pull in all peers. With duplicate ordering values, rangeBetween rows share identical frames.

Common mistake: Using rangeBetween on a timestamp and expecting exactly three rows in the frame.

Follow-up: Which one do you need for a true 7-day rolling sum over daily data with missing days?

Lesson · Simulation

senior

Your window job has one task running for 40 minutes while 199 finished in seconds. What did you write?

Window functions: keep the row, add the answer · tap to open the answer

Short: A Window.orderBy(...) with no partitionBy — every row was shuffled into a single partition.

Detailed: Without partitionBy, Spark must order the whole dataset globally, so it plans an Exchange SinglePartition beneath WindowExec and one task processes everything. The stage shows exactly 1 task carrying the entire Shuffle Read plus a large Spill (Disk) while the rest of the cluster idles.

Senior: Reduce first, then order: filter or pre-aggregate to candidates, or rank within real partitions and stitch them with a small per-partition offset table. A truly global ordering over billions of rows is a design problem, not a tuning knob.

Common mistake: Diagnosing it as skew and salting a key that the window does not even partition by.

Follow-up: How do you compute a global row number without a single-task stage?

Lesson · Simulation

architect

Dedup with row_number() = 1 versus dropDuplicates — which goes in the pipeline and why?

Window functions: keep the row, add the answer · tap to open the answer

Short: row_number() when you must keep a specific survivor; dropDuplicates when any row per key is acceptable.

Detailed: dropDuplicates(keys) dedups via a shuffle and aggregate and hands you an arbitrary row per key. partitionBy(keys).orderBy(updated_at.desc) with a filter on row_number() = 1 pays for a Sort inside WindowExec but deterministically keeps the newest record — which is exactly what CDC and SCD merges require.

Senior: Add a deterministic tiebreaker such as a sequence number or commit ordering. Without one, reruns pick different survivors and downstream diffs look like data loss.

Common mistake: Running dropDuplicates on a CDC feed and shipping a random version of each record.

Follow-up: What breaks if the ordering column has ties?

Lesson · Simulation