TT Lab
Get started
Learn Learning paths Courses

Apache Spark — The answer to a slow job is in the plan and the event log

Reduce the files and tasks produced by 200 shuffle partitions

Continue in TT Lab

Goal

You run the per-customer revenue aggregation three times: with AQE off and the default (200 shuffle partitions), with the value reduced by hand (8), and with AQE, and measure with the event log how the number of post-shuffle tasks and the number of result files differ. You also confirm how repartition and coalesce differ in whether they shuffle and in the number of files.

Why it matters

A wide transformation such as groupBy or join has to gather the same key into one task, so it creates a shuffle. The earlier stage writes its results as per-partition files, and the later stage gathers and reads those pieces. The number of tasks in the later stage is determined by spark.sql.shuffle.partitions, which defaults to 200. 200 is a number based on a large cluster. If you use 200 on a few MB of data, 200 tasks each do a little work, and if you write the result to files, 200 small files are created. Conversely, with 200 on several hundred GB, one task takes on several GB and spills to disk. So this number has to fit the data size, and AQE does this work for you by looking at the actual sizes after the shuffle has finished and merging small pieces. There are two ways to change the number of partitions. repartition(n) creates new pieces with a shuffle, and coalesce(n) attaches adjacent pieces without a shuffle. coalesce is cheap, but it can only reduce: you cannot turn two pieces into four.

Steps

  1. Put a function that builds the per-customer revenue DataFrame in /root/spk/shuffle/common.py, and with /root/spk/shuffle/noaqe.py (application spk-shuffle-noaqe, spark.sql.adaptive.enabled=false), write the result as Parquet to /root/spk/shuffle/out/by_customer_200.
  2. Count the data files (those starting with part-) in that result folder and write the count as an integer to /root/spk/shuffle/out/files_200.txt.
  3. With /root/spk/shuffle/eight.py (application spk-shuffle-8, AQE off, spark.sql.shuffle.partitions=8), write the same result to /root/spk/shuffle/out/by_customer_8.
  4. With /root/spk/shuffle/aqe.py (application spk-shuffle-aqe, default settings), write the same result to /root/spk/shuffle/out/by_customer_aqe, and write the number of tasks in that application that read the shuffle to /root/spk/shuffle/out/aqe_tasks.txt as an integer.
  5. Write the click original with repartition(4) using /root/spk/shuffle/repart.py (application spk-shuffle-repart) to /root/spk/shuffle/out/clicks_repart, and with coalesce(4) using /root/spk/shuffle/coalesce.py (application spk-shuffle-coalesce) to /root/spk/shuffle/out/clicks_coalesce, both as Parquet, and write the number of data files in the two folders to /root/spk/shuffle/out/repart.json as {"repartition_files": 정수, "coalesce_files": 정수} (both values are integers).
  6. From the log, find the sum of the bytes that the step 3 application wrote through the shuffle, and write it to /root/spk/shuffle/out/shuffle_bytes.json as {"app": "spk-shuffle-8", "shuffle_write_bytes": 정수} (the second value is an integer).
  7. With /root/spk/shuffle/narrow.py (application spk-shuffle-narrow), add a day column to the paid orders and write only order_id,customer_id,qty,day to /root/spk/shuffle/out/narrow. This application must have no shuffle.
  8. In /root/spk/shuffle/report.md, write three sections: ## 200 개의 파티션, ## AQE 가 합친 것 and ## repartition 과 coalesce (use exactly these Korean headings in this order; they mean "200 partitions", "What AQE merged" and "repartition and coalesce"). Put the numbers from steps 2, 4 and 5 in the respective sections.

Notes

Turn AQE off and use the default of 200

In /root/spk/shuffle/common.py, put a function that builds a DataFrame of revenue per customer (customer_id, revenue = the sum of qty×price) for paid orders, and create /root/spk/shuffle/noaqe.py with the application name spk-shuffle-noaqe and the setting spark.sql.adaptive.enabled=false, and write the result as Parquet to /root/spk/shuffle/out/by_customer_200.

With AQE off, the stage after the shuffle runs with exactly spark.sql.shuffle.partitions tasks. The grader checks from the log whether the number of tasks of that stage is 200 and whether the result equals the per-customer revenue computed from the original.

