Apache Spark — The answer to a slow job is in the plan and the event log
Avoid computing the same thing twice — how a cache shows up in the plan
Goal
When you call two actions on the same DataFrame, you confirm from the plan that without a cache it recomputes from the source, and see how cache, persist, unpersist, checkpoints and CACHE TABLE appear in the plan and the event log. You also measure with block update events how much memory the cache actually eats.
Why it matters
A DataFrame is not a result but a way of making it (the lineage). Every time you call an action, Spark runs that method again from the beginning. Code that uses the same intermediate result twice reads the source twice and does the join twice, without an error or a warning.
cache() is a mark telling Spark to keep that intermediate result on the executors. It is only a mark, so the first action fills it, and subsequent actions read the cache with InMemoryTableScan in the plan instead of reading the source. The storage level (persist) chooses memory, disk and serialization, and decides what to give up when memory runs short. A cache you are done with must be put down with unpersist so that it becomes memory for other work.
If the lineage gets too long (iterative computation), the plan itself grows and the driver slows down. A checkpoint writes the result to files and cuts the lineage, making the plan after it a single line: "read this file".
Steps
- Put a function that attaches an amount to the paid orders and a function that calls two actions (
count,sum(amount)) in /root/spk/cache/common.py, and with /root/spk/cache/none.py (applicationspk-cache-none), write the results of the two actions without a cache to /root/spk/cache/out/none.json. - In /root/spk/cache/cached.py (application
spk-cache-on), applycache()to the same DataFrame and write the results of the two actions to /root/spk/cache/out/cached.json and the storage level (str(df.storageLevel)) to /root/spk/cache/out/cached_level.txt. - In /root/spk/cache/disk.py (application
spk-cache-disk), do the same withpersist(StorageLevel.DISK_ONLY)and write the result to /root/spk/cache/out/disk.json and the storage level to /root/spk/cache/out/level.txt. - In /root/spk/cache/unpersist.py (application
spk-cache-un), run it in the order cache → action →unpersist()→ action, and write the before and after of whether it is cached (df.is_cached) to /root/spk/cache/out/unpersist.json as{"before": 참거짓, "after": 참거짓}(both values are booleans). - In /root/spk/cache/ckpt.py (application
spk-cache-ckpt), set the checkpoint folder to /root/spk/cache/ckpt, and with the result ofcheckpoint()after tenwithColumncalls that add 1 to 10 to the amount in turn, write the sum of the amount per channel as Parquet (channel,amount) to /root/spk/cache/out/ckpt_result. - In /root/spk/cache/cache_sql.py (application
spk-cache-sql), create the temporary viewpaid, runCACHE TABLE paid_cached AS SELECT channel, amount FROM paid, compute the sum per channel, and write to /root/spk/cache/out/cache_table.json as{"is_cached": 참거짓, "by_channel": {채널: 합}}(is_cached is a boolean, and by_channel maps each channel name to its sum). - With /root/spk/cache/size.py (application
spk-cache-size,spark.eventLog.logBlockUpdates.enabled=true), fill the cache, then add up therdd_block sizes in that log and write them to /root/spk/cache/out/cache_size.json as{"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수}(all three values are integers). - In /root/spk/cache/report.md, write three sections:
## 다시 계산하는 비용,## 저장 수준and## 캐시의 크기(use exactly these Korean headings in this order; they mean "Cost of recomputing", "Storage level" and "Size of the cache"). In the third section, put the number of blocks and the memory bytes from step 7.
Notes
- Put the scripts in
/root/spk/cacheand run them from there (from common import …). The original data is/data/shop/orders.csvand/data/shop/products.csv, and amount = qty × price. - The plan of an execution that reads the cache shows
InMemoryTableScanand, beneath it,InMemoryRelation … StorageLevel(…). The source scan (Scan csv) remains inside the InMemoryRelation description but is not executed. - The same block may appear several times in block update events (
SparkListenerBlockUpdated). Keep only the last one for each block ID and add them up. - Common mistakes: doing only
cache()without calling an action, so the cache is empty; not putting down the cache after you are done with it; and caching both the original and a derived DataFrame within one application so that it becomes unclear which cache was read. - Official documentation: RDD Persistence · Caching Data · CACHE TABLE · Configuration — spark.eventLog.logBlockUpdates.enabled
Two actions without a cache: the source is read twice
In /root/spk/cache/common.py, put a function that joins the paid orders with the products and attaches amount (qty×price), and a function that calls the two actions count() and sum(amount) on a DataFrame and writes {"count": 정수, "sum": 정수} (two integers) to a file. Create /root/spk/cache/none.py with the application name spk-cache-none and write the result without a cache to /root/spk/cache/out/none.json.
Two actions mean two SQL executions, and both have Scan csv at the very bottom of the plan. This means the same file is read twice and the same join is done twice. The grader checks whether there are two or more executions that read the source without a cache.
cache(): the first action fills it and the next one reads it
Create /root/spk/cache/cached.py with the application name spk-cache-on, apply cache() to the step 1 DataFrame, and write the results of the same two actions to /root/spk/cache/out/cached.json and str(df.storageLevel) to /root/spk/cache/out/cached_level.txt.
InMemoryTableScan appears in the plan of the second execution. Also notice that the default storage level of a DataFrame differs from the cache of an RDD: if memory runs short, it spills to disk. The result numbers must be the same as in step 1.
persist(DISK_ONLY): choosing the storage level
Create /root/spk/cache/disk.py with the application name spk-cache-disk, keep the step 1 DataFrame with persist(StorageLevel.DISK_ONLY), and write the results of the same two actions to /root/spk/cache/out/disk.json and str(df.storageLevel) to /root/spk/cache/out/level.txt.
Keeping it only on disk saves executor memory, but it needs deserialization every time you read it. See that the StorageLevel printed on the InMemoryRelation line of the plan has changed. Even with the name InMemory, it can be on disk.
unpersist: if you put the cache down, it is recomputed
Create /root/spk/cache/unpersist.py with the application name spk-cache-un, run the step 1 DataFrame in the order cache() → count() → unpersist() → sum(amount), and write df.is_cached before and after unpersist() to /root/spk/cache/out/unpersist.json as {"before": 참거짓, "after": 참거짓} (both values are booleans).
A cache is not free. It takes executor memory and reduces the space for other work. Put down a cache you are done with. The grader checks from the log whether an execution that does not read the cache follows the execution that read the cache.
Cut the lineage with a checkpoint
Create /root/spk/cache/ckpt.py with the application name spk-cache-ckpt, set the checkpoint folder to /root/spk/cache/ckpt, and with the result of checkpoint() after ten withColumn calls that add 1, 2, …, 10 in turn to amount of the step 1 DataFrame, write the sum of amount per channel as Parquet (channel,amount) to /root/spk/cache/out/ckpt_result.
checkpoint() writes the result to the checkpoint folder on the spot, and the plan of the DataFrame it returns becomes a single line that reads that file (Scan ExistingRDD) rather than the source. A cache keeps the lineage (if lost, it is recomputed), and a checkpoint cuts the lineage (the file is the source).
SQL's CACHE TABLE is not lazy
Create /root/spk/cache/cache_sql.py with the application name spk-cache-sql, register the step 1 DataFrame as the temporary view paid, run CACHE TABLE paid_cached AS SELECT channel, amount FROM paid, then compute the sum of amount per channel from paid_cached and write to /root/spk/cache/out/cache_table.json as {"is_cached": spark.catalog.isCached("paid_cached"), "by_channel": {채널: 합}} (by_channel maps each channel name to its sum).
The cache() of a DataFrame only marks, but CACHE TABLE fills it on the spot (with CACHE LAZY TABLE, it defers). So the CACHE TABLE statement itself launches a job. The plan that reads a named cache prints Scan In-memory table paid_cached. The grader checks whether such an execution exists and the sums.
Measure the memory the cache actually eats
Create /root/spk/cache/size.py with the application name spk-cache-size and the setting spark.eventLog.logBlockUpdates.enabled=true, apply cache() to the step 1 DataFrame and fill it with count(). From that application's log, among the SparkListenerBlockUpdated events, keep only the last one per block for those whose block ID starts with rdd_, and write the count and the sums of Memory Size and Disk Size to /root/spk/cache/out/cache_size.json as {"blocks": 정수, "memory_bytes": 정수, "disk_bytes": 정수} (all three values are integers).
There is one cache block per partition (rdd_<RDD 번호>_<파티션>, where the placeholders are the RDD number and the partition). A cache unpacked in memory (deserialized) can be larger or smaller than the original CSV, depending on the chosen columns and types. The habit of looking at this number before growing a cache protects executor memory.
Where to put a cache and where to put it down
In /root/spk/cache/report.md, write three sections: ## 다시 계산하는 비용, ## 저장 수준 and ## 캐시의 크기 (use exactly these Korean headings in this order; they mean "Cost of recomputing", "Storage level" and "Size of the cache"). Put the storage level strings from steps 2 and 3 in the second section and the number of blocks and the memory bytes from step 7 in the third section.
In the first section, write what was done twice when there was no cache; in the second, what the two storage levels trade off; and in the third, what percent of the executor heap (1g) that memory is.