A job that grows from 20 minutes to 3 hours over six months is almost always doing work proportional to total history instead of to new data. Start by finding which step got slow, then ask what grew.
Step 1: find the slow step
Look at the run time history of the job, not only today's. Plot duration per day. A gradual slope points to data growth. A sudden jump points to a change (a deploy, a new source, a configuration change). Then break the job into stages or tasks and see which one holds the extra time, using the Spark UI, query profile or task durations.
The usual suspects
- Full scans of growing tables. A job that reads the whole
eventstable every day was fine at 50 GB and is slow at 5 TB. The fix is incremental processing: read only new or changed partitions. - Missing partition pruning. A filter that stops using the partition column (a function on it, or a type change) makes the engine read all partitions. Check the plan or the bytes scanned.
- Small files piling up. Daily appends of many tiny files make reads slow. Compact them, and write fewer, larger files.
- Joins against history. Joining a daily batch with a growing full history table (for a lookup or a dedup check) costs more every month. Limit the lookup to a recent window, or use a MERGE with a pruned target.
- Skew. One customer grew from 1 percent to 30 percent of the data, and now one task does most of the work.
- State growth in streaming or SCD logic, which keeps more history each day.
- Stale statistics, so the planner picks a poor join.
- A cluster or warehouse sized for last year's data.
Fix by cause
Switch to incremental loads, add or repair partition filters, compact files, bound the lookup windows, handle the hot key, refresh stats, and then size resources. Measure after each change, so you know what helped.
Catch it earlier
Store the runtime, rows read and rows written of every run, and alert when the duration is more than, say, 50 percent above the 30-day average. A trend chart in a weekly review catches a slope long before it threatens the SLA. Set the alert on rows scanned per row written too, since that ratio shows inefficiency as the data grows.