Count the files that 200 made

Count the data files (those starting with part-) in /root/spk/shuffle/out/by_customer_200 and write the count as a single integer to /root/spk/shuffle/out/files_200.txt.

One post-shuffle task writes one file (an empty partition may not write a file). The revenue of 20,000 customers is a few hundred KB, yet look at how many files there are: this is the most common source of the small files problem. Do not count _SUCCESS and .crc.

Reduce to 8 by hand

Create /root/spk/shuffle/eight.py with the application name spk-shuffle-8 and the settings spark.sql.adaptive.enabled=false and spark.sql.shuffle.partitions=8, and write the same result as Parquet to /root/spk/shuffle/out/by_customer_8.

This time there are 8 post-shuffle tasks and 8 or fewer files. The result must not differ from step 1 by a single row: the number of partitions is a way of dividing the work, not the answer.

AQE merges after the shuffle

Create /root/spk/shuffle/aqe.py with the application name spk-shuffle-aqe (settings left at their defaults), write the same result to /root/spk/shuffle/out/by_customer_aqe, then count from the log the number of tasks in that application that read the shuffle and write it as an integer to /root/spk/shuffle/out/aqe_tasks.txt.

After the shuffle map side has finished, AQE looks at the actual size of each partition and merges adjacent pieces that fall short of the target size (spark.sql.adaptive.advisoryPartitionSizeInBytes). The shuffle partitions are still 200, but the tasks that read become far fewer. In the final plan you see AQEShuffleRead … coalesced.

repartition and coalesce: with and without a shuffle

With /root/spk/shuffle/repart.py (application spk-shuffle-repart), apply repartition(4) to /data/clicks/clicks.jsonl and write it to /root/spk/shuffle/out/clicks_repart, and with /root/spk/shuffle/coalesce.py (application spk-shuffle-coalesce), apply coalesce(4) to the same original and write it to /root/spk/shuffle/out/clicks_coalesce, both as Parquet. Write the number of data files in the two folders to /root/spk/shuffle/out/repart.json as {"repartition_files": 정수, "coalesce_files": 정수} (both values are integers).

First see into how many pieces the 18MB original is read on local[2] (rdd.getNumPartitions()). repartition makes exactly four with a shuffle, and coalesce only attaches neighbors without a shuffle, so it does not go above the original number of pieces. The grader also checks from the log that only the repart application has shuffle writes.

Bytes the shuffle wrote to disk

From the most recent log of the step 3 application (spk-shuffle-8), add up Shuffle Write Metrics.Shuffle Bytes Written of all the tasks and write it to /root/spk/shuffle/out/shuffle_bytes.json as {"app": "spk-shuffle-8", "shuffle_write_bytes": 정수} (the second value is an integer).

A shuffle write is the map-side tasks dividing their results by partition and writing them to local disk. Thanks to partial aggregation (the partial of HashAggregate), only about 20,000 customers × the number of map tasks rows are passed on, so it is much smaller than the original. This number is the actual cost of the shuffle.

With only narrow transformations there is one stage

Create /root/spk/shuffle/narrow.py with the application name spk-shuffle-narrow, add day = to_date(order_ts) to the paid orders, select only order_id,customer_id,qty,day, and write it as Parquet to /root/spk/shuffle/out/narrow. This application must have no shuffle at all.

where, withColumn and select look at one row and produce one row. They do not need data from other tasks, so they run continuously within one stage (pipelining). If you add even a single orderBy or distinct, a shuffle appears.

Record the basis for the number of partitions

In /root/spk/shuffle/report.md, write three sections: ## 200 개의 파티션, ## AQE 가 합친 것 and ## repartition 과 coalesce (use exactly these Korean headings in this order; they mean "200 partitions", "What AQE merged" and "repartition and coalesce"). Put the number of files from step 2 in the first section, the number of tasks from step 4 in the second, and the two file counts from step 5 in the third, as numbers.

Add one line each on how many shuffle partitions you would set for this data size, and whether you need to set it by hand at all given that AQE exists. What serves as the basis for this decision in production is the number of result files and the shuffle bytes.