TT Lab
Get started
Learn Learning paths Courses

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

Run the same join with four strategies and confirm each in the plan

Continue in TT Lab

Goal

You run the join of orders and products with the three strategies (broadcast hash, sort merge and shuffle hash), compare the plans and results, and confirm with the event log that AQE changes a sort merge into a broadcast during execution. You also do a join with a table whose keys overlap, which inflates rows, and a join that selects rows with no match on the other side.

Why it matters

A join is often the most expensive operation in a Spark job, and its cost is decided by the strategy. If one side is small, you copy the small side to every task (broadcast) and join without shuffling the large side. If both are large, you shuffle both sides by key, sort them, and merge them while matching (sort merge). A shuffle hash, which builds a hash table from one side after the shuffle and joins, uses memory in exchange for skipping the sort. Spark estimates the size from statistics and chooses the strategy. The criterion for considering something small is spark.sql.autoBroadcastJoinThreshold (default 10MB). If the estimate is wrong, the strategy is wrong too: after a filter the table is actually small, but Spark estimates from the original size and chooses a sort merge; or conversely it broadcasts a large table and driver memory overflows. AQE fixes this by judging again with the actual size after the shuffle has finished. The number of rows in a join result is determined by key duplication. If you believed a key was unique on one side and it is not, the join multiplies rows without an error and every total after it inflates.

Steps

  1. Put a function that reads the four tables and a per-category revenue function in /root/spk/join/common.py, and with /root/spk/join/auto.py (application spk-join-auto, default settings), write the per-category revenue as a CSV with a header row (category,revenue) to /root/spk/join/out/by_category.
  2. With /root/spk/join/smj.py (application spk-join-smj, spark.sql.autoBroadcastJoinThreshold=-1 and spark.sql.adaptive.enabled=false), write the same result to /root/spk/join/out/by_category_smj.
  3. In /root/spk/join/hint.py (application spk-join-hint, the same settings as step 2), give a broadcast hint to the products side and write to /root/spk/join/out/by_category_hint.
  4. In /root/spk/join/shash.py (application spk-join-shash, the same settings as step 2), give a shuffle_hash hint to the products side and write to /root/spk/join/out/by_category_shash.
  5. With /root/spk/join/aqe.py (application spk-join-aqe, spark.sql.autoBroadcastJoinThreshold=100k, AQE on), join the paid orders with the customers where tier == 'vip' and write the number of orders per city as a CSV (city,orders) to /root/spk/join/out/vip_by_city.
  6. With /root/spk/join/dup.py (application spk-join-dup), write the number of rows from simply joining the paid orders with the promotions table and the number of rows from joining with left_semi to /root/spk/join/out/dup.json as {"naive": 정수, "semi": 정수} (both values are integers).
  7. With /root/spk/join/anti.py (application spk-join-anti), write the customer_id of customers who have never placed an order as a CSV to /root/spk/join/out/no_orders.
  8. In /root/spk/join/report.md, write three sections: ## 네 가지 전략, ## AQE 의 전환 and ## 키 중복 (use exactly these Korean headings in this order; they mean "Four strategies", "AQE's switch" and "Duplicate keys"). Put the two numbers from step 6 in the third section.

Notes

The small side is broadcast automatically

In /root/spk/join/common.py, put a function that reads the four tables (orders, products, customers, promotions) with a schema and a per-category revenue function (paid orders × products, category, revenue = the sum of qty×price), and create /root/spk/join/auto.py with the application name spk-join-auto (default settings) and write the result as a CSV with a header row to /root/spk/join/out/by_category.

The products table is a few KB, far smaller than the threshold (10MB). Spark copies it to every task and does not shuffle the orders side. The grader checks whether the plan in the event log has a BroadcastHashJoin and whether the per-category revenue matches the original.

Turn off the threshold and you get a sort merge

Create /root/spk/join/smj.py with the application name spk-join-smj and the settings spark.sql.autoBroadcastJoinThreshold=-1 and spark.sql.adaptive.enabled=false, and write the same result as in step 1 to /root/spk/join/out/by_category_smj.

A threshold of -1 means "no automatic broadcast". Now both the orders and the products are shuffled by product_id, sorted and then merged. The result must not differ from step 1 by a single row. The reason for turning AQE off is to keep the strategy from changing during execution.

Force a broadcast with a hint

Create /root/spk/join/hint.py with the application name spk-join-hint and the same settings as step 2 (threshold -1, AQE off), but wrap the products side with F.broadcast(p), and write to /root/spk/join/out/by_category_hint.

A hint takes precedence over statistics. Even if you have turned the threshold off, it broadcasts when there is a hint. Put the other way, a broadcast hint carelessly attached to a large table eats the memory of the driver and every executor as is.

Shuffle hash join: in exchange for skipping the sort

Create /root/spk/join/shash.py with the application name spk-join-shash and the same settings as step 2, but give the products side p.hint("shuffle_hash"), and write to /root/spk/join/out/by_category_shash.

A shuffle hash join shuffles both sides by key, then builds a hash table from the smaller side's partition and streams the larger side through it. It can be fast because there is no sort, but the hash table has to fit in memory. Confirm in the plan that the Sort operator is gone.

AQE changes the strategy during execution

Create /root/spk/join/aqe.py with the application name spk-join-aqe and the setting spark.sql.autoBroadcastJoinThreshold=100k (leaving AQE on), join the paid orders with the customers where tier == 'vip' by customer_id, and write the number of orders per city as a CSV with a header row (city,orders) to /root/spk/join/out/vip_by_city.

The customers table file is larger than 100KB, so the first plan is a sort merge. But the actual size after filtering by vip is a few dozen KB. After the shuffle map stage has finished, AQE looks at that size and changes the rest of the plan into a broadcast. The grader compares the first plan and the final plan of the same run.

If keys overlap, a join inflates rows

Create /root/spk/join/dup.py with the application name spk-join-dup, and write the number of rows from simply joining the paid orders with the promotions table (/data/shop/promos.csv) by product_id and the number of rows from joining with left_semi to /root/spk/join/out/dup.json as {"naive": 정수, "semi": 정수} (both values are integers).

The promotions table has products with two codes. The orders of such a product become two rows when you join simply. left_semi looks only at "is there a match", so it does not increase the left rows. Think about which one you should use when splitting revenue by whether there was a promotion.

Select the side with no match: left_anti

Create /root/spk/join/anti.py with the application name spk-join-anti and write the customer_id of customers who have never placed an order (regardless of status) as a CSV with a header row to /root/spk/join/out/no_orders.

A not in subquery or filtering nulls after a left join would also work, but left_anti shows the meaning directly and does not get confusing with keys that contain nulls. See whether the join type is printed as LeftAnti in the plan.

Record which strategy fits when

In /root/spk/join/report.md, write three sections: ## 네 가지 전략, ## AQE 의 전환 and ## 키 중복 (use exactly these Korean headings in this order; they mean "Four strategies", "AQE's switch" and "Duplicate keys"). Put the join operator names you saw in steps 1–4 in the first section and the two numbers from step 6 in the third section.

In the first section, write what each strategy shuffles and what it loads into memory; in the second, how the first plan and the final plan differed; and in the third, how many rows were inflated.