If tasks are even in size, skew is not the problem, so the time is going elsewhere. Look at three things: how much work Spark does, how it is spread over the cluster, and what each task spends time on.
Check in this order
First, read the plan with df.explain() and the SQL tab in the UI. Is Spark reading more than needed?
PartitionFiltersandPushedFiltersin the scan node show what was pruned. An empty list when you filter by date means you read every partition.- Columns: is it reading 80 columns when you use 5? Selecting early can cut I/O a lot on Parquet.
Second, check parallelism in the stage page.
- Too few partitions: 20 tasks on 400 cores leaves most idle. Increase partitions, or lower
maxPartitionBytes. - Too many tiny tasks: 100,000 tasks of 50 ms each means scheduler overhead dominates. Coalesce, or fix small files at the source.
Third, find what the tasks are doing.
- Python UDFs, shown as
BatchEvalPythonorArrowEvalPythonin the plan. Replace with built-in functions. - Repeated recomputation. A DataFrame used three times without
cache()is recomputed three times, and the UI shows the same scan appearing in several jobs. - Heavy spill or GC time columns.
- Reading small files: listing and opening thousands of files can take longer than reading them.
- An external bottleneck: a JDBC source that cannot serve more than a few connections, or throttled object storage requests.
Fourth, check the cluster. Are executors being allocated at all (dynamic allocation ramping slowly)? Is the job waiting in a queue? Are a lot of cores busy with another job?
The stage timeline
The event timeline in the stage view shows gaps. Large gaps between tasks mean scheduling delay or waiting for executors. Long green bars mean computation. Long blue bars show shuffle read waiting on the network.
How to say it
Do not guess. Say you would go from the plan (what is read), to the stage (how it is parallelised), to the tasks (what they spend time on), and change one thing at a time. That order finds most problems quickly.