Quick plain-English version: the database is doing more work than it needs to because it can't tell in advance which rows actually match. An index is basically a shortcut list so it doesn't have to check every single row.
A running total using Window.partitionBy("customer_id").orderBy("date") worked fine in dev on 1M rows, then fell over on 500M rows in prod with massive shuffle spill.
w = Window.partitionBy("customer_id").orderBy("date")
df.withColumn("running_total", F.sum("amount").over(w))Anything to do besides get a bigger cluster?
Quick plain-English version: the database is doing more work than it needs to because it can't tell in advance which rows actually match. An index is basically a shortcut list so it doesn't have to check every single row.
We replaced custom sensors with data contracts and row count checks.
Check whether AQE is disabled in your Spark conf, skew join handling helped us a lot here.
Careful with NULL in join keys, they'll drop rows in an inner join.
Check partition skew first. If a handful of customer_ids have millions of rows each, the window computation for those partitions runs on a single task no matter how big the cluster is.
Turns out one test or bot customer had 40M synthetic rows. Filtering it out before the window fixed 90% of the spill.
Careful with NULL in join keys, they'll drop rows in an inner join.
Careful with NULL in join keys, they'll drop rows in an inner join.
Idempotent writes with merge keys saved us during backfills.
For anyone who wants the two-sentence version: data arrived out of order, and the code assumed it wouldn't. Fix the assumption or sort the data, pick one.
For an unbounded-preceding running total specifically, consider a two-pass aggregate plus join for just the skewed keys, so you avoid a full window scan for those partitions.
Prefer a staging table plus validation gate before promoting to prod tables.
Sign in to reply.
© 2026 Lakebench, operated by Hunnurji Rao. Bengaluru, Karnataka, India.
No cluster. No install. Just the tab.