Wide rows after an explode or array join earlier in the DAG can inflate shuffle a lot. Check the sparkPlan for a Generate or Explode before the aggregate.
Spark stage shows 120GB shuffle read on a 12GB input DataFrame after groupBy(customer_id) and count.
A single skewed customer_id isn't the whole story here, the top key is only 2% of rows. What else should I check?
Wide rows after an explode or array join earlier in the DAG can inflate shuffle a lot. Check the sparkPlan for a Generate or Explode before the aggregate.
Another path: push the compute to the warehouse if the data's already there.
Aggregate before exploding where you can, or filter line items before the explode to cut shuffle volume way down.
Prefer a staging table plus validation gate before promoting to prod tables.
There was an explode on line_items before the groupBy, row count went 8x. That explains it.
Sign in to reply.
© 2026 Lakebench, operated by Hunnurji Rao. Bengaluru, Karnataka, India.
No cluster. No install. Just the tab.