Tasks and Task Scheduling

A task is one stage's code applied to one partition on one core. The TaskScheduler offers tasks to free executor slots and prefers locality: a task that reads a cached partition or a local HDFS block waits up to spark.locality.wait (3 seconds) for a slot on the right machine before running elsewhere. For file scans, the partition count follows file sizes and spark.sql.files.maxPartitionBytes (128 MB by default). The 20 MB order_lines table is four Parquet 129 files, and lines.rdd.getNumPartitions() returns 4: one scan task per core. Setting maxPartitionBytes to 4m raised it to 8.

Aim for tasks of at least a second or two: thousands of tiny tasks spend their time on scheduling, while a few giant ones leave cores idle and make each retry expensive. spark.speculation=true relaunches straggling tasks elsewhere, which helps with a sick machine but not with data skew (Diagnosing Data Skew).